Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
115 changes: 87 additions & 28 deletions reader.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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
}
Expand Down Expand Up @@ -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) {
Expand All @@ -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
Expand Down
140 changes: 138 additions & 2 deletions reader_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Expand Down Expand Up @@ -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 {
Expand Down