rust: Add a MessageReader implementation - #1854
Open
clalancette wants to merge 11 commits into
Open
clalancette wants to merge 11 commits into
clalancette wants to merge 11 commits into
Conversation
MessageStream only reads from a byte slice, so consumers following the streaming-reader guidance had to hand-roll schema/channel linking; the CLI's cat and merge and the crate's examples each carried a copy. Add sans_io::MessageReader, which composes LinearReader with the existing ChannelAccumulator so it applies the same validation as MessageStream (schema id 0, conflicting redefinitions, unknown channels), validates chunk CRCs by default, yields owned messages, and stops after the first error. Add io::MessageReader<R: Read> as the blocking Iterator adapter, the counterpart of MessageStream for a file or socket. Expose get_channel and channels() so callers can see channel definitions, including message-less ones. Tests cover parity with MessageStream on chunked and unchunked input, short reads, truncation ending the iterator even under flatten, source I/O errors surfacing once, CRC validation on by default and off by option, schema id 0, and unknown channels. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
The example drove LinearReader by hand and printed raw records, losing the linked messages the old mmap example showed. Use the new streaming MessageReader instead: shorter, and message-level again. Point the read module header at it as MessageStream's streaming counterpart. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
The CLI and the examples each carried copies of the crate's schema/channel bookkeeping. Expose read::ChannelAccumulator with add_schema, add_channel, get, and channels(), plus from_summary so callers reading chunks after the summary can check in-chunk definitions against it. sans_io::MessageReader delegates to it. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
The examples carried their own message iterator in a shared helper because the crate had none. Now that mcap::io::MessageReader exists, mcapcat, mcapcopy, and recover iterate it directly, which also gives them the crate's full schema/channel validation. The helper shrinks to the summary reader and is renamed common/summary.rs to match. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
clalancette
requested review from
bennetthardwick,
gasmith and
james-rms
as code owners
September 30, 2026 20:50
The reader kept going through the summary section, re-checking the repeated schema and channel records and failing on a summary that was damaged or cut off, while its docs said it stopped at the end of the data section. Return end-of-stream at DataEnd. Test that a file cut inside its summary still yields every message without an error and that an intact file's summary bytes are never requested. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Read::read may return ErrorKind::Interrupted, which the Read contract says to retry; the adapter ended the iteration on it instead, so a signal during a read from a pipe or socket stopped a healthy stream. Retry before notify_read, since notifying zero bytes would signal EOF to the linear reader. Test with a source that interrupts every other read. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
new() turned on chunk CRC validation while new_with_options(default()) did not, so the two constructors disagreed about what "default" means. Every other reader in the crate has new() == new_with_options(default()). Make this one match, and update the docs and test accordingly. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Nothing called ChannelAccumulator::channels() or the reader's delegating channels(), so remove both. Add a test for from_summary, which the CLI relies on: a chunk channel that matches the summary's is accepted and one that conflicts is rejected. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
No code changes. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
gasmith
approved these changes
Sep 30, 2026
| /// Streams linked [`Message`]s from any [`Read`] source in file order; memory scales with the | ||
| /// largest record or chunk, not the file. | ||
| /// | ||
| /// The streaming counterpart of [`crate::MessageStream`]: same schema/channel validation, |
Contributor
There was a problem hiding this comment.
Nit: I had a hard time understanding how this could be "the streaming counterpart of MessageStream", until I realized that we're referring to the Read source as the stream here. We might want to reword this a bit.
Might also want to duplicate MessageStream's explanation of when we take copies of message data.
This branch was successfully deployed
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Changelog
Rust: add
mcap::io::MessageReader, a structure that reads linked messages from anystd::io::Read, andmcap::sans_io::MessageReader, the sans-io core it is built on.Docs
None.
Description
#1853 pointed users at the sans-io readers, but those only yield raw
records: there was no way to get linked
Messages except from a byteslice via
MessageStream. Every consumer following the new guidance hadto reimplement schema/channel linking, and the examples and the CLI each
carried a copy.
sans_io::MessageReadercomposesLinearReaderwith the crate'sexisting
ChannelAccumulator, so it applies the same validation asMessageStream(schema ID 0, conflicting redefinitions, unknownchannels), validates chunk CRCs by default, yields owned messages, and
ends after the first error. It reads nothing itself; it emits read
requests like the other sans-io readers.
io::MessageReader<R: Read>owns aReadsource, drives that core, andis an
Iterator<Item = McapResult<Message>>. It ends after the firsterror, including a source I/O error, so
.flatten()on a truncated fileterminates.
read::ChannelAccumulatoris now public, withadd_schema,add_channel,get,channels(), andfrom_summary(for callers thatread chunks after the summary). Callers that need record-level control
can link messages without copying the rules; the CLI switches to it in
cli: replace mmap I/O with sans-io ByteSource #1834.
Split out of #1834, which will be rebased on top of this once it lands. @gasmith This is the follow-up to #1853 I alluded to.