From 7886d9dd6e27086898ebe446930804dd6fae5d99 Mon Sep 17 00:00:00 2001 From: Khaliq Date: Thu, 10 Sep 2026 06:47:45 +0200 Subject: [PATCH 1/5] fix(core): align Relaycast broker base URL --- .../src/__tests__/workflow-runner.test.ts | 150 ++++++++++++++++++ packages/core/src/runner.ts | 85 ++++++++-- 2 files changed, 222 insertions(+), 13 deletions(-) diff --git a/packages/core/src/__tests__/workflow-runner.test.ts b/packages/core/src/__tests__/workflow-runner.test.ts index bdb2228..a6fc463 100644 --- a/packages/core/src/__tests__/workflow-runner.test.ts +++ b/packages/core/src/__tests__/workflow-runner.test.ts @@ -154,6 +154,7 @@ const mockRelayInstance = { onBrokerExit: vi.fn(() => () => {}), connectEvents: vi.fn(), getStatus: vi.fn().mockResolvedValue({ status: 'running' }), + getSession: vi.fn().mockResolvedValue({ relay_base_url: 'https://api.relaycast.dev' }), listAgents: vi.fn().mockResolvedValue([]), release: vi.fn().mockResolvedValue({ name: '' }), sendMessage: vi.fn().mockResolvedValue({ event_id: 'evt', targets: [] }), @@ -497,6 +498,155 @@ agents: // ── Execution ────────────────────────────────────────────────────────── + describe('Relaycast base URL consistency', () => { + it('uses the default origin for workspace creation, observer minting, broker, and child env', async () => { + const tmpDir = mkdtempSync(path.join(os.tmpdir(), 'relayflows-base-url-')); + const fetchSpy = vi.spyOn(globalThis, 'fetch').mockImplementation((async (url: string) => { + if (url.includes('/v1/observer-tokens')) { + return { + ok: true, + json: async () => ({ data: { token: 'ot_live_test', id: 'ot_test' } }), + } as Response; + } + return { + ok: true, + json: async () => ({ data: { api_key: 'rk_live_test' } }), + text: async () => '', + } as Response; + }) as unknown as typeof fetch); + let localRunner: InstanceType | undefined; + + try { + vi.stubEnv('RELAYCAST_BASE_URL', ''); + vi.stubEnv('RELAY_BASE_URL', ''); + localRunner = new WorkflowRunner({ db, cwd: tmpDir }); + + await (localRunner as any).ensureRelaycastApiKey('wf-default'); + await expect((localRunner as any).mintRunObserverUrl('wf-default')).resolves.toContain( + 'ot_live_test' + ); + await (localRunner as any).startOrReuseSharedBroker('run-default', 'wf-default', false); + + const engineCalls = fetchSpy.mock.calls + .map(([url]) => String(url)) + .filter((url) => url.includes('/v1/')); + expect(engineCalls[0]).toBe('https://api.relaycast.dev/v1/workspaces'); + expect(engineCalls[1]).toBe('https://api.relaycast.dev/v1/observer-tokens'); + + const relayEnv = (localRunner as any).getRelayEnv(); + expect(relayEnv).toEqual( + expect.objectContaining({ + RELAYCAST_BASE_URL: 'https://api.relaycast.dev', + RELAY_BASE_URL: 'https://api.relaycast.dev', + }) + ); + const brokerOptions = mockHarnessDriverSpawn.mock.calls.at(-1)?.[0] as any; + expect(brokerOptions.env).toEqual( + expect.objectContaining({ + RELAYCAST_BASE_URL: 'https://api.relaycast.dev', + RELAY_BASE_URL: 'https://api.relaycast.dev', + }) + ); + } finally { + await localRunner?.shutdownRelay().catch(() => undefined); + fetchSpy.mockRestore(); + vi.unstubAllEnvs(); + rmSync(tmpDir, { recursive: true, force: true }); + } + }); + + it('normalizes an explicit override across workspace creation, observer minting, broker, and children', async () => { + const tmpDir = mkdtempSync(path.join(os.tmpdir(), 'relayflows-base-url-explicit-')); + const fetchSpy = vi.spyOn(globalThis, 'fetch').mockImplementation((async () => ({ + ok: true, + json: async () => ({ data: { api_key: 'rk_live_test', token: 'ot_live_test', id: 'ot_test' } }), + text: async () => '', + })) as unknown as typeof fetch); + const localRunner = new WorkflowRunner({ + db, + cwd: tmpDir, + relay: { env: { RELAYCAST_BASE_URL: 'https://engine.example.test///' } }, + }); + + try { + await (localRunner as any).ensureRelaycastApiKey('wf-explicit'); + await (localRunner as any).mintRunObserverUrl('wf-explicit'); + await (localRunner as any).startOrReuseSharedBroker('run-explicit', 'wf-explicit', false); + expect((localRunner as any).getRelaycastBaseUrl()).toBe('https://engine.example.test'); + expect( + fetchSpy.mock.calls + .map(([url]) => String(url)) + .filter((url) => url.includes('/v1/')) + ).toEqual([ + 'https://engine.example.test/v1/workspaces', + 'https://engine.example.test/v1/observer-tokens', + ]); + expect((localRunner as any).getRelayEnv()).toEqual( + expect.objectContaining({ + RELAYCAST_BASE_URL: 'https://engine.example.test', + RELAY_BASE_URL: 'https://engine.example.test', + }) + ); + const brokerOptions = mockHarnessDriverSpawn.mock.calls.at(-1)?.[0] as any; + expect(brokerOptions.env).toEqual( + expect.objectContaining({ + RELAYCAST_BASE_URL: 'https://engine.example.test', + RELAY_BASE_URL: 'https://engine.example.test', + }) + ); + } finally { + await localRunner.shutdownRelay().catch(() => undefined); + fetchSpy.mockRestore(); + vi.unstubAllEnvs(); + rmSync(tmpDir, { recursive: true, force: true }); + } + }); + + it('fails closed on an invalid explicit origin before creating a workspace', async () => { + const fetchSpy = vi.spyOn(globalThis, 'fetch'); + const localRunner = new WorkflowRunner({ + db, + relay: { env: { RELAYCAST_BASE_URL: 'javascript:alert(1)' } }, + }); + + try { + await expect((localRunner as any).ensureRelaycastApiKey('wf-invalid')).rejects.toThrow( + 'Relaycast base URL must use http or https' + ); + expect(fetchSpy).not.toHaveBeenCalled(); + } finally { + fetchSpy.mockRestore(); + } + }); + + it('does not reuse a shared broker whose reported origin differs from the runner origin', async () => { + const tmpDir = mkdtempSync(path.join(os.tmpdir(), 'relayflows-base-url-reuse-')); + const stateDir = path.join(tmpDir, '.agentworkforce', 'relay'); + mkdirSync(stateDir, { recursive: true }); + writeFileSync( + path.join(stateDir, 'connection.json'), + JSON.stringify({ url: 'http://127.0.0.1:3889', api_key: 'br_test', pid: process.pid }), + 'utf-8' + ); + mockRelayInstance.getSession.mockResolvedValue({ relay_base_url: 'https://cast.agentrelay.com' }); + const localRunner = new WorkflowRunner({ + db, + cwd: tmpDir, + relay: { env: { RELAY_API_KEY: 'rk_live_test', RELAYCAST_BASE_URL: 'https://api.relaycast.dev' } }, + }); + + try { + await (localRunner as any).startOrReuseSharedBroker('run-mismatch', 'wf-mismatch', false); + expect(mockRelayInstance.disconnect).toHaveBeenCalled(); + expect(mockHarnessDriverSpawn).toHaveBeenCalled(); + } finally { + await localRunner.shutdownRelay().catch(() => undefined); + mockRelayInstance.getSession.mockResolvedValue({ relay_base_url: 'https://api.relaycast.dev' }); + rmSync(tmpDir, { recursive: true, force: true }); + } + }); + }); + describe('execute', () => { // `ensureRelaycastApiKey` promises "each run gets full isolation" by creating // a fresh workspace per run. It early-returns when `relayApiKey` is already diff --git a/packages/core/src/runner.ts b/packages/core/src/runner.ts index 304bd14..4d723a1 100644 --- a/packages/core/src/runner.ts +++ b/packages/core/src/runner.ts @@ -152,6 +152,9 @@ import { WorkflowCompletionError, } from './verification.js'; +/** One engine origin for workspace, observer, broker, and child-agent calls. */ +export const DEFAULT_RELAYCAST_BASE_URL = 'https://api.relaycast.dev'; + // ── Broker client / messaging imports ─────────────────────────────────────── // Broker / PTY / lifecycle is driven by the harness-driver client; messaging @@ -182,6 +185,7 @@ const ENV_ALLOWLIST = new Set([ 'RUST_LOG', 'RUST_BACKTRACE', 'RELAY_API_KEY', + 'RELAY_BASE_URL', 'RELAYCAST_BASE_URL', 'RELAY_LLM_PROXY', 'RELAY_LLM_PROXY_URL', @@ -240,6 +244,39 @@ function filteredEnv(extra?: Record): Record): string { + const configured = [ + env.RELAYCAST_BASE_URL, + env.RELAY_BASE_URL, + ] + .map((value) => value?.trim()) + .find((value): value is string => Boolean(value)); + const value = configured ?? DEFAULT_RELAYCAST_BASE_URL; + + let parsed: URL; + try { + parsed = new URL(value); + } catch { + throw new Error(`Invalid Relaycast base URL: ${value}`); + } + if (parsed.protocol !== 'http:' && parsed.protocol !== 'https:') { + throw new Error(`Relaycast base URL must use http or https: ${value}`); + } + if (!parsed.hostname) { + throw new Error(`Relaycast base URL must include a hostname: ${value}`); + } + + return value.replace(/\/+$/, ''); +} + // ── Shared broker coordination ────────────────────────────────────────────── const BROKER_CONNECTION_FILENAME = 'connection.json'; @@ -2433,10 +2470,7 @@ export class WorkflowRunner { // Always create a fresh workspace — each run gets full isolation. const workspaceName = `relay-${channel}-${randomBytes(4).toString('hex')}`; - const baseUrl = - this.relayOptions.env?.RELAYCAST_BASE_URL ?? - process.env.RELAYCAST_BASE_URL ?? - 'https://api.relaycast.dev'; + const baseUrl = this.getRelaycastBaseUrl(); const res = await fetch(`${baseUrl}/v1/workspaces`, { method: 'POST', headers: { 'content-type': 'application/json' }, @@ -2705,10 +2739,19 @@ export class WorkflowRunner { } private getMergedRelayEnvSource(): NodeJS.ProcessEnv { + const baseUrl = this.relayApiKey ? this.getRelaycastBaseUrl() : undefined; return { ...process.env, ...(this.relayOptions.env ?? {}), - ...(this.relayApiKey ? { RELAY_API_KEY: this.relayApiKey } : {}), + ...(this.relayApiKey + ? { + RELAY_API_KEY: this.relayApiKey, + // Keep both names aligned: the broker reads RELAYCAST_BASE_URL, + // while its spawned MCP clients use the legacy RELAY_BASE_URL. + RELAYCAST_BASE_URL: baseUrl, + RELAY_BASE_URL: baseUrl, + } + : {}), }; } @@ -2757,7 +2800,8 @@ export class WorkflowRunner { private async tryConnectSharedBroker( connectionPath: string, - brokerCwd: string + brokerCwd: string, + expectedBaseUrl?: string ): Promise { const conn = readBrokerConnectionFile(connectionPath); if (!conn) { @@ -2772,6 +2816,15 @@ export class WorkflowRunner { try { const client = HarnessDriverClient.connect({ cwd: brokerCwd, connectionPath }); await client.getStatus(); + if (expectedBaseUrl !== undefined) { + // Older harness-driver type declarations omit this field even though + // current brokers return it from /api/session. + const session = (await client.getSession()) as { relay_base_url?: string }; + if (session.relay_base_url !== expectedBaseUrl) { + this.disconnectRelayClient(client); + return null; + } + } return client; } catch { return null; @@ -2917,10 +2970,11 @@ export class WorkflowRunner { const connectionPath = path.join(stateDir, BROKER_CONNECTION_FILENAME); const startupTimeoutMs = this.relayOptions.startupTimeoutMs ?? SHARED_BROKER_DEFAULT_STARTUP_TIMEOUT_MS; + const expectedBaseUrl = relaycastDisabled ? undefined : this.getRelaycastBaseUrl(); const lease = this.createSharedBrokerLease(stateDir, connectionPath, runId, false); this.sharedBrokerLease = lease; - const existing = await this.tryConnectSharedBroker(connectionPath, brokerCwd); + const existing = await this.tryConnectSharedBroker(connectionPath, brokerCwd, expectedBaseUrl); if (existing) { this.log('Reusing shared broker...'); this.relay = existing; @@ -2929,7 +2983,7 @@ export class WorkflowRunner { const releaseLock = await this.acquireSharedBrokerStartLock(stateDir, startupTimeoutMs); try { - const lockedExisting = await this.tryConnectSharedBroker(connectionPath, brokerCwd); + const lockedExisting = await this.tryConnectSharedBroker(connectionPath, brokerCwd, expectedBaseUrl); if (lockedExisting) { this.log('Reusing shared broker...'); this.relay = lockedExisting; @@ -2944,6 +2998,12 @@ export class WorkflowRunner { const relayEnv = { ...(this.getRelayEnv() ?? filteredEnv()), AGENT_RELAY_STATE_DIR: stateDir, + // Harness-driver/broker has its own hosted default. Always pass the + // runner's resolved origin so auto-created keys and broker auth land + // on the same engine, including the default path. + ...(expectedBaseUrl + ? { RELAYCAST_BASE_URL: expectedBaseUrl, RELAY_BASE_URL: expectedBaseUrl } + : {}), }; this.relay = await HarnessDriverClient.spawn({ ...this.relayOptions, @@ -3060,11 +3120,10 @@ export class WorkflowRunner { } private getRelaycastBaseUrl(): string { - return ( - this.relayOptions.env?.RELAYCAST_BASE_URL ?? - process.env.RELAYCAST_BASE_URL ?? - 'https://api.relaycast.dev' - ); + return resolveRelaycastBaseUrl({ + ...process.env, + ...(this.relayOptions.env ?? {}), + }); } private getRelaycastClient(): RelayCast { From 3e652009800480afaed621bba9e15a25b621fcd8 Mon Sep 17 00:00:00 2001 From: Khaliq Date: Thu, 10 Sep 2026 07:01:46 +0200 Subject: [PATCH 2/5] fix(core): handle legacy broker origin capability --- .../src/__tests__/workflow-runner.test.ts | 75 +++++++++++++++++++ packages/core/src/runner.ts | 37 +++++++-- 2 files changed, 105 insertions(+), 7 deletions(-) diff --git a/packages/core/src/__tests__/workflow-runner.test.ts b/packages/core/src/__tests__/workflow-runner.test.ts index a6fc463..fb178a6 100644 --- a/packages/core/src/__tests__/workflow-runner.test.ts +++ b/packages/core/src/__tests__/workflow-runner.test.ts @@ -619,6 +619,23 @@ agents: } }); + it.each([ + 'https://engine.example.test?tenant=one', + 'https://engine.example.test#fragment', + ])('fails closed when an explicit origin contains a query or fragment: %s', async (baseUrl) => { + const fetchSpy = vi.spyOn(globalThis, 'fetch'); + const localRunner = new WorkflowRunner({ db, relay: { env: { RELAYCAST_BASE_URL: baseUrl } } }); + + try { + await expect((localRunner as any).ensureRelaycastApiKey('wf-invalid')).rejects.toThrow( + 'Relaycast base URL must not include a query or fragment' + ); + expect(fetchSpy).not.toHaveBeenCalled(); + } finally { + fetchSpy.mockRestore(); + } + }); + it('does not reuse a shared broker whose reported origin differs from the runner origin', async () => { const tmpDir = mkdtempSync(path.join(os.tmpdir(), 'relayflows-base-url-reuse-')); const stateDir = path.join(tmpDir, '.agentworkforce', 'relay'); @@ -645,6 +662,64 @@ agents: rmSync(tmpDir, { recursive: true, force: true }); } }); + + it('reuses a legacy shared broker when origin reporting is unavailable', async () => { + const tmpDir = mkdtempSync(path.join(os.tmpdir(), 'relayflows-base-url-legacy-')); + const stateDir = path.join(tmpDir, '.agentworkforce', 'relay'); + mkdirSync(stateDir, { recursive: true }); + writeFileSync( + path.join(stateDir, 'connection.json'), + JSON.stringify({ url: 'http://127.0.0.1:3889', api_key: 'br_test', pid: process.pid }), + 'utf-8' + ); + mockRelayInstance.getSession.mockResolvedValue({}); + const localRunner = new WorkflowRunner({ + db, + cwd: tmpDir, + relay: { env: { RELAY_API_KEY: 'rk_live_test', RELAYCAST_BASE_URL: 'https://api.relaycast.dev' } }, + }); + + try { + await (localRunner as any).startOrReuseSharedBroker('run-legacy', 'wf-legacy', false); + expect(mockRelayInstance.disconnect).not.toHaveBeenCalled(); + expect(mockHarnessDriverSpawn).not.toHaveBeenCalled(); + expect((localRunner as any).relay).toBe(mockRelayInstance); + } finally { + await localRunner.shutdownRelay().catch(() => undefined); + mockRelayInstance.getSession.mockResolvedValue({ relay_base_url: 'https://api.relaycast.dev' }); + rmSync(tmpDir, { recursive: true, force: true }); + } + }); + + it('reuses a shared broker when reported origin differs only by casing, default port, and slash', async () => { + const tmpDir = mkdtempSync(path.join(os.tmpdir(), 'relayflows-base-url-canonical-')); + const stateDir = path.join(tmpDir, '.agentworkforce', 'relay'); + mkdirSync(stateDir, { recursive: true }); + writeFileSync( + path.join(stateDir, 'connection.json'), + JSON.stringify({ url: 'http://127.0.0.1:3889', api_key: 'br_test', pid: process.pid }), + 'utf-8' + ); + mockRelayInstance.getSession.mockResolvedValue({ + relay_base_url: 'https://API.RELAYCAST.DEV:443///', + }); + const localRunner = new WorkflowRunner({ + db, + cwd: tmpDir, + relay: { env: { RELAY_API_KEY: 'rk_live_test', RELAYCAST_BASE_URL: 'https://api.relaycast.dev' } }, + }); + + try { + await (localRunner as any).startOrReuseSharedBroker('run-canonical', 'wf-canonical', false); + expect(mockRelayInstance.disconnect).not.toHaveBeenCalled(); + expect(mockHarnessDriverSpawn).not.toHaveBeenCalled(); + expect((localRunner as any).relay).toBe(mockRelayInstance); + } finally { + await localRunner.shutdownRelay().catch(() => undefined); + mockRelayInstance.getSession.mockResolvedValue({ relay_base_url: 'https://api.relaycast.dev' }); + rmSync(tmpDir, { recursive: true, force: true }); + } + }); }); describe('execute', () => { diff --git a/packages/core/src/runner.ts b/packages/core/src/runner.ts index 4d723a1..13071af 100644 --- a/packages/core/src/runner.ts +++ b/packages/core/src/runner.ts @@ -273,8 +273,15 @@ function resolveRelaycastBaseUrl(env: Record): strin if (!parsed.hostname) { throw new Error(`Relaycast base URL must include a hostname: ${value}`); } + if (parsed.search || parsed.hash) { + throw new Error(`Relaycast base URL must not include a query or fragment: ${value}`); + } - return value.replace(/\/+$/, ''); + // URL#origin normalizes hostname casing and removes default ports. Keep a + // configured path (some self-hosted deployments mount the API below one), + // while removing insignificant trailing slashes. + const pathname = parsed.pathname.replace(/\/+$/, ''); + return `${parsed.origin}${pathname}`; } // ── Shared broker coordination ────────────────────────────────────────────── @@ -2817,12 +2824,28 @@ export class WorkflowRunner { const client = HarnessDriverClient.connect({ cwd: brokerCwd, connectionPath }); await client.getStatus(); if (expectedBaseUrl !== undefined) { - // Older harness-driver type declarations omit this field even though - // current brokers return it from /api/session. - const session = (await client.getSession()) as { relay_base_url?: string }; - if (session.relay_base_url !== expectedBaseUrl) { - this.disconnectRelayClient(client); - return null; + // `relay_base_url` was added to the broker session after the pinned + // harness-driver release. Treat its absence as an older broker's + // missing capability; brokers that report an origin are checked + // strictly so a known mismatch is never reused. + const getSession = (client as { getSession?: () => Promise }).getSession; + if (typeof getSession === 'function') { + const session = (await getSession.call(client)) as { relay_base_url?: unknown }; + if (typeof session.relay_base_url === 'string' && session.relay_base_url.trim()) { + let reportedBaseUrl: string; + try { + reportedBaseUrl = resolveRelaycastBaseUrl({ + RELAYCAST_BASE_URL: session.relay_base_url, + }); + } catch { + this.disconnectRelayClient(client); + return null; + } + if (reportedBaseUrl !== expectedBaseUrl) { + this.disconnectRelayClient(client); + return null; + } + } } } return client; From e11bf644b14cbfb8d731badc4cf5abdfc38296bd Mon Sep 17 00:00:00 2001 From: Khaliq Date: Thu, 10 Sep 2026 07:22:04 +0200 Subject: [PATCH 3/5] fix(core): coordinate shared broker origin changes --- .../src/__tests__/workflow-runner.test.ts | 84 +++++++++++- packages/core/src/runner.ts | 125 +++++++++++++----- 2 files changed, 176 insertions(+), 33 deletions(-) diff --git a/packages/core/src/__tests__/workflow-runner.test.ts b/packages/core/src/__tests__/workflow-runner.test.ts index fb178a6..3f0dfb4 100644 --- a/packages/core/src/__tests__/workflow-runner.test.ts +++ b/packages/core/src/__tests__/workflow-runner.test.ts @@ -602,6 +602,21 @@ agents: } }); + it('prioritizes runner-level aliases over process-level aliases', () => { + vi.stubEnv('RELAYCAST_BASE_URL', 'https://process-canonical.example.test'); + vi.stubEnv('RELAY_BASE_URL', 'https://process-legacy.example.test'); + const localRunner = new WorkflowRunner({ + db, + relay: { env: { RELAY_BASE_URL: 'https://runner-legacy.example.test' } }, + }); + + try { + expect((localRunner as any).getRelaycastBaseUrl()).toBe('https://runner-legacy.example.test'); + } finally { + vi.unstubAllEnvs(); + } + }); + it('fails closed on an invalid explicit origin before creating a workspace', async () => { const fetchSpy = vi.spyOn(globalThis, 'fetch'); const localRunner = new WorkflowRunner({ @@ -619,6 +634,28 @@ agents: } }); + it('allows loopback HTTP for local Relaycast engines but rejects remote HTTP', async () => { + const fetchSpy = vi.spyOn(globalThis, 'fetch'); + const localRunner = new WorkflowRunner({ + db, + relay: { env: { RELAYCAST_BASE_URL: 'http://localhost:4000///' } }, + }); + const remoteRunner = new WorkflowRunner({ + db, + relay: { env: { RELAYCAST_BASE_URL: 'http://engine.example.test' } }, + }); + + try { + expect((localRunner as any).getRelaycastBaseUrl()).toBe('http://localhost:4000'); + await expect((remoteRunner as any).ensureRelaycastApiKey('wf-http')).rejects.toThrow( + 'Relaycast base URL must use https except for loopback hosts' + ); + expect(fetchSpy).not.toHaveBeenCalled(); + } finally { + fetchSpy.mockRestore(); + } + }); + it.each([ 'https://engine.example.test?tenant=one', 'https://engine.example.test#fragment', @@ -645,6 +682,11 @@ agents: JSON.stringify({ url: 'http://127.0.0.1:3889', api_key: 'br_test', pid: process.pid }), 'utf-8' ); + writeFileSync( + path.join(stateDir, 'relayflows-owner.json'), + JSON.stringify({ pid: process.pid }), + 'utf-8' + ); mockRelayInstance.getSession.mockResolvedValue({ relay_base_url: 'https://cast.agentrelay.com' }); const localRunner = new WorkflowRunner({ db, @@ -654,7 +696,10 @@ agents: try { await (localRunner as any).startOrReuseSharedBroker('run-mismatch', 'wf-mismatch', false); - expect(mockRelayInstance.disconnect).toHaveBeenCalled(); + expect(mockRelayInstance.shutdown).toHaveBeenCalled(); + // The unlocked probe disconnects its transport; the locked probe then + // owns and shuts down the incompatible broker before replacement. + expect(mockRelayInstance.disconnect).toHaveBeenCalledTimes(1); expect(mockHarnessDriverSpawn).toHaveBeenCalled(); } finally { await localRunner.shutdownRelay().catch(() => undefined); @@ -663,6 +708,43 @@ agents: } }); + it('fails closed without replacing a mismatched broker used by another workflow', async () => { + const tmpDir = mkdtempSync(path.join(os.tmpdir(), 'relayflows-base-url-conflict-')); + const stateDir = path.join(tmpDir, '.agentworkforce', 'relay'); + const leaseDir = path.join(stateDir, 'relayflows-runs'); + mkdirSync(leaseDir, { recursive: true }); + writeFileSync( + path.join(stateDir, 'connection.json'), + JSON.stringify({ url: 'http://127.0.0.1:3889', api_key: 'br_test', pid: process.pid }), + 'utf-8' + ); + writeFileSync(path.join(stateDir, 'relayflows-owner.json'), JSON.stringify({ pid: process.pid }), 'utf-8'); + writeFileSync( + path.join(leaseDir, 'other-workflow.json'), + JSON.stringify({ pid: process.pid }), + 'utf-8' + ); + mockRelayInstance.getSession.mockResolvedValue({ relay_base_url: 'https://cast.agentrelay.com' }); + const localRunner = new WorkflowRunner({ + db, + cwd: tmpDir, + relay: { env: { RELAY_API_KEY: 'rk_live_test', RELAYCAST_BASE_URL: 'https://api.relaycast.dev' } }, + }); + + try { + await expect( + (localRunner as any).startOrReuseSharedBroker('run-conflict', 'wf-conflict', false) + ).rejects.toThrow('Cannot replace shared broker'); + expect(mockRelayInstance.disconnect).toHaveBeenCalled(); + expect(mockRelayInstance.shutdown).not.toHaveBeenCalled(); + expect(mockHarnessDriverSpawn).not.toHaveBeenCalled(); + } finally { + await localRunner.shutdownRelay().catch(() => undefined); + mockRelayInstance.getSession.mockResolvedValue({ relay_base_url: 'https://api.relaycast.dev' }); + rmSync(tmpDir, { recursive: true, force: true }); + } + }); + it('reuses a legacy shared broker when origin reporting is unavailable', async () => { const tmpDir = mkdtempSync(path.join(os.tmpdir(), 'relayflows-base-url-legacy-')); const stateDir = path.join(tmpDir, '.agentworkforce', 'relay'); diff --git a/packages/core/src/runner.ts b/packages/core/src/runner.ts index 13071af..79977a7 100644 --- a/packages/core/src/runner.ts +++ b/packages/core/src/runner.ts @@ -252,10 +252,15 @@ function filteredEnv(extra?: Record): Record): string { +function resolveRelaycastBaseUrl( + optionEnv: Record, + processEnv: Record = {} +): string { const configured = [ - env.RELAYCAST_BASE_URL, - env.RELAY_BASE_URL, + optionEnv.RELAYCAST_BASE_URL, + optionEnv.RELAY_BASE_URL, + processEnv.RELAYCAST_BASE_URL, + processEnv.RELAY_BASE_URL, ] .map((value) => value?.trim()) .find((value): value is string => Boolean(value)); @@ -273,6 +278,14 @@ function resolveRelaycastBaseUrl(env: Record): strin if (!parsed.hostname) { throw new Error(`Relaycast base URL must include a hostname: ${value}`); } + const isLoopbackHost = + parsed.hostname === 'localhost' || + parsed.hostname === '127.0.0.1' || + parsed.hostname === '[::1]' || + parsed.hostname.endsWith('.localhost'); + if (parsed.protocol === 'http:' && !isLoopbackHost) { + throw new Error(`Relaycast base URL must use https except for loopback hosts: ${value}`); + } if (parsed.search || parsed.hash) { throw new Error(`Relaycast base URL must not include a query or fragment: ${value}`); } @@ -2808,7 +2821,8 @@ export class WorkflowRunner { private async tryConnectSharedBroker( connectionPath: string, brokerCwd: string, - expectedBaseUrl?: string + expectedBaseUrl?: string, + allowRetireMismatch = false ): Promise { const conn = readBrokerConnectionFile(connectionPath); if (!conn) { @@ -2820,38 +2834,80 @@ export class WorkflowRunner { return null; } + let client: HarnessDriverClient; try { - const client = HarnessDriverClient.connect({ cwd: brokerCwd, connectionPath }); + client = HarnessDriverClient.connect({ cwd: brokerCwd, connectionPath }); await client.getStatus(); - if (expectedBaseUrl !== undefined) { - // `relay_base_url` was added to the broker session after the pinned - // harness-driver release. Treat its absence as an older broker's - // missing capability; brokers that report an origin are checked - // strictly so a known mismatch is never reused. - const getSession = (client as { getSession?: () => Promise }).getSession; - if (typeof getSession === 'function') { - const session = (await getSession.call(client)) as { relay_base_url?: unknown }; - if (typeof session.relay_base_url === 'string' && session.relay_base_url.trim()) { - let reportedBaseUrl: string; - try { - reportedBaseUrl = resolveRelaycastBaseUrl({ - RELAYCAST_BASE_URL: session.relay_base_url, - }); - } catch { - this.disconnectRelayClient(client); + } catch { + return null; + } + + if (expectedBaseUrl !== undefined) { + // `relay_base_url` was added to the broker session after the pinned + // harness-driver release. Treat its absence as an older broker's + // missing capability; brokers that report an origin are checked + // strictly so a known mismatch is never reused. + const getSession = (client as { getSession?: () => Promise }).getSession; + if (typeof getSession === 'function') { + let session: { relay_base_url?: unknown }; + try { + session = (await getSession.call(client)) as { relay_base_url?: unknown }; + } catch { + this.disconnectRelayClient(client); + return null; + } + if (session && typeof session.relay_base_url === 'string' && session.relay_base_url.trim()) { + let reportedBaseUrl: string; + try { + reportedBaseUrl = resolveRelaycastBaseUrl({ + RELAYCAST_BASE_URL: session.relay_base_url, + }); + } catch { + if (allowRetireMismatch) { + await this.retireMismatchedSharedBroker(client, connectionPath, expectedBaseUrl); return null; } - if (reportedBaseUrl !== expectedBaseUrl) { - this.disconnectRelayClient(client); + this.disconnectRelayClient(client); + return null; + } + if (reportedBaseUrl !== expectedBaseUrl) { + if (allowRetireMismatch) { + await this.retireMismatchedSharedBroker(client, connectionPath, expectedBaseUrl); return null; } + this.disconnectRelayClient(client); + return null; } } } - return client; - } catch { - return null; } + return client; + } + + private async retireMismatchedSharedBroker( + relay: HarnessDriverClient, + connectionPath: string, + expectedBaseUrl: string + ): Promise { + const lease = this.sharedBrokerLease; + const otherLiveLeases = lease + ? this.countLiveSharedBrokerLeases(lease.stateDir, lease.leasePath) + : 0; + const workflowOwned = lease ? this.isWorkflowOwnedSharedBroker(lease) : false; + if (!lease || otherLiveLeases > 0 || !workflowOwned) { + throw new Error( + `Cannot replace shared broker for ${expectedBaseUrl}: it is still in use or not workflow-owned` + ); + } + + try { + await relay.shutdown(); + } catch (error) { + this.disconnectRelayClient(relay); + throw error; + } + safeUnlinkSync(connectionPath); + safeUnlinkSync(lease.ownerPath); } private async acquireSharedBrokerStartLock( @@ -2956,7 +3012,7 @@ export class WorkflowRunner { } } - private countLiveSharedBrokerLeases(stateDir: string): number { + private countLiveSharedBrokerLeases(stateDir: string, excludeLeasePath?: string): number { const leaseDir = path.join(stateDir, SHARED_BROKER_LEASE_DIRNAME); let entries: Dirent[]; try { @@ -2969,6 +3025,9 @@ export class WorkflowRunner { for (const entry of entries) { if (!entry.isFile()) continue; const leasePath = path.join(leaseDir, entry.name); + if (excludeLeasePath && path.resolve(leasePath) === path.resolve(excludeLeasePath)) { + continue; + } try { const lease = JSON.parse(readFileSync(leasePath, 'utf-8')) as { pid?: unknown }; if (typeof lease.pid === 'number' && lease.pid > 0 && isPidRunning(lease.pid)) { @@ -3006,7 +3065,12 @@ export class WorkflowRunner { const releaseLock = await this.acquireSharedBrokerStartLock(stateDir, startupTimeoutMs); try { - const lockedExisting = await this.tryConnectSharedBroker(connectionPath, brokerCwd, expectedBaseUrl); + const lockedExisting = await this.tryConnectSharedBroker( + connectionPath, + brokerCwd, + expectedBaseUrl, + true + ); if (lockedExisting) { this.log('Reusing shared broker...'); this.relay = lockedExisting; @@ -3143,10 +3207,7 @@ export class WorkflowRunner { } private getRelaycastBaseUrl(): string { - return resolveRelaycastBaseUrl({ - ...process.env, - ...(this.relayOptions.env ?? {}), - }); + return resolveRelaycastBaseUrl(this.relayOptions.env ?? {}, process.env); } private getRelaycastClient(): RelayCast { From 751252578739bbcf64851d247fa6407968ff70e0 Mon Sep 17 00:00:00 2001 From: Khaliq Date: Thu, 10 Sep 2026 07:23:24 +0200 Subject: [PATCH 4/5] fix(core): harden Relaycast origin selection --- .../src/__tests__/workflow-runner.test.ts | 29 +++++++++++++++++++ packages/core/src/runner.ts | 12 +++++--- 2 files changed, 37 insertions(+), 4 deletions(-) diff --git a/packages/core/src/__tests__/workflow-runner.test.ts b/packages/core/src/__tests__/workflow-runner.test.ts index 3f0dfb4..d013a0a 100644 --- a/packages/core/src/__tests__/workflow-runner.test.ts +++ b/packages/core/src/__tests__/workflow-runner.test.ts @@ -745,6 +745,35 @@ agents: } }); + it('disconnects a shared broker when its session capability check fails', async () => { + const tmpDir = mkdtempSync(path.join(os.tmpdir(), 'relayflows-base-url-session-error-')); + const stateDir = path.join(tmpDir, '.agentworkforce', 'relay'); + mkdirSync(stateDir, { recursive: true }); + writeFileSync( + path.join(stateDir, 'connection.json'), + JSON.stringify({ url: 'http://127.0.0.1:3889', api_key: 'br_test', pid: process.pid }), + 'utf-8' + ); + writeFileSync(path.join(stateDir, 'relayflows-owner.json'), JSON.stringify({ pid: process.pid }), 'utf-8'); + mockRelayInstance.getSession.mockRejectedValue(new Error('session unavailable')); + const localRunner = new WorkflowRunner({ + db, + cwd: tmpDir, + relay: { env: { RELAY_API_KEY: 'rk_live_test', RELAYCAST_BASE_URL: 'https://api.relaycast.dev' } }, + }); + + try { + await (localRunner as any).startOrReuseSharedBroker('run-session-error', 'wf-session-error', false); + expect(mockRelayInstance.disconnect).toHaveBeenCalledTimes(2); + expect(mockRelayInstance.shutdown).not.toHaveBeenCalled(); + expect(mockHarnessDriverSpawn).toHaveBeenCalled(); + } finally { + await localRunner.shutdownRelay().catch(() => undefined); + mockRelayInstance.getSession.mockResolvedValue({ relay_base_url: 'https://api.relaycast.dev' }); + rmSync(tmpDir, { recursive: true, force: true }); + } + }); + it('reuses a legacy shared broker when origin reporting is unavailable', async () => { const tmpDir = mkdtempSync(path.join(os.tmpdir(), 'relayflows-base-url-legacy-')); const stateDir = path.join(tmpDir, '.agentworkforce', 'relay'); diff --git a/packages/core/src/runner.ts b/packages/core/src/runner.ts index 79977a7..740325e 100644 --- a/packages/core/src/runner.ts +++ b/packages/core/src/runner.ts @@ -2849,18 +2849,22 @@ export class WorkflowRunner { // strictly so a known mismatch is never reused. const getSession = (client as { getSession?: () => Promise }).getSession; if (typeof getSession === 'function') { - let session: { relay_base_url?: unknown }; + let session: unknown; try { - session = (await getSession.call(client)) as { relay_base_url?: unknown }; + session = await getSession.call(client); } catch { this.disconnectRelayClient(client); return null; } - if (session && typeof session.relay_base_url === 'string' && session.relay_base_url.trim()) { + const reportedRelayBaseUrl = + session && typeof session === 'object' && 'relay_base_url' in session + ? (session as { relay_base_url?: unknown }).relay_base_url + : undefined; + if (typeof reportedRelayBaseUrl === 'string' && reportedRelayBaseUrl.trim()) { let reportedBaseUrl: string; try { reportedBaseUrl = resolveRelaycastBaseUrl({ - RELAYCAST_BASE_URL: session.relay_base_url, + RELAYCAST_BASE_URL: reportedRelayBaseUrl, }); } catch { if (allowRetireMismatch) { From 0050d0ef64faae38d9f7f9de5e445a078d70bfbc Mon Sep 17 00:00:00 2001 From: Khaliq Date: Thu, 10 Sep 2026 07:39:12 +0200 Subject: [PATCH 5/5] fix(core): disconnect rejected broker probes --- packages/core/src/__tests__/workflow-runner.test.ts | 4 +++- packages/core/src/runner.ts | 3 +++ 2 files changed, 6 insertions(+), 1 deletion(-) diff --git a/packages/core/src/__tests__/workflow-runner.test.ts b/packages/core/src/__tests__/workflow-runner.test.ts index d013a0a..76054cc 100644 --- a/packages/core/src/__tests__/workflow-runner.test.ts +++ b/packages/core/src/__tests__/workflow-runner.test.ts @@ -735,7 +735,9 @@ agents: await expect( (localRunner as any).startOrReuseSharedBroker('run-conflict', 'wf-conflict', false) ).rejects.toThrow('Cannot replace shared broker'); - expect(mockRelayInstance.disconnect).toHaveBeenCalled(); + // Both the unlocked probe and the lock-protected probe connected to + // the incompatible broker and must release their transports. + expect(mockRelayInstance.disconnect).toHaveBeenCalledTimes(2); expect(mockRelayInstance.shutdown).not.toHaveBeenCalled(); expect(mockHarnessDriverSpawn).not.toHaveBeenCalled(); } finally { diff --git a/packages/core/src/runner.ts b/packages/core/src/runner.ts index 740325e..781a273 100644 --- a/packages/core/src/runner.ts +++ b/packages/core/src/runner.ts @@ -2899,6 +2899,9 @@ export class WorkflowRunner { : 0; const workflowOwned = lease ? this.isWorkflowOwnedSharedBroker(lease) : false; if (!lease || otherLiveLeases > 0 || !workflowOwned) { + // The locked probe owns a live transport even when replacement is + // refused; release it before propagating the fail-closed error. + this.disconnectRelayClient(relay); throw new Error( `Cannot replace shared broker for ${expectedBaseUrl}: it is still in use or not workflow-owned` );