diff --git a/apps/server/src/ws.test.ts b/apps/server/src/ws.test.ts index e280f0750cef..fcef5b78810e 100644 --- a/apps/server/src/ws.test.ts +++ b/apps/server/src/ws.test.ts @@ -1,5 +1,14 @@ +import * as NodeSqliteClient from "@t3tools/shared/nodeSqliteClient"; +import * as Layer from "effect/Layer"; +import * as PubSub from "effect/PubSub"; +import * as Stream from "effect/Stream"; +import * as ThreadManagementService from "./orchestration-v2/ThreadManagementService.ts"; +import * as ProjectStore from "./orchestration-v2/ProjectStore.ts"; +import * as ProjectService from "./project/ProjectService.ts"; +import * as ProjectEnrichmentService from "./project/ProjectEnrichmentService.ts"; +import * as OrchestrationEventStore from "./persistence/Services/OrchestrationEventStore.ts"; import { assert, it } from "@effect/vitest"; -import { ORCHESTRATION_PROTOCOL_VERSION } from "@t3tools/contracts"; +import { ORCHESTRATION_PROTOCOL_VERSION, ProjectId } from "@t3tools/contracts"; import * as Deferred from "effect/Deferred"; import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; @@ -10,6 +19,7 @@ import { hasCompatibleOrchestrationProtocol, resolveAvailableEditorsForConfig, shouldUseBoundedThreadSnapshot, + subscribeOrchestrationV2Shell, } from "./ws.ts"; it("accepts only the current orchestration protocol before websocket RPC setup", () => { @@ -48,3 +58,98 @@ it.effect("does not block server config when editor discovery never resolves", ( assert.deepEqual(availableEditors, []); }), ); + +it.effect( + "enrichment refreshes reuse changed identities without re-requesting project metadata", + () => + Effect.scoped( + Effect.gen(function* () { + const changes = yield* PubSub.unbounded(); + const synchronized = yield* Deferred.make(); + const requestedRoots: string[] = []; + const projects = ["changed", "unchanged"].map((name) => ({ + id: ProjectId.make(name), + title: name, + workspaceRoot: `/workspace/${name}`, + defaultModelSelection: null, + scripts: [], + repositoryIdentity: null, + createdAt: "2026-10-02T00:00:00.000Z", + updatedAt: "2026-10-02T00:00:00.000Z", + })); + const identity = { + canonicalKey: "github.com/example/changed", + locator: { + source: "git-remote" as const, + remoteName: "origin", + remoteUrl: "https://github.com/example/changed.git", + }, + }; + const layer = Layer.mergeAll( + NodeSqliteClient.layer({ filename: ":memory:" }), + Layer.mock(ThreadManagementService.ThreadManagementService)({}), + Layer.mock(ProjectService.ProjectService)({}), + Layer.mock(ProjectStore.ProjectStoreV2)({ listShells: () => Effect.succeed(projects) }), + Layer.mock(OrchestrationEventStore.OrchestrationEventStore)({ + latestApplicationSequence: Effect.succeed(7), + getReplayStats: () => Effect.succeed({ eventCount: 0, rawPayloadBytes: 0 }), + readApplicationEvents: () => Stream.empty, + streamProjectedApplicationEvents: () => Stream.never, + }), + Layer.mock(ProjectEnrichmentService.ProjectEnrichmentService)({ + subscribeChanges: PubSub.subscribe(changes), + getAvailable: (root) => + Effect.sync(() => { + requestedRoots.push(root); + return { + repositoryIdentity: null, + faviconPath: null, + repositoryIdentityResolved: false, + }; + }), + }), + ); + const result = yield* subscribeOrchestrationV2Shell({ + afterSequence: 7, + requestCompletionMarker: true, + }).pipe( + Effect.flatMap((stream) => + stream.pipe( + Stream.tap((item) => + item.kind === "synchronized" + ? Deferred.succeed(synchronized, undefined) + : Effect.void, + ), + Stream.take(2), + Stream.runCollect, + ), + ), + Effect.provide(layer), + Effect.forkChild, + ); + yield* Deferred.await(synchronized); + assert.sameMembers( + requestedRoots, + projects.map((project) => project.workspaceRoot), + ); + // Fill the coalescing batch so this proof does not depend on a timer. + for (let index = 0; index < 64; index++) { + yield* PubSub.publish(changes, { + workspaceRoot: projects[0]!.workspaceRoot, + repositoryIdentityResolved: true, + enrichment: { + repositoryIdentity: identity, + faviconPath: null, + repositoryIdentityResolved: true, + }, + }); + } + const items = yield* Fiber.join(result); + const refresh = items.find((item) => item.kind === "snapshot"); + assert.deepEqual(refresh?.snapshot.projects, [ + { ...projects[0]!, repositoryIdentity: identity }, + ]); + assert.lengthOf(requestedRoots, 2); + }), + ), +); diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index 58bbcb2dbede..2a0006ddb24b 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -948,16 +948,35 @@ export const subscribeOrchestrationV2Shell = Effect.fn("ws.orchestrationV2.subsc const enrichmentRefreshes = Stream.fromSubscription(enrichmentChanges).pipe( Stream.filter((change) => change.repositoryIdentityResolved), Stream.groupedWithin(64, Duration.millis(25)), + // Build the refresh from the identities the changes carry. Re-enriching + // every project here re-requested each expired root, whose resolution + // published again, so one expiry kept every subscriber reloading every + // project's metadata once a minute. Stream.mapEffect((changes) => - applicationEvents.latestApplicationSequence.pipe( - Effect.flatMap(loadProjectMetadataSnapshot), - Effect.map(({ snapshot }) => - shellStreamItemFromEnrichmentRefresh({ - snapshot, - changes: Array.from(changes), - }), - ), - ), + Effect.gen(function* () { + const identities = new Map( + Array.from(changes, (change) => [ + change.workspaceRoot, + change.enrichment.repositoryIdentity, + ]), + ); + const snapshotSequence = yield* applicationEvents.latestApplicationSequence; + const changedProjects = (yield* projects.listShells()).flatMap((project) => + identities.has(project.workspaceRoot) + ? [{ ...project, repositoryIdentity: identities.get(project.workspaceRoot) ?? null }] + : [], + ); + return shellStreamItemFromEnrichmentRefresh({ + snapshot: { + schemaVersion: ORCHESTRATION_V2_PROJECTION_SCHEMA_VERSION, + snapshotSequence, + projects: changedProjects, + threads: [], + archivedThreads: [], + } as OrchestrationV2ShellSnapshot, + changes: Array.from(changes), + }); + }), ), );