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
21 changes: 12 additions & 9 deletions apps/mobile/src/features/threads/thread-work-log.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,7 @@ import { scopeThreadRef } from "@t3tools/client-runtime/environment";
import { environmentThreadDetails, threadEnvironment } from "../../state/threads";
import { useAtomCommand } from "../../state/use-atom-command";
import { toolActivityFaviconUrl } from "@t3tools/shared/favicon";
import { runInterruptSenderThreadId } from "@t3tools/shared/orchestrationV2Timeline";

import { AppText as Text } from "../../components/AppText";
import { T3Wordmark } from "../../components/T3Wordmark";
Expand Down Expand Up @@ -933,7 +934,9 @@ const ThreadWorkLogRow = memo(function ThreadWorkLogRow(
row.projectedItem.item.type === "notification"
? notificationChildThreadId(row.projectedItem.item.source)
: undefined;
const canExpand = row.canExpand && notifiedSubagentThreadId === undefined;
const interruptSenderThreadId = runInterruptSenderThreadId(row.projectedItem.item);
const linkedThreadId = notifiedSubagentThreadId ?? interruptSenderThreadId;
const canExpand = row.canExpand && linkedThreadId === undefined;
const reasoning = row.projectedItem.item.type === "reasoning" ? row.projectedItem.item : null;
const fetchedItem = fetchedDetail.data?.item ?? null;
// Reads keep their path list; the fetched file contents show as output.
Expand Down Expand Up @@ -1000,26 +1003,26 @@ const ThreadWorkLogRow = memo(function ThreadWorkLogRow(
{...(isFreshRow(row.createdAt) ? { entering: FadeIn.duration(200) } : {})}
>
<WorkLogPressable
accessibilityRole={
notifiedSubagentThreadId !== undefined ? "link" : canExpand ? "button" : undefined
}
accessibilityRole={linkedThreadId !== undefined ? "link" : canExpand ? "button" : undefined}
accessibilityLabel={failed ? `${accessiblePreview}, tool call failed` : accessiblePreview}
accessibilityHint={
notifiedSubagentThreadId !== undefined
? "Opens this agent's thread. Long press to copy."
: canExpand
? `Double tap to ${expanded ? "hide" : "show"} full details. Long press to copy.`
: "Long press to copy."
: interruptSenderThreadId !== undefined
? "Opens the thread that stopped this run. Long press to copy."
: canExpand
? `Double tap to ${expanded ? "hide" : "show"} full details. Long press to copy.`
: "Long press to copy."
}
accessibilityState={canExpand ? { expanded } : undefined}
onPress={() => {
if (notifiedSubagentThreadId !== undefined) {
if (linkedThreadId !== undefined) {
// Push, not navigate: navigate reuses this Thread route, so back
// would skip the parent thread and land on Home (matches #15068).
navigation.dispatch(
StackActions.push("Thread", {
environmentId: String(props.environmentId),
threadId: String(notifiedSubagentThreadId),
threadId: String(linkedThreadId),
}),
);
return;
Expand Down
4 changes: 3 additions & 1 deletion apps/server/src/mcp/OrchestratorMcpService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -499,10 +499,12 @@ describe("OrchestratorMcpService", () => {
clientRequestId: "cancel-dispose-failed-task",
});
assert.equal(result.status, "cancel_requested");
const commands = yield* Ref.get(dispatched);
assert.deepEqual(
(yield* Ref.get(dispatched)).map((command) => (command as { type: string }).type),
commands.map((command) => (command as { type: string }).type),
["thread.stop", "delegated_task.completion-delivery.dispose"],
);
assert.deepInclude(commands[0], { createdBy: "agent", senderThreadId: parentThreadId });
}).pipe(Effect.provide(OrchestratorMcpService.layer.pipe(Layer.provide(layerDependencies))));
}),
);
Expand Down
14 changes: 13 additions & 1 deletion apps/server/src/mcp/OrchestratorMcpService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2054,12 +2054,17 @@ const make = Effect.gen(function* () {
const stopChild = Effect.gen(function* () {
const commandId = stableCommandId({ scope, requestKey: key, operation: "cancel-task" });
const reason = input.reason;
const attribution = {
createdBy: "agent",
senderThreadId: scope.thread.threadId,
} as const;
yield* threadManagement
.dispatch({
type: "thread.stop",
commandId,
threadId: current.childThreadId,
...(reason === undefined ? {} : { reason }),
...attribution,
})
.pipe(
Effect.mapError((error) =>
Expand All @@ -2071,7 +2076,12 @@ const make = Effect.gen(function* () {
);
// A retry with the same clientRequestId repeats only the stops that failed.
yield* threadManagement
.stopDelegatedTasks({ threadId: current.childThreadId, commandId, reason })
.stopDelegatedTasks({
threadId: current.childThreadId,
commandId,
reason,
...attribution,
})
.pipe(
Effect.mapError((error) =>
failure(
Expand Down Expand Up @@ -2518,6 +2528,8 @@ const make = Effect.gen(function* () {
threadId: input.threadId,
...(input.runId === undefined ? {} : { runId: input.runId }),
...(input.reason === undefined ? {} : { reason: input.reason }),
createdBy: "agent",
...(parent === undefined ? {} : { senderThreadId: parent.thread.id }),
})
.pipe(
Effect.mapError((error) =>
Expand Down
18 changes: 18 additions & 0 deletions apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2055,6 +2055,15 @@ describe("orchestrator MCP toolkit", () => {
cancelledStatusCall.structuredContent,
).pipe(Effect.orDie);
expect(cancelledStatus.status).toBe("interrupted");
expect(
(yield* orchestrator.getThreadProjection(cancellable.childThreadId)).turnItems.find(
(item) => item.type === "run_interrupt_result",
),
).toMatchObject({
createdBy: "agent",
senderThreadId: parentThreadId,
message: "Run interrupted by an agent",
});

// Explicit task_cancel disposes automatic delivery after the child
// interrupt succeeds. The interrupted result remains readable,
Expand Down Expand Up @@ -2571,6 +2580,15 @@ describe("orchestrator MCP toolkit", () => {
interruptedWaitCall.structuredContent,
).pipe(Effect.orDie);
expect(interruptedWait.status).toBe("interrupted");
expect(
(yield* orchestrator.getThreadProjection(activeThread.threadId)).turnItems.find(
(item) => item.type === "run_interrupt_result" && item.runId === activeRun.id,
),
).toMatchObject({
createdBy: "agent",
senderThreadId: parentThreadId,
message: "Run interrupted by an agent before provider start",
});
const repeatedInterruptCall = yield* invoke("t3_thread_interrupt", {
threadId: activeThread.threadId,
runId: activeRun.id,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -351,6 +351,7 @@ const stopEarlierBackgroundWork = ({
type: "thread.stop",
commandId: CommandId.make("stop-stalled-run"),
threadId,
createdBy: "agent",
},
);
yield* worker.drain();
Expand Down Expand Up @@ -381,10 +382,16 @@ const stopEarlierBackgroundWork = ({
after.turnItems.find((candidate) => candidate.type === "assistant_message")?.status,
interrupted ? "interrupted" : "running",
);
assert.equal(
after.turnItems.filter((candidate) => candidate.type === "run_interrupt_result").length,
interrupted ? 1 : 0,
const results = after.turnItems.filter(
(candidate) => candidate.type === "run_interrupt_result",
);
assert.equal(results.length, interrupted ? 1 : 0);
if (interrupted) {
assert.deepInclude(results[0], {
createdBy: "agent",
message: "Run interrupted by an agent",
});
}
assert.isEmpty(after.runs.filter((candidate) => candidate.status === "waiting"));
return;
}
Expand Down
40 changes: 39 additions & 1 deletion apps/server/src/orchestration-v2/Orchestrator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -536,6 +536,19 @@ function hasLiveRun(projection: Pick<OrchestrationV2ThreadProjection, "runs">):
);
}

/** A stop with no `createdBy` came from a client's Stop, so the user. */
function runInterruptAttribution(
command: Extract<
OrchestrationV2ServerCommand,
{ readonly type: "run.interrupt" | "thread.stop" }
>,
) {
return {
createdBy: command.createdBy ?? "user",
...(command.senderThreadId === undefined ? {} : { senderThreadId: command.senderThreadId }),
} as const;
}

/** The link with its watch replaced, or removed when `watch` is undefined. */
function withPullRequestWatch(
link: ThreadPullRequestLink,
Expand Down Expand Up @@ -8530,6 +8543,20 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
});
}
if (providerThread !== undefined) {
const requestId = idAllocator.derive.runSignalTurnItem({
runId: run.id,
signal: "interrupt-request",
});
// A dead session's Stop writes its request in this same command, before the store has it.
const pendingRequest = (yield* Ref.get(input.events)).findLast(
(event) => event.type === "turn-item.updated" && event.payload.id === requestId,
);
const request =
pendingRequest?.type === "turn-item.updated"
? pendingRequest.payload
: yield* projectionStore
.getTurnItem({ threadId: run.threadId, itemId: requestId })
.pipe(Effect.catchCause(() => Effect.succeed(null)));
yield* emitEvent({
...base,
type: "turn-item.updated",
Expand All @@ -8538,6 +8565,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
run,
rootNode,
providerThread,
request: request?.type === "run_interrupt_request" ? request : null,
completedAt: input.now,
}),
});
Expand Down Expand Up @@ -8638,6 +8666,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
completedAt: now,
updatedAt: now,
type: "run_interrupt_request",
...runInterruptAttribution(command),
message: command.reason ?? "Interrupt requested",
},
});
Expand Down Expand Up @@ -8909,6 +8938,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
});

const emitEvent = emit(events, command);
const interruptAttribution = runInterruptAttribution(command);
const interruptRequestItem: OrchestrationV2TurnItem = {
id: idAllocator.derive.runSignalTurnItem({
runId: run.id,
Expand All @@ -8928,6 +8958,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
completedAt: now,
updatedAt: now,
type: "run_interrupt_request",
...interruptAttribution,
message: command.reason ?? "Interrupt requested",
};

Expand Down Expand Up @@ -9022,7 +9053,13 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
completedAt: now,
updatedAt: now,
type: "run_interrupt_result",
message: "Run interrupted before provider start",
...interruptAttribution,
message:
interruptAttribution.createdBy === "agent"
? "Run interrupted by an agent before provider start"
: interruptAttribution.createdBy === "user"
? "Run interrupted by user before provider start"
: "Run interrupted before provider start",
};
yield* emitEvent({
type: "turn-item.updated",
Expand Down Expand Up @@ -9252,6 +9289,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
runId: target.id,
holdQueue: true,
...(command.reason === undefined ? {} : { reason: command.reason }),
...runInterruptAttribution(command),
},
interruptEvents,
interruptEffects,
Expand Down
14 changes: 14 additions & 0 deletions apps/server/src/orchestration-v2/ProviderTurnStartService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -175,6 +175,19 @@ export const layer: Layer.Layer<
}),
)
.pipe(Effect.catchCause(() => Effect.succeed(false))),
loadRunInterruptRequest: () =>
projectionStore
.getTurnItem({
threadId: input.threadId,
itemId: idAllocator.derive.runSignalTurnItem({
runId: input.runId,
signal: "interrupt-request",
}),
})
.pipe(
Effect.map((item) => (item?.type === "run_interrupt_request" ? item : null)),
Effect.catchCause(() => Effect.succeed(null)),
),
};
};

Expand Down Expand Up @@ -1257,6 +1270,7 @@ export const layer: Layer.Layer<
shouldStartProviderTurn: runControls.shouldStartProviderTurn,
shouldFinalizeRun: runControls.shouldFinalizeRun,
hasUnpairedRunInterruptRequest: runControls.hasUnpairedRunInterruptRequest,
loadRunInterruptRequest: runControls.loadRunInterruptRequest,
message: {
messageId: message.id,
text: userText,
Expand Down
65 changes: 49 additions & 16 deletions apps/server/src/orchestration-v2/RunExecutionService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3352,22 +3352,51 @@ it.effect("does not overwrite Stop when ownership changes after the finalization
}),
);

it.effect("emits run_interrupt_result when hard-stop finalizes the active attempt", () =>
Effect.gen(function* () {
const { written, observed, committedEffects } = yield* captureRootRunTermination({
key: "hard-stop",
shouldFinalizeRun: () => Effect.succeed(true),
});
assert.deepEqual(
written.map((item) => item.type),
["run_interrupt_result"],
);
assert.deepEqual(observed, ["run:interrupted", "pull-requests-refreshed"]);
assert.deepEqual(
committedEffects.map((effect) => effect.request.type),
["checkpoint.capture"],
);
}),
const interruptingThreadId = ThreadId.make("thread:hard-stop:interrupting");

it.effect.each([
{
name: "an agent's request",
request: { createdBy: "agent", senderThreadId: interruptingThreadId },
message: "Run interrupted by an agent",
},
{ name: "the user's Stop", request: { createdBy: "user" }, message: "Run interrupted by user" },
{ name: "no request", request: null, message: "" },
] as const)(
"emits run_interrupt_result attributed to $name when hard-stop finalizes the active attempt",
({ name, request, message }) =>
Effect.gen(function* () {
const { written, observed, committedEffects } = yield* captureRootRunTermination({
key: `hard-stop:${name}`,
shouldFinalizeRun: () => Effect.succeed(true),
interruptRequest:
request === null
? null
: ({
type: "run_interrupt_request",
...request,
} as RunExecutionService.RunInterruptRequestTurnItem),
});
assert.deepEqual(
written.map((item) => item.type),
["run_interrupt_result"],
);
const result = written[0];
if (result?.type !== "run_interrupt_result") {
return assert.fail("expected run_interrupt_result");
}
assert.equal(result.message, message);
assert.equal(result.createdBy, request?.createdBy);
assert.equal(
result.senderThreadId,
request !== null && "senderThreadId" in request ? request.senderThreadId : undefined,
);
assert.deepEqual(observed, ["run:interrupted", "pull-requests-refreshed"]);
assert.deepEqual(
committedEffects.map((effect) => effect.request.type),
["checkpoint.capture"],
);
}),
);

it.effect.each(["completed", "interrupted", "cancelled", "failed"] as const)(
Expand Down Expand Up @@ -3497,6 +3526,7 @@ function captureRootRunTermination(input: {
readonly shouldFinalizeRun: () => Effect.Effect<boolean, ProjectionStore.ProjectionStoreV2Error>;
readonly rejectTerminalWrite?: boolean;
readonly hasUnpairedRunInterruptRequest?: () => Effect.Effect<boolean, never>;
readonly interruptRequest?: RunExecutionService.RunInterruptRequestTurnItem | null;
readonly seedOpenSubagent?: boolean;
readonly events?: (
ids: BackgroundScenarioIds,
Expand Down Expand Up @@ -3653,6 +3683,9 @@ function captureRootRunTermination(input: {
: {
hasUnpairedRunInterruptRequest: input.hasUnpairedRunInterruptRequest,
}),
...(input.interruptRequest === undefined
? {}
: { loadRunInterruptRequest: () => Effect.succeed(input.interruptRequest ?? null) }),
message: {
messageId: MessageId.make(`message:${input.key}`),
text: "interrupt projection",
Expand Down
Loading
Loading