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
29 changes: 23 additions & 6 deletions apps/server/src/orchestration-v2/EffectOutbox.ts
Original file line number Diff line number Diff line change
Expand Up @@ -123,6 +123,24 @@ export const PROCESS_BOUND_EFFECT_TYPES = [
"runtime-request.respond",
] as const satisfies ReadonlyArray<OrchestrationEffectRequestV2["type"]>;

export const SESSION_BOUND_EFFECT_TYPES = [
"provider-turn.interrupt",
"provider-turn.steer",
"provider-turn.restart",
"runtime-request.respond",
] as const satisfies ReadonlyArray<OrchestrationEffectRequestV2["type"]>;

/**
* Which unsettled effects of a thread to cancel. `providerSessionId` narrows
* it to effects addressed to that provider session, leaving effects for other
* sessions on the same thread alone.
*/
export interface UnsettledEffectCancellation {
readonly effectTypes: ReadonlyArray<OrchestrationEffectRequestV2["type"]>;
readonly reason: string;
readonly providerSessionId?: ProviderSessionId;
}

export const OrchestrationEffectStatusV2 = Schema.Literals([
"pending",
"running",
Expand Down Expand Up @@ -184,11 +202,9 @@ export interface EffectOutboxV2Shape {
readonly listByCommandId: (
commandId: CommandId,
) => Effect.Effect<ReadonlyArray<OrchestrationEffectV2>, EffectOutboxError>;
readonly cancelUnsettled: (input: {
readonly threadId: ThreadId;
readonly effectTypes: ReadonlyArray<OrchestrationEffectRequestV2["type"]>;
readonly reason: string;
}) => Effect.Effect<ReadonlyArray<string>, EffectOutboxError>;
readonly cancelUnsettled: (
input: UnsettledEffectCancellation & { readonly threadId: ThreadId },
) => Effect.Effect<ReadonlyArray<string>, EffectOutboxError>;
readonly signalCancellations: (effectIds: ReadonlyArray<string>) => Effect.Effect<void>;
readonly awaitCancellation: (effectId: string) => Effect.Effect<void>;
readonly clearCancellation: (effectId: string) => Effect.Effect<void>;
Expand Down Expand Up @@ -391,7 +407,7 @@ export const layer: Layer.Layer<EffectOutboxV2, never, SqlClient.SqlClient> = La
: new EffectOutboxError({ operation: "list", cause }),
),
),
cancelUnsettled: ({ threadId, effectTypes, reason }) =>
cancelUnsettled: ({ threadId, effectTypes, reason, providerSessionId }) =>
Effect.gen(function* () {
if (effectTypes.length === 0) return [];
const now = DateTime.formatIso(yield* DateTime.now);
Expand All @@ -407,6 +423,7 @@ export const layer: Layer.Layer<EffectOutboxV2, never, SqlClient.SqlClient> = La
WHERE thread_id = ${threadId}
AND status IN ('pending', 'running')
AND effect_type IN ${sql.in(effectTypes)}
${providerSessionId === undefined ? sql`` : sql`AND json_extract(payload_json, '$.providerSessionId') = ${providerSessionId}`}
RETURNING effect_id
`;
return rows.map(({ effect_id }) => effect_id);
Expand Down
5 changes: 1 addition & 4 deletions apps/server/src/orchestration-v2/EventSink.ts
Original file line number Diff line number Diff line change
Expand Up @@ -124,10 +124,7 @@ export interface EventSinkV2Shape {
readonly acceptedAt: DateTime.Utc;
readonly events: ReadonlyArray<OrchestrationV2DomainEvent>;
readonly effects: ReadonlyArray<EffectOutbox.PendingOrchestrationEffectV2>;
readonly cancelUnsettledEffects?: {
readonly effectTypes: ReadonlyArray<EffectOutbox.OrchestrationEffectRequestV2["type"]>;
readonly reason: string;
};
readonly cancelUnsettledEffects?: EffectOutbox.UnsettledEffectCancellation;
}) => Effect.Effect<
{
readonly receipt: CommandReceiptStore.CommandReceiptV2;
Expand Down
148 changes: 124 additions & 24 deletions apps/server/src/orchestration-v2/Orchestrator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -77,7 +77,12 @@ import { notificationTurnItem } from "./Notification.ts";
import { isRestartNoteSource } from "./RestartBackgroundNote.ts";
import { isUndeliveredMailboxSteer } from "./NotificationMailbox.ts";
import { EventSinkV2 } from "./EventSink.ts";
import type { OrchestrationEffectRequestV2, PendingOrchestrationEffectV2 } from "./EffectOutbox.ts";
import {
EffectOutboxV2,
SESSION_BOUND_EFFECT_TYPES,
type PendingOrchestrationEffectV2,
type UnsettledEffectCancellation,
} from "./EffectOutbox.ts";
import { IdAllocatorV2 } from "./IdAllocator.ts";
import {
ThreadCommandExecutor,
Expand All @@ -99,6 +104,7 @@ import type { ProviderAdapterV2Shape } from "./ProviderAdapter.ts";
import { ProviderAdapterRegistryV2 } from "./ProviderAdapterRegistry.ts";
import { ProviderContinuationRequests } from "./ProviderContinuationRequests.ts";
import { makeProviderFailure } from "./ProviderFailure.ts";
import { planProcessLossReconciliation } from "./ProviderRuntimeRecoveryService.ts";
import {
metricAttributes,
orchestrationCommandAckDuration,
Expand Down Expand Up @@ -402,6 +408,15 @@ function pendingThreadTitleGenerationEffect(

const WORKSPACE_PREPARATION_INPUT = "Preparing workspace";

// Turn item types a live provider process owns and reports the end of itself
// (RunExecutionService/ProviderRuntimeRecoveryService background-work sweeps).
// Interrupt finalization must not re-settle these outside that sweep.
const BACKGROUND_CAPABLE_TURN_ITEM_TYPES: ReadonlySet<OrchestrationV2TurnItem["type"]> = new Set([
"command_execution",
"dynamic_tool",
"subagent",
]);

function isBlockingRun(run: OrchestrationV2Run): boolean {
return (
run.status === "preparing" ||
Expand Down Expand Up @@ -661,6 +676,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
const eventSink = yield* EventSinkV2;
const commandReceipts = yield* CommandReceiptStoreV2;
const idAllocator = yield* IdAllocatorV2;
const outbox = yield* EffectOutboxV2;
const projects = yield* ProjectStore.ProjectStoreV2;
const projectionStore = yield* ProjectionStoreV2;
const nextTurnItemOrdinal = (
Expand Down Expand Up @@ -7716,7 +7732,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
// Thread-wide, like the Waiting strip: a settled thread's background
// work can belong to an earlier run than the one Stop targets.
{
turnItemTypes: ["command_execution", "dynamic_tool", "subagent"],
turnItemTypes: [...BACKGROUND_CAPABLE_TURN_ITEM_TYPES],
turnItemStatuses: ["pending", "running", "waiting"],
},
);
Expand Down Expand Up @@ -7920,10 +7936,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
return {
effectTypes: ["provider-turn.start", "provider-turn.restart"],
reason: `Run ${run.id} was interrupted before its provider turn started.`,
} satisfies {
readonly effectTypes: ReadonlyArray<OrchestrationEffectRequestV2["type"]>;
readonly reason: string;
};
} satisfies UnsettledEffectCancellation;
}

if (providerTurn === undefined) {
Expand All @@ -7936,14 +7949,14 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
// Stop on a settled thread's background work. Its process may be gone
// (released, restarted) and only the projection still shows the work;
// the settle follow-up ends whatever no provider reports ending.
const settleOnly =
providerTurn.status !== "running" &&
(providerThread.providerSessionId === null ||
Option.isNone(
yield* providerSessions
.get(providerThread.providerSessionId)
.pipe(Effect.orElseSucceed(() => Option.none())),
));
const sessionIsDead =
providerThread.providerSessionId === null ||
Option.isNone(
yield* providerSessions
.get(providerThread.providerSessionId)
.pipe(Effect.orElseSucceed(() => Option.none())),
);
const settleOnly = providerTurn.status !== "running" && sessionIsDead;
if (settleOnly) {
yield* emitEvent({
type: "turn-item.updated",
Expand All @@ -7966,6 +7979,100 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
});
return undefined;
}
// The projection still shows this turn running, but its provider
// session is already gone (e.g. released on idle timeout): no live
// process will ever report a terminal for it. Settle the run the same
// way as an interrupt before provider start, plus the stuck turn
// itself, instead of erroring and leaving the run wedged forever.
if (providerTurn.status === "running" && sessionIsDead) {
if (!projection.attempts.some((candidate) => candidate.id === run.activeAttemptId)) {
return yield* new OrchestratorDispatchError({
commandId: command.commandId,
commandType: command.type,
cause: `Run ${command.runId} has no active attempt to interrupt.`,
});
}
yield* emitEvent({
type: "turn-item.updated",
threadId: command.threadId,
runId: run.id,
nodeId: rootNode.id,
providerInstanceId: run.providerInstanceId,
occurredAt: now,
payload: interruptRequestItem,
});
// No live process remains to settle this run, so finalize it exactly
// the way ProviderRuntimeRecoveryService reconciles a dead provider
// process: scoped to this run and its own provider thread, so sibling
// live provider threads on the same orchestration thread (e.g. a
// concurrently-running native subagent) are left untouched.
const recoveryProjection = yield* projectionStore
.getRuntimeRecoveryProjection(command.threadId)
.pipe(
Effect.mapError(
(cause) => new OrchestratorProjectionError({ threadId: command.threadId, cause }),
),
);
const plan = yield* planProcessLossReconciliation({
projection: recoveryProjection,
runId: run.id,
providerThreadId: providerThread.id,
commandId: command.commandId,
now,
ids: idAllocator,
outbox,
}).pipe(
Effect.mapError(
(cause) =>
new OrchestratorDispatchError({
commandId: command.commandId,
commandType: command.type,
cause,
}),
),
);
yield* Ref.update(events, (existing) => [...existing, ...plan.events]);
yield* Ref.update(effects, (existing) => [...existing, ...plan.effects]);
Comment thread
yordis marked this conversation as resolved.
// The workspace may have changed before the session died, so the run
// still gets its end-of-turn checkpoint like any interrupted turn.
const checkpointScopeId = rootNode.checkpointScopeId;
if (checkpointScopeId !== null) {
yield* Ref.update(effects, (existing) => [
...existing,
{
id: `effect:checkpoint.capture:${run.id}`,
commandId: CommandId.make(`command:effect:checkpoint.capture:${run.id}`),
threadId: run.threadId,
request: {
type: "checkpoint.capture" as const,
runId: run.id,
scopeId: checkpointScopeId,
},
},
]);
}
if (command.holdQueue === true) yield* holdQueuedRuns;
yield* stopCompletionCohort();
// Read past the plan's events so a stale provider-thread row cannot
// overwrite the idle state the plan just wrote.
yield* settleBackgroundWork({
command,
events,
projection: yield* getProjectionWithPendingEvents(command.threadId, events),
stoppedProviderThreadId: providerThread.id,
throughRunOrdinal: run.ordinal,
now,
});
Comment thread
yordis marked this conversation as resolved.
// Only effects addressed to the dead session lost their process; a
// thread's other sessions keep theirs.
return {
effectTypes: providerThread.providerSessionId === null ? [] : SESSION_BOUND_EFFECT_TYPES,
reason: plan.detail,
...(providerThread.providerSessionId === null
? {}
: { providerSessionId: providerThread.providerSessionId }),
} satisfies UnsettledEffectCancellation;
}
Comment thread
cursor[bot] marked this conversation as resolved.
Comment thread
coderabbitai[bot] marked this conversation as resolved.
if (providerThread.providerSessionId === null) {
return yield* new OrchestratorDispatchError({
commandId: command.commandId,
Expand Down Expand Up @@ -9021,10 +9128,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
{
readonly events: ReadonlyArray<OrchestrationV2DomainEvent>;
readonly effects: ReadonlyArray<PendingOrchestrationEffectV2>;
readonly cancelUnsettledEffects?: {
readonly effectTypes: ReadonlyArray<OrchestrationEffectRequestV2["type"]>;
readonly reason: string;
};
readonly cancelUnsettledEffects?: UnsettledEffectCancellation;
},
OrchestratorV2Error
> {
Expand All @@ -9036,12 +9140,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio

const events = yield* Ref.make<Array<OrchestrationV2DomainEvent>>([]);
const effects = yield* Ref.make<Array<PendingOrchestrationEffectV2>>([]);
let cancelUnsettledEffects:
| {
readonly effectTypes: ReadonlyArray<OrchestrationEffectRequestV2["type"]>;
readonly reason: string;
}
| undefined;
let cancelUnsettledEffects: UnsettledEffectCancellation | undefined;
switch (command.type) {
case "thread.create":
yield* dispatchThreadCreate(command, events);
Expand Down Expand Up @@ -9740,6 +9839,7 @@ export const layer: Layer.Layer<
| CommandPolicyV2
| CommandReceiptStoreV2
| ContextHandoffServiceV2
| EffectOutboxV2
| EventSinkV2
| IdAllocatorV2
| ProjectStore.ProjectStoreV2
Expand Down
Loading
Loading