Skip to content
61 changes: 61 additions & 0 deletions apps/server/src/orchestration-v2/LiveStreamBudget.test.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import { it } from "@effect/vitest";
import type * as Cause from "effect/Cause";
import * as Deferred from "effect/Deferred";
import * as Effect from "effect/Effect";
import * as Fiber from "effect/Fiber";
Expand All @@ -9,6 +10,7 @@ import { describe, expect } from "vite-plus/test";

import {
bufferLiveStream,
bufferLatestLiveStream,
makeLiveStreamBudget,
replayAndBufferLiveEvents,
type RetainedLiveItem,
Expand Down Expand Up @@ -77,6 +79,65 @@ it.effect("stops draining a slow subscriber when its unacknowledged tail fills",

type Event = { readonly sequence: number; readonly text: string; readonly threadId?: string };

describe("bufferLatestLiveStream", () => {
it.effect(
"replaces queued payloads across batches without replacing the unacknowledged batch",
() =>
Effect.scoped(
Effect.gen(function* () {
const input = yield* Queue.unbounded<Event, Cause.Done>();
const drained = yield* Deferred.make<void>();
const first = { sequence: 1, threadId: "a", text: "x".repeat(1000) };
const pull = yield* Stream.toPull(
bufferLatestLiveStream(
Stream.fromQueue(input).pipe(Stream.ensuring(Deferred.succeed(drained, undefined))),
(event) => event.threadId!,
{ maxItems: 3, maxSerializedBytes: 3300 },
),
);
yield* Queue.offer(input, first);
expect(yield* pull).toEqual([first]);
const updates = Array.from({ length: 300 }, (_, index) => ({
sequence: index + 2,
threadId: index % 2 === 0 ? "a" : "b",
text: "x".repeat(1000),
}));
yield* Queue.offerAll(input, updates);
yield* Queue.end(input);
// No ACK/pull while all 300 updates pass through the producer.
yield* Deferred.await(drained);
expect(yield* pull).toEqual(updates.slice(-2));
}),
),
);

it.effect.each(["items", "bytes"] as const)(
"keeps an in-flight update charged to the %s budget when the same key updates",
(limit) =>
Effect.scoped(
Effect.gen(function* () {
const input = yield* Queue.unbounded<Event>();
const closed = yield* Deferred.make<void>();
const first = { sequence: 1, text: "first" };
const pull = yield* Stream.toPull(
bufferLatestLiveStream(
Stream.fromQueue(input).pipe(Stream.ensuring(Deferred.succeed(closed, undefined))),
() => "same-key",
limit === "items" ? { maxItems: 1 } : { maxSerializedBytes: 32 },
),
);
yield* Queue.offer(input, first);
expect(yield* pull).toEqual([first]);
yield* Queue.offer(input, { sequence: 2, text: "next" });
yield* Deferred.await(closed);
const result = yield* pull.pipe(Effect.result);
expect(result._tag).toBe("Failure");
if (result._tag === "Failure") expect(result.failure._tag).toBe("LiveStreamBufferError");
}),
),
);
});

describe("replayAndBufferLiveEvents", () => {
it.effect.each(["high-water", "replay"] as const)(
"unsubscribes and cancels a blocked %s read on live overflow",
Expand Down
79 changes: 79 additions & 0 deletions apps/server/src/orchestration-v2/LiveStreamBudget.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import * as Exit from "effect/Exit";
import * as Queue from "effect/Queue";
import * as PubSub from "effect/PubSub";
import * as Scope from "effect/Scope";
import * as Semaphore from "effect/Semaphore";
import * as Stream from "effect/Stream";

export class LiveStreamBufferError extends Schema.TaggedError<LiveStreamBufferError>()(
Expand Down Expand Up @@ -262,6 +263,84 @@ export const bufferLiveStream = <A extends object, E, R>(
}),
);

/** For ordered full-state sources, keep only the latest queued state per aggregate while awaiting ACK. */
export const bufferLatestLiveStream = <A extends { readonly sequence: number }, E, R>(
source: Stream.Stream<A, E, R>,
key: (value: A) => string,
limits?: LiveStreamLimits,
) =>
Stream.unwrap(
Effect.gen(function* () {
const budget = yield* makeLiveStreamBudget(limits);
const ready = yield* Queue.unbounded<void, E | LiveStreamBufferError | Cause.Done>();
const mutex = yield* Semaphore.make(1);
const pending = new Map<string, RetainedLiveItem<A>>();
let closed = false;
const close = (error?: LiveStreamBufferError) =>
mutex.withPermits(1)(
Effect.gen(function* () {
if (closed) return;
closed = true;
budget.release(pending.values());
pending.clear();
if (error) yield* Queue.fail(ready, error);
yield* Queue.shutdown(ready);
}),
);
yield* Effect.addFinalizer(() => close());
yield* budget.failed.pipe(
Effect.catchTags({ LiveStreamBufferError: close }),
Effect.forkScoped,
);
yield* source.pipe(
Stream.runForEach((value) =>
mutex.withPermits(1)(
Effect.gen(function* () {
yield* budget.check;
if (closed) return;
const identity = key(value);
const previous = pending.get(identity);
if (previous && previous.value.sequence >= value.sequence) return;
const wasEmpty = pending.size === 0;
const [item] = yield* budget.replace(previous ? [previous] : [], [value]);
pending.set(identity, item!);
// One wakeup per pending batch, not per update. Replacements
// cannot leave an unbounded queue of obsolete keys behind.
if (wasEmpty) yield* Queue.offer(ready, undefined);
}).pipe(Effect.uninterruptible),
),
),
Effect.raceFirst(budget.failed),
Effect.exit,
Effect.flatMap((exit) =>
Exit.isFailure(exit) ? Queue.failCause(ready, exit.cause) : Queue.end(ready),
),
Effect.forkScoped({ startImmediately: true }),
);
return budget.deliver(
Stream.fromPull(
Effect.succeed(
Queue.take(ready).pipe(
Effect.andThen(
mutex.withPermits(1)(
Effect.sync(() => {
const items = Array.from(pending.values()).sort(
(left, right) => left.value.sequence - right.value.sequence,
);
pending.clear();
// The wakeup exists only for a nonempty pending batch.
// Delivery now owns these items until the next ACK.
return items as Arr.NonEmptyArray<RetainedLiveItem<A>>;
}),
),
),
),
),
),
);
}),
);

/**
* Subscribe and start draining before reading the high-water mark. A blocked
* catch-up query or an unacknowledged replay batch must not strand an unbounded
Expand Down
61 changes: 61 additions & 0 deletions apps/server/src/orchestration-v2/ShellStream.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,17 +8,21 @@ import type {
} from "@t3tools/contracts";
import { ProjectId, ProviderInstanceId, ThreadId } from "@t3tools/contracts";
import * as NodeSqliteClient from "@t3tools/shared/nodeSqliteClient";
import type * as Cause from "effect/Cause";
import * as DateTime from "effect/DateTime";
import * as Deferred from "effect/Deferred";
import * as Duration from "effect/Duration";
import * as Effect from "effect/Effect";
import * as Option from "effect/Option";
import * as Queue from "effect/Queue";
import * as SqlClient from "effect/sql/SqlClient";
import * as Stream from "effect/Stream";
import * as TestClock from "effect/testing/TestClock";

import {
archivedShellStreamItemFromThreadShell,
buildActiveShellSnapshot,
bufferShellLiveStream,
coalesceShellApplicationEvents,
coalesceStoredThreadEvents,
composeShellStreamWithEnrichment,
Expand Down Expand Up @@ -54,6 +58,63 @@ const emptyShellSnapshot = {
archivedThreads: [],
} as OrchestrationV2ShellSnapshot;

it.effect("buffers the latest shell state in sequence order across deletion and recreation", () =>
Effect.scoped(
Effect.gen(function* () {
const input = yield* Queue.unbounded<
Extract<OrchestrationV2ShellStreamItem, { readonly sequence: number }>,
Cause.Done
>();
const drained = yield* Deferred.make<void>();
const pull = yield* Stream.toPull(
bufferShellLiveStream(
Stream.fromQueue(input).pipe(Stream.ensuring(Deferred.succeed(drained, undefined))),
{ maxItems: 5 },
),
);
const initial = {
kind: "project.removed" as const,
sequence: 1,
projectId: ProjectId.make("initial"),
};
yield* Queue.offer(input, initial);
expect(yield* pull).toEqual([initial]);
const shell = shellFixture({ id: ThreadId.make("same") });
const updates = [
{
kind: "thread.updated" as const,
sequence: 2,
location: "active" as const,
thread: shell,
},
{
kind: "thread.removed" as const,
sequence: 3,
location: "active" as const,
threadId: shell.id,
},
{ kind: "project.removed" as const, sequence: 4, projectId: ProjectId.make("same") },
{
kind: "thread.updated" as const,
sequence: 5,
location: "active" as const,
thread: shell,
},
{
kind: "thread.removed" as const,
sequence: 6,
location: "archive" as const,
threadId: shell.id,
},
];
yield* Queue.offerAll(input, updates);
yield* Queue.end(input);
yield* Deferred.await(drained);
expect(yield* pull).toEqual(updates.slice(2));
}),
),
);

describe("buildActiveShellSnapshot", () => {
it("never duplicates archived rows into the regular shell", () => {
const active = shellFixture({ archivedAt: null });
Expand Down
31 changes: 30 additions & 1 deletion apps/server/src/orchestration-v2/ShellStream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,32 @@ import * as Schema from "effect/Schema";
import type * as SqlClient from "effect/sql/SqlClient";
import * as Stream from "effect/Stream";

import { bufferLatestLiveStream } from "./LiveStreamBudget.ts";

type ShellDelta = Extract<OrchestrationV2ShellStreamItem, { readonly sequence: number }>;

/** Full shell deltas replace earlier queued state; snapshots and sync markers stay outside this buffer. */
export const bufferShellLiveStream = <E, R>(
source: Stream.Stream<ShellDelta, E, R>,
limits?: { readonly maxItems?: number; readonly maxSerializedBytes?: number },
) =>
bufferLatestLiveStream(
source,
(item) => {
switch (item.kind) {
case "project.updated":
return `project:${item.project.id}`;
case "project.removed":
return `project:${item.projectId}`;
case "thread.updated":
return `thread:${item.location}:${item.thread.id}`;
case "thread.removed":
return `thread:${item.location}:${item.threadId}`;
}
},
limits,
);

/** Build the regular navigation shell without duplicating the archive dataset. */
export function buildActiveShellSnapshot(input: {
readonly projects: ReadonlyArray<OrchestrationProjectShell>;
Expand Down Expand Up @@ -292,7 +318,10 @@ export function coalesceStoredThreadEvents(
export function shellStreamItemFromThreadShell(input: {
readonly stored: Extract<ShellApplicationEvent, { readonly event: unknown }>;
readonly shell: OrchestrationV2ThreadShell | null;
}): Exclude<OrchestrationV2ShellStreamItem, { readonly kind: "snapshot" }> {
}): Extract<
OrchestrationV2ShellStreamItem,
{ readonly kind: "thread.updated" | "thread.removed" }
> {
if (input.shell !== null) {
if (input.shell.archivedAt !== null) {
return {
Expand Down
3 changes: 2 additions & 1 deletion apps/server/src/ws.ts
Original file line number Diff line number Diff line change
Expand Up @@ -127,6 +127,7 @@ import * as SecretRequests from "./secrets/SecretRequests.ts";
import {
archivedShellStreamItemFromThreadShell,
buildActiveShellSnapshot,
bufferShellLiveStream,
coalesceShellApplicationEvents,
coalesceStoredThreadEvents,
composeShellStreamWithEnrichment,
Expand Down Expand Up @@ -1034,7 +1035,7 @@ export const subscribeOrchestrationV2Shell = Effect.fn("ws.orchestrationV2.subsc
);

const liveFrom = (afterSequence: number) =>
bufferLiveStream(
bufferShellLiveStream(
toShellStream(
applicationEvents.streamProjectedApplicationEvents({
afterSequence,
Expand Down
Loading