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
193 changes: 193 additions & 0 deletions apps/server/src/orchestration-v2/Adapters/OpenCode2AdapterV2.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -881,6 +881,199 @@ describe("OpenCode2 adapter", () => {
}).pipe(Effect.scoped),
);

/** A `subagent` call's result while its subagent runs on (`session.tool.success` is at version 2). */
const returnedRunning = (id: string): ProviderReplayEntry => ({
type: "emit_inbound",
frame: {
type: "sdk.event",
event: {
id: `evt_toolsuccess_${id}`,
created: 1,
type: "session.tool.success",
data: {
sessionID: SESSION,
assistantMessageID: "msg_assistant",
id,
content: [{ type: "text", text: "The subagent is working in the background." }],
metadata: { sessionID: CHILD, status: "running" },
executed: true,
},
durable: { aggregateID: SESSION, seq: 1, version: 2 },
},
},
});
/** The model calls the subagent again by its session; the tool returns while it runs on. */
const callAgain = (input: Record<string, unknown>): ReadonlyArray<ProviderReplayEntry> => {
const tool = { sessionID: SESSION, assistantMessageID: "msg_assistant", id: "call-again" };
return [
event("session.tool.input.started", { ...tool, name: "subagent" }),
event("session.tool.called", {
...tool,
name: "subagent",
input: { description: "Again", prompt: "also this", sessionID: CHILD, ...input },
executed: false,
}),
event("session.tool.progress", {
...tool,
metadata: { sessionID: CHILD, status: "running" },
}),
// A session called again gets the thread's rules.
out("session.update", { sessionID: CHILD, permissions: "<any>" }),
reply("session.update", null),
returnedRunning(tool.id),
];
};
const childReport = (inboxID: string, text: string) =>
event("session.inbox.enqueued", {
inboxID,
sessionID: SESSION,
item: {
type: "synthetic",
payload: {
text: `<subagent sessionID="${CHILD}" state="completed" description="Sleep">\n${text}\n</subagent>`,
description: "Sleep",
metadata: { source: "subagent", childID: CHILD, agent: "General", state: "completed" },
},
delivery: "steer",
},
});
/** The subagent rows by title, the thread's background roster and the child's turns, as emitted. */
const watchCalls = (
runtime: ProviderAdapterV2SessionRuntime,
thread: OrchestrationV2ProviderThread,
) =>
Effect.gen(function* () {
const calls = new Map<string, { status: string; result: string | null }>();
const childTurns = new Map<string, string>();
const callsAtChildTurnStart = new Map<
string,
Record<string, { status: string; result: string | null }>
>();
const state: { roster: ReadonlyArray<unknown> | undefined } = { roster: undefined };
let wake = yield* Deferred.make<void>();
yield* runtime.events.pipe(
Stream.tap((event) =>
Effect.gen(function* () {
if (event.type === "subagent.updated" && event.subagent.title !== null) {
calls.set(event.subagent.title, {
status: event.subagent.status,
result: event.subagent.result,
});
}
if (event.type === "provider_thread.updated" && event.providerThread.id === thread.id) {
state.roster = event.providerThread.pendingBackgroundTasks;
}
const id =
event.type === "provider_turn.updated"
? event.providerTurn.nativeTurnRef?.nativeId
: undefined;
if (
event.type === "provider_turn.updated" &&
id?.startsWith(`${CHILD}:turn:`) === true
) {
if (!childTurns.has(id)) {
callsAtChildTurnStart.set(id, Object.fromEntries(calls));
}
childTurns.set(id, event.providerTurn.status);
}
yield* Deferred.succeed(wake, undefined);
}),
),
Stream.runDrain,
Effect.forkScoped,
);
const until = (check: () => boolean): Effect.Effect<void> =>
Effect.suspend(() => {
if (check()) return Effect.void;
return Effect.gen(function* () {
wake = yield* Deferred.make<void>();
if (check()) return;
yield* Deferred.await(wake);
yield* until(check);
});
});
return { calls, callsAtChildTurnStart, childTurns, state, until };
});
const launchedAndRunning = [
...backgroundLaunch(CHILD),
returnedRunning("call-background"),
event("session.execution.started", { sessionID: CHILD }),
];

it.effect("settles every call that joined a running subagent with its one report", () =>
Effect.gen(function* () {
const { runtime, thread } = yield* resumed([
...launchedAndRunning,
// OpenCode steers the running subagent and joins its run: one report answers both calls.
...callAgain({ background: true }),
event("session.execution.succeeded", { sessionID: SESSION }),
event("session.execution.succeeded", { sessionID: CHILD }),
childReport("msg_report", "CHILD_OK"),
]);
const watch = yield* watchCalls(runtime, thread);
yield* runtime.startTurn(withLineage(thread));
yield* watch.until(
() =>
watch.calls.get("Again")?.status === "completed" &&
watch.calls.get("Sleep")?.status === "completed" &&
watch.state.roster?.length === 0,
);
assert.deepEqual(Object.fromEntries(watch.calls), {
Sleep: { status: "completed", result: "CHILD_OK" },
Again: { status: "completed", result: "CHILD_OK" },
});
}).pipe(Effect.scoped),
);

it.effect("keeps a foreground call that joined a background run for that run's report", () =>
Effect.gen(function* () {
const { runtime, thread } = yield* resumed([
...launchedAndRunning,
// Without `background`, the call joins the run in the background all the same.
...callAgain({}),
event("session.execution.succeeded", { sessionID: SESSION }),
event("session.execution.succeeded", { sessionID: CHILD }),
childReport("msg_report", "CHILD_OK"),
]);
const watch = yield* watchCalls(runtime, thread);
yield* runtime.startTurn(withLineage(thread));
yield* watch.until(
() =>
watch.calls.get("Again")?.status === "completed" &&
watch.calls.get("Sleep")?.status === "completed" &&
watch.state.roster?.length === 0,
);
// Not settled empty when the turn that made it ended.
assert.deepEqual(watch.calls.get("Again"), { status: "completed", result: "CHILD_OK" });
}).pipe(Effect.scoped),
);

it.effect("keeps a call to a finished subagent for its own run's report", () =>
Effect.gen(function* () {
const { runtime, thread } = yield* resumed([
...launchedAndRunning,
event("session.execution.succeeded", { sessionID: CHILD }),
// Called again after it ended but before its report: the call starts a new run.
...callAgain({ background: true }),
childReport("msg_report_1", "FIRST_OK"),
event("session.execution.started", { sessionID: CHILD }),
event("session.execution.succeeded", { sessionID: CHILD }),
event("session.execution.succeeded", { sessionID: SESSION }),
childReport("msg_report_2", "SECOND_OK"),
]);
const watch = yield* watchCalls(runtime, thread);
yield* runtime.startTurn(withLineage(thread));
// The new run's turn starts after the first report was handled.
yield* watch.until(() => watch.childTurns.has(`${CHILD}:turn:2`));
assert.deepEqual(watch.callsAtChildTurnStart.get(`${CHILD}:turn:2`), {
Sleep: { status: "completed", result: "FIRST_OK" },
Again: { status: "running", result: null },
});
yield* watch.until(() => watch.calls.get("Again")?.status === "completed");
assert.deepEqual(watch.calls.get("Again"), { status: "completed", result: "SECOND_OK" });
}).pipe(Effect.scoped),
);

it.effect("leaves no background work pending once a reconnect finds its subagent gone", () =>
Effect.gen(function* () {
// The subagent ends while the stream is down: its end, its report and
Expand Down
43 changes: 29 additions & 14 deletions packages/provider-opencode/src/server/v2/adapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -291,6 +291,8 @@ interface SubagentCall {
model: string | null;
/** The tool returned while the subagent runs on; the report OpenCode gives its parent settles it. */
background: boolean;
/** The running call whose subagent run this call joined, so that run's one report settles both. */
joined: SubagentCall | undefined;
child: ThreadState | undefined;
status: OrchestrationV2Subagent["status"];
result: string | null;
Expand Down Expand Up @@ -1489,6 +1491,14 @@ export const make = Effect.fn("OpenCode2Adapter.make")(function* (instanceId: Pr
// ran (each one's native id is its own), its grants and rules, and the
// background work an earlier call left running, which Stop must reach.
const previous = threads.get(childId);
// Called again while it runs, the subagent takes the prompt as a steer
// and OpenCode reports that run once, for every call that joined it.
if (previous?.active !== undefined) {
call.joined = [...call.state.calls.values()].findLast(
(candidate) =>
candidate !== call && candidate.child === previous && candidate.joined === undefined,
);
}
const subagent = {
call,
appThread,
Expand Down Expand Up @@ -2175,6 +2185,7 @@ export const make = Effect.fn("OpenCode2Adapter.make")(function* (instanceId: Pr
agent: undefined,
model: null,
background: false,
joined: undefined,
child: undefined,
status: "running",
result: null,
Expand Down Expand Up @@ -2240,8 +2251,13 @@ export const make = Effect.fn("OpenCode2Adapter.make")(function* (instanceId: Pr
if (call !== undefined) {
const childId = stringField(event.data.metadata, "sessionID");
if (childId !== undefined) yield* attachChild(call, childId);
// A background call returns at launch; its report settles it.
if (event.data.metadata?.["status"] === "running") return;
// A background call returns at launch, and so does one that joined a
// background run; the report settles it.
if (event.data.metadata?.["status"] === "running") {
if (call.background) return;
call.background = true;
return yield* emitSubagent(call);
}
return yield* settleCall(call, "completed", subagentOutput(textOf(event.data.content)));
}
const output = textOf(event.data.content);
Expand Down Expand Up @@ -2379,6 +2395,9 @@ export const make = Effect.fn("OpenCode2Adapter.make")(function* (instanceId: Pr
const call = [...state.calls.values()].find(
(candidate) => candidate.child?.sessionId === childId,
);
const calls = [...state.calls.values()].filter(
(candidate) => candidate === call || (call !== undefined && candidate.joined === call),
);
const outcome = reportOutcome(stringField(payload.metadata, "state"));
state.reports.set(inboxId, {
inboxId,
Expand All @@ -2391,18 +2410,14 @@ export const make = Effect.fn("OpenCode2Adapter.make")(function* (instanceId: Pr
outcome,
},
});
if (call === undefined) return;
yield* settleCall(
call,
state.stoppedChildren.has(childId)
? "interrupted"
: outcome === "failed"
? "failed"
: outcome === "cancelled"
? "cancelled"
: "completed",
subagentOutput(payload.text),
);
const status = state.stoppedChildren.has(childId)
? "interrupted"
: outcome === "failed"
? "failed"
: outcome === "cancelled"
? "cancelled"
: "completed";
for (const each of calls) yield* settleCall(each, status, subagentOutput(payload.text));
});

/**
Expand Down
Loading