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
Original file line number Diff line number Diff line change
Expand Up @@ -310,6 +310,7 @@ it("taskStatus returns task.providerInstanceId rather than the driver kind", asy
],
subagents: [],
providerThreads: [],
turnItems: [],
updatedAt: now,
} as unknown as OrchestrationV2ThreadProjection;

Expand Down
285 changes: 160 additions & 125 deletions apps/server/src/mcp/OrchestratorMcpService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ import {
ProviderInstanceId,
RunId,
ThreadId,
TurnItemId,
type OrchestrationV2ThreadProjection,
type ServerProvider,
} from "@t3tools/contracts";
Expand Down Expand Up @@ -41,135 +42,165 @@ import { idleThreadProjection, liveThreadShell } from "./McpToolAccess.testkit.t
import * as OrchestratorMcpService from "./OrchestratorMcpService.ts";

describe("OrchestratorMcpService", () => {
it.effect("retries terminal acknowledgement with a fresh command id", () =>
Effect.gen(function* () {
const parentThreadId = ThreadId.make("thread:mcp-ack-parent");
const childThreadId = ThreadId.make("thread:mcp-ack-child");
const childRunId = RunId.make("run:mcp-ack-child");
const taskId = NodeId.make("node:mcp-ack-task");
const acknowledgementCommandIds = yield* Ref.make<ReadonlyArray<string>>([]);
const acknowledgementAttempts = yield* Ref.make(0);
const parentProjection = {
thread: { id: parentThreadId },
runs: [],
contextTransfers: [],
subagents: [
{
id: taskId,
it.effect.each(["nested task", "background command"] as const)(
"waits for a %s before acknowledging the result and retries failed acknowledgement",
(backgroundWork) =>
Effect.gen(function* () {
const parentThreadId = ThreadId.make("thread:mcp-ack-parent");
const childThreadId = ThreadId.make("thread:mcp-ack-child");
const childRunId = RunId.make("run:mcp-ack-child");
const taskId = NodeId.make("node:mcp-ack-task");
const acknowledgementCommandIds = yield* Ref.make<ReadonlyArray<string>>([]);
const acknowledgementAttempts = yield* Ref.make(0);
const parentProjection = {
thread: { id: parentThreadId },
runs: [],
contextTransfers: [],
subagents: [
{
id: taskId,
threadId: parentThreadId,
origin: "app_owned",
childThreadId,
driver: "codex",
model: "gpt-5.6-terra",
result: "terminal result",
completionDelivery: { state: "pending" },
},
],
} as unknown as OrchestrationV2ThreadProjection;
const childProjection = {
thread: { id: childThreadId },
runs: [
{
id: childRunId,
ordinal: 1,
status: "completed",
startedAt: DateTime.makeUnsafe("2026-10-03T10:00:00Z"),
completedAt: DateTime.makeUnsafe("2026-10-03T10:10:00Z"),
},
{
id: RunId.make("run:mcp-ack-continuation"),
ordinal: 2,
status: "failed",
startedAt: DateTime.makeUnsafe("2026-10-03T10:02:00Z"),
completedAt: DateTime.makeUnsafe("2026-10-03T10:05:00Z"),
},
],
contextTransfers: [],
messages: [],
subagents: [],
providerThreads: [],
turnItems: [],
} as unknown as OrchestrationV2ThreadProjection;
let hasNestedWork = true;
const layerDependencies = Layer.mergeAll(
NodeServices.layer,
Layer.mock(ThreadManagementService.ThreadManagementService)({
getThreadRecords: (threadId, records) =>
Effect.succeed(
threadId === parentThreadId
? hasNestedWork
? {
...parentProjection,
subagents: parentProjection.subagents.map((task) => ({
...task,
result: null,
status: "running" as const,
})),
}
: parentProjection
: hasNestedWork
? {
...childProjection,
subagents:
backgroundWork === "nested task"
? [{ ...parentProjection.subagents[0]!, status: "running" as const }]
: [],
turnItems:
backgroundWork === "background command" &&
records.some((record) => record === "turnItems")
? [
{
id: TurnItemId.make("background-command"),
threadId: childThreadId,
runId: childRunId,
nodeId: null,
providerThreadId: null,
providerTurnId: null,
nativeItemRef: null,
parentItemId: null,
ordinal: 1,
type: "command_execution" as const,
status: "running" as const,
title: null,
startedAt: null,
completedAt: null,
updatedAt: DateTime.makeUnsafe("2026-10-03T10:10:00Z"),
input: "test --watch",
},
]
: [],
}
: childProjection,
),
dispatch: (command) =>
Ref.update(acknowledgementCommandIds, (commandIds) => [
...commandIds,
String(command.commandId),
]).pipe(
Effect.andThen(Ref.updateAndGet(acknowledgementAttempts, (count) => count + 1)),
Effect.flatMap((attempt) =>
attempt === 1
? Effect.fail(new Error("simulated acknowledgement failure") as never)
: Effect.succeed({} as never),
),
),
}),
Layer.mock(ProviderRegistry.ProviderRegistry)({ getProviders: Effect.succeed([]) }),
Layer.mock(ProviderAdapterRegistry.ProviderAdapterRegistryV2)({
list: () => Effect.succeed([]),
}),
Layer.mock(ProjectService.ProjectService)({}),
Layer.mock(SecretRequests.SecretRequests)({}),
Layer.mock(ScheduledTaskService.ScheduledTaskService)({}),
);
const scope: McpInvocationScope = {
environmentId: EnvironmentId.make("environment:mcp-ack"),
requestNamespace: "provider-session:mcp-ack",
thread: {
threadId: parentThreadId,
origin: "app_owned",
childThreadId,
driver: "codex",
model: "gpt-5.6-terra",
result: "terminal result",
completionDelivery: { state: "pending" },
providerSessionId: "provider-session:mcp-ack",
providerInstanceId: ProviderInstanceId.make("codex"),
},
],
} as unknown as OrchestrationV2ThreadProjection;
const childProjection = {
thread: { id: childThreadId },
runs: [
{
id: childRunId,
ordinal: 1,
status: "completed",
startedAt: DateTime.makeUnsafe("2026-10-03T10:00:00Z"),
completedAt: DateTime.makeUnsafe("2026-10-03T10:10:00Z"),
},
{
id: RunId.make("run:mcp-ack-continuation"),
ordinal: 2,
status: "failed",
startedAt: DateTime.makeUnsafe("2026-10-03T10:02:00Z"),
completedAt: DateTime.makeUnsafe("2026-10-03T10:05:00Z"),
},
],
contextTransfers: [],
messages: [],
subagents: [],
providerThreads: [],
} as unknown as OrchestrationV2ThreadProjection;
let hasNestedWork = true;
const layerDependencies = Layer.mergeAll(
NodeServices.layer,
Layer.mock(ThreadManagementService.ThreadManagementService)({
getThreadRecords: (threadId) =>
Effect.succeed(
threadId === parentThreadId
? hasNestedWork
? {
...parentProjection,
subagents: parentProjection.subagents.map((task) => ({
...task,
result: null,
status: "running" as const,
})),
}
: parentProjection
: hasNestedWork
? {
...childProjection,
subagents: [
{ ...parentProjection.subagents[0]!, status: "running" as const },
],
}
: childProjection,
),
dispatch: (command) =>
Ref.update(acknowledgementCommandIds, (commandIds) => [
...commandIds,
String(command.commandId),
]).pipe(
Effect.andThen(Ref.updateAndGet(acknowledgementAttempts, (count) => count + 1)),
Effect.flatMap((attempt) =>
attempt === 1
? Effect.fail(new Error("simulated acknowledgement failure") as never)
: Effect.succeed({} as never),
),
),
}),
Layer.mock(ProviderRegistry.ProviderRegistry)({ getProviders: Effect.succeed([]) }),
Layer.mock(ProviderAdapterRegistry.ProviderAdapterRegistryV2)({
list: () => Effect.succeed([]),
}),
Layer.mock(ProjectService.ProjectService)({}),
Layer.mock(SecretRequests.SecretRequests)({}),
Layer.mock(ScheduledTaskService.ScheduledTaskService)({}),
);
const scope: McpInvocationScope = {
environmentId: EnvironmentId.make("environment:mcp-ack"),
requestNamespace: "provider-session:mcp-ack",
thread: {
threadId: parentThreadId,
providerSessionId: "provider-session:mcp-ack",
providerInstanceId: ProviderInstanceId.make("codex"),
},
client: undefined,
capabilities: new Set(["orchestration"]),
issuedAt: 1,
};
client: undefined,
capabilities: new Set(["orchestration"]),
issuedAt: 1,
};

yield* Effect.gen(function* () {
const service = yield* OrchestratorMcpService.OrchestratorMcpService;
const pending = yield* service.taskStatus(scope, taskId);
assert.equal(pending.status, "running");
assert.equal(pending.workState, "waiting_for_children");
assert.isNull(pending.summary);
assert.equal(yield* Ref.get(acknowledgementAttempts), 0);
hasNestedWork = false;
const error = yield* service.taskStatus(scope, taskId).pipe(Effect.flip);
assert.equal(error.code, "orchestration_error");
yield* Effect.gen(function* () {
const service = yield* OrchestratorMcpService.OrchestratorMcpService;
const pending = yield* service.taskStatus(scope, taskId);
assert.equal(pending.status, "running");
assert.equal(pending.workState, "waiting_for_children");
assert.isNull(pending.summary);
assert.equal(yield* Ref.get(acknowledgementAttempts), 0);
hasNestedWork = false;
const error = yield* service.taskStatus(scope, taskId).pipe(Effect.flip);
assert.equal(error.code, "orchestration_error");

const result = yield* service.taskStatus(scope, taskId);
assert.equal(result.status, "completed");
assert.equal(result.summary, "terminal result");
assert.equal(result.latestTerminalRunId, childRunId);
assert.equal(result.latestTerminalStatus, "completed");
const commandIds = yield* Ref.get(acknowledgementCommandIds);
assert.equal(commandIds.length, 2);
assert.notEqual(commandIds[0], commandIds[1]);
}).pipe(Effect.provide(OrchestratorMcpService.layer.pipe(Layer.provide(layerDependencies))));
}),
const result = yield* service.taskStatus(scope, taskId);
assert.equal(result.status, "completed");
assert.equal(result.summary, "terminal result");
assert.equal(result.latestTerminalRunId, childRunId);
assert.equal(result.latestTerminalStatus, "completed");
const commandIds = yield* Ref.get(acknowledgementCommandIds);
assert.equal(commandIds.length, 2);
assert.notEqual(commandIds[0], commandIds[1]);
}).pipe(
Effect.provide(OrchestratorMcpService.layer.pipe(Layer.provide(layerDependencies))),
);
}),
);

it.effect("reports a restart-cut child as working until its continuation settles", () =>
Expand Down Expand Up @@ -289,6 +320,7 @@ describe("OrchestratorMcpService", () => {
messages: [],
subagents: [],
providerThreads: [],
turnItems: [],
} as unknown as OrchestrationV2ThreadProjection;
const layerDependencies = Layer.mergeAll(
NodeServices.layer,
Expand Down Expand Up @@ -365,6 +397,7 @@ describe("OrchestratorMcpService", () => {
messages: [],
subagents: [],
providerThreads: [],
turnItems: [],
} as unknown as OrchestrationV2ThreadProjection;
const layerDependencies = Layer.mergeAll(
NodeServices.layer,
Expand Down Expand Up @@ -444,6 +477,7 @@ describe("OrchestratorMcpService", () => {
messages: [],
subagents: [],
providerThreads: [],
turnItems: [],
} as unknown as OrchestrationV2ThreadProjection;
const layerDependencies = Layer.mergeAll(
NodeServices.layer,
Expand Down Expand Up @@ -530,6 +564,7 @@ describe("OrchestratorMcpService", () => {
messages: [],
subagents: [],
providerThreads: [],
turnItems: [],
} as unknown as OrchestrationV2ThreadProjection;
const layerDependencies = Layer.mergeAll(
NodeServices.layer,
Expand Down
8 changes: 6 additions & 2 deletions apps/server/src/mcp/OrchestratorMcpService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1210,8 +1210,12 @@ const make = Effect.gen(function* () {
const childControls = yield* threadManagement
.getThreadRecords(
task.childThreadId,
["runs", "messages", "contextTransfers", "subagents", "providerThreads"],
{ messageRoles: ["user"] },
["runs", "messages", "contextTransfers", "subagents", "providerThreads", "turnItems"],
{
messageRoles: ["user"],
turnItemTypes: ["command_execution", "dynamic_tool", "subagent"],
turnItemStatuses: ["pending", "running", "waiting"],
},
)
.pipe(Effect.mapError(threadManagementFailure));
const childRun = delegatedTaskRun(childControls, task);
Expand Down
Loading
Loading