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
338 changes: 338 additions & 0 deletions packages/core/src/__tests__/workflow-runner.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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: [] }),
Expand Down Expand Up @@ -497,6 +498,343 @@ 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<typeof WorkflowRunner> | 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('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({
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('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',
])('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');
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.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.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);
mockRelayInstance.getSession.mockResolvedValue({ relay_base_url: 'https://api.relaycast.dev' });
rmSync(tmpDir, { recursive: true, force: true });
}
});

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');
// 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 {
await localRunner.shutdownRelay().catch(() => undefined);
mockRelayInstance.getSession.mockResolvedValue({ relay_base_url: 'https://api.relaycast.dev' });
rmSync(tmpDir, { recursive: true, force: true });
}
});

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');
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', () => {
// `ensureRelaycastApiKey` promises "each run gets full isolation" by creating
// a fresh workspace per run. It early-returns when `relayApiKey` is already
Expand Down
Loading
Loading