import {
CommandId,
type OrchestrationV2Run,
OrchestrationV2DomainEvent,
OrchestrationV2StoredEvent,
ProviderThreadId,
RunAttemptId,
RunId,
RuntimeRequestId,
NodeId,
type ProjectId,
ThreadId,
} from "@t3tools/contracts";
import * as Context from "effect/Context";
import * as DateTime from "effect/DateTime";
import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as PubSub from "effect/PubSub";
import * as Semaphore from "effect/Semaphore";
import * as Schema from "effect/Schema";
import * as Stream from "effect/Stream";
import * as SqlClient from "effect/sql/SqlClient";
import { replayAndBufferProjectedLiveEvents } from "./LiveStreamBudget.ts";
import type { UnsequencedProjectEvent } from "../persistence/Services/OrchestrationEventStore.ts";
import { projectDomainEventForWire } from "./WireProjection.ts";
import * as CommandReceiptStore from "./CommandReceiptStore.ts";
import * as EffectOutbox from "./EffectOutbox.ts";
import * as EventStore from "./EventStore.ts";
import * as ProjectionStore from "./ProjectionStore.ts";
import * as ProjectStore from "./ProjectStore.ts";
import * as TurnItemPositionStore from "./TurnItemPositionStore.ts";
/**
* ERRORS
*/
export class EventSinkWriteError extends Schema.TaggedError<EventSinkWriteError>()(
"EventSinkWriteError",
{
eventCount: Schema.Number,
commandId: Schema.optional(CommandId),
cause: Schema.optional(Schema.Defect()),
},
) {
override get message(): string {
return `Failed to write ${this.eventCount} orchestration V2 event(s).`;
}
}
export class EventSinkStreamError extends Schema.TaggedError<EventSinkStreamError>()(
"EventSinkStreamError",
{
threadId: Schema.optional(ThreadId),
afterSequence: Schema.optional(Schema.Number),
cause: Schema.optional(Schema.Defect()),
},
) {
override get message(): string {
return this.threadId === undefined
? "Failed to stream orchestration V2 events."
: `Failed to stream orchestration V2 events for thread ${this.threadId}.`;
}
}
export const EventSinkV2Error = Schema.Union([EventSinkWriteError, EventSinkStreamError]);
export type EventSinkV2Error = typeof EventSinkV2Error.Type;
/**
* SERVICE DEFINITION
*/
export interface EventSinkV2Shape {
readonly write: (input: {
readonly guardPendingUserInputCancellations?: boolean;
readonly commandId?: CommandId;
readonly events: ReadonlyArray<OrchestrationV2DomainEvent>;
}) => Effect.Effect<ReadonlyArray<OrchestrationV2StoredEvent>, EventSinkV2Error>;
readonly writeWithEffects: (input: {
readonly guardPendingUserInputCancellations?: boolean;
readonly commandId?: CommandId;
readonly events: ReadonlyArray<OrchestrationV2DomainEvent>;
readonly effects: ReadonlyArray<EffectOutbox.PendingOrchestrationEffectV2>;
}) => Effect.Effect<ReadonlyArray<OrchestrationV2StoredEvent>, EventSinkV2Error>;
readonly writeIfRunCurrent: (input: {
readonly guardPendingUserInputCancellations?: boolean;
readonly commandId?: CommandId;
readonly threadId: ThreadId;
readonly runId: RunId;
readonly activeAttemptId: RunAttemptId;
readonly expectedStatus: OrchestrationV2Run["status"];
readonly events: ReadonlyArray<OrchestrationV2DomainEvent>;
}) => Effect.Effect<
{
readonly committed: boolean;
readonly storedEvents: ReadonlyArray<OrchestrationV2StoredEvent>;
},
EventSinkV2Error
>;
/**
* Atomically commit only when the provider thread is still owned by the
* expected run attempt and ordinal. Used for late post-terminal
* provider_thread updates so a completed or superseded attempt cannot clobber
* a newer attempt that already claimed the thread.
*/
readonly writeIfProviderThreadOwner: (input: {
readonly guardPendingUserInputCancellations?: boolean;
readonly commandId?: CommandId;
readonly providerThreadId: ProviderThreadId;
readonly runId: RunId;
readonly activeAttemptId: RunAttemptId;
readonly expectedLastRunOrdinal: number;
readonly events: ReadonlyArray<OrchestrationV2DomainEvent>;
}) => Effect.Effect<
{
readonly committed: boolean;
readonly storedEvents: ReadonlyArray<OrchestrationV2StoredEvent>;
},
EventSinkV2Error
>;
readonly commitCommand: (input: {
readonly commandId: CommandId;
readonly threadId: ThreadId;
readonly commandType: string;
readonly acceptedAt: DateTime.Utc;
readonly events: ReadonlyArray<OrchestrationV2DomainEvent>;
readonly effects: ReadonlyArray<EffectOutbox.PendingOrchestrationEffectV2>;
readonly cancelUnsettledEffects?: {
readonly effectTypes: ReadonlyArray<EffectOutbox.OrchestrationEffectRequestV2["type"]>;
readonly reason: string;
};
}) => Effect.Effect<
{
readonly receipt: CommandReceiptStore.CommandReceiptV2;
readonly storedEvents: ReadonlyArray<OrchestrationV2StoredEvent>;
readonly committed: boolean;
readonly cancelledEffectCount: number;
},
EventSinkV2Error
>;
readonly commitRejectedCommand: (input: {
readonly commandId: CommandId;
readonly threadId: ThreadId;
readonly commandType: string;
readonly rejectedAt: DateTime.Utc;
readonly error: string;
}) => Effect.Effect<CommandReceiptStore.CommandReceiptV2, EventSinkV2Error>;
/**
* Append a project event, fold it into its row and record the receipt in one
* transaction. A reused command id commits nothing and returns its receipt.
*/
readonly commitProjectCommand: (input: {
readonly commandId: CommandId;
readonly projectId: ProjectId;
readonly commandType: string;
readonly acceptedAt: DateTime.Utc;
readonly event: UnsequencedProjectEvent;
}) => Effect.Effect<
{ readonly receipt: CommandReceiptStore.ProjectCommandReceiptV2; readonly committed: boolean },
EventSinkV2Error
>;
/** Record a rejected project command, or return the receipt its command id already has. */
readonly commitRejectedProjectCommand: (input: {
readonly commandId: CommandId;
readonly projectId: ProjectId;
readonly commandType: string;
readonly rejectedAt: DateTime.Utc;
readonly error: string;
}) => Effect.Effect<CommandReceiptStore.ProjectCommandReceiptV2, EventSinkV2Error>;
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>;
readonly latestSequence: (input?: {
readonly threadId?: ThreadId;
}) => Effect.Effect<number, EventSinkV2Error>;
readonly readByCommandId: (input: {
readonly commandId: CommandId;
}) => Stream.Stream<OrchestrationV2StoredEvent, EventSinkV2Error>;
}
export class EventSinkV2 extends Context.Service<EventSinkV2, EventSinkV2Shape>()(
"t3/orchestration-v2/EventSink/EventSinkV2",
) {}
/**
* IMPLEMENTATIONS
*/
const layerBase: Layer.Layer<
EventSinkV2,
never,
| CommandReceiptStore.CommandReceiptStoreV2
| EffectOutbox.EffectOutboxV2
| EventStore.EventStoreV2
| ProjectionStore.ProjectionStoreV2
| ProjectStore.ProjectStoreV2
| SqlClient.SqlClient
| TurnItemPositionStore.TurnItemPositionStoreV2
> = Layer.effect(
EventSinkV2,
Effect.gen(function* () {
const sql = yield* SqlClient.SqlClient;
const commandReceipts = yield* CommandReceiptStore.CommandReceiptStoreV2;
const effectOutbox = yield* EffectOutbox.EffectOutboxV2;
const eventStore = yield* EventStore.EventStoreV2;
const projectionStore = yield* ProjectionStore.ProjectionStoreV2;
const projectStore = yield* ProjectStore.ProjectStoreV2;
const turnItemPositions = yield* TurnItemPositionStore.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),
);
}
});
const publishStoredEvents = (events: ReadonlyArray<OrchestrationV2StoredEvent>) =>
eventStore.publishCommitted(events).pipe(Effect.andThen(publishLiveEvents(events)));
// Transactions commit one at a time, but each writer publishes after its
// commit. If a writer is descheduled in between, a later commit reaches
// subscribers first, and clients drop any event at or below the newest
// sequence they have applied. So a writer takes this lane as the last step
// of its transaction and holds it until it has published. Publishing never
// waits, so a writer that holds the transaction while it waits for the
// lane is not blocked for long.
const publishLane = yield* Semaphore.make(1);
const commitThenPublish = <A, E, R>(
transaction: Effect.Effect<A, E, R>,
publish: (committed: A) => Effect.Effect<void>,
) =>
Effect.suspend(() => {
let holdsLane = false;
const takeLane = publishLane.take(1).pipe(
Effect.andThen(
Effect.sync(() => {
holdsLane = true;
}),
),
Effect.uninterruptible,
);
return sql
.withTransaction(Effect.tap(transaction, () => takeLane))
.pipe(
Effect.tap(publish),
Effect.ensuring(
Effect.suspend(() => (holdsLane ? publishLane.release(1) : Effect.void)),
),
);
});
// 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>) =>
Effect.gen(function* () {
const staleRequests = new Set<RuntimeRequestId>();
const staleNodes = new Set<NodeId>();
for (const event of events) {
if (
event.type !== "runtime-request.updated" ||
event.payload.kind !== "user_input" ||
event.payload.status !== "cancelled"
)
continue;
const current = yield* projectionStore.getRuntimeRequest(
event.threadId,
event.payload.id,
);
if (
current?.status !== "pending" ||
current.kind !== "user_input" ||
current.providerTurnId !== event.payload.providerTurnId ||
current.responseCapability.type === "message"
) {
staleRequests.add(event.payload.id);
staleNodes.add(event.payload.nodeId);
}
}
return events.filter((event) => {
switch (event.type) {
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);
case "turn-item.updated":
return (
event.payload.type !== "user_input_request" ||
event.payload.status !== "cancelled" ||
!staleRequests.has(event.payload.requestId)
);
default:
return true;
}
});
});
const normalizeEvents = (events: ReadonlyArray<OrchestrationV2DomainEvent>) => {
const runOrdinals = new Map(
events.flatMap((event) =>
event.type === "run.created" || event.type === "run.updated"
? [[event.payload.id, event.payload.ordinal] as const]
: [],
),
);
return Effect.forEach(
events,
(event): Effect.Effect<OrchestrationV2DomainEvent, unknown> =>
event.type === "turn-item.updated"
? turnItemPositions
.normalize(
event.payload,
event.payload.runId === null ? undefined : runOrdinals.get(event.payload.runId),
)
.pipe(Effect.map((payload) => ({ ...event, payload })))
: Effect.succeed(event),
{ concurrency: 1 },
);
};
const applyStoredEvents = (storedEvents: ReadonlyArray<OrchestrationV2StoredEvent>) =>
Effect.gen(function* () {
yield* Effect.forEach(storedEvents, (stored) => projectionStore.apply(stored.event), {
concurrency: 1,
});
const sequence = storedEvents.at(-1)?.sequence;
if (sequence !== undefined) {
const now = DateTime.formatIso(yield* DateTime.now);
yield* sql`
INSERT INTO orchestration_v2_projection_metadata (
projection_name,
schema_version,
last_sequence,
updated_at
)
VALUES (
'thread-projections',
${ProjectionStore.ORCHESTRATION_V2_PROJECTION_SCHEMA_VERSION},
${sequence},
${now}
)
ON CONFLICT(projection_name)
DO UPDATE SET
schema_version = excluded.schema_version,
last_sequence = excluded.last_sequence,
updated_at = excluded.updated_at
`;
}
});
const writeEffect = Effect.fn("orchestrationV2.EventSink.write")(function* (
input: Parameters<EventSinkV2Shape["writeWithEffects"]>[0],
) {
yield* Effect.annotateCurrentSpan({
"orchestration_v2.command_id": input.commandId ?? null,
"orchestration_v2.event_count": input.events.length,
"orchestration_v2.thread_id": input.events[0]?.threadId ?? null,
});
return yield* commitThenPublish(
Effect.gen(function* () {
const normalized = yield* normalizeEvents(
input.guardPendingUserInputCancellations === true
? yield* guardUserInputCancellations(input.events)
: input.events,
);
const committed = yield* eventStore.append({
...(input.commandId === undefined ? {} : { commandId: input.commandId }),
events: normalized,
});
yield* applyStoredEvents(committed);
yield* effectOutbox.enqueue(input.effects);
return committed;
}),
(storedEvents) =>
Effect.gen(function* () {
if (input.effects.length > 0) {
yield* effectOutbox.notifyAvailable(input.effects.length);
}
yield* publishStoredEvents(storedEvents);
}),
);
});
const writeIfRunCurrentEffect = Effect.fn("orchestrationV2.EventSink.writeIfRunCurrent")(
function* (input: Parameters<EventSinkV2Shape["writeIfRunCurrent"]>[0]) {
yield* Effect.annotateCurrentSpan({
"orchestration_v2.command_id": input.commandId ?? null,
"orchestration_v2.event_count": input.events.length,
"orchestration_v2.run_id": input.runId,
"orchestration_v2.thread_id": input.threadId,
});
return yield* commitThenPublish(
Effect.gen(function* () {
const rows = yield* sql<{
readonly status: string;
readonly active_attempt_id: string | null;
}>`
SELECT
status,
json_extract(payload_json, '$.activeAttemptId') AS active_attempt_id
FROM orchestration_v2_projection_runs
WHERE run_id = ${input.runId}
AND thread_id = ${input.threadId}
LIMIT 1
`;
const current = rows[0];
if (
current === undefined ||
current.status !== input.expectedStatus ||
current.active_attempt_id !== input.activeAttemptId
) {
return {
committed: false as const,
storedEvents: [] as ReadonlyArray<OrchestrationV2StoredEvent>,
};
}
const normalized = yield* normalizeEvents(
input.guardPendingUserInputCancellations === true
? yield* guardUserInputCancellations(input.events)
: input.events,
);
const storedEvents = yield* eventStore.append({
...(input.commandId === undefined ? {} : { commandId: input.commandId }),
events: normalized,
});
yield* applyStoredEvents(storedEvents);
return { committed: true as const, storedEvents };
}),
(result) => (result.committed ? publishStoredEvents(result.storedEvents) : Effect.void),
);
},
);
const writeIfProviderThreadOwnerEffect = Effect.fn(
"orchestrationV2.EventSink.writeIfProviderThreadOwner",
)(function* (input: Parameters<EventSinkV2Shape["writeIfProviderThreadOwner"]>[0]) {
yield* Effect.annotateCurrentSpan({
"orchestration_v2.command_id": input.commandId ?? null,
"orchestration_v2.event_count": input.events.length,
"orchestration_v2.provider_thread_id": input.providerThreadId,
"orchestration_v2.run_id": input.runId,
"orchestration_v2.active_attempt_id": input.activeAttemptId,
"orchestration_v2.expected_last_run_ordinal": input.expectedLastRunOrdinal,
});
return yield* commitThenPublish(
Effect.gen(function* () {
const rows = yield* sql<{
readonly active_attempt_id: string | null;
readonly last_run_ordinal: number | null;
}>`
SELECT
json_extract(r.payload_json, '$.activeAttemptId') AS active_attempt_id,
p.last_run_ordinal
FROM orchestration_v2_projection_provider_threads p
JOIN orchestration_v2_projection_runs r
ON r.run_id = ${input.runId}
AND r.thread_id = p.thread_id
WHERE p.provider_thread_id = ${input.providerThreadId}
LIMIT 1
`;
const current = rows[0];
if (
current === undefined ||
current.active_attempt_id !== input.activeAttemptId ||
current.last_run_ordinal !== input.expectedLastRunOrdinal
) {
return {
committed: false as const,
storedEvents: [] as ReadonlyArray<OrchestrationV2StoredEvent>,
};
}
const normalized = yield* normalizeEvents(
input.guardPendingUserInputCancellations === true
? yield* guardUserInputCancellations(input.events)
: input.events,
);
const storedEvents = yield* eventStore.append({
...(input.commandId === undefined ? {} : { commandId: input.commandId }),
events: normalized,
});
yield* applyStoredEvents(storedEvents);
return { committed: true as const, storedEvents };
}),
(result) => (result.committed ? publishStoredEvents(result.storedEvents) : Effect.void),
);
});
const existingCommandResult = (commandId: CommandId) =>
Effect.gen(function* () {
const existing = yield* commandReceipts.getByCommandId(commandId);
if (Option.isNone(existing)) {
return yield* Effect.die(
new Error(`Command receipt ${commandId} disappeared during its transaction.`),
);
}
const storedEvents = yield* eventStore.readByCommandId({ commandId }).pipe(
Stream.runCollect,
Effect.map((events): ReadonlyArray<OrchestrationV2StoredEvent> => Array.from(events)),
);
return { receipt: existing.value, storedEvents };
});
const commitCommandEffect = Effect.fn("orchestrationV2.EventSink.commitCommand")(function* (
input: Parameters<EventSinkV2Shape["commitCommand"]>[0],
) {
const result = yield* commitThenPublish(
Effect.gen(function* () {
const reserved = yield* commandReceipts.insertIfAbsent({
commandId: input.commandId,
threadId: input.threadId,
commandType: input.commandType,
acceptedAt: input.acceptedAt,
resultSequence: 0,
status: "accepted",
error: null,
});
if (!reserved) {
const existing = yield* existingCommandResult(input.commandId);
return { ...existing, committed: false as const, cancelledEffectIds: [] };
}
const normalized = yield* normalizeEvents(input.events);
const storedEvents = yield* eventStore.append({
commandId: input.commandId,
events: normalized,
});
const sequence = storedEvents.at(-1)?.sequence;
if (sequence === undefined) {
return yield* Effect.die(
new Error(`Command ${input.commandId} produced no orchestration events.`),
);
}
yield* applyStoredEvents(storedEvents);
yield* effectOutbox.enqueue(input.effects);
const receipt: CommandReceiptStore.CommandReceiptV2 = {
commandId: input.commandId,
threadId: input.threadId,
commandType: input.commandType,
acceptedAt: input.acceptedAt,
resultSequence: sequence,
status: "accepted",
error: null,
};
yield* commandReceipts.upsert(receipt);
const cancelledEffectIds =
input.cancelUnsettledEffects === undefined
? []
: yield* effectOutbox.cancelUnsettled({
threadId: input.threadId,
...input.cancelUnsettledEffects,
});
return { receipt, storedEvents, committed: true as const, cancelledEffectIds };
}),
(result) =>
Effect.gen(function* () {
yield* effectOutbox.signalCancellations(result.cancelledEffectIds);
if (result.committed && input.effects.length > 0) {
yield* effectOutbox.notifyAvailable(input.effects.length);
}
if (result.committed) yield* publishStoredEvents(result.storedEvents);
}),
);
return {
receipt: result.receipt,
storedEvents: result.storedEvents,
committed: result.committed,
cancelledEffectCount: result.cancelledEffectIds.length,
};
});
const commitRejectedCommandEffect = Effect.fn(
"orchestrationV2.EventSink.commitRejectedCommand",
)(function* (input: Parameters<EventSinkV2Shape["commitRejectedCommand"]>[0]) {
return yield* sql.withTransaction(
Effect.gen(function* () {
const sequence = yield* eventStore.latestSequence({ threadId: input.threadId });
const receipt: CommandReceiptStore.CommandReceiptV2 = {
commandId: input.commandId,
threadId: input.threadId,
commandType: input.commandType,
acceptedAt: input.rejectedAt,
resultSequence: sequence,
status: "rejected",
error: input.error,
};
const inserted = yield* commandReceipts.insertIfAbsent(receipt);
if (inserted) {
return receipt;
}
const existing = yield* commandReceipts.getByCommandId(input.commandId);
return Option.getOrElse(existing, () => receipt);
}),
);
});
const existingProjectReceipt = (commandId: CommandId) =>
commandReceipts.getProjectByCommandId(commandId).pipe(
Effect.flatMap(
Option.match({
onNone: () =>
Effect.fail(`Command ${commandId} was already used by a thread command.` as const),
onSome: Effect.succeed,
}),
),
);
const commitProjectCommandEffect = Effect.fn("orchestrationV2.EventSink.commitProjectCommand")(
function* (input: Parameters<EventSinkV2Shape["commitProjectCommand"]>[0]) {
const result = yield* commitThenPublish(
Effect.gen(function* () {
const reserved: CommandReceiptStore.ProjectCommandReceiptV2 = {
commandId: input.commandId,
projectId: input.projectId,
commandType: input.commandType,
acceptedAt: input.acceptedAt,
resultSequence: 0,
status: "accepted",
error: null,
};
if (!(yield* commandReceipts.insertIfAbsent(reserved))) {
return { receipt: yield* existingProjectReceipt(input.commandId), event: undefined };
}
const event = yield* eventStore.appendProjectEvent(input.event);
yield* projectStore.apply(event);
const receipt = { ...reserved, resultSequence: event.sequence };
yield* commandReceipts.upsert(receipt);
return { receipt, event };
}),
(result) =>
result.event === undefined ? Effect.void : eventStore.publishCommitted([result.event]),
);
return { receipt: result.receipt, committed: result.event !== undefined };
},
);
const commitRejectedProjectCommandEffect = Effect.fn(
"orchestrationV2.EventSink.commitRejectedProjectCommand",
)(function* (input: Parameters<EventSinkV2Shape["commitRejectedProjectCommand"]>[0]) {
return yield* sql.withTransaction(
Effect.gen(function* () {
const receipt: CommandReceiptStore.ProjectCommandReceiptV2 = {
commandId: input.commandId,
projectId: input.projectId,
commandType: input.commandType,
acceptedAt: input.rejectedAt,
resultSequence: yield* eventStore.latestApplicationSequence,
status: "rejected",
error: input.error,
};
return (yield* commandReceipts.insertIfAbsent(receipt))
? receipt
: yield* existingProjectReceipt(input.commandId);
}),
);
});
const catchUp = (input: {
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> =>
Stream.unwrap(
eventStore
.read({
afterSequence,
throughSequence: input.throughSequence,
...(input.threadId === undefined ? {} : { threadId: input.threadId }),
...(input.eventType === undefined ? {} : { eventType: input.eventType }),
limit: pageSize,
})
.pipe(
Stream.runCollect,
Effect.map((chunk) => Array.from(chunk)),
Effect.map((events) => {
if (events.length === 0) {
return Stream.empty;
}
const current = Stream.fromIterable(events);
const last = events.at(-1)?.sequence ?? input.throughSequence;
return events.length < pageSize || last >= input.throughSequence
? current
: Stream.concat(current, loop(last));
}),
),
);
return loop(input.afterSequence);
};
const stream = (input?: Parameters<EventSinkV2Shape["stream"]>[0]) => {
const afterSequence = input?.afterSequence ?? 0;
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 }),
...(input?.eventType === undefined ? {} : { eventType: input.eventType }),
}).pipe(Stream.filter(matches));
return Stream.unwrap(
Effect.gen(function* () {
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(matches),
);
return Stream.concat(replay(highWater), live);
}),
);
};
return EventSinkV2.of({
write: (input) =>
writeEffect({ ...input, effects: [] }).pipe(
Effect.mapError(
(cause) =>
new EventSinkWriteError({
eventCount: input.events.length,
...(input.commandId === undefined ? {} : { commandId: input.commandId }),
cause,
}),
),
),
writeWithEffects: (input) =>
writeEffect(input).pipe(
Effect.mapError(
(cause) =>
new EventSinkWriteError({
eventCount: input.events.length,
...(input.commandId === undefined ? {} : { commandId: input.commandId }),
cause,
}),
),
),
writeIfRunCurrent: (input) =>
writeIfRunCurrentEffect(input).pipe(
Effect.mapError(
(cause) =>
new EventSinkWriteError({
eventCount: input.events.length,
...(input.commandId === undefined ? {} : { commandId: input.commandId }),
cause,
}),
),
),
writeIfProviderThreadOwner: (input) =>
writeIfProviderThreadOwnerEffect(input).pipe(
Effect.mapError(
(cause) =>
new EventSinkWriteError({
eventCount: input.events.length,
...(input.commandId === undefined ? {} : { commandId: input.commandId }),
cause,
}),
),
),
commitCommand: (input) =>
commitCommandEffect(input).pipe(
Effect.mapError(
(cause) =>
new EventSinkWriteError({
commandId: input.commandId,
eventCount: input.events.length,
cause,
}),
),
),
commitRejectedCommand: (input) =>
commitRejectedCommandEffect(input).pipe(
Effect.mapError(
(cause) =>
new EventSinkWriteError({
commandId: input.commandId,
eventCount: 0,
cause,
}),
),
),
commitProjectCommand: (input) =>
commitProjectCommandEffect(input).pipe(
Effect.mapError(
(cause) =>
new EventSinkWriteError({ commandId: input.commandId, eventCount: 1, cause }),
),
),
commitRejectedProjectCommand: (input) =>
commitRejectedProjectCommandEffect(input).pipe(
Effect.mapError(
(cause) =>
new EventSinkWriteError({ commandId: input.commandId, eventCount: 0, cause }),
),
),
stream: (input) =>
stream(input).pipe(
Stream.mapError(
(cause) =>
new EventSinkStreamError({
...(input?.threadId === undefined ? {} : { threadId: input.threadId }),
...(input?.afterSequence === undefined
? {}
: { afterSequence: input.afterSequence }),
cause,
}),
),
),
latestSequence: (input) =>
eventStore.latestSequence(input).pipe(
Effect.mapError(
(cause) =>
new EventSinkStreamError({
...(input?.threadId === undefined ? {} : { threadId: input.threadId }),
cause,
}),
),
),
readByCommandId: (input) =>
eventStore.readByCommandId(input).pipe(
Stream.mapError(
(cause) =>
new EventSinkStreamError({
cause,
}),
),
),
} satisfies EventSinkV2Shape);
}),
);
/**
* Event sink layer for application compositions that already own the
* persistence services. Keeping the outbox instance shared with the worker is
* important because enqueue notifications are in-memory wakeups backed by the
* durable SQL queue.
*/
export const layerFromStores = layerBase;
export const layer: Layer.Layer<
EventSinkV2,
never,
EventStore.EventStoreV2 | ProjectionStore.ProjectionStoreV2 | SqlClient.SqlClient
> = layerBase.pipe(
Layer.provide(
Layer.mergeAll(
CommandReceiptStore.layer,
EffectOutbox.layer,
ProjectStore.layer,
TurnItemPositionStore.layer,
),
),
);