diff --git a/apps/server/src/workspace/WorkspaceEntries.test.ts b/apps/server/src/workspace/WorkspaceEntries.test.ts index 7aa2b6d0f7a4..ff6c62de26bc 100644 --- a/apps/server/src/workspace/WorkspaceEntries.test.ts +++ b/apps/server/src/workspace/WorkspaceEntries.test.ts @@ -3,12 +3,16 @@ import * as NodeFSP from "node:fs/promises"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { FileFinder } from "@ff-labs/fff-node"; import { it, afterEach, describe, expect } from "@effect/vitest"; +import * as Deferred from "effect/Deferred"; +import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as FileSystem from "effect/FileSystem"; import * as Layer from "effect/Layer"; import * as Path from "effect/Path"; import * as Option from "effect/Option"; import * as PlatformError from "effect/PlatformError"; +import * as Ref from "effect/Ref"; +import * as TestClock from "effect/testing/TestClock"; import * as Stream from "effect/Stream"; import { vi } from "vite-plus/test"; @@ -764,3 +768,133 @@ it.layer(TestLayer, { excludeTestServices: true })("WorkspaceEntries", (it) => { ); }); }); + +/** + * A `FileFinder` whose entry count is whatever `readEntryCount` returns, so a + * test can move the count under the index the way an external write would. + */ +function mockEntryCountFinder(readEntryCount: () => number) { + const mixedSearch = vi.fn(() => { + const entryCount = readEntryCount(); + return { + ok: true as const, + value: { + items: [], + scores: [], + totalMatched: entryCount, + totalFiles: entryCount, + totalDirs: 0, + }, + }; + }); + const finder = { + destroy: vi.fn(), + waitForIndexReady: vi.fn(async () => ({ ok: true as const, value: true })), + getScanProgress: vi.fn(() => ({ + ok: true as const, + value: { + scannedFilesCount: readEntryCount(), + isScanning: false, + isWatcherReady: true, + isWarmupComplete: true, + }, + })), + mixedSearch, + } as unknown as FileFinder; + vi.spyOn(FileFinder, "create").mockReturnValueOnce({ ok: true, value: finder }); + return { mixedSearch }; +} + +/** + * The same stack as `TestLayer` without excluding the test services: the poll + * fibre below is built inside the test's context, so it runs on the TestClock + * the test advances instead of on wall-clock ticks. + */ +const EntryCountTestLayer = Layer.empty.pipe( + Layer.provideMerge(WorkspaceEntries.layer.pipe(Layer.provide(WorkspacePaths.layer))), + Layer.provideMerge(WorkspacePaths.layer), + Layer.provideMerge(VcsProcess.layer), + Layer.provide( + ServerConfig.ServerConfig.layerTest(process.cwd(), { + prefix: "t3-workspace-entry-count-test-", + }), + ), + Layer.provideMerge(NodeServices.layer), +); + +it.effect( + "signals a watcher when the index entry count moves, and stays quiet while it does not", + () => + Effect.gen(function* () { + const cwd = yield* makeTempDir({ prefix: "t3code-workspace-watch-entry-count-" }); + // The index is mocked so the entry count is the test's to move: the real + // one would tie the assertion to filesystem watcher timing. + let entryCount = 3; + const { mixedSearch } = mockEntryCountFinder(() => entryCount); + + const workspaceEntries = yield* WorkspaceEntries.WorkspaceEntries; + const watched = yield* workspaceEntries.watchEntries({ cwd }); + const signals = yield* Ref.make(0); + const firstSignal = yield* Deferred.make(); + yield* Effect.forkScoped( + Stream.runForEach(watched, () => + Ref.update(signals, (count) => count + 1).pipe( + Effect.andThen(Deferred.succeed(firstSignal, undefined)), + ), + ), + ); + + // The seeding read plus two intervals of an unchanged count. Asserting + // the reads happened is what makes the silence below mean something. + yield* TestClock.adjust(Duration.seconds(2)); + expect(mixedSearch).toHaveBeenCalledTimes(3); + expect(yield* Ref.get(signals)).toBe(0); + + entryCount = 4; + yield* TestClock.adjust(Duration.seconds(1)); + yield* Deferred.await(firstSignal); + expect(yield* Ref.get(signals)).toBe(1); + }).pipe(Effect.scoped, Effect.provide(EntryCountTestLayer)), +); + +it.effect("resumes watching after a failed read of the index", () => + Effect.gen(function* () { + const cwd = yield* makeTempDir({ prefix: "t3code-workspace-watch-entry-count-retry-" }); + let entryCount = 3; + let failNextRead = true; + mockEntryCountFinder(() => { + if (failNextRead) { + failNextRead = false; + throw new Error("native search failed"); + } + return entryCount; + }); + + const workspaceEntries = yield* WorkspaceEntries.WorkspaceEntries; + const watched = yield* workspaceEntries.watchEntries({ cwd }); + const signals = yield* Ref.make(0); + const firstSignal = yield* Deferred.make(); + yield* Effect.forkScoped( + Stream.runForEach(watched, () => + Ref.update(signals, (count) => count + 1).pipe( + Effect.andThen(Deferred.succeed(firstSignal, undefined)), + ), + ), + ); + + // The seeding read failed, so this workspace is unwatched until the retry. + yield* TestClock.adjust(Duration.seconds(1)); + expect(yield* Ref.get(signals)).toBe(0); + + // Whatever moved during the outage is the retry's new baseline, not a + // change to report against a count from before it. + entryCount = 4; + yield* TestClock.adjust(Duration.seconds(30)); + expect(yield* Ref.get(signals)).toBe(0); + + entryCount = 5; + yield* TestClock.adjust(Duration.seconds(1)); + yield* Deferred.await(firstSignal); + expect(yield* Ref.get(signals)).toBe(1); + }).pipe(Effect.scoped, Effect.provide(EntryCountTestLayer)), +); diff --git a/apps/server/src/workspace/WorkspaceEntries.ts b/apps/server/src/workspace/WorkspaceEntries.ts index e9908e90ba41..207bee8b5ad9 100644 --- a/apps/server/src/workspace/WorkspaceEntries.ts +++ b/apps/server/src/workspace/WorkspaceEntries.ts @@ -3,11 +3,13 @@ import * as NodeFSP from "node:fs/promises"; import * as NodeOS from "node:os"; import * as Context from "effect/Context"; +import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as Path from "effect/Path"; import * as PubSub from "effect/PubSub"; import * as RcMap from "effect/RcMap"; +import * as Schedule from "effect/Schedule"; import * as Schema from "effect/Schema"; import type * as Scope from "effect/Scope"; import * as Stream from "effect/Stream"; @@ -30,6 +32,21 @@ import { normalizeSearchQuery } from "@t3tools/shared/searchRanking"; import * as WorkspacePaths from "./WorkspacePaths.ts"; import * as WorkspaceSearchIndex from "./WorkspaceSearchIndex.ts"; +/** + * How often a subscribed workspace re-reads its index entry count. The index + * picks up external changes on its own within ~100ms; this is only how long a + * client waits to hear that it did. One read costs ~0.5ms, so a second is + * already a rounding error against a core. + */ +const WORKSPACE_ENTRY_COUNT_POLL_INTERVAL = Duration.seconds(1); + +/** + * How long a workspace whose index could not be created or read waits before + * watching it again. Slow enough that a root which keeps failing costs + * nothing, quick enough that a subscriber outlives the outage. + */ +const WORKSPACE_ENTRY_COUNT_RETRY_INTERVAL = Duration.seconds(30); + export class WorkspaceEntriesWindowsPathUnsupportedError extends Schema.TaggedErrorClass()( "WorkspaceEntriesWindowsPathUnsupportedError", { @@ -105,8 +122,9 @@ export class WorkspaceEntries extends Context.Service< ) => Effect.Effect; readonly refresh: (cwd: string) => Effect.Effect; /** - * Emits once per `refresh` of the same workspace root. It is a signal, not - * a payload: subscribers re-read through `list`. + * Emits once per `refresh` of the same workspace root, and once per change + * the index picks up on its own while at least one subscriber is watching. + * It is a signal, not a payload: subscribers re-read through `list`. * * Acquiring the subscription is the effect and consuming it is the stream, * so a caller that has awaited this call is already listening and cannot @@ -180,6 +198,70 @@ export const make = Effect.gen(function* () { */ const entryChanges = yield* PubSub.unbounded(); + /** + * The index already tracks external creates and deletes through its own + * filesystem watcher, it just has no callback to tell us it moved. Reading + * the entry count is the cheapest way to notice, so a subscribed workspace + * polls it and publishes the same coarse signal `refresh` does. + * + * A rename, or a balanced add and delete inside one interval, leaves the + * count where it was and publishes nothing; the client's own staleness + * window still catches those. Content-only edits deliberately publish + * nothing - the tree does not render contents. + */ + const publishWhenEntryCountChanges = Effect.fn("WorkspaceEntries.publishWhenEntryCountChanges")( + function* (normalizedCwd: string): Effect.fn.Return { + // The count this fibre owns: read and written by nothing else, and + // undefined only until the first tick establishes a baseline. + let lastCount: number | undefined; + const watch = Effect.gen(function* () { + const searchIndex = yield* WorkspaceSearchIndex.WorkspaceSearchIndex; + yield* Effect.gen(function* () { + const count = yield* searchIndex.entryCount(); + const previousCount = lastCount; + lastCount = count; + if (previousCount === undefined || previousCount === count) return; + yield* PubSub.publish(entryChanges, normalizedCwd); + }).pipe(Effect.repeat({ schedule: Schedule.spaced(WORKSPACE_ENTRY_COUNT_POLL_INTERVAL) })); + }).pipe( + Effect.provide( + workspaceSearchIndexes.get( + WorkspaceSearchIndex.workspaceSearchIndexKey(normalizedCwd, "paths"), + ), + ), + // An index that cannot be created or read is already surfacing through + // `list`, so the watch degrades to the client's staleness window + // instead of taking the subscription down with it. The baseline is + // dropped with it: the next attempt re-seeds rather than publishing a + // change against a count from before the outage. + Effect.catch((cause) => + Effect.gen(function* () { + lastCount = undefined; + yield* Effect.logWarning("Stopped watching the workspace index for external changes", { + cwd: normalizedCwd, + cause, + }); + }), + ), + ); + + // `watch` only returns when it has failed, so this retries the outage + // rather than repeating a healthy poll loop. + yield* watch.pipe( + Effect.repeat({ schedule: Schedule.spaced(WORKSPACE_ENTRY_COUNT_RETRY_INTERVAL) }), + ); + }, + ); + + /** + * One poll fibre per workspace root, shared by every subscriber of that root + * and released when the last of them goes away. + */ + const entryCountPollers = yield* RcMap.make({ + lookup: (normalizedCwd: string) => + Effect.forkScoped(publishWhenEntryCountChanges(normalizedCwd)), + }); + const refresh: WorkspaceEntries["Service"]["refresh"] = Effect.fn("WorkspaceEntries.refresh")( function* (cwd) { const normalizedCwd = yield* workspaceRootKey(cwd); @@ -227,6 +309,7 @@ export const make = Effect.gen(function* () { // is being normalised is buffered rather than dropped. const subscription = yield* PubSub.subscribe(entryChanges); const normalizedCwd = yield* workspaceRootKey(input.cwd); + yield* RcMap.get(entryCountPollers, normalizedCwd); return Stream.fromSubscription(subscription).pipe( Stream.filter((changedCwd) => changedCwd === normalizedCwd), Stream.map(() => ({})), diff --git a/apps/server/src/workspace/WorkspaceSearchIndex.test.ts b/apps/server/src/workspace/WorkspaceSearchIndex.test.ts index 078278242942..dcd933233b24 100644 --- a/apps/server/src/workspace/WorkspaceSearchIndex.test.ts +++ b/apps/server/src/workspace/WorkspaceSearchIndex.test.ts @@ -508,3 +508,32 @@ it.live("indexes a file created externally right after construction, without a r ); }), ); + +it.effect("counts index entries without materialising a page", () => + Effect.scoped( + Effect.gen(function* () { + const mixedSearch = vi.fn(() => ({ + ok: true as const, + value: { + items: [], + scores: [], + totalMatched: 42, + totalFiles: 30, + totalDirs: 12, + }, + })); + const finder = { + destroy: vi.fn(), + waitForIndexReady: vi.fn(async () => ({ ok: true as const, value: true })), + getScanProgress: watcherReadyProgress(), + mixedSearch, + } as unknown as FileFinder; + vi.spyOn(FileFinder, "create").mockReturnValueOnce({ ok: true, value: finder }); + + const searchIndex = yield* WorkspaceSearchIndex.make("/workspace/project"); + + expect(yield* searchIndex.entryCount()).toBe(42); + expect(mixedSearch).toHaveBeenCalledWith("", { pageSize: 1 }); + }), + ), +); diff --git a/apps/server/src/workspace/WorkspaceSearchIndex.ts b/apps/server/src/workspace/WorkspaceSearchIndex.ts index 1ac91c35393a..26c8d599a1c3 100644 --- a/apps/server/src/workspace/WorkspaceSearchIndex.ts +++ b/apps/server/src/workspace/WorkspaceSearchIndex.ts @@ -30,6 +30,7 @@ import { isWorkspaceImagePreviewPath } from "@t3tools/shared/filePreview"; const WORKSPACE_INDEX_MAX_ENTRIES = 25_000; const WORKSPACE_INDEX_PAGE_SIZE = WORKSPACE_INDEX_MAX_ENTRIES + 2; +const WORKSPACE_INDEX_COUNT_PAGE_SIZE = 1; const WORKSPACE_INDEX_SCAN_TIMEOUT = "15 seconds"; const WORKSPACE_INDEX_SCAN_TIMEOUT_MS = 15_000; const WORKSPACE_INDEX_IDLE_TTL = "15 minutes"; @@ -113,6 +114,12 @@ export class WorkspaceSearchIndex extends Context.Service< WorkspaceSearchIndex, { readonly list: () => Effect.Effect; + /** + * How many entries the index currently holds. Cheap enough to poll: the + * page size of 1 is load-bearing, it returns the total without + * materialising the page. + */ + readonly entryCount: () => Effect.Effect; readonly search: ( query: string, limit: number, @@ -527,6 +534,15 @@ export const make = Effect.fn("WorkspaceSearchIndex.make")(function* ( }, ); + const entryCount: WorkspaceSearchIndex["Service"]["entryCount"] = Effect.fn( + "WorkspaceSearchIndex.entryCount", + )(function* () { + const result = yield* runSearch("", WORKSPACE_INDEX_COUNT_PAGE_SIZE, "mixedSearch", () => + finder.mixedSearch("", { pageSize: WORKSPACE_INDEX_COUNT_PAGE_SIZE }), + ); + return result.totalMatched; + }); + const search: WorkspaceSearchIndex["Service"]["search"] = Effect.fn( "WorkspaceSearchIndex.search", )(function* (query, limit, kind, imageOnly) { @@ -600,7 +616,7 @@ export const make = Effect.fn("WorkspaceSearchIndex.make")(function* ( }; }); - return WorkspaceSearchIndex.of({ list, refresh, search, searchContents }); + return WorkspaceSearchIndex.of({ entryCount, list, refresh, search, searchContents }); }); export const WORKSPACE_SEARCH_INDEX_VARIANTS = ["paths", "content"] as const; diff --git a/packages/client-runtime/src/state/projectCommands.ts b/packages/client-runtime/src/state/projectCommands.ts index 73c3b4fbcfd1..7bbfa4ec80d5 100644 --- a/packages/client-runtime/src/state/projectCommands.ts +++ b/packages/client-runtime/src/state/projectCommands.ts @@ -61,9 +61,10 @@ export function createProjectEnvironmentAtoms( tag: WS_METHODS.projectsSearchEntries, staleTimeMs: 15_000, }), - // The server says when its own entry index moved, so files an agent creates - // or deletes during a turn land in the tree without waiting out the stale - // window. Changes made outside the app still need a remount. + // The server says when its own entry index moved, so files created or + // deleted during a turn - by an agent, or by anything else on disk - land + // in the tree without waiting out the stale window. A rename outside the + // app leaves the server's entry count where it was, so those still wait. listEntries: createEnvironmentRpcQueryAtomFamily(runtime, { label: "environment-data:projects:list-entries", tag: WS_METHODS.projectsListEntries, diff --git a/packages/contracts/src/project.ts b/packages/contracts/src/project.ts index d23fd73af432..2167a186dc43 100644 --- a/packages/contracts/src/project.ts +++ b/packages/contracts/src/project.ts @@ -83,7 +83,8 @@ export type ProjectListEntriesResult = typeof ProjectListEntriesResult.Type; /** * Emitted by `subscribeProjectEntryChanges` whenever the server's entry index - * for a workspace is refreshed. Like `ProjectFileChangedEvent` it carries no + * for a workspace changes: either the app refreshed it, or the index picked up + * a create or delete made outside the app. Like `ProjectFileChangedEvent` it carries no * payload: clients re-read the listing through `projects.listEntries` so * limits, ordering and error mapping keep living in one place. */