Skip to content
This repository was archived by the owner on Oct 1, 2026. It is now read-only.
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
15 changes: 15 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,21 @@ and this project adheres to no released version yet.

### Added

- **Inbound request continuation (CTX-0084, bitty#1482)**: live and headless
requests above one `256 KiB` frame, up to the `1 MiB` inbound limit, are sent
as devtools-rfc Amendment A4 continuation fragments. Each fragment carries a
16-byte `\0BC1` header with a nonzero continuation id, a sequence, a FINAL
flag, and the declared total. Every non-final fragment is a full frame, so a
request is at most five fragments. The core reassembles them into one
exchange with one admission. This replaces the headerless split and the live
`FrameTooLarge` fail-closed. The TypeScript live socket now continues a
partial write on `drain` within a bounded outbound queue, instead of failing
the connection. A fragmented request needs an idle connection and blocks
other requests until it settles, because the server rejects any frame that
interleaves a reassembly. The Rust `devtools-client` sends the same fragments
with byte-identical headers. Requests above `1 MiB` still fail closed before
any write.

- **Test-automation drivers (CTX-0024, candidate)**: typed, fail-closed client
bindings in `src/automation.ts` for `bitty.debug/synthesizeInput`,
`bitty.debug/captureFrame`, and `bitty.debug/frameHash` — the headless
Expand Down
35 changes: 21 additions & 14 deletions crates/devtools-client/src/ipc_socket.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,9 @@ use crate::auth::{AuthError, MAX_SOCKET_PATH_BYTES, resolve_socket_path};
use crate::auth::{DIR_MODE, SOCKET_MODE};
use crate::transport::TransportError;
#[cfg(target_os = "linux")]
use crate::transport::{MAX_FRAME_BYTES, decode_frame, encode_frame};
use crate::transport::{
MAX_FRAME_BYTES, decode_frame, encode_request_frames, next_continuation_id,
};

/// Per-dial and per-response timeout (matches the TS seam).
pub const LIVE_SOCKET_TIMEOUT_SECS: u64 = 5;
Expand Down Expand Up @@ -232,6 +234,8 @@ pub struct LiveSocketConnection {
stream: std::os::unix::net::UnixStream,
socket_path: String,
identity: LiveSocketIdentity,
/// Next Amendment A4 continuation id (nonzero; wraps to 1).
continuation_id: u32,
}

#[cfg(target_os = "linux")]
Expand All @@ -247,10 +251,12 @@ impl LiveSocketConnection {
self.identity
}

/// Write one framed request and read the next framed response payload
/// (raw bytes, still to be JSON-decoded by the caller). Bounded at one
/// 256 KiB frame each way, matching the `bitty-ipc` framing. Times out
/// instead of blocking forever.
/// Write one request and read the next framed response payload (raw
/// bytes, still to be JSON-decoded by the caller). A request above one
/// 256 KiB frame, up to the 1 MiB inbound limit, is written as
/// Amendment A4 continuation fragments (bitty#1482); the round trip is
/// synchronous, so no other request can interleave them. The response
/// stays bounded at one frame. Times out instead of blocking forever.
pub fn request_response(
&mut self,
request_json: &[u8],
Expand All @@ -259,22 +265,22 @@ impl LiveSocketConnection {
use std::io::{Read, Write};
use std::time::Duration;

if request_json.len() > MAX_FRAME_BYTES {
return Err(TransportError::FrameTooLarge {
actual: request_json.len(),
limit: MAX_FRAME_BYTES,
});
let frames = encode_request_frames(request_json, self.continuation_id)?;
if frames.len() > 1 {
self.continuation_id = next_continuation_id(self.continuation_id);
}
self.stream
.set_read_timeout(Some(Duration::from_secs(LIVE_SOCKET_TIMEOUT_SECS)))
.map_err(|_| TransportError::TransportClosed)?;
self.stream
.set_write_timeout(Some(Duration::from_secs(LIVE_SOCKET_TIMEOUT_SECS)))
.map_err(|_| TransportError::TransportClosed)?;
let wire = encode_frame(request_json)?;
self.stream
.write_all(&wire)
.map_err(|_| TransportError::TransportClosed)?;
// `write_all` retries partial writes; fragments go out back to back.
for frame in &frames {
self.stream
.write_all(frame)
.map_err(|_| TransportError::TransportClosed)?;
}
self.stream
.flush()
.map_err(|_| TransportError::TransportClosed)?;
Expand Down Expand Up @@ -326,6 +332,7 @@ pub fn connect_live_socket(
.map_err(|_| TransportError::TransportClosed)?;
Ok(LiveSocketConnection {
stream,
continuation_id: 1,
socket_path: endpoint.socket_path.clone(),
identity: LiveSocketIdentity {
runtime_uid: endpoint.runtime_uid,
Expand Down
95 changes: 89 additions & 6 deletions crates/devtools-client/src/transport.rs
Original file line number Diff line number Diff line change
Expand Up @@ -114,6 +114,77 @@ pub fn encode_frame(payload: &[u8]) -> Result<Vec<u8>, TransportError> {
Ok(out)
}

// ── inbound request continuation (devtools-rfc Amendment A4, bitty#1482) ──

/// Magic prefix of every continuation fragment payload (`\0BC1`).
pub const CONTINUATION_MAGIC: [u8; 4] = [0x00, b'B', b'C', b'1'];
/// Fragment header bytes before each chunk.
pub const CONTINUATION_HEADER_BYTES: usize = 16;
/// Chunk bytes in every non-final fragment (one full physical frame).
pub const CONTINUATION_CHUNK_BYTES: usize = MAX_FRAME_BYTES - CONTINUATION_HEADER_BYTES;
/// `flags` bit of the fragment that completes the logical request.
pub const CONTINUATION_FLAG_FINAL: u8 = 0b0000_0001;

/// Encode one logical request as wire frames (length prefix included).
///
/// A request of at most [`MAX_FRAME_BYTES`] stays one plain frame. A larger
/// one, up to the 1 MiB inbound limit, becomes canonical Amendment A4
/// fragments tagged `continuation_id`: every non-final fragment is a full
/// frame and the final one carries the remainder (at most five fragments).
///
/// # Errors
///
/// `PayloadTooLarge` above [`RC9_PAYLOAD_CAP_BYTES`] and `InvalidFrame` for
/// a zero `continuation_id` when fragmentation is needed. No frame is
/// produced on error.
pub fn encode_request_frames(
request: &[u8],
continuation_id: u32,
) -> Result<Vec<Vec<u8>>, TransportError> {
check_payload_cap(request.len())?;
if request.len() <= MAX_FRAME_BYTES {
return Ok(vec![encode_frame(request)?]);
}
if continuation_id == 0 {
return Err(TransportError::InvalidFrame(
"continuation id must be nonzero".into(),
));
}
// The payload cap above keeps the total within u32 (1 MiB).
let total = u32::try_from(request.len()).map_err(|_| TransportError::PayloadTooLarge {
field: "request".into(),
limit: RC9_PAYLOAD_CAP_BYTES,
actual: request.len(),
})?;
let chunks: Vec<&[u8]> = request.chunks(CONTINUATION_CHUNK_BYTES).collect();
let last = chunks.len() - 1;
let mut frames = Vec::with_capacity(chunks.len());
for (index, chunk) in chunks.into_iter().enumerate() {
let sequence = u16::try_from(index)
.map_err(|_| TransportError::InvalidFrame("continuation sequence overflow".into()))?;
let mut payload = Vec::with_capacity(CONTINUATION_HEADER_BYTES + chunk.len());
payload.extend_from_slice(&CONTINUATION_MAGIC);
payload.extend_from_slice(&continuation_id.to_be_bytes());
payload.extend_from_slice(&sequence.to_be_bytes());
payload.push(if index == last {
CONTINUATION_FLAG_FINAL
} else {
0
});
payload.push(0);
payload.extend_from_slice(&total.to_be_bytes());
payload.extend_from_slice(chunk);
frames.push(encode_frame(&payload)?);
}
Ok(frames)
}

/// The continuation id after `id` (wraps to 1; 0 is never used).
#[must_use]
pub fn next_continuation_id(id: u32) -> u32 {
id.checked_add(1).unwrap_or(1)
}

pub fn decode_frame(buf: &[u8]) -> Result<(Frame, usize), TransportError> {
if buf.len() < 4 {
return Err(TransportError::FrameTruncated {
Expand Down Expand Up @@ -482,6 +553,8 @@ pub struct IpcTransport {
limiter: RateLimiter,
active_connections: usize,
requests: usize,
/// Next Amendment A4 continuation id (nonzero; wraps to 1).
continuation_id: u32,
peer: Option<PeerCredentials>,
runtime_uid: u32,
socket_path: String,
Expand Down Expand Up @@ -514,6 +587,7 @@ impl IpcTransport {
limiter: RateLimiter::rc9_default(),
active_connections: 0,
requests: 0,
continuation_id: 1,
peer,
runtime_uid,
socket_path,
Expand Down Expand Up @@ -681,12 +755,21 @@ impl IpcTransport {
self.limiter.check(now_ms)?;
let bytes = json.as_bytes();
check_payload_cap(bytes.len())?;
if bytes.len() <= MAX_FRAME_BYTES {
self.stub.try_send_payload(bytes)?;
} else {
for chunk in bytes.chunks(MAX_FRAME_BYTES) {
self.stub.try_send_payload(chunk)?;
}
// Amendment A4 (bitty#1482): above one frame the stub carries the
// continuation fragment payloads, never a headerless split.
let frames = encode_request_frames(bytes, self.continuation_id)?;
// A fragmented request is queued whole or not at all: a partial
// enqueue would leave the peer an unfinished reassembly that the next
// request would interleave.
let capacity = self.stub.capacity();
if self.stub.outgoing_len() + frames.len() > capacity {
return Err(TransportError::TransportFull { capacity });
}
if frames.len() > 1 {
self.continuation_id = next_continuation_id(self.continuation_id);
}
for frame in &frames {
self.stub.try_send_payload(&frame[4..])?;
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}
self.requests += 1;
Ok(())
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -299,7 +299,7 @@
"jsonLen": 1048576
},
"expect": {
"chunks": 4
"chunks": 5
},
"id": "Q03",
"kind": "req-1mib-exact",
Expand Down
19 changes: 14 additions & 5 deletions crates/devtools-client/tests/framer_fuzz_smoke.rs
Original file line number Diff line number Diff line change
@@ -1,8 +1,8 @@
use bitty_devtools_client::protocol::chunk_text;
use bitty_devtools_client::transport::{
Framer, IpcTransport, MAX_BUFFERED_BYTES, MAX_FRAME_BYTES, RC9_PAYLOAD_CAP_BYTES,
RC10_CHUNK_CEILING, StdioTransportStub, TransportError, check_payload_cap, decode_frame,
encode_frame,
CONTINUATION_HEADER_BYTES, CONTINUATION_MAGIC, Framer, IpcTransport, MAX_BUFFERED_BYTES,
MAX_FRAME_BYTES, RC9_PAYLOAD_CAP_BYTES, RC10_CHUNK_CEILING, StdioTransportStub, TransportError,
check_payload_cap, decode_frame, encode_frame,
};
use std::time::Instant;

Expand Down Expand Up @@ -341,8 +341,11 @@ fn t3_chunk_boundary_parity_with_oracle() {

#[test]
fn t4_chunk_counts_for_sized_requests() {
// Amendment A4 (bitty#1482): above one frame the stub carries
// continuation fragments; stripping each 16-byte header reproduces the
// request, and a 1 MiB request is five fragments.
let start = Instant::now();
for (target, chunks) in [(262_144usize, 1), (262_145, 2), (1_048_576, 4)] {
for (target, chunks) in [(262_144usize, 1), (262_145, 2), (1_048_576, 5)] {
let json = request_json_of_len(target);
let mut transport = connected_transport();
transport.send_request(&json, 0).unwrap();
Expand All @@ -351,7 +354,13 @@ fn t4_chunk_counts_for_sized_requests() {
for _ in 0..chunks {
let frame = transport.stub_mut().recv_outgoing().unwrap();
assert!(frame.payload().len() <= MAX_FRAME_BYTES);
got.extend_from_slice(frame.payload());
let payload = if chunks == 1 {
frame.payload()
} else {
assert_eq!(frame.payload()[..4], CONTINUATION_MAGIC);
&frame.payload()[CONTINUATION_HEADER_BYTES..]
};
got.extend_from_slice(payload);
}
assert_eq!(got, json.as_bytes());
}
Expand Down
Loading
Loading