Skip to content
Closed
164 changes: 161 additions & 3 deletions apps/server/src/orchestration-v2/EventSink.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
import {
CommandId,
type OrchestrationV2Run,
OrchestrationV2ConversationMessageJson,
OrchestrationV2TurnItemJson,
OrchestrationV2DomainEvent,
OrchestrationV2StoredEvent,
ProviderThreadId,
Expand Down Expand Up @@ -33,6 +35,13 @@ import * as ProjectionStore from "./ProjectionStore.ts";
import * as ProjectStore from "./ProjectStore.ts";
import * as TurnItemPositionStore from "./TurnItemPositionStore.ts";

const decodeMessagePayload = Schema.decodeUnknownEffect(
Schema.fromJsonString(OrchestrationV2ConversationMessageJson),
);
const decodeTurnItemPayload = Schema.decodeUnknownEffect(
Schema.fromJsonString(OrchestrationV2TurnItemJson),
);

/**
* ERRORS
*/
Expand Down Expand Up @@ -78,6 +87,8 @@ export interface EventSinkV2Shape {
}) => Effect.Effect<ReadonlyArray<OrchestrationV2StoredEvent>, EventSinkV2Error>;
readonly writeWithEffects: (input: {
readonly guardPendingUserInputCancellations?: boolean;
/** Skip artifact cleanup when its node now belongs to a different run attempt. */
readonly guardNodeRootId?: NodeId;
readonly commandId?: CommandId;
readonly events: ReadonlyArray<OrchestrationV2DomainEvent>;
readonly effects: ReadonlyArray<EffectOutbox.PendingOrchestrationEffectV2>;
Expand All @@ -90,6 +101,7 @@ export interface EventSinkV2Shape {
readonly activeAttemptId: RunAttemptId;
readonly expectedStatus: OrchestrationV2Run["status"];
readonly events: ReadonlyArray<OrchestrationV2DomainEvent>;
readonly effects?: ReadonlyArray<EffectOutbox.PendingOrchestrationEffectV2>;
}) => Effect.Effect<
{
readonly committed: boolean;
Expand Down Expand Up @@ -305,6 +317,139 @@ const baseLayer: Layer.Layer<
});
});

const guardNodeRootArtifacts = Effect.fnUntraced(function* (
events: ReadonlyArray<OrchestrationV2DomainEvent>,
rootNodeId: NodeId,
) {
const offeredNodes = new Map(
events.flatMap((event) =>
event.type === "node.updated" && event.payload.rootNodeId === rootNodeId
? [[event.payload.id, event.payload.providerTurnId] as const]
: [],
),
);
const accepted: Array<OrchestrationV2DomainEvent> = [];
for (const event of events) {
if (
event.type !== "node.updated" &&
event.type !== "message.updated" &&
event.type !== "turn-item.updated"
) {
accepted.push(event);
continue;
}
const nodeId = event.type === "node.updated" ? event.payload.id : event.payload.nodeId;
if (nodeId === null) continue;
if (event.type === "node.updated" && event.payload.rootNodeId !== rootNodeId) continue;
const nodes = yield* sql<{
readonly root_node_id: string;
readonly status: string;
readonly provider_turn_id: string | null;
}>`
SELECT root_node_id, status, provider_turn_id FROM orchestration_v2_projection_nodes
WHERE thread_id = ${event.threadId} AND node_id = ${nodeId}`;
const node = nodes[0];
if (node === undefined ? !offeredNodes.has(nodeId) : node.root_node_id !== rootNodeId)
continue;
// Assistant streaming may suppress the successor's running node update.
// Check item provenance too, even when the projected node still looks old.
// Root-linked reasoning shares a node with the user message. Its
// turn identity belongs to the item, and sibling items may have none.
const isRootTurnItem =
event.type === "turn-item.updated" &&
nodeId === rootNodeId &&
event.payload.providerTurnId !== null;
const expectedTurnId = isRootTurnItem
? event.payload.providerTurnId
: (offeredNodes.get(nodeId) ?? node?.provider_turn_id);
const items = isRootTurnItem
? []
: yield* sql<{ readonly provider_turn_id: string | null }>`
SELECT provider_turn_id FROM orchestration_v2_projection_turn_items
WHERE thread_id = ${event.threadId} AND node_id = ${nodeId}`;
if (items.some((item) => item.provider_turn_id !== expectedTurnId)) continue;
if (expectedTurnId == null && (node !== undefined || items.length > 0)) continue;
if (event.type === "node.updated") {
if (
node !== undefined &&
node.status !== "pending" &&
node.status !== "running" &&
node.status !== "waiting"
)
continue;
} else if (event.type === "message.updated") {
const rows = yield* sql<{
readonly node_id: string | null;
readonly streaming: number;
readonly payload_json: string;
}>`
SELECT node_id, streaming, payload_json FROM orchestration_v2_projection_messages
WHERE thread_id = ${event.threadId} AND message_id = ${event.payload.id}`;
const current = rows[0];
if (current !== undefined && (current.node_id !== nodeId || current.streaming === 0))
continue;
if (current !== undefined) {
if (node === undefined && items.length === 0) continue;
const message = yield* decodeMessagePayload(current.payload_json);
accepted.push({
...event,
payload: {
...message,
text:
message.text.length > event.payload.text.length
? message.text
: event.payload.text,
streaming: false,
updatedAt: event.payload.updatedAt,
},
});
continue;
}
} else {
const rows = yield* sql<{
readonly node_id: string | null;
readonly status: string;
readonly provider_turn_id: string | null;
readonly payload_json: string;
}>`
SELECT node_id, status, provider_turn_id, payload_json FROM orchestration_v2_projection_turn_items
WHERE thread_id = ${event.threadId} AND turn_item_id = ${event.payload.id}`;
const current = rows[0];
if (
current !== undefined &&
(current.node_id !== nodeId ||
current.provider_turn_id !== event.payload.providerTurnId ||
(current.status !== "pending" &&
current.status !== "running" &&
current.status !== "waiting"))
)
continue;
if (current !== undefined) {
const item = yield* decodeTurnItemPayload(current.payload_json);
const content =
"text" in item &&
"text" in event.payload &&
item.text.length < event.payload.text.length
? event.payload
: item;
accepted.push({
...event,
payload: {
...content,
...("streaming" in content ? { streaming: false } : {}),
status: event.payload.status,
completedAt: event.payload.completedAt,
updatedAt: event.payload.updatedAt,
},
});
continue;
}
}
accepted.push(event);
}
return accepted;
});

const normalizeEvents = (events: ReadonlyArray<OrchestrationV2DomainEvent>) => {
const runOrdinals = new Map(
events.flatMap((event) =>
Expand Down Expand Up @@ -369,10 +514,14 @@ const baseLayer: Layer.Layer<

return yield* commitThenPublish(
Effect.gen(function* () {
const guardedEvents =
input.guardNodeRootId === undefined
? input.events
: yield* guardNodeRootArtifacts(input.events, input.guardNodeRootId);
const normalized = yield* normalizeEvents(
input.guardPendingUserInputCancellations === true
? yield* guardUserInputCancellations(input.events)
: input.events,
? yield* guardUserInputCancellations(guardedEvents)
: guardedEvents,
);
const committed = yield* eventStore.append({
...(input.commandId === undefined ? {} : { commandId: input.commandId }),
Expand Down Expand Up @@ -437,9 +586,18 @@ const baseLayer: Layer.Layer<
events: normalized,
});
yield* applyStoredEvents(storedEvents);
yield* effectOutbox.enqueue(input.effects ?? []);
return { committed: true as const, storedEvents };
}),
(result) => (result.committed ? publishStoredEvents(result.storedEvents) : Effect.void),
(result) =>
result.committed
? Effect.gen(function* () {
if (input.effects !== undefined && input.effects.length > 0) {
yield* effectOutbox.notifyAvailable(input.effects.length);
}
yield* publishStoredEvents(result.storedEvents);
})
: Effect.void,
);
},
);
Expand Down
Loading
Loading