Skip to content
Closed
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
5 changes: 5 additions & 0 deletions .changeset/ws-transport-global-registry.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
'@workflow/world-vercel': patch
---

Hold the WebSocket events transport's channel registry on `globalThis` instead of at module scope, so an app that ends up with two copies of the bundled world in one process (for example a Next.js app whose `instrumentation.ts` warms the world) still resolves the channel its queue consumer opened instead of silently writing every event over HTTP.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

too verboe

51 changes: 51 additions & 0 deletions packages/world-vercel/src/ws-transport.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1048,6 +1048,57 @@ describe('transport selection', () => {
});
});

/**
* Two copies of this module can share a process, and the open and the lookup
* run through different Worlds, so they can land in different copies. With
* the registry at module scope that made an open socket carry nothing.
*/
describe('the channel registry is process-wide', () => {
const WsEventsStateKey = Symbol.for(
'@workflow/world-vercel//wsEventsTransports/v1'
);

/** The registry as a second copy of this module in the process sees it. */
const sharedTransports = () =>
(
globalThis as typeof globalThis & {
[WsEventsStateKey]?: { transports: Map<string, unknown> };
}
)[WsEventsStateKey]?.transports;

it('registers an opened channel there, and evicts it on release', () => {
process.env.WORKFLOW_EVENTS_TRANSPORT = 'ws';
const release = openWsChannel('wrun_1', directConfig);
const wsUrl = resolveWsTransport('wrun_1', directConfig)?.wsUrl;

expect(wsUrl).toBeDefined();
expect(sharedTransports()?.has(String(wsUrl))).toBe(true);

release?.();

expect(sharedTransports()?.has(String(wsUrl))).toBe(false);
});

it('resolves a channel another copy of this module opened', () => {
process.env.WORKFLOW_EVENTS_TRANSPORT = 'ws';
// Borrow the URL the open path derives, then stand in for that other
// copy by registering under it directly.
const release = openWsChannel('wrun_1', directConfig);
const wsUrl = String(resolveWsTransport('wrun_1', directConfig)?.wsUrl);
release?.();
expect(resolveWsTransport('wrun_1', directConfig)).toBeNull();

// Stands in for a transport built by that other copy: structurally used
// by the write path, and `close`able by the reset seam.
const openedElsewhere = { copy: 'other', close: () => {} };
sharedTransports()?.set(wsUrl, openedElsewhere);
const resolved = resolveWsTransport('wrun_1', directConfig);

expect(resolved?.wsUrl).toBe(wsUrl);
expect(resolved?.transport as unknown).toBe(openedElsewhere);
});
});

/**
* The flow route calls these on every message, for every World — including the
* HTTP default and the proxy World that can't speak WS at all. So "does
Expand Down
73 changes: 53 additions & 20 deletions packages/world-vercel/src/ws-transport.ts
Original file line number Diff line number Diff line change
Expand Up @@ -290,7 +290,9 @@ class WsEventsTransport {
clearTimeout(this.reconnectTimer);
this.reconnectTimer = null;
}
if (transports.get(this.wsUrl) === this) transports.delete(this.wsUrl);
if (wsEvents.transports.get(this.wsUrl) === this) {
wsEvents.transports.delete(this.wsUrl);
}
const conn = this.connection;
this.connection = null;
// Normal closure: a clean client-side release, not an aborted run.
Expand Down Expand Up @@ -660,7 +662,43 @@ class WsEventsTransport {
}
}

const transports = new Map<string, WsEventsTransport>();
/**
* On `globalThis`, not at module scope: a bundler can put two copies of this
* module in one process — Next.js gives the `instrumentation.ts` entry its own
* copy of the world now that it is bundled rather than external
* (vercel/workflow#3493) — and the open and the lookup run through different
* Worlds (`getWorldHandlers()` vs `getWorld()`, two `createWorld()` calls under
* two `globalThis` keys in core), so they can land in different copies. With the
* registry at module scope, opens went into one Map and every write read the
* other, empty one: socket open, every event silently on HTTP.
*
* `/v1` so a later change to `WsEventsTransport`'s shape gets a fresh registry

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

nice

* rather than a structurally-incompatible hit from an older copy.
*/
const WsEventsStateKey = Symbol.for(
'@workflow/world-vercel//wsEventsTransports/v1'
);

interface WsEventsState {
/** Open channels by `wsUrl`; see `getWsEventsTransport` below. */
transports: Map<string, WsEventsTransport>;
/** Log-once latches; per-copy ones would log once *per copy*. */
loggedWsInUse: boolean;
loggedWsProxyFallback: boolean;
}

const globalWsEventsState = globalThis as typeof globalThis & {
[WsEventsStateKey]?: WsEventsState;
};

// First copy to evaluate seeds the state; later copies adopt that same object,
// so nothing here may ever *replace* the slot.
const wsEvents: WsEventsState = globalWsEventsState[WsEventsStateKey] ?? {
transports: new Map(),
loggedWsInUse: false,
loggedWsProxyFallback: false,
};
globalWsEventsState[WsEventsStateKey] = wsEvents;

/**
* Get (or lazily create) the shared WS transport for `wsUrl`. `getHeaders` runs
Expand All @@ -678,26 +716,26 @@ export function getWsEventsTransport(
forceRefresh: boolean;
}) => Promise<Record<string, string>>
): WsEventsTransport {
let transport = transports.get(wsUrl);
let transport = wsEvents.transports.get(wsUrl);
if (!transport) {
transport = new WsEventsTransport(wsUrl, getHeaders);
transports.set(wsUrl, transport);
wsEvents.transports.set(wsUrl, transport);
}
return transport;
}

/**
* Test seam: close and drop every cached transport, and re-arm the
* once-per-process log latches below so a test asserting on either message
* isn't silenced by an earlier one having already logged it.
* once-per-process log latches so a test asserting on either message isn't
* silenced by an earlier one having already logged it.
*/
export function resetWsEventsTransportsForTest(): void {
for (const transport of [...transports.values()]) {
for (const transport of [...wsEvents.transports.values()]) {
transport.close('test reset');
}
transports.clear();
loggedWsProxyFallback = false;
loggedWsInUse = false;
wsEvents.transports.clear();
wsEvents.loggedWsProxyFallback = false;
wsEvents.loggedWsInUse = false;
}

/**
Expand Down Expand Up @@ -772,8 +810,8 @@ export function openWsChannel(
if (!isWsEventsTransportEnabled()) return undefined;
const resolved = resolveChannelUrl(runId, config);
if (!resolved) return undefined;
if (!loggedWsInUse) {
loggedWsInUse = true;
if (!wsEvents.loggedWsInUse) {
wsEvents.loggedWsInUse = true;
console.log(`world-vercel: using ws events transport (${resolved}).`);
}
// Cheap: a URL plus a map lookup, no token mint and no I/O. The socket work
Expand Down Expand Up @@ -829,11 +867,6 @@ async function refreshOidcTokenBestEffort(): Promise<void> {
}
}

// Each logged at most once per process — both branches below are expected
// to repeat (every event), and a per-request log would just be noise.
let loggedWsProxyFallback = false;
let loggedWsInUse = false;

/**
* Resolve this run's channel URL, or `null` when this World can't hold a socket
* at all and every caller must use HTTP. Says nothing about whether a channel is
Expand All @@ -855,8 +888,8 @@ function resolveChannelUrl(
// platform-level upgrade path, which is what surfaces as
// "experimental_upgradeWebSocket is not available in the current runtime
// environment". Fall back rather than fail a connection it can't serve.
if (!loggedWsProxyFallback) {
loggedWsProxyFallback = true;
if (!wsEvents.loggedWsProxyFallback) {
wsEvents.loggedWsProxyFallback = true;
console.warn(
`world-vercel: ws events transport requested but a World with projectConfig ` +
`(api-workflow proxy, resolved baseUrl: ${baseUrl}) is active — falling back.`
Expand Down Expand Up @@ -885,6 +918,6 @@ export function resolveWsTransport(
} | null {
const wsUrl = resolveChannelUrl(runId, config);
if (!wsUrl) return null;
const transport = transports.get(wsUrl);
const transport = wsEvents.transports.get(wsUrl);
return transport ? { transport, wsUrl } : null;
}
Loading