Skip to content
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
130 changes: 123 additions & 7 deletions crates/devtools-client/src/transport.rs
Original file line number Diff line number Diff line change
Expand Up @@ -480,7 +480,7 @@ impl Default for StdioTransportStub {
pub struct IpcTransport {
stub: StdioTransportStub,
limiter: RateLimiter,
connected: bool,
active_connections: usize,
requests: usize,
peer: Option<PeerCredentials>,
runtime_uid: u32,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -571,7 +571,7 @@ impl IpcTransport {

#[must_use]
pub fn is_connected(&self) -> bool {
self.connected
self.active_connections == 1
}

#[must_use]
Expand All @@ -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)
{
Expand Down Expand Up @@ -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();
}

Expand Down Expand Up @@ -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()?;
Expand Down Expand Up @@ -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);
Expand Down
22 changes: 16 additions & 6 deletions src/transport.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<
Expand Down Expand Up @@ -509,7 +509,7 @@ export class IpcTransport {
}

isConnected(): boolean {
return this.connected;
return this.activeConnections === 1;
}

getSocketPath(): string {
Expand All @@ -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(
Expand Down Expand Up @@ -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();
}

Expand Down Expand Up @@ -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);
Expand Down
93 changes: 93 additions & 0 deletions tests/transport.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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({
Expand Down