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
162 changes: 128 additions & 34 deletions packages/server/src/ws-server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -257,12 +257,13 @@ export interface ReadCookiesResultError {
/**
* The MCP-facing handle for the fetchproxy bridge.
*
* On `listen()`, the server races the configured port: if it binds,
* the instance becomes the concentrator (role `'host'`) the extension
* dials. If the port is already taken by another fetchproxy host, the
* instance becomes a peer (role `'peer'`) and tunnels through that
* host's existing WebSocket. Either way, callers issue `fetch()` (or
* one of the verb shortcuts) and get the response from the user's
* `listen()` loads identity and reserves nothing. The first verb call
* (or an explicit `connect()`) races the configured port: if the bind
* succeeds, the instance becomes the concentrator (role `'host'`) the
* extension dials. If the port is already taken by another fetchproxy
* host, the instance becomes a peer (role `'peer'`) and tunnels through
* that host's existing WebSocket. Either way, callers issue `fetch()`
* (or one of the verb shortcuts) and get the response from the user's
* signed-in browser tab as if they'd run `window.fetch` there
* themselves.
*
Expand All @@ -271,7 +272,11 @@ export interface ReadCookiesResultError {
* branch on it.
*/
export class FetchproxyServer {
/** Set after `listen()` succeeds. Null while not listening. */
/**
* Bridge role. `null` until the first verb call (or an explicit
* `connect()`) — `listen()` no longer triggers the role election
* as of 0.5.3+. Reset to `null` on `close()`.
*/
public role: 'host' | 'peer' | null = null;

private opts: ResolvedOpts;
Expand Down Expand Up @@ -306,6 +311,12 @@ export class FetchproxyServer {
>();
private mcpId: string | null = null;
private identity: Identity | null = null;
// 0.5.3+: in-flight role-election / handle-start promise. Set the
// first time a verb call runs `ensureConnected`, awaited by concurrent
// callers, cleared once the connection is up. Single source of truth
// for "we're connecting right now" so two parallel first-calls don't
// race the port bind.
private connectingPromise: Promise<void> | null = null;

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Nit — two stale docstrings elsewhere in the class

Not in the diff but worth flagging here since this PR introduced the semantic change that makes them wrong:

  1. Class-level JSDoc (line ~260): "On listen(), the server races the configured port…" — this is the old eager behavior. Should say the port race happens lazily on the first verb call (or explicit connect()).
  2. role field JSDoc (line ~274): "Set after listen() succeeds. Null while not listening." — role is now null after listen() and only set after connect() / first verb call.

Fix this →


constructor(opts: FetchproxyServerOpts) {
if (!Array.isArray(opts.domains) || opts.domains.length === 0) {
Expand Down Expand Up @@ -370,23 +381,95 @@ export class FetchproxyServer {
}

/**
* Start the WebSocket bridge. Loads the long-term identity keypair
* from disk (creating it on first call), elects the host-vs-peer
* role by attempting to bind the configured port, and stands up the
* matching handshake machinery. Idempotent only insofar as it leaves
* `role` non-null on success; calling `listen()` twice without an
* intervening `close()` is a programming error.
* Prepare the bridge for use. Loads the long-term identity keypair
* from disk (creating it on first call) and computes this instance's
* `mcpId`. Does NOT bind the bridge port or dial any WebSocket — the
* connection is established lazily on the first verb call (see
* `ensureConnected` / `getOrConnect`).
*
* Pre-0.5.3 behavior: `listen()` also did role election and started
* the host/peer immediately, which meant every configured-but-unused
* MCP claimed bridge resources at MCP-client boot. Several MCPs
* starting in parallel under Claude Desktop also produced noisy
* `ERR_CONNECTION_REFUSED` errors in the extension if it raced ahead
* of the first MCP's port bind. Deferring keeps boot quiet and
* leaves the port unowned until something actually needs it.
*
* Calling `listen()` twice without an intervening `close()` is a
* no-op (the second call's identity load is idempotent).
*/
async listen(): Promise<void> {
this.identity = await loadOrCreateIdentity(this.opts.serverName, this.opts.identityDir);
this.mcpId = generateMcpId(this.opts.serverName, this.opts.version);
if (!this.identity) {
this.identity = await loadOrCreateIdentity(this.opts.serverName, this.opts.identityDir);
}
if (!this.mcpId) {
this.mcpId = generateMcpId(this.opts.serverName, this.opts.version);
}
}

/**
* Force an eager bridge connection (role-election + host/peer handle
* start + listener wiring) without waiting for the first verb call.
* Useful for callers that want to surface the role / connection
* outcome at boot, or for tests whose harness dials a mock extension
* immediately after server construction. Production MCPs that just
* answer tool calls should NOT call this — the lazy connect via
* `ensureConnected` will do the right thing on first use, keeping
* boot cheap and avoiding port-bind contention for MCPs that never
* actually get invoked.
*
* Idempotent: a second call after the first has resolved is a no-op
* (the existing handle is reused). Throws if `listen()` was never
* called.
*/
async connect(): Promise<void> {
await this.ensureConnected();
}

/**
* Establish the bridge connection (role-election + host/peer handle
* start + listener wiring) the first time a verb is invoked.
* Idempotent after the connection is up; concurrent first-callers
* share the same in-flight promise so only one election happens.
*
* Throws if `listen()` was never called — the contract is that the
* MCP author still must wire `transport.start()` at boot to load
* identity / set mcpId, even though the WS doesn't open until a
* verb runs.
*/
private async ensureConnected(): Promise<void> {
if (this.hostHandle || this.peerHandle) return;
if (this.connectingPromise) {
await this.connectingPromise;
return;
}
if (!this.identity || !this.mcpId) {
throw new Error(
'FetchproxyServer: ensureConnected called before listen() — call listen() at MCP boot to load identity',
);
}
this.connectingPromise = this.doConnect();
try {
await this.connectingPromise;
} finally {
// Always clear so a transient connect failure can be retried by
// the next verb call. Successful path: the handle is now set, so
// the next `ensureConnected` short-circuits on the first branch.
this.connectingPromise = null;

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Nit — close() called while doConnect() is in flight can leak a handle

The finally correctly clears connectingPromise so a failed connection is retryable, but close() doesn't set connectingPromise = null before it exits. If a verb call triggers doConnect() and close() races it, close() will see hostHandle == null (the handle hasn't been set yet), so hostHandle.close() never runs. doConnect() then finishes and writes this.hostHandle = <live handle> after close() has already returned — the handle leaks.

One simple fix: have close() also clear connectingPromise (and optionally await it first so any in-flight election gets shut down cleanly):

async close(): Promise<void> {
  this.rejectAllPending();
  // Wait for any in-flight election before tearing down handles.
  if (this.connectingPromise) {
    await this.connectingPromise.catch(() => {});
  }
  if (this.hostHandle) await this.hostHandle.close();
  if (this.peerHandle) this.peerHandle.close();
  this.hostHandle = null;
  this.peerHandle = null;
  this.role = null;
  this.connectingPromise = null;
}

This is a narrow race (requires concurrent close() + verb call) so not blocking, but worth a follow-up.

}
}

private async doConnect(): Promise<void> {
// Identity / mcpId are guaranteed by `ensureConnected`'s precondition.
const identity = this.identity!;
const mcpId = this.mcpId!;
const el = await electRole({ host: this.opts.host, port: this.opts.port });
if (el.role === 'host') {
this.role = 'host';
this.hostHandle = await startHost({
httpServer: el.server,
ownIdentity: this.identity,
ownMcpId: this.mcpId,
ownIdentity: identity,
ownMcpId: mcpId,
ownServerName: this.opts.serverName,
ownVersion: this.opts.version,
ownDomains: this.opts.domains,
Expand All @@ -413,8 +496,8 @@ export class FetchproxyServer {
this.peerHandle = await startPeer({
host: this.opts.host,
port: this.opts.port,
identity: this.identity,
mcpId: this.mcpId,
identity,
mcpId,
serverName: this.opts.serverName,
version: this.opts.version,
domains: this.opts.domains,
Expand Down Expand Up @@ -472,9 +555,10 @@ export class FetchproxyServer {
* offline, etc.).
*/
async fetch(init: FetchInit): Promise<FetchResult | FetchResultError> {
if (!this.hostHandle && !this.peerHandle) {
throw new Error('FetchproxyServer.fetch called before listen() — not listening');
}
// 0.5.3+: connect lazily on first verb call. `listen()` only loads
// identity now; the bridge port bind / WS dial happens here so a
// configured-but-unused MCP doesn't tie up resources at boot.
await this.ensureConnected();
// 0.5.2+: if the extension has us queued for pairing, fail this call
// immediately with the actionable pair-code error rather than sealing
// a frame the extension can't process until the user approves.
Expand Down Expand Up @@ -672,9 +756,8 @@ export class FetchproxyServer {
'FetchproxyServer.readCookies(): MCP did not declare "read_cookies" in capabilities — add it to FetchproxyServerOpts.capabilities to enable this verb',
);
}
if (!this.hostHandle && !this.peerHandle) {
throw new Error('FetchproxyServer.readCookies called before listen() — not listening');
}
// 0.5.3+: lazy connect — see the doc comment on `ensureConnected`.
await this.ensureConnected();
this.throwIfPendingPair();
if (opts.subdomain !== undefined) assertSubdomainLabel(opts.subdomain);
const baseDomain = this.resolveBaseDomain(opts.domain);
Expand Down Expand Up @@ -778,9 +861,8 @@ export class FetchproxyServer {
`FetchproxyServer.${op === 'read_local_storage' ? 'readLocalStorage' : 'readSessionStorage'}(): MCP did not declare ${JSON.stringify(op)} in capabilities`,
);
}
if (!this.hostHandle && !this.peerHandle) {
throw new Error(`FetchproxyServer.${op} called before listen() — not listening`);
}
// 0.5.3+: lazy connect — see the doc comment on `ensureConnected`.
await this.ensureConnected();
this.throwIfPendingPair();
if (!Array.isArray(opts.keys) || opts.keys.length === 0) {
throw new Error(`FetchproxyServer.${op}: opts.keys must be a non-empty array`);
Expand Down Expand Up @@ -848,9 +930,8 @@ export class FetchproxyServer {
'FetchproxyServer.captureRequestHeader(): MCP did not declare "capture_request_header" in capabilities',
);
}
if (!this.hostHandle && !this.peerHandle) {
throw new Error('FetchproxyServer.captureRequestHeader called before listen() — not listening');
}
// 0.5.3+: lazy connect — see the doc comment on `ensureConnected`.
await this.ensureConnected();
this.throwIfPendingPair();
const declared = this.opts.captureHeaders.find(
(d) => d.urlPattern === opts.urlPattern && d.headerName === opts.headerName,
Expand Down Expand Up @@ -905,9 +986,8 @@ export class FetchproxyServer {
'FetchproxyServer.readIndexedDb(): MCP did not declare "read_indexed_db" in capabilities',
);
}
if (!this.hostHandle && !this.peerHandle) {
throw new Error('FetchproxyServer.readIndexedDb called before listen() — not listening');
}
// 0.5.3+: lazy connect — see the doc comment on `ensureConnected`.
await this.ensureConnected();
this.throwIfPendingPair();
if (!Array.isArray(opts.keys) || opts.keys.length === 0) {
throw new Error('FetchproxyServer.readIndexedDb: opts.keys must be a non-empty array');
Expand Down Expand Up @@ -1166,10 +1246,24 @@ export class FetchproxyServer {
*/
async close(): Promise<void> {
this.rejectAllPending();
// 0.5.3+: if a verb call has triggered `doConnect()` but the handle
// hasn't been written yet, wait it out before we tear things down.
// Otherwise `close()` would observe `hostHandle == null`, return
// without calling its `.close()`, and `doConnect()` would then
// write `this.hostHandle = <live handle>` after we'd already
// returned — leaking the socket. Swallow the rejection: if the
// election itself failed, there's nothing to tear down.
if (this.connectingPromise) {
await this.connectingPromise.catch(() => undefined);
}
if (this.hostHandle) await this.hostHandle.close();
if (this.peerHandle) this.peerHandle.close();
this.hostHandle = null;
this.peerHandle = null;
this.role = null;
// Match the `finally` in `ensureConnected` so a subsequent
// `listen()` + `connect()` after `close()` doesn't observe a stale
// promise reference.
this.connectingPromise = null;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,7 @@ describe('integration: all 0.3.0 bootstrap verbs', () => {
identityDir: idDir,
});
await server.listen();
await server.connect();

extWs = new WebSocket(`ws://127.0.0.1:${port}`);
await new Promise<void>((resolve, reject) => {
Expand Down
2 changes: 2 additions & 0 deletions packages/server/tests/integration/mutual-auth-idb.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,7 @@ describe('integration: 0.4.0 mutual auth + read_indexed_db', () => {
},
});
await server.listen();
await server.connect();

extWs = new WebSocket(`ws://127.0.0.1:${port}`);
await new Promise<void>((resolve, reject) => {
Expand Down Expand Up @@ -190,6 +191,7 @@ describe('integration: 0.4.0 mutual auth + read_indexed_db', () => {
identityDir: idDir,
});
await server.listen();
await server.connect();

extWs = new WebSocket(`ws://127.0.0.1:${port}`);
await new Promise<void>((resolve, reject) => {
Expand Down
4 changes: 4 additions & 0 deletions packages/server/tests/integration/pair-pending.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -113,6 +113,7 @@ describe('pair-pending surfaces the pair code to MCP-side callers', () => {
identityDir: idDir,
});
await host.listen();
await host.connect();
expect(host.role).toBe('host');

const ext = await connectMockExtensionThatNeverApproves(port, {
Expand Down Expand Up @@ -147,6 +148,7 @@ describe('pair-pending surfaces the pair code to MCP-side callers', () => {
identityDir: idDir,
});
await host.listen();
await host.connect();
peer = new FetchproxyServer({
port,
serverName: 'peer-mcp',
Expand All @@ -155,6 +157,7 @@ describe('pair-pending surfaces the pair code to MCP-side callers', () => {
identityDir: idDir,
});
await peer.listen();
await peer.connect();
expect(peer.role).toBe('peer');

const ext = await connectMockExtensionThatNeverApproves(port, {
Expand Down Expand Up @@ -193,6 +196,7 @@ describe('pair-pending surfaces the pair code to MCP-side callers', () => {
onPairCode: (c) => codes.push(c),
});
await host.listen();
await host.connect();

const ext = await connectMockExtensionThatNeverApproves(port, {
'host-mcp': '567-890',
Expand Down
2 changes: 2 additions & 0 deletions packages/server/tests/integration/read-cookies.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,7 @@ describe('integration: readCookies() round-trip', () => {
identityDir: idDir,
});
await server.listen();
await server.connect();
expect(server.role).toBe('host');

extWs = new WebSocket(`ws://127.0.0.1:${port}`);
Expand Down Expand Up @@ -183,6 +184,7 @@ describe('integration: readCookies() round-trip', () => {
identityDir: idDir,
});
await server.listen();
await server.connect();

// listen() succeeded; the developer-facing guard should bite without
// any extension being connected, because it predates the wire send.
Expand Down
8 changes: 8 additions & 0 deletions packages/server/tests/integration/reconnect.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -122,6 +122,7 @@ describe('extension reconnect (session-key renegotiation)', () => {
identityDir: idDir,
});
await server.listen();
await server.connect();

// First connection — complete handshake and verify a fetch round-trips.
const ext1 = await connectMockExtension(port);
Expand Down Expand Up @@ -216,6 +217,7 @@ describe('extension reconnect (session-key renegotiation)', () => {
identityDir: idDir,
});
await server.listen();
await server.connect();

const ext = await connectMockExtension(port);

Expand Down Expand Up @@ -251,6 +253,7 @@ describe('extension reconnect (session-key renegotiation)', () => {
identityDir: idDir,
});
await server.listen();
await server.connect();

const ext = await connectMockExtension(port);

Expand All @@ -274,6 +277,7 @@ describe('extension reconnect (session-key renegotiation)', () => {
identityDir: idDir,
});
await server.listen();
await server.connect();

const ext1 = await connectMockExtension(port);
const ext1Closed = new Promise<void>((r) => ext1.ws.on('close', () => r()));
Expand Down Expand Up @@ -438,6 +442,7 @@ describe('extension reconnect (peer MCPs)', () => {
identityDir: idDir,
});
await host.listen();
await host.connect();
expect(host.role).toBe('host');

peer = new FetchproxyServer({
Expand All @@ -448,6 +453,7 @@ describe('extension reconnect (peer MCPs)', () => {
identityDir: idDir,
});
await peer.listen();
await peer.connect();
expect(peer.role).toBe('peer');

// Mock extension responder: echoes the request body back to the caller,
Expand Down Expand Up @@ -549,6 +555,7 @@ describe('extension reconnect (peer MCPs)', () => {
identityDir: idDir,
});
await host.listen();
await host.connect();

peer = new FetchproxyServer({
port,
Expand All @@ -558,6 +565,7 @@ describe('extension reconnect (peer MCPs)', () => {
identityDir: idDir,
});
await peer.listen();
await peer.connect();

// Connect extension but never answer the fetch — we want to observe
// the in-flight rejection path, not an upstream response.
Expand Down
Loading