Skip to content
Open
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
55 changes: 55 additions & 0 deletions packages/core/src/__tests__/broker-transport.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,55 @@
import { describe, expect, it, vi } from 'vitest';
import {
HarnessBrokerTransport,
resolveBrokerTransportMode,
type BrokerTransportPort,
} from '../broker-transport.js';
import { WorkflowRunner } from '../runner.js';

describe('broker transport selection', () => {
it('defaults to legacy mode', () => {
// WorkflowRunner resolves the mode from process.env when no selector is
// passed; pin the variable so the assertion cannot depend on the shell.
const saved = process.env.RELAYFLOWS_INTEGRATION_TRANSPORT;
delete process.env.RELAYFLOWS_INTEGRATION_TRANSPORT;
try {
expect(resolveBrokerTransportMode(undefined, {})).toBe('legacy');
const runner = new WorkflowRunner();
expect((runner as any).brokerTransport).toBeInstanceOf(HarnessBrokerTransport);
expect((runner as any).brokerTransport.mode).toBe('legacy');
} finally {
if (saved !== undefined) process.env.RELAYFLOWS_INTEGRATION_TRANSPORT = saved;
}
});

it('prefers the per-run selector over the environment selector', () => {
const runner = new WorkflowRunner({
brokerTransportMode: 'shadow',
relay: { env: { RELAYFLOWS_INTEGRATION_TRANSPORT: 'adapter' } },
});
expect((runner as any).brokerTransport.mode).toBe('shadow');
});

it('uses the rollout environment selector when no per-run selector is supplied', () => {
const runner = new WorkflowRunner({
relay: { env: { RELAYFLOWS_INTEGRATION_TRANSPORT: 'adapter' } },
});
expect((runner as any).brokerTransport.mode).toBe('adapter');
});

it('prefers an explicitly injected port over both selectors', () => {
const explicit = { mode: 'legacy', shutdown: vi.fn() } as unknown as BrokerTransportPort;
const runner = new WorkflowRunner({
brokerTransport: explicit,
brokerTransportMode: 'shadow',
relay: { env: { RELAYFLOWS_INTEGRATION_TRANSPORT: 'adapter' } },
});
expect((runner as any).brokerTransport).toBe(explicit);
});

it('rejects an invalid rollout selector', () => {
expect(() =>
resolveBrokerTransportMode(undefined, { RELAYFLOWS_INTEGRATION_TRANSPORT: 'invalid' })
).toThrow('Invalid broker transport mode');
});
});
2 changes: 1 addition & 1 deletion packages/core/src/__tests__/idle-nudge.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -225,7 +225,7 @@ describe('Idle Nudge Detection', () => {
const agentDef = { name: 'worker', cli: 'claude' };

(runner as any).currentConfig = config;
(runner as any).relay = { sendMessage: mockSendMessage };
(runner as any).brokerTransport.relay = { sendMessage: mockSendMessage };
const result = await (runner as any).waitForExitWithIdleNudging(
wrappedMockAgent(),
agentDef,
Expand Down
29 changes: 24 additions & 5 deletions packages/core/src/__tests__/workflow-runner.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -650,14 +650,28 @@ agents:

await Promise.all(
runners.map((candidate, index) =>
(candidate as any).startOrReuseSharedBroker(`run-${index}`, `wf-${index}`, true)
(candidate as any).brokerTransport.start(
{
runId: `run-${index}`,
brokerName: `relayflows-run-${index}`,
channel: `wf-${index}`,
relaycastDisabled: true,
},
{
onEvent: vi.fn(),
onLog: vi.fn(),
getActiveAgentNames: () => [],
}
)
)
);

expect(HarnessDriverClient.spawn).toHaveBeenCalledTimes(1);
expect(HarnessDriverClient.connect).toHaveBeenCalledTimes(3);

const starter = runners.find((candidate) => (candidate as any).sharedBrokerLease?.startedBroker);
const starter = runners.find(
(candidate) => (candidate as any).brokerTransport.sharedBrokerLease?.startedBroker
);
expect(starter).toBeDefined();
const attached = runners.filter((candidate) => candidate !== starter);

Expand Down Expand Up @@ -810,15 +824,20 @@ agents:

it('refuses broker recovery while workflow agents are still active', async () => {
const candidate = new WorkflowRunner({ db, workspaceId: 'ws-test' });
(candidate as any).currentBrokerContext = {
const transport = (candidate as any).brokerTransport;
transport.context = {
runId: 'run-active',
brokerName: 'relayflows-run-active',
channel: 'wf-active',
relaycastDisabled: true,
};
(candidate as any).activeAgentHandles.set('active-agent', { name: 'active-agent' });
transport.hooks = {
onEvent: vi.fn(),
onLog: vi.fn(),
getActiveAgentNames: () => ['active-agent'],
};

await expect((candidate as any).recoverBroker('spawn failed')).rejects.toThrow(
await expect(transport.recoverBroker('spawn failed')).rejects.toThrow(
'Broker recovery is unsafe while 1 agent is still active: active-agent'
);
expect(mockHarnessDriverSpawn).not.toHaveBeenCalled();
Expand Down
6 changes: 3 additions & 3 deletions packages/core/src/agent-handle.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,16 +8,16 @@
* string contract the workflow runner has always used (`'exited'` / `'timeout'`
* / `'idle'`), so `runner.ts` keeps consuming handles unchanged.
*/
import type { SpawnedAgentHandle } from '@agent-relay/harness-driver';
import type { BrokerAgentHandle } from './broker-transport.js';

export class WorkflowAgentHandle {
constructor(private readonly inner: SpawnedAgentHandle) {}
constructor(private readonly inner: BrokerAgentHandle) {}

get name(): string {
return this.inner.name;
}

get runtime(): SpawnedAgentHandle['runtime'] {
get runtime(): BrokerAgentHandle['runtime'] {
return this.inner.runtime;
}

Expand Down
Loading
Loading