diff --git a/apps/server/src/orchestration-v2/LiveStreamBudget.test.ts b/apps/server/src/orchestration-v2/LiveStreamBudget.test.ts index ad50b410c906..737e042e7bd3 100644 --- a/apps/server/src/orchestration-v2/LiveStreamBudget.test.ts +++ b/apps/server/src/orchestration-v2/LiveStreamBudget.test.ts @@ -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"; @@ -9,6 +10,7 @@ import { describe, expect } from "vite-plus/test"; import { bufferLiveStream, + bufferLatestLiveStream, makeLiveStreamBudget, replayAndBufferLiveEvents, type RetainedLiveItem, @@ -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(); + const drained = yield* Deferred.make(); + 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(); + const closed = yield* Deferred.make(); + 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", diff --git a/apps/server/src/orchestration-v2/LiveStreamBudget.ts b/apps/server/src/orchestration-v2/LiveStreamBudget.ts index 6353b4da9065..7db11573f6c3 100644 --- a/apps/server/src/orchestration-v2/LiveStreamBudget.ts +++ b/apps/server/src/orchestration-v2/LiveStreamBudget.ts @@ -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()( @@ -262,6 +263,84 @@ export const bufferLiveStream = ( }), ); +/** For ordered full-state sources, keep only the latest queued state per aggregate while awaiting ACK. */ +export const bufferLatestLiveStream = ( + source: Stream.Stream, + key: (value: A) => string, + limits?: LiveStreamLimits, +) => + Stream.unwrap( + Effect.gen(function* () { + const budget = yield* makeLiveStreamBudget(limits); + const ready = yield* Queue.unbounded(); + const mutex = yield* Semaphore.make(1); + const pending = new Map>(); + 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>; + }), + ), + ), + ), + ), + ), + ); + }), + ); + /** * 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 diff --git a/apps/server/src/orchestration-v2/ShellStream.test.ts b/apps/server/src/orchestration-v2/ShellStream.test.ts index fccff41a5d46..4b7eaadbb2f7 100644 --- a/apps/server/src/orchestration-v2/ShellStream.test.ts +++ b/apps/server/src/orchestration-v2/ShellStream.test.ts @@ -8,10 +8,13 @@ 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"; @@ -19,6 +22,7 @@ import * as TestClock from "effect/testing/TestClock"; import { archivedShellStreamItemFromThreadShell, buildActiveShellSnapshot, + bufferShellLiveStream, coalesceShellApplicationEvents, coalesceStoredThreadEvents, composeShellStreamWithEnrichment, @@ -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, + Cause.Done + >(); + const drained = yield* Deferred.make(); + 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 }); diff --git a/apps/server/src/orchestration-v2/ShellStream.ts b/apps/server/src/orchestration-v2/ShellStream.ts index cec809f45ad4..1e9c5a188889 100644 --- a/apps/server/src/orchestration-v2/ShellStream.ts +++ b/apps/server/src/orchestration-v2/ShellStream.ts @@ -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; + +/** Full shell deltas replace earlier queued state; snapshots and sync markers stay outside this buffer. */ +export const bufferShellLiveStream = ( + source: Stream.Stream, + 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; @@ -292,7 +318,10 @@ export function coalesceStoredThreadEvents( export function shellStreamItemFromThreadShell(input: { readonly stored: Extract; readonly shell: OrchestrationV2ThreadShell | null; -}): Exclude { +}): Extract< + OrchestrationV2ShellStreamItem, + { readonly kind: "thread.updated" | "thread.removed" } +> { if (input.shell !== null) { if (input.shell.archivedAt !== null) { return { diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index 8e7a1210a507..ea90f3798321 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -127,6 +127,7 @@ import * as SecretRequests from "./secrets/SecretRequests.ts"; import { archivedShellStreamItemFromThreadShell, buildActiveShellSnapshot, + bufferShellLiveStream, coalesceShellApplicationEvents, coalesceStoredThreadEvents, composeShellStreamWithEnrichment, @@ -1034,7 +1035,7 @@ export const subscribeOrchestrationV2Shell = Effect.fn("ws.orchestrationV2.subsc ); const liveFrom = (afterSequence: number) => - bufferLiveStream( + bufferShellLiveStream( toShellStream( applicationEvents.streamProjectedApplicationEvents({ afterSequence,