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
69 changes: 50 additions & 19 deletions apps/server/src/orchestration-v2/EventSink.ts
Original file line number Diff line number Diff line change
Expand Up @@ -159,6 +159,8 @@ export interface EventSinkV2Shape {
readonly stream: (input?: {
readonly threadId?: ThreadId;
readonly afterSequence?: number;
/** Filter before queuing live events so a busy worker retains only the events it handles. */
readonly eventType?: OrchestrationV2DomainEvent["type"];
/** Bound RPC subscribers; internal workers must not drop their subscription under load. */
readonly bounded?: boolean;
}) => Stream.Stream<OrchestrationV2StoredEvent, EventSinkV2Error>;
Expand Down Expand Up @@ -196,6 +198,20 @@ const baseLayer: Layer.Layer<
const projectionStore = yield* ProjectionStoreV2;
const turnItemPositions = yield* TurnItemPositionStoreV2;
const liveEvents = yield* PubSub.unbounded<OrchestrationV2StoredEvent>();
const liveEventsByType = new Map<
OrchestrationV2DomainEvent["type"],
PubSub.PubSub<OrchestrationV2StoredEvent>
>();
const publishLiveEvents = (events: ReadonlyArray<OrchestrationV2StoredEvent>) =>
Effect.gen(function* () {
yield* PubSub.publishAll(liveEvents, events);
for (const [type, pubsub] of liveEventsByType) {
yield* PubSub.publishAll(
pubsub,
events.filter((stored) => stored.event.type === type),
);
}
});

// A user can answer after terminal normalization reads the pending request.
// Recheck inside the write transaction so stale cleanup cannot erase answers.
Expand Down Expand Up @@ -324,7 +340,7 @@ const baseLayer: Layer.Layer<
yield* effectOutbox.notifyAvailable(input.effects.length);
}
yield* eventStore.publishCommitted(storedEvents);
yield* PubSub.publishAll(liveEvents, storedEvents);
yield* publishLiveEvents(storedEvents);
return storedEvents;
});

Expand Down Expand Up @@ -378,7 +394,7 @@ const baseLayer: Layer.Layer<
);
if (result.committed) {
yield* eventStore.publishCommitted(result.storedEvents);
yield* PubSub.publishAll(liveEvents, result.storedEvents);
yield* publishLiveEvents(result.storedEvents);
}
return result;
},
Expand Down Expand Up @@ -439,7 +455,7 @@ const baseLayer: Layer.Layer<
);
if (result.committed) {
yield* eventStore.publishCommitted(result.storedEvents);
yield* PubSub.publishAll(liveEvents, result.storedEvents);
yield* publishLiveEvents(result.storedEvents);
}
return result;
});
Expand Down Expand Up @@ -517,7 +533,7 @@ const baseLayer: Layer.Layer<
}
if (result.committed) {
yield* eventStore.publishCommitted(result.storedEvents);
yield* PubSub.publishAll(liveEvents, result.storedEvents);
yield* publishLiveEvents(result.storedEvents);
}
return {
receipt: result.receipt,
Expand Down Expand Up @@ -556,6 +572,7 @@ const baseLayer: Layer.Layer<
readonly afterSequence: number;
readonly throughSequence: number;
readonly threadId?: ThreadId;
readonly eventType?: OrchestrationV2DomainEvent["type"];
}): Stream.Stream<OrchestrationV2StoredEvent, unknown> => {
const pageSize = 256;
const loop = (afterSequence: number): Stream.Stream<OrchestrationV2StoredEvent, unknown> =>
Expand All @@ -565,6 +582,7 @@ const baseLayer: Layer.Layer<
afterSequence,
throughSequence: input.throughSequence,
...(input.threadId === undefined ? {} : { threadId: input.threadId }),
...(input.eventType === undefined ? {} : { eventType: input.eventType }),
limit: pageSize,
})
.pipe(
Expand All @@ -587,31 +605,44 @@ const baseLayer: Layer.Layer<

const stream = (input?: Parameters<EventSinkV2Shape["stream"]>[0]) => {
const afterSequence = input?.afterSequence ?? 0;
const matchesThread = (stored: OrchestrationV2StoredEvent) =>
input?.threadId === undefined || stored.event.threadId === input.threadId;
const matches = (stored: OrchestrationV2StoredEvent) =>
(input?.threadId === undefined || stored.event.threadId === input.threadId) &&
(input?.eventType === undefined || stored.event.type === input.eventType);
const replay = (throughSequence: number) =>
catchUp({
afterSequence,
throughSequence,
...(input?.threadId === undefined ? {} : { threadId: input.threadId }),
});
if (input?.bounded === true) {
return replayAndBufferProjectedLiveEvents({
subscribe: PubSub.subscribe(liveEvents),
latestSequence: eventStore.latestSequence(),
afterSequence,
filter: matchesThread,
replay,
project: (stored) => ({ ...stored, event: projectDomainEventForWire(stored.event) }),
});
}
...(input?.eventType === undefined ? {} : { eventType: input.eventType }),
}).pipe(Stream.filter(matches));
return Stream.unwrap(
Effect.gen(function* () {
const subscription = yield* PubSub.subscribe(liveEvents);
let pubsub = liveEvents;
if (input?.eventType !== undefined) {
const existing = liveEventsByType.get(input.eventType);
if (existing !== undefined) {
pubsub = existing;
} else {
const created = yield* PubSub.unbounded<OrchestrationV2StoredEvent>();
pubsub = liveEventsByType.get(input.eventType) ?? created;
liveEventsByType.set(input.eventType, pubsub);
}
}
if (input?.bounded === true) {
return replayAndBufferProjectedLiveEvents({
subscribe: PubSub.subscribe(pubsub),
latestSequence: eventStore.latestSequence(),
afterSequence,
filter: matches,
replay,
project: (stored) => ({ ...stored, event: projectDomainEventForWire(stored.event) }),
});
}
const subscription = yield* PubSub.subscribe(pubsub);
const highWater = yield* eventStore.latestSequence();
const live = Stream.fromSubscription(subscription).pipe(
Stream.filter((stored) => stored.sequence > Math.max(highWater, afterSequence)),
Stream.filter(matchesThread),
Stream.filter(matches),
);
return Stream.concat(replay(highWater), live);
}),
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/orchestration-v2/EventStore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,7 @@ export interface EventStoreV2Shape {
readonly afterSequence?: number;
readonly throughSequence?: number;
readonly threadId?: ThreadId;
readonly eventType?: OrchestrationV2DomainEvent["type"];
readonly limit?: number;
}) => Stream.Stream<OrchestrationV2StoredEvent, EventStoreV2Error>;
readonly readByCommandId: (input: {
Expand Down Expand Up @@ -86,6 +87,7 @@ const baseLayer: Layer.Layer<EventStoreV2, never, OrchestrationEventStore> = Lay
? {}
: { throughSequence: input.throughSequence }),
...(input?.threadId === undefined ? {} : { threadId: input.threadId }),
...(input?.eventType === undefined ? {} : { eventType: input.eventType }),
...(input?.limit === undefined ? {} : { limit: input.limit }),
})
.pipe(
Expand Down
95 changes: 95 additions & 0 deletions apps/server/src/orchestration-v2/FoundationPersistence.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -550,6 +550,101 @@ it.layer(TestLayer)("orchestration V2 foundation persistence", (it) => {
),
);

it.effect("filters worker replay and live queues without losing matching events", () =>
Effect.scoped(
Effect.gen(function* () {
const sink = yield* EventSinkV2;
const sql = yield* SqlClient.SqlClient;
const now = yield* DateTime.now;
const thread = makeThread(ThreadId.make("thread:filtered-worker"), now);
const other = makeThread(ThreadId.make("thread:filtered-worker-other"), now);
const run: OrchestrationV2Run = {
id: RunId.make("run:filtered-worker"),
threadId: thread.id,
ordinal: 1,
providerInstanceId,
modelSelection,
providerThreadId: null,
userMessageId: MessageId.make("message:filtered-worker"),
rootNodeId: null,
activeAttemptId: null,
status: "completed",
queuePosition: null,
requestedAt: now,
startedAt: now,
completedAt: now,
checkpointId: null,
contextHandoffId: null,
};
const runEvent = (id: string, payload: OrchestrationV2Run): OrchestrationV2DomainEvent => ({
id: EventId.make(id),
type: "run.updated",
threadId: payload.threadId,
runId: payload.id,
occurredAt: now,
payload,
});
const history = yield* sink.write({
events: [
threadCreatedEvent({ id: "event:filtered-worker:created", thread, now }),
threadCreatedEvent({ id: "event:filtered-worker:other", thread: other, now }),
runEvent("event:filtered-worker:history", run),
runEvent("event:filtered-worker:other-history", {
...run,
id: RunId.make("run:filtered-worker-other"),
threadId: other.id,
}),
],
});
// Unrelated payloads must be skipped in SQL, before decoding or
// allocating their bodies, even when a retained row is unreadable.
const original = yield* sql<{ readonly payload_json: string }>`
SELECT payload_json FROM orchestration_events
WHERE sequence = ${history[0]!.sequence}
`;
yield* Effect.acquireRelease(
sql`
UPDATE orchestration_events SET payload_json = 'unreadable unrelated payload'
WHERE sequence = ${history[0]!.sequence}
`,
() =>
sql`
UPDATE orchestration_events SET payload_json = ${original[0]!.payload_json}
WHERE sequence = ${history[0]!.sequence}
`.pipe(Effect.orDie),
);
const pull = yield* Stream.toPull(
sink.stream({ threadId: thread.id, eventType: "run.updated" }),
);
assert.deepEqual(
(yield* pull).map((stored) => stored.sequence),
[history[2]!.sequence],
);
// The worker is occupied with the previous batch while the thread
// publishes output. Its live queue must receive just run updates.
yield* sink.write({
events: Array.from({ length: LIVE_STREAM_MAX_ITEMS + 1 }, (_, index) => ({
id: EventId.make(`event:filtered-worker:output:${index}`),
type: "thread.metadata-updated" as const,
threadId: thread.id,
occurredAt: now,
payload: { ...thread, title: `Output ${index}` },
})),
});
const live = yield* sink.write({
events: [
runEvent("event:filtered-worker:live:1", { ...run, status: "interrupted" }),
runEvent("event:filtered-worker:live:2", { ...run, status: "failed" }),
],
});
assert.deepEqual(
(yield* pull).map((stored) => stored.sequence),
live.map((stored) => stored.sequence),
);
}),
),
);

it.effect("paginates catch-up beyond the event-store read limit", () =>
Effect.gen(function* () {
const eventSink = yield* EventSinkV2;
Expand Down
32 changes: 18 additions & 14 deletions apps/server/src/orchestration-v2/Orchestrator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8549,20 +8549,24 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
// below. Replaying the full event table on every server start delays live
// queue promotion in proportion to the lifetime size of the database.
const terminalEventsAfterSequence = yield* eventSink.latestSequence().pipe(Effect.orDie);
yield* eventSink.stream({ afterSequence: terminalEventsAfterSequence }).pipe(
Stream.filter(
(stored) =>
stored.event.type === "run.updated" &&
!String(stored.commandId).startsWith("command:runtime-reconcile:") &&
(stored.event.payload.status === "completed" ||
stored.event.payload.status === "interrupted" ||
stored.event.payload.status === "failed" ||
stored.event.payload.status === "cancelled" ||
stored.event.payload.status === "rolled_back"),
),
Stream.runForEach(handleTerminalRun),
Effect.forkDetach,
);
// Queue promotion can wait on a provider or a thread lock. Subscribe to run
// updates before buffering so that wait never retains unrelated tool bodies.
yield* eventSink
.stream({ afterSequence: terminalEventsAfterSequence, eventType: "run.updated" })
.pipe(
Stream.filter(
(stored) =>
stored.event.type === "run.updated" &&
!String(stored.commandId).startsWith("command:runtime-reconcile:") &&
(stored.event.payload.status === "completed" ||
stored.event.payload.status === "interrupted" ||
stored.event.payload.status === "failed" ||
stored.event.payload.status === "cancelled" ||
stored.event.payload.status === "rolled_back"),
),
Stream.runForEach(handleTerminalRun),
Effect.forkDetach,
);

// The high-water subscription deliberately skips history, so recover the
// two terminal side effects from current projections instead: one queued
Expand Down
3 changes: 3 additions & 0 deletions apps/server/src/persistence/Layers/OrchestrationEventStore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -372,6 +372,7 @@ const makeEventStore = Effect.gen(function* () {
readonly throughSequence?: number;
readonly threadId?: ThreadId;
readonly commandId?: CommandId;
readonly eventType?: OrchestrationV2DomainEvent["type"];
readonly onlyAgentEvents?: boolean;
readonly limit: number;
}) =>
Expand Down Expand Up @@ -408,6 +409,7 @@ const makeEventStore = Effect.gen(function* () {
AND ${sql.and([
...(input.threadId === undefined ? [] : [sql`stream_id = ${input.threadId}`]),
...(input.commandId === undefined ? [] : [sql`command_id = ${input.commandId}`]),
...(input.eventType === undefined ? [] : [sql`event_type = ${input.eventType}`]),
])}
ORDER BY sequence ASC
LIMIT ${input.limit}
Expand Down Expand Up @@ -492,6 +494,7 @@ const makeEventStore = Effect.gen(function* () {
: { throughSequence: input.throughSequence }),
...(input?.threadId === undefined ? {} : { threadId: input.threadId }),
...(input?.commandId === undefined ? {} : { commandId: input.commandId }),
...(input?.eventType === undefined ? {} : { eventType: input.eventType }),
onlyAgentEvents: true,
limit: pageLimit,
}).pipe(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,7 @@ export interface OrchestrationEventStoreShape {
readonly throughSequence?: number;
readonly threadId?: ThreadId;
readonly commandId?: CommandId;
readonly eventType?: OrchestrationV2DomainEvent["type"];
readonly limit?: number;
}) => Stream.Stream<OrchestrationV2StoredEvent, OrchestrationEventStoreError>;

Expand Down
Loading