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
156 changes: 135 additions & 21 deletions apps/server/src/orchestration-v2/EventSink.ts
Original file line number Diff line number Diff line change
Expand Up @@ -84,6 +84,18 @@ export interface EventSinkV2Shape {
}) => Effect.Effect<ReadonlyArray<OrchestrationV2StoredEvent>, EventSinkV2Error>;
readonly writeIfRunCurrent: (input: {
readonly guardPendingUserInputCancellations?: boolean;
/**
* Guards approvals and every other request kind too. For a writer that
* cancels requests an earlier projection read showed pending; a
* provider's own approval cancellation stays authoritative.
*/
readonly guardPendingRequestCancellations?: boolean;
/**
* Drops subagent, node and turn-item updates for rows that ended after an
* earlier projection read showed them open. For a writer that settles
* work it no longer receives events for.
*/
readonly guardSettledWork?: boolean;
readonly commandId?: CommandId;
readonly threadId: ThreadId;
readonly runId: RunId;
Expand Down Expand Up @@ -263,14 +275,65 @@ const layerBase: Layer.Layer<

// A user can answer after terminal normalization reads the pending request.
// Recheck inside the write transaction so stale cleanup cannot erase answers.
const guardUserInputCancellations = (events: ReadonlyArray<OrchestrationV2DomainEvent>) =>
// `allKinds` extends this from questions to approvals and every other kind.
const guardRequestCancellations = (
events: ReadonlyArray<OrchestrationV2DomainEvent>,
allKinds: boolean,
) =>
Effect.gen(function* () {
const staleRequests = new Set<RuntimeRequestId>();
const staleNodes = new Set<NodeId>();
// A run's cleanup settles a request's node and item with the run's own
// status, which is not always `cancelled`.
const settles = (status: string) =>
status === "cancelled" || (allKinds && (status === "failed" || status === "interrupted"));
// Whether a request was answered after the cleanup read it. A missing
// request leaves nothing to protect.
const answered = new Map<RuntimeRequestId, boolean>();
const wasAnswered = (threadId: ThreadId, requestId: RuntimeRequestId) =>
Effect.gen(function* () {
const known = answered.get(requestId);
if (known !== undefined) return known;
const current = yield* projectionStore.getRuntimeRequest(threadId, requestId);
const result =
current !== undefined &&
(current.status !== "pending" || current.responseCapability.type === "message");
answered.set(requestId, result);
return result;
});
for (const event of events) {
// The cleanup reads nodes, requests and items at different moments,
// so it can settle a request's node or item after the request was
// answered and left its pending read. Each one is checked against
// its own request, with or without a cancellation beside it.
if (allKinds && event.type === "node.updated") {
const requestId = event.payload.runtimeRequestId;
if (
requestId !== null &&
settles(event.payload.status) &&
(yield* wasAnswered(event.payload.threadId, requestId))
) {
staleNodes.add(event.payload.id);
}
continue;
}
if (
allKinds &&
event.type === "turn-item.updated" &&
(event.payload.type === "user_input_request" ||
event.payload.type === "approval_request")
) {
if (
settles(event.payload.status) &&
(yield* wasAnswered(event.payload.threadId, event.payload.requestId))
) {
staleRequests.add(event.payload.requestId);
}
continue;
}
if (
event.type !== "runtime-request.updated" ||
event.payload.kind !== "user_input" ||
(!allKinds && event.payload.kind !== "user_input") ||
event.payload.status !== "cancelled"
)
continue;
Expand All @@ -280,7 +343,7 @@ const layerBase: Layer.Layer<
);
if (
current?.status !== "pending" ||
current.kind !== "user_input" ||
current.kind !== event.payload.kind ||
current.providerTurnId !== event.payload.providerTurnId ||
current.responseCapability.type === "message"
) {
Expand All @@ -293,18 +356,81 @@ const layerBase: Layer.Layer<
case "runtime-request.updated":
return event.payload.status !== "cancelled" || !staleRequests.has(event.payload.id);
case "node.updated":
return event.payload.status !== "cancelled" || !staleNodes.has(event.payload.id);
return !settles(event.payload.status) || !staleNodes.has(event.payload.id);
case "turn-item.updated":
return (
event.payload.type !== "user_input_request" ||
event.payload.status !== "cancelled" ||
(event.payload.type !== "user_input_request" &&
event.payload.type !== "approval_request") ||
!settles(event.payload.status) ||
!staleRequests.has(event.payload.requestId)
);
default:
return true;
}
});
});
// A cleanup settles subagents, nodes and items it read as open. One that
// ended on its own before the commit keeps its own outcome and result.
const guardSettledWork = (events: ReadonlyArray<OrchestrationV2DomainEvent>) =>
Effect.gen(function* () {
const ended = (status: string) =>
status === "completed" ||
status === "failed" ||
status === "cancelled" ||
status === "interrupted";
const threads = new Map<
ThreadId,
ProjectionStore.ProjectionRecords<"nodes" | "subagents"> | undefined
>();
const records = (threadId: ThreadId) =>
Effect.gen(function* () {
if (threads.has(threadId)) return threads.get(threadId);
const read = yield* projectionStore

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟠 High orchestration-v2/EventSink.ts:388

A persistent getThreadRecords read/decode error aborts settleRunBeforeStart before eventStore.append, leaving the run and attempt starting instead of recording the provider-start failure. guardSettledWork catches only ProjectionStoreThreadNotFoundError, so its reread of nodes and subagents can fail even after inherited-work collection already failed; keep this guard’s read failure from blocking mandatory run/attempt/root settlement.

🚀 Reply "fix it for me" or copy this AI Prompt for your agent:
In file @apps/server/src/orchestration-v2/EventSink.ts around line 388:

A persistent `getThreadRecords` read/decode error aborts `settleRunBeforeStart` before `eventStore.append`, leaving the run and attempt `starting` instead of recording the provider-start failure. `guardSettledWork` catches only `ProjectionStoreThreadNotFoundError`, so its reread of nodes and subagents can fail even after inherited-work collection already failed; keep this guard’s read failure from blocking mandatory run/attempt/root settlement.

.getThreadRecords(threadId, ["nodes", "subagents"])
.pipe(
Effect.catchTags({
ProjectionStoreThreadNotFoundError: () => Effect.succeed(undefined),
}),
);
threads.set(threadId, read);
return read;
});
const kept: Array<OrchestrationV2DomainEvent> = [];
for (const event of events) {
const current =
event.type === "subagent.updated"
? (yield* records(event.payload.threadId))?.subagents.find(
(row) => row.id === event.payload.id,
)
: event.type === "node.updated"
? (yield* records(event.payload.threadId))?.nodes.find(
(row) => row.id === event.payload.id,
)
: event.type === "turn-item.updated"
? yield* projectionStore.getTurnItem({
threadId: event.payload.threadId,
itemId: event.payload.id,
})
: undefined;
if (current == null || !ended(current.status)) kept.push(event);
}
return kept;
});
const guardCancellations = (input: {
readonly guardPendingUserInputCancellations?: boolean;
readonly guardPendingRequestCancellations?: boolean;
readonly guardSettledWork?: boolean;
readonly events: ReadonlyArray<OrchestrationV2DomainEvent>;
}) =>
Effect.gen(function* () {
const events: ReadonlyArray<OrchestrationV2DomainEvent> =
input.guardPendingRequestCancellations === true
? yield* guardRequestCancellations(input.events, true)
: input.guardPendingUserInputCancellations === true
? yield* guardRequestCancellations(input.events, false)
: input.events;
return input.guardSettledWork === true ? yield* guardSettledWork(events) : events;
});

const normalizeEvents = (events: ReadonlyArray<OrchestrationV2DomainEvent>) => {
const runOrdinals = new Map(
Expand Down Expand Up @@ -370,11 +496,7 @@ const layerBase: Layer.Layer<

return yield* commitThenPublish(
Effect.gen(function* () {
const normalized = yield* normalizeEvents(
input.guardPendingUserInputCancellations === true
? yield* guardUserInputCancellations(input.events)
: input.events,
);
const normalized = yield* normalizeEvents(yield* guardCancellations(input));
const committed = yield* eventStore.append({
...(input.commandId === undefined ? {} : { commandId: input.commandId }),
events: normalized,
Expand Down Expand Up @@ -428,11 +550,7 @@ const layerBase: Layer.Layer<
};
}

const normalized = yield* normalizeEvents(
input.guardPendingUserInputCancellations === true
? yield* guardUserInputCancellations(input.events)
: input.events,
);
const normalized = yield* normalizeEvents(yield* guardCancellations(input));
const storedEvents = yield* eventStore.append({
...(input.commandId === undefined ? {} : { commandId: input.commandId }),
events: normalized,
Expand Down Expand Up @@ -493,11 +611,7 @@ const layerBase: Layer.Layer<
};
}

const normalized = yield* normalizeEvents(
input.guardPendingUserInputCancellations === true
? yield* guardUserInputCancellations(input.events)
: input.events,
);
const normalized = yield* normalizeEvents(yield* guardCancellations(input));
const storedEvents = yield* eventStore.append({
...(input.commandId === undefined ? {} : { commandId: input.commandId }),
events: normalized,
Expand Down
Loading
Loading