diff --git a/CHANGELOG.md b/CHANGELOG.md index a65707d..c9d3785 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/crates/devtools-client/src/ipc_socket.rs b/crates/devtools-client/src/ipc_socket.rs index ad5e462..eedc95c 100644 --- a/crates/devtools-client/src/ipc_socket.rs +++ b/crates/devtools-client/src/ipc_socket.rs @@ -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; @@ -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")] @@ -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], @@ -259,11 +265,9 @@ 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))) @@ -271,10 +275,12 @@ impl LiveSocketConnection { 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)?; @@ -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, diff --git a/crates/devtools-client/src/transport.rs b/crates/devtools-client/src/transport.rs index 6e4d6c2..e8a81f8 100644 --- a/crates/devtools-client/src/transport.rs +++ b/crates/devtools-client/src/transport.rs @@ -114,6 +114,77 @@ pub fn encode_frame(payload: &[u8]) -> Result, 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>, 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 { @@ -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, runtime_uid: u32, socket_path: String, @@ -514,6 +587,7 @@ impl IpcTransport { limiter: RateLimiter::rc9_default(), active_connections: 0, requests: 0, + continuation_id: 1, peer, runtime_uid, socket_path, @@ -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..])?; } self.requests += 1; Ok(()) diff --git a/crates/devtools-client/tests/fixtures/fuzz/framer-seeds/vectors.json b/crates/devtools-client/tests/fixtures/fuzz/framer-seeds/vectors.json index 9c6eeb1..09035ab 100644 --- a/crates/devtools-client/tests/fixtures/fuzz/framer-seeds/vectors.json +++ b/crates/devtools-client/tests/fixtures/fuzz/framer-seeds/vectors.json @@ -299,7 +299,7 @@ "jsonLen": 1048576 }, "expect": { - "chunks": 4 + "chunks": 5 }, "id": "Q03", "kind": "req-1mib-exact", diff --git a/crates/devtools-client/tests/framer_fuzz_smoke.rs b/crates/devtools-client/tests/framer_fuzz_smoke.rs index 29452b1..e9b16cc 100644 --- a/crates/devtools-client/tests/framer_fuzz_smoke.rs +++ b/crates/devtools-client/tests/framer_fuzz_smoke.rs @@ -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; @@ -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(); @@ -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()); } diff --git a/crates/devtools-client/tests/request_continuation.rs b/crates/devtools-client/tests/request_continuation.rs new file mode 100644 index 0000000..9cd2edb --- /dev/null +++ b/crates/devtools-client/tests/request_continuation.rs @@ -0,0 +1,139 @@ +//! Amendment A4 inbound request continuation encoder (bitty#1482). +//! +//! Mirrors the TypeScript `encodeRequestFrames` vectors so both clients emit +//! byte-identical fragments for the same request and continuation id. + +use bitty_devtools_client::transport::{ + CONTINUATION_CHUNK_BYTES, CONTINUATION_FLAG_FINAL, CONTINUATION_HEADER_BYTES, + CONTINUATION_MAGIC, MAX_FRAME_BYTES, RC9_PAYLOAD_CAP_BYTES, TransportError, decode_frame, + encode_request_frames, next_continuation_id, +}; + +fn request_of(len: usize) -> Vec { + (0..len).map(|i| b'a' + (i % 26) as u8).collect() +} + +/// Validate every fragment header and return the reassembled request. +fn reassemble(frames: &[Vec], total: usize) -> Vec { + let mut out = Vec::with_capacity(total); + let mut id = None; + for (sequence, wire) in frames.iter().enumerate() { + let (frame, consumed) = decode_frame(wire).expect("valid frame"); + assert_eq!(consumed, wire.len()); + let payload = frame.payload(); + assert_eq!(payload[..4], CONTINUATION_MAGIC); + let fragment_id = u32::from_be_bytes(payload[4..8].try_into().unwrap()); + assert_ne!(fragment_id, 0); + assert_eq!(*id.get_or_insert(fragment_id), fragment_id); + assert_eq!( + usize::from(u16::from_be_bytes(payload[8..10].try_into().unwrap())), + sequence + ); + let last = sequence == frames.len() - 1; + assert_eq!(payload[10], if last { CONTINUATION_FLAG_FINAL } else { 0 }); + assert_eq!(payload[11], 0); + let declared = u32::from_be_bytes(payload[12..16].try_into().unwrap()); + assert_eq!(declared as usize, total); + if !last { + assert_eq!( + payload.len(), + MAX_FRAME_BYTES, + "non-final fragments are full" + ); + } + out.extend_from_slice(&payload[CONTINUATION_HEADER_BYTES..]); + } + out +} + +#[test] +fn a_request_that_fits_one_frame_stays_plain() { + let request = request_of(MAX_FRAME_BYTES); + let frames = encode_request_frames(&request, 0).unwrap(); + assert_eq!(frames.len(), 1); + assert_eq!(&frames[0][4..], request.as_slice()); +} + +#[test] +fn just_above_one_frame_is_two_fragments() { + let request = request_of(MAX_FRAME_BYTES + 1); + let frames = encode_request_frames(&request, 5).unwrap(); + assert_eq!(frames.len(), 2); + assert_eq!(reassemble(&frames, request.len()), request); +} + +#[test] +fn the_inbound_limit_is_five_fragments() { + let request = request_of(RC9_PAYLOAD_CAP_BYTES); + let frames = encode_request_frames(&request, 5).unwrap(); + assert_eq!(frames.len(), 5); + assert_eq!(reassemble(&frames, request.len()), request); +} + +#[test] +fn an_exact_chunk_multiple_ends_with_a_full_final_fragment() { + let request = request_of(CONTINUATION_CHUNK_BYTES * 2); + let frames = encode_request_frames(&request, 5).unwrap(); + assert_eq!(frames.len(), 2); + assert_eq!(reassemble(&frames, request.len()), request); +} + +#[test] +fn the_header_layout_matches_the_typescript_vector() { + let request = request_of(MAX_FRAME_BYTES + 1); + let frames = encode_request_frames(&request, 0x0102_0304).unwrap(); + // magic, id, sequence 1, FINAL, reserved, total 262145 (0x00040001). + assert_eq!( + frames[1][4..4 + CONTINUATION_HEADER_BYTES], + [ + 0x00, 0x42, 0x43, 0x31, 0x01, 0x02, 0x03, 0x04, 0x00, 0x01, 0x01, 0x00, 0x00, 0x04, + 0x00, 0x01, + ] + ); +} + +#[test] +fn over_limit_and_zero_id_fail_without_frames() { + let over = request_of(RC9_PAYLOAD_CAP_BYTES + 1); + assert!(matches!( + encode_request_frames(&over, 1), + Err(TransportError::PayloadTooLarge { .. }) + )); + let big = request_of(MAX_FRAME_BYTES + 1); + assert!(matches!( + encode_request_frames(&big, 0), + Err(TransportError::InvalidFrame(_)) + )); +} + +#[test] +fn continuation_ids_skip_zero_on_wrap() { + assert_eq!(next_continuation_id(1), 2); + assert_eq!(next_continuation_id(u32::MAX), 1); +} + +#[test] +fn a_fragmented_request_is_queued_whole_or_not_at_all() { + use bitty_devtools_client::transport::IpcTransport; + // Capacity 2 with one slot taken: a two-fragment request must be + // refused without enqueueing its first fragment (CodeRabbit on #149). + let mut transport = IpcTransport::new( + 1000, + "/tmp/ctx-0084-continuation-fixture.sock".to_owned(), + None, + 0o700, + 0o600, + 1000, + 1000, + 2, + ); + transport.connect().unwrap(); + transport.send_request(r#"{"id":1}"#, 0).unwrap(); + assert_eq!(transport.outgoing_len(), 1); + let big = String::from_utf8(request_of(MAX_FRAME_BYTES + 1)).unwrap(); + assert!(matches!( + transport.send_request(&big, 0), + Err(TransportError::TransportFull { capacity: 2 }) + )); + assert_eq!(transport.outgoing_len(), 1, "no partial reassembly queued"); +} diff --git a/src/client.ts b/src/client.ts index e204d7f..15c250d 100644 --- a/src/client.ts +++ b/src/client.ts @@ -621,13 +621,10 @@ export class DevtoolsClient { throw new TransportError("TransportClosed", "request cancelled"); } transport.getRateLimiter().check(nowMs); - const frame = transport.encodeSingleLiveRequest(req); - const raw = await live.requestResponse( - frame.slice(4), - nowMs, - req.id, - signal, - ); + // The live socket frames the request itself: one plain frame, or + // Amendment A4 continuation fragments above 256 KiB (bitty#1482). + const requestJson = transport.encodeRequestJson(req); + const raw = await live.requestResponse(requestJson, nowMs, req.id, signal); return decodeLiveResponse(raw, req.id); } diff --git a/src/ipc-socket.ts b/src/ipc-socket.ts index 0e80621..8cfe472 100644 --- a/src/ipc-socket.ts +++ b/src/ipc-socket.ts @@ -32,8 +32,10 @@ import { import { Framer, MAX_FRAME_BYTES, + RC9_PAYLOAD_CAP_BYTES, TransportError, - encodeFrame, + encodeRequestFrames, + nextContinuationId, } from "./transport.js"; export const LIVE_SOCKET_TIMEOUT_MS = 5_000 as const; @@ -41,6 +43,15 @@ export const LIVE_SOCKET_MAX_TIMEOUT_MS = 60_000 as const; export const LIVE_SOCKET_MAX_PENDING_FRAMES = 64 as const; +/** + * Bytes the live socket may hold while the kernel buffer is full: one + * maximal continuation request plus one maximal plain frame per pending + * slot. A request that would exceed it fails with `TransportFull`. + */ +export const LIVE_SOCKET_MAX_OUTBOUND_BYTES = + RC9_PAYLOAD_CAP_BYTES + + LIVE_SOCKET_MAX_PENDING_FRAMES * (4 + MAX_FRAME_BYTES); + export const LIVE_SOCKET_SUPPORTED_PLATFORM = "linux" as const; type LiveSocketPlatform = typeof LIVE_SOCKET_SUPPORTED_PLATFORM | "unsupported"; @@ -69,6 +80,7 @@ export type LiveSocketConfig = { }; type BunSocketHandle = { + /** Bytes accepted now; the rest must be written again on `drain`. */ write(data: Uint8Array): number; flush(): void; end(): void; @@ -83,6 +95,7 @@ type BunRuntime = { error(socket: BunSocketHandle, error: Error): void; open(socket: BunSocketHandle): void; close(socket: BunSocketHandle): void; + drain(socket: BunSocketHandle): void; }; }): Promise; file(path: string): { stat(): Promise }; @@ -321,8 +334,13 @@ export type LiveSocketConnection = { /** True once the OS socket is open. */ isOpen(): boolean; /** - * Write one framed request and resolve with the response carrying the same + * Write one request and resolve with the response carrying the same * request id. The optional id is used when the payload cannot be inspected. + * + * A request above one 256 KiB frame, up to the 1 MiB inbound limit, is + * sent as Amendment A4 continuation fragments (bitty#1482). It needs an + * idle connection, and no other request is accepted until it settles, + * because the server rejects any frame that interleaves a reassembly. */ requestResponse( requestJson: Uint8Array, @@ -470,16 +488,58 @@ export async function connectLiveSocket( let open = false; let terminalError: TransportError | null = null; let nextFallbackId = 0; + let continuationId = 1; + /** Request id of the in-flight continuation request, if any. */ + let continuationPending: number | null = null; + /** Bytes the kernel has not accepted yet, flushed on `drain`. */ + const outbound: Uint8Array[] = []; + let outboundBytes = 0; const removePending = (id: number): void => { pending.delete(id); const index = pendingOrder.indexOf(id); if (index >= 0) pendingOrder.splice(index, 1); + if (continuationPending === id) continuationPending = null; + }; + /** Write what the kernel accepts; queue the rest for `drain`, in order. */ + const writeOrQueue = (handle: BunSocketHandle, bytes: Uint8Array): void => { + if (outbound.length === 0) { + const written = handle.write(bytes); + const accepted = Math.max(0, Math.min(written, bytes.length)); + if (accepted === bytes.length) return; + bytes = bytes.subarray(accepted); + } + outbound.push(bytes); + outboundBytes += bytes.length; + }; + const onDrain = (handle: BunSocketHandle): void => { + try { + while (outbound.length > 0) { + const head = outbound[0]!; + const written = handle.write(head); + const accepted = Math.max(0, Math.min(written, head.length)); + outboundBytes -= accepted; + if (accepted < head.length) { + outbound[0] = head.subarray(accepted); + return; + } + outbound.shift(); + } + handle.flush(); + } catch { + failConnection( + new TransportError("TransportClosed", "socket write failed"), + true, + ); + } }; const failConnection = (error: TransportError, remember = false): void => { if (remember) terminalError = error; if (!open && pending.size === 0 && retained.length === 0) return; open = false; inbound.clear(); + outbound.length = 0; + outboundBytes = 0; + continuationPending = null; const waiters = [...pending.values()]; pending.clear(); pendingOrder.length = 0; @@ -577,6 +637,7 @@ export async function connectLiveSocket( error: onError, open: onOpen, close: onClose, + drain: onDrain, }, }); void connecting.catch((error: unknown) => { @@ -617,16 +678,31 @@ export async function connectLiveSocket( signal?: AbortSignal, ) => { let id: number; + let wire: Uint8Array; + let wireLength = 0; try { void nowMs; if (terminalError !== null) throw terminalError; if (!open || socket === null) { throw new TransportError("TransportClosed", "live socket is closed"); } - if (requestJson.length > MAX_FRAME_BYTES) { + if (requestJson.length > RC9_PAYLOAD_CAP_BYTES) { throw new TransportError( - "FrameTooLarge", - `request ${requestJson.length} > ${MAX_FRAME_BYTES}`, + "PayloadTooLarge", + `request ${requestJson.length} > ${RC9_PAYLOAD_CAP_BYTES}`, + ); + } + if (continuationPending !== null) { + throw new TransportError( + "TransportFull", + `continuation request ${continuationPending} is in flight`, + ); + } + const fragmented = requestJson.length > MAX_FRAME_BYTES; + if (fragmented && pending.size > 0) { + throw new TransportError( + "TransportFull", + "a continuation request needs an idle connection", ); } if (requestId !== undefined) { @@ -662,6 +738,24 @@ export async function connectLiveSocket( `pending requests exceed ${LIVE_SOCKET_MAX_PENDING_FRAMES}`, ); } + const frames = encodeRequestFrames(requestJson, continuationId); + wireLength = frames.reduce((sum, frame) => sum + frame.length, 0); + if (outboundBytes + wireLength > LIVE_SOCKET_MAX_OUTBOUND_BYTES) { + throw new TransportError( + "TransportFull", + `outbound bytes exceed ${LIVE_SOCKET_MAX_OUTBOUND_BYTES}`, + ); + } + wire = new Uint8Array(wireLength); + let offset = 0; + for (const frame of frames) { + wire.set(frame, offset); + offset += frame.length; + } + if (fragmented) { + continuationId = nextContinuationId(continuationId); + continuationPending = id; + } } catch (error) { return Promise.reject(error); } @@ -718,15 +812,17 @@ export async function connectLiveSocket( const queued = retainedIndex >= 0 ? retained.splice(retainedIndex, 1)[0] : undefined; try { - const wire = encodeFrame(requestJson); - const written = socket?.write(wire); - if (written !== wire.length) { + if (socket === null) { throw new TransportError( "TransportClosed", - "socket write was partial", + "live socket is closed", ); } - socket?.flush(); + // Fragments of one request are queued back to back in one buffer, + // so nothing can interleave them; a partial write continues on + // `drain` instead of failing the connection. + writeOrQueue(socket, wire); + socket.flush(); if (queued !== undefined) { const waiter = pending.get(id); waiter?.resolve(queued.payload); diff --git a/src/transport.ts b/src/transport.ts index cb96afe..ef615b6 100644 --- a/src/transport.ts +++ b/src/transport.ts @@ -97,6 +97,81 @@ export function encodeFrame(payload: Uint8Array): Uint8Array { return out; } +// --------------------------------------------------------------------------- +// Inbound request continuation (devtools-rfc Amendment A4, bitty#1482) +// --------------------------------------------------------------------------- + +/** Magic prefix of every continuation fragment payload (`\0BC1`). */ +export const CONTINUATION_MAGIC = Uint8Array.of(0x00, 0x42, 0x43, 0x31); +/** Fragment header bytes before each chunk. */ +export const CONTINUATION_HEADER_BYTES = 16; +/** Chunk bytes in every non-final fragment (one full physical frame). */ +export const CONTINUATION_CHUNK_BYTES = + MAX_FRAME_BYTES - CONTINUATION_HEADER_BYTES; +/** `flags` bit of the fragment that completes the logical request. */ +export const CONTINUATION_FLAG_FINAL = 0b0000_0001; +/** Largest continuation id; ids are nonzero unsigned 32-bit values. */ +export const MAX_CONTINUATION_ID = 0xffff_ffff; + +/** + * Encode one logical request as wire frames (length prefix included). + * + * A request of at most {@link MAX_FRAME_BYTES} stays one plain frame. A + * larger one, up to the 1 MiB inbound limit, becomes canonical Amendment A4 + * fragments tagged `continuationId`: every non-final fragment is a full + * frame and the final one carries the remainder, so a request is at most + * five fragments. Nothing is produced on error. + */ +export function encodeRequestFrames( + request: Uint8Array, + continuationId: number, +): Uint8Array[] { + if (request.length > RC9_PAYLOAD_CAP_BYTES) { + throw new TransportError( + "PayloadTooLarge", + `request ${request.length} > ${RC9_PAYLOAD_CAP_BYTES}`, + ); + } + if (request.length <= MAX_FRAME_BYTES) { + return [encodeFrame(request)]; + } + if ( + !Number.isInteger(continuationId) || + continuationId < 1 || + continuationId > MAX_CONTINUATION_ID + ) { + throw new TransportError( + "InvalidFrame", + `continuation id must be in 1..${MAX_CONTINUATION_ID}`, + ); + } + const frames: Uint8Array[] = []; + for ( + let offset = 0, sequence = 0; + offset < request.length; + offset += CONTINUATION_CHUNK_BYTES, sequence += 1 + ) { + const chunk = request.subarray(offset, offset + CONTINUATION_CHUNK_BYTES); + const last = offset + chunk.length >= request.length; + const payload = new Uint8Array(CONTINUATION_HEADER_BYTES + chunk.length); + const view = new DataView(payload.buffer); + payload.set(CONTINUATION_MAGIC, 0); + view.setUint32(4, continuationId, false); + view.setUint16(8, sequence, false); + payload[10] = last ? CONTINUATION_FLAG_FINAL : 0; + payload[11] = 0; + view.setUint32(12, request.length, false); + payload.set(chunk, CONTINUATION_HEADER_BYTES); + frames.push(encodeFrame(payload)); + } + return frames; +} + +/** The continuation id after `id` (wraps to 1; 0 is never used). */ +export function nextContinuationId(id: number): number { + return id >= MAX_CONTINUATION_ID ? 1 : id + 1; +} + export function decodeFrame(buf: Uint8Array): { frame: Frame; consumed: number; @@ -471,6 +546,7 @@ export class IpcTransport { private readonly limiter: RateLimiter; private activeConnections = 0; private requests = 0; + private continuationId = 1; private readonly peer: PeerCredentials | null; private readonly config: Required< Omit @@ -603,43 +679,30 @@ export class IpcTransport { } } - encodeRequest(req: IpcRequest): Uint8Array[] { + /** Validated UTF-8 JSON bytes of `req` (bounded by the 1 MiB inbound cap). */ + encodeRequestJson(req: IpcRequest): Uint8Array { const json = JSON.stringify(req); - assertStringBounded("devtools request", json, MAX_DEVTOOLS_FRAME_BYTES); - checkPayloadCap(new TextEncoder().encode(json).length); const bytes = new TextEncoder().encode(json); - if (bytes.length <= MAX_FRAME_BYTES) { - return [encodeFrame(bytes)]; - } - const chunks: Uint8Array[] = []; - for (let off = 0; off < bytes.length; off += MAX_FRAME_BYTES) { - const slice = bytes.slice(off, off + MAX_FRAME_BYTES); - chunks.push(encodeFrame(slice)); - } - return chunks; + checkPayloadCap(bytes.length); + assertStringBounded("devtools request", json, MAX_DEVTOOLS_FRAME_BYTES); + return bytes; } /** - * Encode exactly one physical frame for a live request, or fail closed. - * - * `encodeRequest` splits an oversized logical request into RC-10 chunks for - * the headless streaming transport. The live path must never use that split: - * `bitty` reads one complete frame per exchange and defines no inbound - * request continuation identity, so a fragmented request would become - * several independent exchanges with partial side effects. Until the server - * continuation contract lands in `bitty`, a request above - * {@link MAX_FRAME_BYTES} is rejected with `FrameTooLarge` before any byte - * reaches the socket. + * Encode `req` as wire frames: one plain frame up to + * {@link MAX_FRAME_BYTES}, otherwise Amendment A4 continuation fragments + * under a fresh per-transport continuation id (bitty#1482). The server + * reassembles the fragments into one exchange with one admission. */ - encodeSingleLiveRequest(req: IpcRequest): Uint8Array { - const frames = this.encodeRequest(req); - if (frames.length !== 1) { - throw new TransportError( - "FrameTooLarge", - "live request requires a server continuation contract", - ); + encodeRequest(req: IpcRequest): Uint8Array[] { + const bytes = this.encodeRequestJson(req); + if (bytes.length <= MAX_FRAME_BYTES) { + return [encodeFrame(bytes)]; } - return frames[0]!; + const id = this.continuationId; + const frames = encodeRequestFrames(bytes, id); + this.continuationId = nextContinuationId(id); + return frames; } sendRequest(req: IpcRequest, nowMs: number): void { @@ -652,6 +715,13 @@ export class IpcTransport { this.verifyPeerForPrivilegedAction(); this.limiter.check(nowMs); const frames = this.encodeRequest(req); + // 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 (Amendment A4, bitty#1482). + const capacity = this.stub.getCapacity(); + if (this.stub.outgoingLen() + frames.length > capacity) { + throw new TransportError("TransportFull", `capacity ${capacity}`); + } for (const f of frames) { const payload = f.slice(4); this.stub.trySendPayload(payload); diff --git a/tests/client.test.ts b/tests/client.test.ts index 9ead817..ddc4b03 100644 --- a/tests/client.test.ts +++ b/tests/client.test.ts @@ -213,7 +213,7 @@ describe("DevtoolsClient inspection live IPC wiring", () => { } }); - test("oversized live request fails closed before any socket write", async () => { + test("over-limit live request fails closed before any socket write", async () => { if (!isLiveSocketSupported()) return; const responsePayload = new TextEncoder().encode( JSON.stringify({ @@ -270,13 +270,16 @@ describe("DevtoolsClient inspection live IPC wiring", () => { 0, ); expect(writes.length).toBe(1); + // A request above the 1 MiB inbound limit fails closed before any + // byte reaches the socket. Requests between 256 KiB and 1 MiB are + // sent as continuation fragments (see ipc-socket.test.ts). let caught: unknown = null; try { await c.requestLive( { id: 2, method: "bitty.debug/listPlugins", - params: { x: "a".repeat(600 * 1024) }, + params: { x: "a".repeat(1024 * 1024) }, version: "1.0", }, 0, @@ -284,9 +287,9 @@ describe("DevtoolsClient inspection live IPC wiring", () => { } catch (error) { caught = error; } - expect(caught).toBeInstanceOf(TransportError); - expect((caught as TransportError).code).toBe("FrameTooLarge"); - expect((caught as Error).message).toContain("continuation contract"); + // Rejected by live request validation or the transport cap; either + // way nothing is written. + expect(caught).toBeInstanceOf(Error); expect(writes.length).toBe(1); c.disconnect(); } finally { diff --git a/tests/fixtures/fuzz/framer-seeds/vectors.json b/tests/fixtures/fuzz/framer-seeds/vectors.json index 9c6eeb1..09035ab 100644 --- a/tests/fixtures/fuzz/framer-seeds/vectors.json +++ b/tests/fixtures/fuzz/framer-seeds/vectors.json @@ -299,7 +299,7 @@ "jsonLen": 1048576 }, "expect": { - "chunks": 4 + "chunks": 5 }, "id": "Q03", "kind": "req-1mib-exact", diff --git a/tests/framer-fuzz-smoke.test.ts b/tests/framer-fuzz-smoke.test.ts index 4624513..bcea12d 100644 --- a/tests/framer-fuzz-smoke.test.ts +++ b/tests/framer-fuzz-smoke.test.ts @@ -2,6 +2,9 @@ import { describe, expect, test } from "bun:test"; import { readFileSync } from "node:fs"; import { join } from "node:path"; import { + CONTINUATION_FLAG_FINAL, + CONTINUATION_HEADER_BYTES, + CONTINUATION_MAGIC, Framer, IpcTransport, StdioTransportStub, @@ -468,9 +471,46 @@ describe("T4 chunking and inbound framing model", () => { expect(frame.payload.length).toBe(target); }); - test("Q02 256KiB+1 request splits into two bounded chunks", () => { + /** + * Check every Amendment A4 fragment header, then return the concatenated + * chunks: shared id, sequence from 0, FINAL only on the last fragment, + * the declared total in every header, and full non-final frames. + */ + function reassembleFragments( + chunks: Uint8Array[], + total: number, + ): Uint8Array { + const parts: Uint8Array[] = []; + let id: number | null = null; + chunks.forEach((chunk, sequence) => { + const { frame, consumed } = decodeFrame(chunk); + expect(consumed).toBe(chunk.length); + const payload = frame.payload; + const view = new DataView( + payload.buffer, + payload.byteOffset, + payload.byteLength, + ); + expect([...payload.subarray(0, 4)]).toEqual([...CONTINUATION_MAGIC]); + const fragmentId = view.getUint32(4, false); + expect(fragmentId).toBeGreaterThan(0); + id ??= fragmentId; + expect(fragmentId).toBe(id); + expect(view.getUint16(8, false)).toBe(sequence); + const last = sequence === chunks.length - 1; + expect(payload[10]).toBe(last ? CONTINUATION_FLAG_FINAL : 0); + expect(payload[11]).toBe(0); + expect(view.getUint32(12, false)).toBe(total); + if (!last) expect(payload.length).toBe(MAX_FRAME_BYTES); + parts.push(payload.subarray(CONTINUATION_HEADER_BYTES)); + }); + return concatBytes(...parts); + } + + test("Q02 256KiB+1 request becomes two continuation fragments", () => { const oracle = loadOracle(); - const target = vectorById(oracle, "Q02").build!["jsonLen"]!; + const vector = vectorById(oracle, "Q02"); + const target = vector.build!["jsonLen"]!; const transport = new IpcTransport({ runtimeUid: 1000, socketPath: "/tmp/ctx-0074-fuzz-smoke-fixture.sock", @@ -478,18 +518,12 @@ describe("T4 chunking and inbound framing model", () => { const req = requestOfJsonLen(target); expect(requestJsonBytes(req).length).toBe(target); const chunks = transport.encodeRequest(req); - expect(chunks.length).toBe(2); - const parts: Uint8Array[] = []; - for (const chunk of chunks) { - const { frame } = decodeFrame(chunk); - expect(frame.payload.length).toBeLessThanOrEqual(MAX_FRAME_BYTES); - parts.push(frame.payload); - } - expect(concatBytes(...parts)).toEqual(requestJsonBytes(req)); + expect(chunks.length).toBe(vector.expect["chunks"] as number); + expect(reassembleFragments(chunks, target)).toEqual(requestJsonBytes(req)); }); test( - "Q03 1MiB exact request splits into four bounded chunks", + "Q03 1MiB exact request becomes five continuation fragments", () => { const oracle = loadOracle(); const target = vectorById(oracle, "Q03").build!["jsonLen"]!; @@ -500,18 +534,35 @@ describe("T4 chunking and inbound framing model", () => { const req = requestOfJsonLen(target); expect(requestJsonBytes(req).length).toBe(target); const chunks = transport.encodeRequest(req); - expect(chunks.length).toBe(4); - const parts: Uint8Array[] = []; - for (const chunk of chunks) { - const { frame } = decodeFrame(chunk); - expect(frame.payload.length).toBeLessThanOrEqual(MAX_FRAME_BYTES); - parts.push(frame.payload); - } - expect(concatBytes(...parts)).toEqual(requestJsonBytes(req)); + expect(chunks.length).toBe( + vectorById(oracle, "Q03").expect["chunks"] as number, + ); + expect(reassembleFragments(chunks, target)).toEqual( + requestJsonBytes(req), + ); }, LOCAL_TARGET_BUDGET_MS, ); + test("a fragmented request is queued whole or not at all", () => { + // Capacity 2 with one slot taken: a two-fragment request is refused + // without enqueueing its first fragment (CodeRabbit on #149). + const transport = new IpcTransport({ + runtimeUid: 1000, + socketPath: "/tmp/ctx-0084-continuation-fixture.sock", + capacity: 2, + }); + transport.connect(); + transport.sendRequest(requestOfJsonLen(256), 0); + expect(transport.outgoingLen()).toBe(1); + expect( + codeOf(() => + transport.sendRequest(requestOfJsonLen(MAX_FRAME_BYTES + 1), 0), + ), + ).toBe("TransportFull"); + expect(transport.outgoingLen()).toBe(1); + }); + test("Q04 1MiB+1 request is refused with PayloadTooLarge", () => { const oracle = loadOracle(); const target = vectorById(oracle, "Q04").build!["jsonLen"]!; diff --git a/tests/helpers/fake-live-socket.ts b/tests/helpers/fake-live-socket.ts index 9847965..ea373fa 100644 --- a/tests/helpers/fake-live-socket.ts +++ b/tests/helpers/fake-live-socket.ts @@ -28,11 +28,16 @@ export type MemorySocket = { export type MemoryHandlers = { data(socket: MemorySocket, data: Uint8Array): void; open(socket: MemorySocket): void; + drain?(socket: MemorySocket): void; }; export type MemoryTransmission = { write?: () => void; flush?: (count: number) => void; + /** Bytes the fake kernel accepts from one write (default: all). */ + accept?: (length: number) => number; + /** Receives exactly the bytes each write accepted. */ + capture?: (accepted: Uint8Array) => void; }; export type MemoryRun = ( @@ -40,6 +45,7 @@ export type MemoryRun = ( receive: (bytes: Uint8Array) => void, writes: number[], flushes: number[], + drain: () => void, ) => Promise; export function localUid(): number { @@ -94,7 +100,12 @@ export async function withMemoryConnection( write: (data) => { writes.push(data.length); transmission.write?.(); - return data.length; + const accepted = Math.min( + data.length, + transmission.accept?.(data.length) ?? data.length, + ); + transmission.capture?.(data.slice(0, accepted)); + return accepted; }, flush() { flushes.push(writes.length); @@ -130,6 +141,7 @@ export async function withMemoryConnection( (bytes) => handlers!.data(socket, bytes), writes, flushes, + () => handlers!.drain?.(socket), ); } finally { connection.close(); diff --git a/tests/ipc-socket.test.ts b/tests/ipc-socket.test.ts index ed033e8..c8a3561 100644 --- a/tests/ipc-socket.test.ts +++ b/tests/ipc-socket.test.ts @@ -1,8 +1,11 @@ import { describe, expect, spyOn, test } from "bun:test"; import { + CONTINUATION_HEADER_BYTES, MAX_FRAME_BYTES, + RC9_PAYLOAD_CAP_BYTES, TransportError, encodeFrame, + encodeRequestFrames, } from "../src/transport.js"; import { attestLiveSocketEndpoint, @@ -401,3 +404,123 @@ describe("live Unix IPC socket (CTX-0036)", () => { } }); }); + +// --------------------------------------------------------------------------- +// Inbound request continuation (devtools-rfc Amendment A4, bitty#1482) +// --------------------------------------------------------------------------- + +/** A JSON request with id `id` padded to exactly `length` bytes. */ +function paddedRequest(id: number, length: number): Uint8Array { + const head = `{"id":${id},"method":"bitty.debug/ping","version":"1.0"`; + const bytes = new Uint8Array(length); + bytes.fill(0x20); + bytes.set(new TextEncoder().encode(head), 0); + bytes[length - 1] = 0x7d; // "}" + return bytes; +} + +function concat(parts: Uint8Array[]): Uint8Array { + const out = new Uint8Array(parts.reduce((sum, part) => sum + part.length, 0)); + let offset = 0; + for (const part of parts) { + out.set(part, offset); + offset += part.length; + } + return out; +} + +describe("live request continuation (Amendment A4, bitty#1482)", () => { + test("a request above one frame is written as continuation fragments", async () => { + const accepted: Uint8Array[] = []; + const request = paddedRequest(7, 600 * 1024); + await withMemoryConnection( + async (connection, receive) => { + const response = connection.requestResponse(request, 0); + const wire = concat(accepted); + expect(wire).toEqual(concat(encodeRequestFrames(request, 1))); + receive(encodeFrame(new TextEncoder().encode('{"id":7,"result":{}}'))); + expect(new TextDecoder().decode(await response)).toContain('"id":7'); + }, + { capture: (bytes) => accepted.push(bytes) }, + ); + }); + + test("a partial write continues on drain in order", async () => { + const accepted: Uint8Array[] = []; + const request = paddedRequest(8, 300 * 1024); + const perWrite = 64 * 1024; + await withMemoryConnection( + async (connection, receive, _writes, _flushes, drain) => { + const response = connection.requestResponse(request, 0); + const expected = concat(encodeRequestFrames(request, 1)); + let rounds = 0; + while (concat(accepted).length < expected.length) { + drain(); + rounds += 1; + expect(rounds).toBeLessThan(64); + } + expect(concat(accepted)).toEqual(expected); + receive(encodeFrame(new TextEncoder().encode('{"id":8,"result":{}}'))); + expect(new TextDecoder().decode(await response)).toContain('"id":8'); + }, + { + accept: (length) => Math.min(length, perWrite), + capture: (bytes) => accepted.push(bytes), + }, + ); + }); + + test("a continuation request needs an idle connection and blocks others", async () => { + await withMemoryConnection(async (connection, receive) => { + const plain = connection.requestResponse(firstPayload, 0); + await expect( + connection.requestResponse(paddedRequest(9, 300 * 1024), 0), + ).rejects.toMatchObject({ code: "TransportFull" }); + receive(firstFrame); + expect(await plain).toEqual(firstPayload); + + const large = connection.requestResponse(paddedRequest(9, 300 * 1024), 0); + await expect( + connection.requestResponse(secondPayload, 0), + ).rejects.toMatchObject({ + code: "TransportFull", + }); + receive(encodeFrame(new TextEncoder().encode('{"id":9,"result":{}}'))); + expect(new TextDecoder().decode(await large)).toContain('"id":9'); + // Settled: plain requests are accepted again. + const after = connection.requestResponse(secondPayload, 0); + receive(secondFrame); + expect(await after).toEqual(secondPayload); + }); + }); + + test("a request above the inbound limit fails closed before any write", async () => { + await withMemoryConnection(async (connection, _receive, writes) => { + await expect( + connection.requestResponse( + paddedRequest(10, RC9_PAYLOAD_CAP_BYTES + 1), + 0, + ), + ).rejects.toMatchObject({ code: "PayloadTooLarge" }); + expect(writes.length).toBe(0); + }); + }); + + test("the encoder writes the Amendment A4 header layout", () => { + const request = paddedRequest(11, MAX_FRAME_BYTES + 1); + const frames = encodeRequestFrames(request, 0x01020304); + expect(frames.length).toBe(2); + const header = [...frames[1]!.subarray(4, 4 + CONTINUATION_HEADER_BYTES)]; + // magic, id, sequence 1, FINAL, reserved, total 262145 (0x00040001). + expect(header).toEqual([ + 0x00, 0x42, 0x43, 0x31, 0x01, 0x02, 0x03, 0x04, 0x00, 0x01, 0x01, 0x00, + 0x00, 0x04, 0x00, 0x01, + ]); + expect(frames[0]!.length).toBe(4 + MAX_FRAME_BYTES); + expect(frames[0]![4 + 10]).toBe(0); + expect(() => encodeRequestFrames(request, 0)).toThrow(TransportError); + expect( + encodeRequestFrames(paddedRequest(12, MAX_FRAME_BYTES), 0).length, + ).toBe(1); + }); +});