diff --git a/crates/devtools-client/src/transport.rs b/crates/devtools-client/src/transport.rs index 3aaffe0..6e4d6c2 100644 --- a/crates/devtools-client/src/transport.rs +++ b/crates/devtools-client/src/transport.rs @@ -480,7 +480,7 @@ impl Default for StdioTransportStub { pub struct IpcTransport { stub: StdioTransportStub, limiter: RateLimiter, - connected: bool, + active_connections: usize, requests: usize, peer: Option, runtime_uid: u32, @@ -512,7 +512,7 @@ impl IpcTransport { Self { stub: StdioTransportStub::new(capacity), limiter: RateLimiter::rc9_default(), - connected: false, + active_connections: 0, requests: 0, peer, runtime_uid, @@ -571,7 +571,7 @@ impl IpcTransport { #[must_use] pub fn is_connected(&self) -> bool { - self.connected + self.active_connections == 1 } #[must_use] @@ -580,6 +580,10 @@ impl IpcTransport { } pub fn connect(&mut self) -> Result<(), TransportError> { + if self.stub.is_closed() { + self.disconnect(); + return Err(TransportError::TransportClosed); + } if let (Some(peer_sid), Some(runtime_sid)) = (self.windows_peer_sid, self.windows_runtime_sid) { @@ -621,13 +625,15 @@ impl IpcTransport { ))); } } - check_connection_cap(self.requests)?; - self.connected = true; + if !self.is_connected() { + check_connection_cap(self.active_connections)?; + self.active_connections += 1; + } Ok(()) } pub fn disconnect(&mut self) { - self.connected = false; + self.active_connections = 0; self.stub.clear(); } @@ -664,7 +670,11 @@ impl IpcTransport { } pub fn send_request(&mut self, json: &str, now_ms: u64) -> Result<(), TransportError> { - if !self.connected { + if self.stub.is_closed() { + self.disconnect(); + return Err(TransportError::TransportClosed); + } + if !self.is_connected() { return Err(TransportError::TransportClosed); } self.verify_peer_for_privileged()?; @@ -924,6 +934,112 @@ mod tests { assert!(t.is_connected()); } + #[test] + fn reconnect_admission_ignores_completed_request_history() { + let mut t = IpcTransport::with_defaults( + 1000, + "/unused/headless.sock".into(), + Some(PeerCredentials::new(1000, 1000, 1)), + ); + for cycle in 0..2 { + t.connect().unwrap(); + for _ in 0..=RC9_MAX_CONNECTIONS { + t.send_request("{}", 0).unwrap(); + assert!(t.stub_mut().recv_outgoing().is_some()); + } + assert!(t.connect().is_ok()); + assert!(t.is_connected()); + t.stub_mut().inject_incoming_payload(b"{}").unwrap(); + t.disconnect(); + t.disconnect(); + assert!(!t.is_connected()); + assert_eq!(t.outgoing_len(), 0); + assert_eq!(t.incoming_len(), 0); + assert_eq!(t.requests, (cycle + 1) * (RC9_MAX_CONNECTIONS + 1)); + assert_eq!(t.limiter_mut().count_in_window(0), t.requests); + } + assert!(t.connect().is_ok()); + t.disconnect(); + } + + #[test] + fn reconnect_preserves_rate_credit_and_releases_closed_ownership() { + let mut t = IpcTransport::new( + 1000, + "/unused/headless.sock".into(), + Some(PeerCredentials::new(1000, 1000, 1)), + 0o700, + 0o600, + 1000, + 1000, + 1, + ); + t.limiter = RateLimiter::new(0, 3); + t.connect().unwrap(); + t.send_request("{}", 0).unwrap(); + assert!(matches!( + t.send_request("{}", 0), + Err(TransportError::TransportFull { .. }) + )); + assert!(t.is_connected()); + t.disconnect(); + assert_eq!(t.outgoing_len(), 0); + t.connect().unwrap(); + t.send_request("{}", 0).unwrap(); + assert_eq!(t.requests, 2); + assert_eq!(t.limiter_mut().count_in_window(0), 3); + t.disconnect(); + t.connect().unwrap(); + assert!(matches!( + t.send_request("{}", 0), + Err(TransportError::RateLimited(_)) + )); + assert!(t.is_connected()); + t.stub_mut().close(); + assert!(matches!( + t.send_request("{}", 0), + Err(TransportError::TransportClosed) + )); + assert!(!t.is_connected()); + assert_eq!(t.outgoing_len(), 0); + assert!(matches!(t.connect(), Err(TransportError::TransportClosed))); + t.disconnect(); + assert!(matches!(t.connect(), Err(TransportError::TransportClosed))); + assert!(!t.is_connected()); + assert_eq!(t.requests, 2); + assert_eq!(t.limiter_mut().count_in_window(0), 3); + } + + #[test] + fn reconnect_failure_preserves_ownership_and_rechecks_peer() { + let mut t = IpcTransport::with_defaults( + 1000, + "/unused/headless.sock".into(), + Some(PeerCredentials::new(1001, 1000, 1)), + ); + assert!(matches!( + t.connect(), + Err(TransportError::Unauthenticated(_)) + )); + assert!(!t.is_connected()); + t.disconnect(); + t.peer = Some(PeerCredentials::new(1000, 1000, 1)); + t.connect().unwrap(); + t.peer = Some(PeerCredentials::new(1001, 1000, 1)); + assert!(matches!( + t.connect(), + Err(TransportError::Unauthenticated(_)) + )); + assert!(matches!( + t.send_request("{}", 0), + Err(TransportError::Unauthenticated(_)) + )); + t.disconnect(); + t.peer = Some(PeerCredentials::new(1000, 1000, 1)); + assert!(t.connect().is_ok()); + t.disconnect(); + } + #[test] fn transport_peer_mismatch_fails() { let peer = PeerCredentials::new(1001, 1000, 1); diff --git a/src/transport.ts b/src/transport.ts index a80c5a6..abc7172 100644 --- a/src/transport.ts +++ b/src/transport.ts @@ -472,7 +472,7 @@ export type IpcResponse = { export class IpcTransport { private readonly stub: StdioTransportStub; private readonly limiter: RateLimiter; - private connected = false; + private activeConnections = 0; private requests = 0; private readonly peer: PeerCredentials | null; private readonly config: Required< @@ -509,7 +509,7 @@ export class IpcTransport { } isConnected(): boolean { - return this.connected; + return this.activeConnections === 1; } getSocketPath(): string { @@ -525,6 +525,10 @@ export class IpcTransport { } connect(nowMs?: number): void { + if (this.stub.isClosed()) { + this.disconnect(); + throw new TransportError("TransportClosed", "stdio transport is closed"); + } if (this.isWindowsPipe()) { // CTX-0043: named-pipe peers carry SIDs, not Unix modes/owners. this.checkWindowsPipe( @@ -566,15 +570,17 @@ export class IpcTransport { ); } } - checkConnectionCap(this.requests); - this.connected = true; + if (!this.isConnected()) { + checkConnectionCap(this.activeConnections); + this.activeConnections += 1; + } if (nowMs !== undefined) { void nowMs; } } disconnect(): void { - this.connected = false; + this.activeConnections = 0; this.stub.clear(); } @@ -619,7 +625,11 @@ export class IpcTransport { } sendRequest(req: IpcRequest, nowMs: number): void { - if (!this.connected) + if (this.stub.isClosed()) { + this.disconnect(); + throw new TransportError("TransportClosed", "stdio transport is closed"); + } + if (!this.isConnected()) throw new TransportError("TransportClosed", "not connected"); this.verifyPeerForPrivilegedAction(); this.limiter.check(nowMs); diff --git a/tests/transport.test.ts b/tests/transport.test.ts index be2c62e..327a7c5 100644 --- a/tests/transport.test.ts +++ b/tests/transport.test.ts @@ -218,6 +218,99 @@ describe("transport framing (phase 2, live runtime, bounded 256 KiB IPC / 1 MiB expect(t.getRateLimiter().countInWindow(0)).toBe(1); }); + test("reconnect admission is independent of completed request history", () => { + const t = new IpcTransport({ + runtimeUid: 1000, + socketPath: "/unused/headless.sock", + peer: peerCredentials(1000, 1000, 1), + }); + for (let cycle = 0; cycle < 2; cycle++) { + t.connect(); + for (let id = 0; id <= RC9_MAX_CONNECTIONS; id++) { + t.injectResponsePayload( + new TextEncoder().encode( + JSON.stringify({ jsonrpc: "2.0", id, result: {}, version: "1.0" }), + ), + ); + expect( + t.request( + { id, method: "bitty.debug/listPlugins", version: "1.0" }, + 0, + ).id, + ).toBe(id); + expect(t.getStub().recvOutgoing()).toBeDefined(); + } + expect(() => t.connect()).not.toThrow(); + expect(t.isConnected()).toBe(true); + t.injectResponsePayload(new TextEncoder().encode("{}")); + t.disconnect(); + t.disconnect(); + expect(t.isConnected()).toBe(false); + expect(t.outgoingLen()).toBe(0); + expect(t.incomingLen()).toBe(0); + expect(t.getRateLimiter().countInWindow(0)).toBe( + (cycle + 1) * (RC9_MAX_CONNECTIONS + 1), + ); + } + expect(() => t.connect()).not.toThrow(); + t.disconnect(); + }); + + test("reconnect preserves rate credit and releases terminally closed ownership", () => { + const t = new IpcTransport({ + runtimeUid: 1000, + socketPath: "/unused/headless.sock", + peer: peerCredentials(1000, 1000, 1), + capacity: 1, + }); + const req = { id: 1, method: "bitty.debug/listPlugins", version: "1.0" }; + t.connect(); + t.sendRequest(req, 0); + expect(() => t.sendRequest(req, 0)).toThrow("capacity"); + expect(t.isConnected()).toBe(true); + t.disconnect(); + expect(t.outgoingLen()).toBe(0); + t.connect(); + t.sendRequest(req, 0); + expect(t.getRateLimiter().countInWindow(0)).toBe(3); + t.getStub().close(); + expect(() => t.sendRequest(req, 0)).toThrow("closed"); + expect(t.isConnected()).toBe(false); + expect(t.outgoingLen()).toBe(0); + expect(() => t.connect()).toThrow("closed"); + t.disconnect(); + expect(() => t.connect()).toThrow("closed"); + expect(t.isConnected()).toBe(false); + expect(t.getRateLimiter().countInWindow(0)).toBe(3); + }); + + test("reconnect failure leaves no ownership and rechecks peer identity", () => { + const peer = peerCredentials(1000, 1000, 1); + const t = new IpcTransport({ + runtimeUid: 1000, + socketPath: "/unused/headless.sock", + peer, + }); + peer.uid = 1001; + expect(() => t.connect()).toThrow("peer uid"); + expect(t.isConnected()).toBe(false); + t.disconnect(); + peer.uid = 1000; + t.connect(); + peer.uid = 1001; + expect(() => t.connect()).toThrow("peer uid"); + expect(() => + t.sendRequest( + { id: 1, method: "bitty.debug/listPlugins", version: "1.0" }, + 0, + ), + ).toThrow("peer uid"); + t.disconnect(); + peer.uid = 1000; + expect(() => t.connect()).not.toThrow(); + t.disconnect(); + }); + test("ipc transport rejects foreign user peer", () => { const peer = peerCredentials(1001, 1000, 1); const t = new IpcTransport({