diff --git a/README.md b/README.md index ca994b4..7cb2952 100644 --- a/README.md +++ b/README.md @@ -42,9 +42,10 @@ The pinned environment, test implementation, fixture, and provenance live togeth - streaming and eager reads - typed and dynamic writing - information and multi-information, parameters and defaults, logging, and dropouts +- reading appended data sections, including sections following an interrupted message - CSV, Arrow record batches, and Parquet output -Appended data sections are rejected rather than silently misread. Multi-information, default-parameter, and tagged-log writing are available at the lower-level `pkg/wire` boundary but do not yet have root-package writer conveniences. +Multi-information, default-parameter, and tagged-log writing are available at the lower-level `pkg/wire` boundary but do not yet have root-package writer conveniences. ## License diff --git a/reader.go b/reader.go index d6c78f4..3f5903f 100644 --- a/reader.go +++ b/reader.go @@ -26,19 +26,21 @@ type Header struct { // advances, it also collects information, multi-information groups, parameters, // logs, and dropouts. type Reader struct { - source io.Reader - header Header - formats map[string]Format - subscriptions map[uint16]subscription - information []KeyValue - multiInfo []MultiInformationGroup - multiInfoLast map[string]int - parameters []KeyValue - defaults []DefaultParameter - logs []LogEntry - dropouts []Dropout - record Record - err error + source io.Reader + offset uint64 + appendedOffsets []uint64 + header Header + formats map[string]Format + subscriptions map[uint16]subscription + information []KeyValue + multiInfo []MultiInformationGroup + multiInfoLast map[string]int + parameters []KeyValue + defaults []DefaultParameter + logs []LogEntry + dropouts []Dropout + record Record + err error } type subscription struct { @@ -79,6 +81,7 @@ func NewReader(source io.Reader) (*Reader, error) { } return &Reader{ source: source, + offset: uint64(len(data)), header: Header{ Version: fileHeader.Version, Timestamp: fileHeader.Timestamp, @@ -106,7 +109,7 @@ func (r *Reader) Next() bool { } for { - messageType, payload, err := readMessage(r.source) + messageType, payload, err := r.readMessage() if errors.Is(err, io.EOF) { return false } @@ -197,22 +200,65 @@ func (r *Reader) Dropouts() []Dropout { return append([]Dropout(nil), r.dropouts...) } -func readMessage(source io.Reader) (wire.MessageType, []byte, error) { - var headerBytes [3]byte - n, err := io.ReadFull(source, headerBytes[:]) - if errors.Is(err, io.EOF) && n == 0 { - return 0, nil, io.EOF - } - if err != nil { - return 0, nil, fmt.Errorf("read ULog message header: %w", err) +func (r *Reader) readMessage() (wire.MessageType, []byte, error) { + for { + if r.atAppendedSection() { + r.appendedOffsets = r.appendedOffsets[1:] + } + + var headerBytes [3]byte + if r.crossesAppendedOffset(uint64(len(headerBytes))) { + if err := r.skipToAppendedSection(); err != nil { + return 0, nil, err + } + continue + } + n, err := r.readFull(headerBytes[:]) + if errors.Is(err, io.EOF) && n == 0 { + return 0, nil, io.EOF + } + if err != nil { + return 0, nil, fmt.Errorf("read ULog message header: %w", err) + } + + size := binary.LittleEndian.Uint16(headerBytes[:2]) + if r.crossesAppendedOffset(uint64(size)) { + if err := r.skipToAppendedSection(); err != nil { + return 0, nil, err + } + continue + } + payload := make([]byte, int(size)) + if _, err := r.readFull(payload); err != nil { + return 0, nil, fmt.Errorf("read ULog %q message payload: %w", headerBytes[2], err) + } + return wire.MessageType(headerBytes[2]), payload, nil } +} + +func (r *Reader) readFull(data []byte) (int, error) { + n, err := io.ReadFull(r.source, data) + r.offset += uint64(n) // #nosec G115 -- io.ReadFull returns a non-negative byte count. + return n, err +} + +func (r *Reader) atAppendedSection() bool { + return len(r.appendedOffsets) > 0 && r.offset == r.appendedOffsets[0] +} - size := binary.LittleEndian.Uint16(headerBytes[:2]) - payload := make([]byte, int(size)) - if _, err := io.ReadFull(source, payload); err != nil { - return 0, nil, fmt.Errorf("read ULog %q message payload: %w", headerBytes[2], err) +func (r *Reader) crossesAppendedOffset(size uint64) bool { + return len(r.appendedOffsets) > 0 && size > r.appendedOffsets[0]-r.offset +} + +func (r *Reader) skipToAppendedSection() error { + remaining := r.appendedOffsets[0] - r.offset + // crossesAppendedOffset limits remaining to a three-byte header or uint16 payload. + if _, err := io.CopyN(io.Discard, r.source, int64(remaining)); err != nil { // #nosec G115 -- bounded above. + return fmt.Errorf("skip interrupted ULog message at appended offset %d: %w", r.appendedOffsets[0], err) } - return wire.MessageType(headerBytes[2]), payload, nil + r.offset += remaining + r.appendedOffsets = r.appendedOffsets[1:] + return nil } func (r *Reader) consume(messageType wire.MessageType, payload []byte) (Record, bool, error) { @@ -231,7 +277,20 @@ func (r *Reader) consume(messageType wire.MessageType, payload []byte) (Record, return Record{}, false, fmt.Errorf("unsupported ULog incompatibility flags %#x", unknown) } if flags.IncompatibilityFlags&wire.IncompatibilityFlagDataAppended != 0 { - return Record{}, false, errors.New("ULog appended data sections are not supported") + previous := r.offset + for _, offset := range flags.AppendedOffsets { + if offset == 0 { + break + } + if offset <= previous { + return Record{}, false, fmt.Errorf("invalid ULog appended offset %d after offset %d", offset, previous) + } + r.appendedOffsets = append(r.appendedOffsets, offset) + previous = offset + } + if len(r.appendedOffsets) == 0 { + return Record{}, false, errors.New("ULog data-appended flag has no appended offsets") + } } case wire.MessageTypeFormat: var message wire.FormatMessage diff --git a/reader_test.go b/reader_test.go index 044e5fa..e31dba8 100644 --- a/reader_test.go +++ b/reader_test.go @@ -222,6 +222,121 @@ func TestReaderAcceptsFutureVersionAndExtendedFlagBits(t *testing.T) { } } +func TestReaderContinuesAtAppendedDataAfterInterruptedMessage(t *testing.T) { + data := newULogFixture(t, 0) + data.message(t, wire.MessageTypeFormat, wire.FormatMessage{Format: "sample:uint64_t timestamp;uint16_t value;"}) + data.message(t, wire.MessageTypeSubscription, wire.SubscriptionMessage{MessageID: 1, MessageName: "sample"}) + + data.rawMessageHeader(t, wire.MessageTypeData, 12) + data.data = append(data.data, 1, 0, 0xff) + appendedOffset := uint64(len(data.data)) + + data.message(t, wire.MessageTypeData, wire.DataMessage{MessageID: 1, Data: samplePayload(99, 42)}) + data.setAppendedOffsets(t, appendedOffset) + + reader, err := NewReader(bytes.NewReader(data.bytes())) + if err != nil { + t.Fatalf("NewReader() error = %v", err) + } + if !reader.Next() { + t.Fatalf("Next() = false, error = %v", reader.Err()) + } + + assertRecordValue(t, reader.Record(), 42) + if reader.Next() { + t.Fatal("second Next() = true, want false") + } + if err := reader.Err(); err != nil { + t.Fatalf("Err() = %v", err) + } +} + +func TestReaderContinuesAcrossMultipleAppendedSections(t *testing.T) { + data := newULogFixture(t, 0) + data.message(t, wire.MessageTypeFormat, wire.FormatMessage{Format: "sample:uint64_t timestamp;uint16_t value;"}) + data.message(t, wire.MessageTypeSubscription, wire.SubscriptionMessage{MessageID: 1, MessageName: "sample"}) + + data.rawMessageHeader(t, wire.MessageTypeData, 12) + data.data = append(data.data, 1, 0) + firstOffset := uint64(len(data.data)) + data.message(t, wire.MessageTypeData, wire.DataMessage{MessageID: 1, Data: samplePayload(10, 1)}) + + data.data = append(data.data, 0xff, 0xff) + secondOffset := uint64(len(data.data)) + data.message(t, wire.MessageTypeData, wire.DataMessage{MessageID: 1, Data: samplePayload(20, 2)}) + data.setAppendedOffsets(t, firstOffset, secondOffset) + + reader, err := NewReader(bytes.NewReader(data.bytes())) + if err != nil { + t.Fatalf("NewReader() error = %v", err) + } + for i, want := range []uint16{1, 2} { + if !reader.Next() { + t.Fatalf("Next() for record %d = false, error = %v", i, reader.Err()) + } + assertRecordValue(t, reader.Record(), want) + } + if reader.Next() { + t.Fatal("third Next() = true, want false") + } + if err := reader.Err(); err != nil { + t.Fatalf("Err() = %v", err) + } +} + +func TestReaderRejectsInvalidAppendedOffsets(t *testing.T) { + tests := map[string]struct { + offsets []uint64 + wantMessage string + }{ + "missing": { + wantMessage: "has no appended offsets", + }, + "before flag bits": { + offsets: []uint64{1}, + wantMessage: "invalid ULog appended offset 1", + }, + "not increasing": { + offsets: []uint64{100, 100}, + wantMessage: "invalid ULog appended offset 100 after offset 100", + }, + } + + for name, tt := range tests { + t.Run(name, func(t *testing.T) { + data := newULogFixture(t, 0) + data.setAppendedOffsets(t, tt.offsets...) + + reader, err := NewReader(bytes.NewReader(data.bytes())) + if err != nil { + t.Fatalf("NewReader() error = %v", err) + } + if reader.Next() { + t.Fatal("Next() = true, want false") + } + if err := reader.Err(); err == nil || !strings.Contains(err.Error(), tt.wantMessage) { + t.Fatalf("Err() = %v, want error containing %q", err, tt.wantMessage) + } + }) + } +} + +func samplePayload(timestamp uint64, value uint16) []byte { + payload := binary.LittleEndian.AppendUint64(nil, timestamp) + return binary.LittleEndian.AppendUint16(payload, value) +} + +func assertRecordValue(t *testing.T, record Record, want uint16) { + t.Helper() + value, err := record.Value("value") + if err != nil { + t.Fatalf("Value(value) error = %v", err) + } + if value != want { + t.Errorf("Value(value) = %#v, want %#v", value, want) + } +} + type ulogFixture struct { data []byte } @@ -263,13 +378,34 @@ func (f *ulogFixture) rawMessage(t *testing.T, messageType wire.MessageType, pay if len(payload) > math.MaxUint16 { t.Fatalf("payload size %d exceeds ULog limit", len(payload)) } - header := wire.MessageHeader{Size: uint16(len(payload)), Type: messageType} // #nosec G115 -- bounded above. + f.rawMessageHeader(t, messageType, uint16(len(payload))) // #nosec G115 -- bounded above. + f.data = append(f.data, payload...) +} + +func (f *ulogFixture) rawMessageHeader(t *testing.T, messageType wire.MessageType, size uint16) { + t.Helper() + header := wire.MessageHeader{Size: size, Type: messageType} var err error f.data, err = binary.Append(f.data, binary.LittleEndian, header) if err != nil { t.Fatalf("append message header: %v", err) } - f.data = append(f.data, payload...) +} + +func (f *ulogFixture) setAppendedOffsets(t *testing.T, offsets ...uint64) { + t.Helper() + if len(offsets) > len(wire.FlagBitsMessage{}.AppendedOffsets) { + t.Fatalf("got %d appended offsets, maximum is %d", len(offsets), len(wire.FlagBitsMessage{}.AppendedOffsets)) + } + + flags := wire.FlagBitsMessage{IncompatibilityFlags: wire.IncompatibilityFlagDataAppended} + copy(flags.AppendedOffsets[:], offsets) + payload, err := binary.Append(nil, binary.LittleEndian, flags) + if err != nil { + t.Fatalf("append flag bits: %v", err) + } + payloadOffset := binary.Size(wire.FileHeader{}) + binary.Size(wire.MessageHeader{}) + copy(f.data[payloadOffset:payloadOffset+len(payload)], payload) } func (f *ulogFixture) bytes() []byte {