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
25 changes: 25 additions & 0 deletions apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1037,6 +1037,19 @@ describe("AcpAdapterV2", () => {
// Production thread 54aeb6d7 split after "(command". Metadata shapes
// below were captured from live Devin sessions showy-mile/fragrant-chamomile.
const updates = [
{
sessionUpdate: "tool_call_update",
toolCallId: "child-a",
status: "in_progress",
_meta: {
"cognition.ai/subagent_started": {
agentId: "child-a",
title: "Map orchestration",
task: "Run pwd, then reply ONE.",
model: " \t ",
},
},
},
{
sessionUpdate: "tool_call_update",
toolCallId: "child-a",
Expand Down Expand Up @@ -1181,6 +1194,18 @@ describe("AcpAdapterV2", () => {
event.type === "subagent.updated" ? [event.subagent] : [],
);
const task = tasks.at(-1);
assert.isNull(tasks[0]?.model);
const childThread = events.find(
(event) =>
event.type === "app_thread.created" && event.appThread.id === task?.childThreadId,
);
assert.equal(
childThread?.type === "app_thread.created"
? childThread.appThread.modelSelection.model
: undefined,
modelSelection.model,
);
assert.equal(task?.model, "SWE-1.7 Medium");
assert.equal(task?.status, "completed");
assert.equal(task?.result, "Final report: ONE");
const childMessages = new Map(
Expand Down
4 changes: 2 additions & 2 deletions apps/server/src/orchestration-v2/Adapters/AcpAdapterV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2670,10 +2670,10 @@ export function makeAcpAdapterV2(
nativeTaskRef: nativeItemRef,
prompt: update.prompt,
title: update.title,
model: update.model,
result: null,
startedAt: now,
}),
model: update.model?.trim() || existing?.task.model || null,
Comment thread
coderabbitai[bot] marked this conversation as resolved.
status: taskStatus,
result: update.result ?? existing?.task.result ?? null,
completedAt: acpSubagentStatusIsTerminal(taskStatus) ? now : null,
Expand Down Expand Up @@ -2709,7 +2709,7 @@ export function makeAcpAdapterV2(
providerInstanceId: context.input.modelSelection.instanceId,
modelSelection: {
...context.input.modelSelection,
model: update.model ?? context.input.modelSelection.model,
model: task.model ?? context.input.modelSelection.model,
},
title: subagentThreadTitle({
parentTitle: context.input.appThread.title,
Expand Down
76 changes: 73 additions & 3 deletions apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@ import { assert, describe, it } from "@effect/vitest";
import { HostProcessEnvironment, HostProcessPlatform } from "@t3tools/shared/hostProcess";
import { SpawnExecutableResolution } from "@t3tools/shared/shell";
import * as CodexClient from "effect-codex-app-server/client";
import * as CodexError from "effect-codex-app-server/errors";
import * as CodexReplay from "effect-codex-app-server/replay";
import * as DateTime from "effect/DateTime";
import * as Deferred from "effect/Deferred";
Expand Down Expand Up @@ -1619,7 +1620,7 @@ describe("CodexAdapterV2 post-settle continuation", () => {
transcript: CodexReplay.CodexAppServerReplayTranscript,
onEvent: (event: ProviderAdapterV2Event) => Effect.Effect<unknown> = () => Effect.void,
onRequest: (method: string, params: unknown) => Effect.Effect<void> = () => Effect.void,
readChildMetadata?: (threadId: string) => Effect.Effect<unknown>,
readChildMetadata?: Parameters<typeof withCodexReplayChildMetadata>[2],
) =>
Effect.gen(function* () {
const fileSystem = yield* FileSystem.FileSystem;
Expand Down Expand Up @@ -6402,6 +6403,8 @@ describe("CodexAdapterV2 post-settle continuation", () => {
});

it.effect.each([
{ name: "current Codex Sol", model: "gpt-6-sol" },
{ name: "current Codex wrong child", model: null },
{ name: "Sol", model: "gpt-5.6-sol" },
{ name: "Fable", model: "gpt-5.6-fable" },
{ name: "Astra", model: "gpt-6-astra" },
Expand All @@ -6413,6 +6416,7 @@ describe("CodexAdapterV2 post-settle continuation", () => {
Effect.gen(function* () {
const metadataRead = yield* Deferred.make<void>();
const modelReported = yield* Deferred.make<void>();
let metadataRequests = 0;
const harness = yield* makeCodexReplayHarness(
resumeSubagentTranscript,
(event) =>
Expand All @@ -6421,14 +6425,22 @@ describe("CodexAdapterV2 post-settle continuation", () => {
: Effect.void,
undefined,
(threadId) => {
metadataRequests++;
assert.equal(threadId, RESUME_CHILD_THREAD);
return Deferred.succeed(metadataRead, undefined).pipe(
Effect.as(
name === "invalid"
? {}
: {
thread: { id: name === "wrong child" ? "other-child" : threadId },
model: name === "wrong child" ? "gpt-5.6-sol" : model,
thread: {
id: name.includes("wrong child") ? "other-child" : threadId,
...(name.startsWith("current Codex") ? { model: "gpt-6-sol" } : {}),
},
model: name.startsWith("current Codex")
? null
: name === "wrong child"
? "gpt-5.6-sol"
: model,
},
),
);
Expand All @@ -6448,10 +6460,68 @@ describe("CodexAdapterV2 post-settle continuation", () => {
yield* TestClock.adjust("100 millis");
yield* harness.firstTerminal;
assert.equal(harness.subagentUpdates().at(-1)?.subagent.model, model);
assert.equal(metadataRequests, name === "current Codex Sol" ? 1 : 2);
}).pipe(Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))),
),
);

it.effect.each(["failed", "malformed", "wrong child", "blank model"] as const)(
"resumes child metadata after a %s read",
(readResult) =>
Effect.scoped(
Effect.gen(function* () {
const modelReported = yield* Deferred.make<void>();
const model = "gpt-6-sol";
const metadataRequests: Array<string> = [];
const harness = yield* makeCodexReplayHarness(
resumeSubagentTranscript,
(event) =>
event.type === "subagent.updated" && event.subagent.model === model
? Deferred.succeed(modelReported, undefined)
: Effect.void,
undefined,
(threadId, method) => {
assert.equal(threadId, RESUME_CHILD_THREAD);
metadataRequests.push(method);
if (method === "thread/resume") {
return Effect.succeed({ thread: { id: threadId }, model });
}
switch (readResult) {
case "failed":
return Effect.fail(
new CodexError.CodexAppServerRequestError({
code: -32000,
errorMessage: "Child metadata unavailable",
method,
}),
);
case "malformed":
return Effect.succeed({});
case "wrong child":
return Effect.succeed({ thread: { id: "other-child", model: "gpt-6-astra" } });
case "blank model":
return Effect.succeed({ thread: { id: threadId, model: " \t " } });
}
},
);
yield* harness.runtime.startTurn(
makeCodexTestTurnInput({
threadId: harness.threadId,
providerThread: harness.providerThread,
now: yield* DateTime.now,
attemptId: RunAttemptId.make("attempt-child-model-fallback"),
text: RESUME_PROMPT,
}),
);
yield* Deferred.await(modelReported);
yield* TestClock.adjust("100 millis");
yield* harness.firstTerminal;
assert.equal(harness.subagentUpdates().at(-1)?.subagent.model, model);
assert.deepEqual(metadataRequests, ["thread/read", "thread/resume"]);
}).pipe(Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))),
),
);

it.effect.each(["thread/settings/updated", "model/rerouted"] as const)(
"keeps %s child metadata when an older lookup finishes later",
(method) =>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ import { DEFAULT_SIGNAL_EXPORT } from "@t3tools/shared/observability";
import * as NodeServices from "@effect/platform-node/NodeServices";
import { type ProviderReplayTranscript } from "@t3tools/contracts";
import * as CodexClient from "effect-codex-app-server/client";
import type * as CodexError from "effect-codex-app-server/errors";
import * as CodexReplay from "effect-codex-app-server/replay";
import * as Effect from "effect/Effect";
import * as FileSystem from "effect/FileSystem";
Expand Down Expand Up @@ -45,7 +46,10 @@ export type CodexOrchestratorReplayHarnessError = typeof CodexOrchestratorReplay
export function withCodexReplayChildMetadata(
client: CodexClient.CodexAppServerClient["Service"],
transcript: CodexReplay.CodexAppServerReplayTranscript,
readMetadata: (threadId: string) => Effect.Effect<unknown> = (threadId) =>
readMetadata: (
threadId: string,
method: "thread/read" | "thread/resume",
) => Effect.Effect<unknown, CodexError.CodexAppServerError> = (threadId) =>
Effect.succeed({ thread: { id: threadId }, model: null }),
): CodexClient.CodexAppServerClient["Service"] {
const childThreadIds = new Set(
Expand All @@ -67,12 +71,12 @@ export function withCodexReplayChildMetadata(
raw: {
...client.raw,
request: (method, params) =>
method === "thread/resume" &&
(method === "thread/read" || method === "thread/resume") &&
Predicate.isObject(params) &&
params.excludeTurns === true &&
(method === "thread/read" ? params.includeTurns === false : params.excludeTurns === true) &&
typeof params.threadId === "string" &&
childThreadIds.has(params.threadId)
? readMetadata(params.threadId)
? readMetadata(params.threadId, method)
: client.raw.request(method, params),
},
};
Expand Down
29 changes: 27 additions & 2 deletions apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1239,6 +1239,15 @@ const decodeCodexChildModel = Schema.decodeUnknownEffect(
}),
);

const decodeCodexChildThread = Schema.decodeUnknownEffect(
Schema.Struct({
thread: Schema.Struct({
id: Schema.String,
model: Schema.optional(Schema.NullOr(Schema.String)),
}),
}),
);

export const makeCodexAppServerSpawnCommand = Effect.fn(
"CodexAdapterV2.makeCodexAppServerSpawnCommand",
)(function* (input: {
Expand Down Expand Up @@ -2661,9 +2670,25 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi
}
if (task.model === null) {
yield* client.raw
.request("thread/resume", { threadId: input.nativeThreadId, excludeTurns: true })
.request("thread/read", { threadId: input.nativeThreadId, includeTurns: false })
.pipe(
Effect.flatMap(decodeCodexChildModel),
Effect.flatMap(decodeCodexChildThread),
Comment thread
coderabbitai[bot] marked this conversation as resolved.
Effect.map((response) =>
response.thread.id === input.nativeThreadId && response.thread.model?.trim()
? { thread: response.thread, model: response.thread.model }
: null,
),
Effect.catch(() => Effect.succeed(null)),
Effect.flatMap((response) =>
response === null
? client.raw
.request("thread/resume", {
threadId: input.nativeThreadId,
excludeTurns: true,
})
.pipe(Effect.flatMap(decodeCodexChildModel))
: Effect.succeed(response),
),
Effect.timeout("5 seconds"),
Effect.flatMap((response) =>
response.thread.id === input.nativeThreadId &&
Expand Down
37 changes: 29 additions & 8 deletions apps/server/src/orchestration-v2/Adapters/CursorAdapterV2.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -40,12 +40,14 @@ const decodeCursorSettings = Schema.decodeEffect(CursorSettings);

describe("CursorAdapterV2", () => {
it.effect.each([
{ status: "finished", model: undefined },
{ status: "cancelled", model: "claude-opus-4-6" },
{ status: "error", model: "custom-fable" },
{ status: "finished", model: undefined, lateModel: undefined },
{ status: "cancelled", model: "claude-opus-4-6", lateModel: undefined },
{ status: "error", model: "custom-fable", lateModel: undefined },
{ status: "finished", model: undefined, lateModel: "gpt-6-sol" },
{ status: "finished", model: "gpt-6-sol", lateModel: null },
] as const)(
"settles missing task completions when the Cursor run is $status",
({ status, model }) =>
"projects Cursor tasks: $status, late model $lateModel",
({ status, model, lateModel }) =>
Effect.gen(function* () {
const fileSystem = yield* FileSystem.FileSystem;
const path = yield* Path.Path;
Expand Down Expand Up @@ -96,12 +98,24 @@ describe("CursorAdapterV2", () => {
mode: "unspecified" as const,
},
};
for (const type of ["partial-tool-call", "tool-call-started"] as const) {
const updates =
lateModel !== undefined
? (["tool-call-started", "tool-call-completed"] as const)
: (["partial-tool-call", "tool-call-started"] as const);
for (const type of updates) {
yield* input.onDelta!({
type,
modelCallId: "model-call",
callId: "task-call",
toolCall: taskToolCall,
toolCall: {
...taskToolCall,
args: {
...taskToolCall.args,
...(type === "tool-call-completed"
? { model: lateModel ?? undefined }
: {}),
},
},
}).pipe(Effect.orDie);
}
return {
Expand Down Expand Up @@ -180,9 +194,16 @@ describe("CursorAdapterV2", () => {
const rows = events.filter((event) => event.type === "subagent.updated");
assert.equal(rows[0]?.subagent.status, "running");
assert.equal(rows[0]?.subagent.model, model ?? null);
assert.equal(rows.at(-1)?.subagent.model, lateModel ?? model ?? null);
assert.equal(
rows.at(-1)?.subagent.status,
status === "finished" ? "idle" : status === "cancelled" ? "cancelled" : "failed",
lateModel !== undefined
? "completed"
: status === "finished"
? "idle"
: status === "cancelled"
? "cancelled"
: "failed",
);
assert.isNotNull(rows.at(-1)?.subagent.completedAt);
}).pipe(Effect.scoped, Effect.provide(Layer.merge(NodeServices.layer, IdAllocator.layer))),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1550,10 +1550,10 @@ export function makeCursorAdapterV2(
},
prompt: args.prompt,
title: args.description,
model: args.model?.trim() || null,
result: null,
startedAt: now,
}),
model: args.model?.trim() || existing?.task.model || null,
Comment thread
Bil0000 marked this conversation as resolved.
nativeTaskRef: {
driver: CursorAgentSdk.CURSOR_PROVIDER,
nativeId: input.callId,
Expand Down
Loading