20
21
// ParseDLQMessages parses a JSONL file containing serialized [tdbg.DLQMessage] objects.
22
>
func ParseDLQMessages[T proto.Message](file io.Reader, newMessage func() T) ([]DLQMessage[T], error) {
output_parsing.go
23
>
var opts temporalproto.CustomJSONUnmarshalOptions
24
>
decodeNext := func(decoder *json.Decoder) (DLQMessage[T], error) {
25
>
var dlqMessage tdbg.DLQMessage
26
>
err := decoder.Decode(&dlqMessage)
27
>
if err != nil {
28
return DLQMessage[T]{}, err
29
}
31
>
b := dlqMessage.Payload.Bytes()
32
>
if err = opts.Unmarshal(b, protoMessage); err != nil {
33
return DLQMessage[T]{}, err
34
}
36
>
MessageID: dlqMessage.MessageID,
37
>
ShardID: dlqMessage.ShardID,
38
>
Payload: protoMessage,
39
>
}, nil
40
}
42
}
43
44
// ParseJSONL parses a JSONL file. We separate this out from [ParseDLQMessages] so that we can reuse it for other JSONL
45
// files that don't contain [tdbg.DLQMessage] objects (i.e. when [tdbg.DLQV1Service] is used).
46
>
func ParseJSONL[T any](file io.Reader, decodeNext func(decoder *json.Decoder) (T, error)) ([]T, error) {
output_parsing.go
47
>
decoder := json.NewDecoder(file)
48
>
var (
49
>
messages []T
50
>
)
51
>
for decoder.More() {
52
>
message, err := decodeNext(decoder)
53
>
if err != nil {
54
return nil, err
55
}
57
}
59
}