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
22 changes: 22 additions & 0 deletions packages/local-surface/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -35,3 +35,25 @@ flushes once more during teardown.
`defineWorkforcePersonaNode` remains the long-lived channel `onMessage` surface.
It composes `@agentworkforce/deploy` for a persona that consumes Relay message
events rather than launching an interactive worker per request.

### Prepared execution ownership

Both `defineWorkforcePersonaSpawnNode` and `workforcePersonaSpawnCapability`
accept `onExecutionPrepared(name, execution)`. The factory awaits this callback
once per prepared launch, including coalesced requests, before asking the broker
to spawn. `execution.handle` is the actual `ExecutionHandle` returned by the
persona executor; `execution.scratchDir` is its factory-owned parent directory.
Hosts can retain this receipt and persist ownership before delegation begins.
Comment thread
cubic-dev-ai[bot] marked this conversation as resolved.
TypeScript hosts can type the stored receipt with
`import type { WorkforcePersonaExecution } from '@agentworkforce/local-surface'`.

If the callback throws or rejects, no spawn is requested. If preparation or
broker delegation fails, the factory disposes its prepared resources and removes
the scratch directory; discard any retained receipt for that failed launch.

After a successful spawn, the host owns the remaining lifetime. Release the
specific launched worker through the broker and verify that it has stopped before
calling `await execution.handle.dispose()`, then remove `execution.scratchDir`.
Disposal stops mount synchronization and restores generated files. Do not dispose
inside the preparation callback or while delegation is pending. The callback is
optional; omitting it preserves the existing successful-execution lifetime.
1 change: 1 addition & 0 deletions packages/local-surface/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ export {
defineWorkforcePersonaSpawnNode,
workforcePersonaSpawnCapability,
type DefineWorkforcePersonaSpawnNodeOptions,
type WorkforcePersonaExecution,
type WorkforcePersonaSpawnInput,
type WorkforcePersonaSpawnOptions,
type WorkforcePersonaSpawnResult
Expand Down
166 changes: 165 additions & 1 deletion packages/local-surface/src/persona-spawn.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import assert from 'node:assert/strict';
import { mkdtemp, rm, stat, writeFile } from 'node:fs/promises';
import { mkdtemp, readFile, rm, stat, writeFile } from 'node:fs/promises';
import { tmpdir } from 'node:os';
import { dirname, join } from 'node:path';
import test from 'node:test';
Expand Down Expand Up @@ -255,3 +255,167 @@ async function removeScratchDirs(scratchDirs: Iterable<string>): Promise<void> {
[...scratchDirs].map((scratchDir) => rm(scratchDir, { recursive: true, force: true }))
);
}

test('awaits ownership of the actual handle before delegation and leaves success cleanup to the host', async () => {
let prepared: import('./persona-spawn.js').WorkforcePersonaExecution | undefined;
let disposed = 0;
let spawnCalls = 0;
let preparedCalls = 0;
let allowDelegation!: () => void;
let signalPrepared!: () => void;
const preparedEntered = new Promise<void>((resolve) => { signalPrepared = resolve; });
const custodyBanked = new Promise<void>((resolve) => { allowDelegation = resolve; });
const handle = { cwd: '/tmp/persona-runtime', dispose: async () => { disposed += 1; } };
__setPersonaSpawnImplementationsForTest({
resolvePersona: () => resolved,
buildPlan: () => plan,
checkFleetCompatibility: () => undefined,
executePlan: async () => handle
});
const node = defineWorkforcePersonaSpawnNode({
nodeName: 'persona-node',
async onExecutionPrepared(name, execution) {
assert.equal(name, 'owned-reviewer');
assert.equal(execution.handle, handle);
assert.ok((await stat(execution.scratchDir)).isDirectory());
prepared = execution;
preparedCalls += 1;
signalPrepared();
await custodyBanked;
}
});
const ctx = {
node: { name: 'persona-node', capabilities: ['spawn:persona'] },
relay: { sendMessage: async () => undefined },
spawnAgent: async (input) => {
spawnCalls += 1;
assert.equal(disposed, 0);
assert.deepEqual(input.agent.channels, ['owned-a', 'owned-b']);
return { ready: true };
}
} satisfies FleetActionContext;
const input = { name: 'owned-reviewer', persona: 'reviewer', channels: ['owned-a', 'owned-b'] };
let launches: Promise<unknown>[] = [];
try {
launches = [invokeNodeHandler(node, 'spawn:persona', input, ctx)];
await preparedEntered;
launches.push(invokeNodeHandler(node, 'spawn:persona', input, ctx));
assert.equal(spawnCalls, 0, 'delegation must wait for durable ownership');
allowDelegation();
await Promise.all(launches);
assert.equal(preparedCalls, 1, 'coalesced launch must publish only one handle');
assert.equal(spawnCalls, 1);
assert.equal(disposed, 0, 'factory must not dispose a running worker mount');
assert.ok(prepared);
// The host has now completed its worker-release contract. Use the exact
// retained executor handle, never a separately constructed cleanup handle.
await prepared.handle.dispose();
await rm(prepared.scratchDir, { recursive: true, force: true });
assert.equal(disposed, 1);
await assert.rejects(stat(prepared.scratchDir), { code: 'ENOENT' });
} finally {
allowDelegation();
await Promise.allSettled(launches);
__setPersonaSpawnImplementationsForTest();
if (prepared) await rm(prepared.scratchDir, { recursive: true, force: true });
}
});

for (const failure of ['ownership', 'delegation', 'dispose'] as const) {
test(`cleans the prepared resources after ${failure} failure`, async () => {
let prepared: import('./persona-spawn.js').WorkforcePersonaExecution | undefined;
let disposed = 0;
let spawnCalls = 0;
const launchError = new Error('launch rejected');
const disposeError = new Error('dispose rejected');
const handle = {
cwd: '/tmp/persona-runtime',
async dispose() {
disposed += 1;
if (failure === 'dispose') throw disposeError;
}
};
__setPersonaSpawnImplementationsForTest({
resolvePersona: () => resolved,
buildPlan: () => plan,
checkFleetCompatibility: () => undefined,
executePlan: async () => handle
});
const node = defineWorkforcePersonaSpawnNode({
nodeName: 'persona-node',
async onExecutionPrepared(name, execution) {
assert.equal(name, 'owned-reviewer');
assert.equal(execution.handle, handle);
prepared = execution;
if (failure !== 'delegation') throw launchError;
}
});
const ctx = {
node: { name: 'persona-node', capabilities: ['spawn:persona'] },
relay: { sendMessage: async () => undefined },
spawnAgent: async () => {
spawnCalls += 1;
assert.ok(prepared);
assert.equal(disposed, 0);
throw launchError;
}
} satisfies FleetActionContext;
try {
await assert.rejects(
invokeNodeHandler(node, 'spawn:persona', { name: 'owned-reviewer', persona: 'reviewer' }, ctx),
failure === 'dispose' ? disposeError : launchError
);
assert.equal(spawnCalls, failure === 'delegation' ? 1 : 0);
assert.equal(disposed, 1);
assert.ok(prepared);
await assert.rejects(stat(prepared.scratchDir), { code: 'ENOENT' });
} finally {
__setPersonaSpawnImplementationsForTest();
if (prepared) await rm(prepared.scratchDir, { recursive: true, force: true });
}
});
}


test('host disposes the real executor mount retained after a successful spawn', async () => {
const project = await mkdtemp(join(tmpdir(), 'persona-spawn-owned-project-'));
await writeFile(join(project, 'input.txt'), 'owned project contents');
let prepared: import('./persona-spawn.js').WorkforcePersonaExecution | undefined;
__setPersonaSpawnImplementationsForTest({
resolvePersona: () => resolved,
buildPlan: () => plan,
checkFleetCompatibility: () => undefined
});
const node = defineWorkforcePersonaSpawnNode({
nodeName: 'persona-node',
cwd: project,
onExecutionPrepared(_name, execution) { prepared = execution; }
});
const ctx = {
node: { name: 'persona-node', capabilities: ['spawn:persona'] },
relay: { sendMessage: async () => undefined },
spawnAgent: async (input) => {
assert.ok(prepared);
assert.equal(input.agent.cwd, prepared.handle.cwd);
assert.equal(await readFile(join(prepared.handle.cwd, 'input.txt'), 'utf8'), 'owned project contents');
return { ready: true };
}
} satisfies FleetActionContext;
try {
await invokeNodeHandler(node, 'spawn:persona', { name: 'owned-real-mount', persona: 'reviewer' }, ctx);
assert.ok(prepared);
assert.ok((await stat(prepared.handle.cwd)).isDirectory());
await prepared.handle.dispose();
await rm(prepared.scratchDir, { recursive: true, force: true });
await assert.rejects(stat(prepared.handle.cwd), { code: 'ENOENT' });
await assert.rejects(stat(prepared.scratchDir), { code: 'ENOENT' });
assert.equal(await readFile(join(project, 'input.txt'), 'utf8'), 'owned project contents');
} finally {
__setPersonaSpawnImplementationsForTest();
if (prepared) {
await prepared.handle.dispose();
await rm(prepared.scratchDir, { recursive: true, force: true });
}
await rm(project, { recursive: true, force: true });
}
});
36 changes: 33 additions & 3 deletions packages/local-surface/src/persona-spawn.ts
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,28 @@ export interface WorkforcePersonaSpawnResult {
result: unknown;
}

/** Resources prepared for one launch; the handle is the executor's actual handle. */
export interface WorkforcePersonaExecution {
readonly handle: ExecutionHandle;
/** Factory-owned parent of handle.cwd; remove after disposing the handle. */
readonly scratchDir: string;
}

export interface WorkforcePersonaSpawnOptions {
/**
* Awaited once per prepared launch, before delegating to the broker. Use this
* to retain the actual resources and bank ownership before a worker can start.
* Throwing/rejecting prevents delegation and the factory cleans up. Delegation
* failure also triggers factory cleanup; discard the retained receipt then.
* After success the host owns cleanup: first release the launched worker and
* confirm it stopped, then await handle.dispose() and remove scratchDir. The
* callback must not dispose resources while preparation/delegation is pending.
* Without this callback, successful resources retain their existing lifetime.
*/
onExecutionPrepared?: (
name: string,
execution: WorkforcePersonaExecution
) => void | Promise<void>;
/** Default project cwd. Request-local `cwd` wins. */
cwd?: string;
/** Registry source configuration shared with `agentworkforce agent`. */
Expand Down Expand Up @@ -121,7 +142,10 @@ export function workforcePersonaSpawnCapability(
const existing = inFlight.get(key);
if (existing) return existing;

const launch = launchResolvedPersona({ input, cwd, resolved, ctx, active });
const launch = launchResolvedPersona({
input, cwd, resolved, ctx, active,
onExecutionPrepared: options.onExecutionPrepared
});
inFlight.set(key, launch);
try {
return await launch;
Expand Down Expand Up @@ -157,6 +181,7 @@ async function launchResolvedPersona(input: {
resolved: ResolvedPersonaReference;
ctx: FleetActionContext;
active: Map<string, PreparedPersonaExecution>;
onExecutionPrepared?: WorkforcePersonaSpawnOptions['onExecutionPrepared'];
}): Promise<WorkforcePersonaSpawnResult> {
const { resolved, ctx } = input;
const scratchDir = await mkdtemp(join(tmpdir(), 'agentworkforce-relay-persona-'));
Expand Down Expand Up @@ -187,6 +212,8 @@ async function launchResolvedPersona(input: {
mount: { mountDir, includeGit: true, autoSync: true }
});

await input.onExecutionPrepared?.(input.input.name, { handle: execution, scratchDir });

const args = plan.initialPrompt ? [...plan.args, plan.initialPrompt] : [...plan.args];
const result = await ctx.spawnAgent({
agent: {
Expand Down Expand Up @@ -233,8 +260,11 @@ async function launchResolvedPersona(input: {
result
};
} catch (error) {
await execution?.dispose();
await rm(scratchDir, { recursive: true, force: true });
try {
await execution?.dispose();
} finally {
await rm(scratchDir, { recursive: true, force: true });
}
throw error;
}
}
Expand Down
Loading