Repository navigation
Streaming Support for msgpack #47
Closed
SeanTAllen
started this conversation in
Research
Replies: 1 comment
|
Went with approach D. |
0 replies
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Uh oh!
There was an error while loading. Please reload this page.
Uh oh!
There was an error while loading. Please reload this page.
Ideas for making the decoder safe for streaming (incremental data
arrival). See issue #14.
Design Constraints
buffered.Readerin ponycare off the table.
package is acceptable.
MessagePackEncoderandMessagePackDecoderprimitives should remain as-is for non-streaminguse. Streaming support should be additive.
Key Facts About
buffered.ReaderThese details about the stdlib Reader inform which approaches are
viable:
u8(),u16_be(),block(), etc. permanentlyadvance the cursor and discard consumed bytes. No undo.
peek_u8(offset),peek_u16_be(offset),peek_u32_be(offset),peek_u64_be(offset),and their little-endian/signed variants. These read without advancing
the cursor.
peek_block: You can peek at fixed-width integers but not atarbitrary byte ranges. This is fine for the size-check strategy
(you only need to peek at format bytes and length fields, not
payloads) but means you can't fully preview variable-length data
without consuming it.
size()returns available bytes: Useful for checking whetherenough data is buffered before attempting a decode.
restore it after a failed partial read.
The Problem
The decoder does multiple consuming reads from the
Readerper value.For example,
MessagePackDecoder.str()reads:If step 1 succeeds but step 3 fails because the payload hasn't arrived
yet, the Reader has already consumed the type and length bytes. They're
gone. When more data arrives and the caller retries, the stream is
corrupted.
The
Readerhas no rollback mechanism — eachu8(),u16_be(),block()call permanently advances the cursor.Cross-Cutting Concern: Error Discrimination
All approaches must solve this: the current API uses
?for both "notenough data" and "invalid data." A streaming layer must distinguish
these because the caller's response is completely different — wait for
more data vs. abort.
Options for expressing this:
(T | NotEnoughData | InvalidData)— explicit,pattern-matchable, no ambiguity. Most idiomatic for Pony.
become method calls too. Natural for push-based streaming.
decoding. Keeps existing
?semantics for actual errors. But createsa check-then-act race if used incorrectly (not really a race in Pony's
single-actor model, but it's two calls that must stay in sync).
Approach A: Peek-and-Check Wrapper
A new primitive (or small class) that peeks at the Reader to calculate
how many bytes the next value needs, checks
reader.size(), and onlydelegates to the existing
MessagePackDecoderif sufficient data isavailable.
How it works: For every MessagePack format, the total byte count is
deterministic from the format byte and (for variable-length types) the
length field:
alone (1-9 bytes total).
bytes using
Reader.peek_u8/peek_u16_be/peek_u32_beto determinetotal size without consuming anything.
bytes). The contained items are the caller's problem — same as today.
If enough bytes: call the corresponding
MessagePackDecodermethod.If not: return
NotEnoughData(no bytes consumed).Pros:
MessagePackDecoderstays untouched.Cons:
streaming_str(),streaming_u32(), etc.). Doesn't solve the"read the next value regardless of type" use case.
decoder. Could drift if the decoder changes.
decoding an array's items and run out of data, you need your own
bookkeeping for where you left off.
Approach B: Custom Reader with Snapshot/Rollback
Create a reader class within the msgpack package that supports
transactional reads. The existing decoder methods expect
buffered.Reader, so this approach also requires new decoder methods(or a new decoder type) that accept the custom reader.
How it works: The custom reader wraps incoming byte chunks and
maintains a cursor. Key addition:
save()returns a snapshot of thecursor position, and
restore(snapshot)resets back to it. Bytesaren't discarded until explicitly committed (or until the snapshot is
dropped).
Pros:
to use the custom reader type).
sizes.
high-level API mentioned in the README).
Cons:
both raise
error. After a rollback, the caller doesn't know whetherto wait for more data or give up. Would need a secondary mechanism
(e.g., a
last_errorfield on the reader, or wrapping the decoderto peek-validate the format byte before attempting the full read).
buffered.Readeris ~800 lines with careful chunk management andendianness handling.
ones, just accepting a different reader type.
Variant — interface extraction: If the custom reader and
buffered.Readershared a common interface, the existing decoder couldpotentially accept either. But
buffered.Readeris a concrete class instdlib and doesn't implement a reading interface. This would require
either structural typing compatibility or stdlib changes (ruled out).
Approach C: State Machine Streaming Decoder
A class that accepts data incrementally and emits complete decoded
values via a notify interface. Internally tracks parse state so it can
pause mid-value and resume when more data arrives.
How it works:
The decoder class holds a
Readerinternally and maintains a stateenum (conceptually):
_AwaitingType— need at least 1 byte_ReadingFixedValue(format, bytes_needed)— have the type, needN more bytes for a fixed-size value
_ReadingLength(format, length_bytes_needed)— need the lengthfield before we know the payload size
_ReadingPayload(format, remaining)— know the full size, waitingfor payload bytes
When
append()is called with new data, the decoder runs a loop:check state, check if enough bytes are available (via
reader.size()),if yes consume and emit via notify, advance to next state. Stop when
insufficient data.
Pros:
callbacks. No corruption possible.
data" (
on_errorcallback).notify interface gives the caller typed values without needing to know
MessagePack format details.
consumes bytes it can fully process.
Cons:
MessagePackDecoderat all.synchronous pull-based API. Users who want "decode the next value"
synchronously need to adapt.
notify callbacks to handle nesting, or the decoder needs an internal
stack. This adds complexity.
Approach D: Peek-Check + Pull API (Hybrid of A and C)
Combine approach A's peek-and-check with a unified "read next value"
method. A class that wraps a
Reader, peeks to determine the type andrequired bytes, and either returns a decoded value or
NotEnoughData.How it works:
The class has a
next(): DecodeResultmethod that:NotEnoughDatalength fields for variable-length types)
NotEnoughDataPros:
"invalid data."
when all bytes are present.
Cons:
MessagePackValueunion type loses some specificity — you getback a union and need to match on it. The current API gives you the
exact type you asked for.
twice (once to peek, once during actual decode). Minor overhead.
contents are still the caller's responsibility to track.
both return U8 but are different formats. The union collapses them.
This might or might not matter — need to decide whether the value
type represents the Pony type or the MessagePack format.
Nested Structures: A Shared Challenge
Arrays and maps in MessagePack are headers followed by N values (or
2*N for maps). The current decoder only reads headers — the caller is
responsible for then reading N items. In a streaming context, you might
get the array header but not all items.
Options:
handles individual values. The caller tracks "I'm 3 items into a
5-item array" themselves. Simple for the library, more work for
users.
tracks an internal stack of containers. It doesn't emit an array
until all its items are decoded. This gives the caller complete
values but means the decoder must buffer an entire array/map in
memory before emitting — which could be arbitrarily large.
on_array_start(5)... five value callbacks ...on_array_end(). SAX-style event stream. The caller gets structureawareness without the decoder needing to buffer.
Recommendation
Not making a recommendation yet — this document is for discussion. But
some observations:
streaming-safe. They work but don't move toward the high-level API
that's already identified as a TODO.
the beginnings of a high-level API.
simplicity and capability. It's streaming-safe, provides a high-level
"read next value" API, and is simpler than a full state machine.
C and D, where the return type itself carries the distinction.
All reactions