Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
134 changes: 134 additions & 0 deletions apps/server/src/workspace/WorkspaceEntries.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";

Expand Down Expand Up @@ -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<void>();
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<void>();
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)),
);
87 changes: 85 additions & 2 deletions apps/server/src/workspace/WorkspaceEntries.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand All @@ -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>()(
"WorkspaceEntriesWindowsPathUnsupportedError",
{
Expand Down Expand Up @@ -105,8 +122,9 @@ export class WorkspaceEntries extends Context.Service<
) => Effect.Effect<ProjectSearchContentsResult, WorkspaceEntriesError>;
readonly refresh: (cwd: string) => Effect.Effect<void>;
/**
* 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
Expand Down Expand Up @@ -180,6 +198,70 @@ export const make = Effect.gen(function* () {
*/
const entryChanges = yield* PubSub.unbounded<string>();

/**
* 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<void> {
// 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);
Expand Down Expand Up @@ -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(() => ({})),
Expand Down
29 changes: 29 additions & 0 deletions apps/server/src/workspace/WorkspaceSearchIndex.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 });
}),
),
);
18 changes: 17 additions & 1 deletion apps/server/src/workspace/WorkspaceSearchIndex.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -113,6 +114,12 @@ export class WorkspaceSearchIndex extends Context.Service<
WorkspaceSearchIndex,
{
readonly list: () => Effect.Effect<ProjectListEntriesResult, WorkspaceSearchIndexSearchFailed>;
/**
* 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<number, WorkspaceSearchIndexSearchFailed>;
readonly search: (
query: string,
limit: number,
Expand Down Expand Up @@ -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) {
Expand Down Expand Up @@ -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;
Expand Down
7 changes: 4 additions & 3 deletions packages/client-runtime/src/state/projectCommands.ts
Original file line number Diff line number Diff line change
Expand Up @@ -61,9 +61,10 @@ export function createProjectEnvironmentAtoms<R, E>(
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,
Expand Down
3 changes: 2 additions & 1 deletion packages/contracts/src/project.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.
*/
Expand Down
Loading