diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index cf5666653a9..a9dc165e390 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -6,6 +6,11 @@ on: - main pull_request: workflow_dispatch: + inputs: + search-profile-base: + description: Commit to compare with this revision using the search profiling fixture + type: string + default: "" env: PRIMARY_NODE_VERSION: "22.x" @@ -20,6 +25,45 @@ concurrency: cancel-in-progress: ${{ github.ref != 'refs/heads/main' }} jobs: + search-profile: + if: github.event_name == 'workflow_dispatch' && inputs.search-profile-base != '' + runs-on: blacksmith-4vcpu-ubuntu-2404 + timeout-minutes: 45 + permissions: + contents: read + steps: + - uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 + with: + fetch-depth: 0 + - uses: ./.github/actions/setup-workspace + with: + node-version: ${{ env.PRIMARY_NODE_VERSION }} + pnpm-version: ${{ env.PNPM_VERSION }} + cache-prefix: search-profile + - name: Compare search timings and collect CPU profiles + env: + SEARCH_PROFILE_BASE: ${{ inputs.search-profile-base }} + SEARCH_PROFILE_HEAD: ${{ github.sha }} + run: | + mkdir -p .search-profile + cp scripts/search-profile.mjs .search-profile/runner.mjs + node .search-profile/runner.mjs compare "$SEARCH_PROFILE_BASE" "$SEARCH_PROFILE_HEAD" + - name: Preserve measurements and profiles + if: always() + uses: actions/upload-artifact@043fb46d1a93c77aae656e7c1c64a875d1fc6a0a + with: + name: search-profile-${{ github.sha }} + include-hidden-files: true + retention-days: 90 + path: | + .search-profile/runner.mjs + .search-profile/manifest.json + .search-profile/results/** + .search-profile/profiles/** + .search-profile/logs/** + .search-profile/*.tgz + .search-profile/summary.md + checks: name: Checks (ubuntu-latest, Node 22.x) runs-on: blacksmith-4vcpu-ubuntu-2404 diff --git a/apps/app/bundle-budget.json b/apps/app/bundle-budget.json index a4af3baff95..e3b114df791 100644 --- a/apps/app/bundle-budget.json +++ b/apps/app/bundle-budget.json @@ -8,6 +8,10 @@ "Chunks behind a dynamic import are deliberately not counted. Moving code", "behind React.lazy or import() is the fix this budget exists to encourage.", "", + "On 2026-09-26 the SplitWorkspaceRoute raw limit rose 1 KiB for", + "independent mention-provider queries (about 0.3 KiB) and the existing", + "0.3 KiB overage. The measured brotli payload shrank; its limit is unchanged.", + "", "On 2026-09-25 the SplitWorkspaceRoute limits rose by 2 KiB raw and", "256 bytes brotli for sanitized HTML video support in thread Markdown.", "", @@ -109,7 +113,7 @@ }, "routeClosures": { "SplitWorkspaceRoute": { - "maxBytes": 2341134, + "maxBytes": 2342158, "maxBrotliBytes": 630625, "forbiddenPackages": [ "@pierre/diffs", diff --git a/apps/app/src/hooks/queries/plugin-contribution-queries.test.tsx b/apps/app/src/hooks/queries/plugin-contribution-queries.test.tsx index 9755af5ace4..5036aa623d5 100644 --- a/apps/app/src/hooks/queries/plugin-contribution-queries.test.tsx +++ b/apps/app/src/hooks/queries/plugin-contribution-queries.test.tsx @@ -1,6 +1,7 @@ // @vitest-environment jsdom -import { cleanup, renderHook, waitFor } from "@testing-library/react"; +import { act, cleanup, renderHook, waitFor } from "@testing-library/react"; +import { createDeferredPromise } from "@bb/test-helpers"; import { afterEach, describe, expect, it, vi } from "vitest"; import { sdk } from "@/lib/sdk"; import { createQueryClientTestHarness } from "@/test/queryClientTestHarness"; @@ -110,6 +111,67 @@ describe("usePluginContributions", () => { }); describe("usePluginMentionSearch", () => { + it("publishes fast provider results while a slow provider is still pending", async () => { + const slow = createDeferredPromise(); + const group = (providerId: string) => ({ + pluginId: "fixture", + providerId, + label: providerId, + items: [ + { + itemId: `${providerId}:one`, + title: providerId, + subtitle: null, + icon: null, + }, + ], + }); + vi.stubGlobal( + "fetch", + vi.fn((input: string) => + input.includes("providerId=slow") + ? slow.promise + : Promise.resolve( + new Response(JSON.stringify({ groups: [group("fast")] })), + ), + ), + ); + const { wrapper } = createQueryClientTestHarness(); + const { result } = renderHook( + () => + usePluginMentionSearch( + { trigger: "@", query: "one", projectId: null, threadId: null }, + { + enabled: true, + providers: [ + { + pluginId: "fixture", + id: "slow", + label: "Slow", + triggers: ["@"], + }, + { + pluginId: "fixture", + id: "fast", + label: "Fast", + triggers: ["@"], + }, + ], + }, + ), + { wrapper }, + ); + await waitFor(() => expect(result.current.data).toEqual([group("fast")])); + expect(result.current.isFetching).toBe(true); + await act(async () => + slow.resolve(new Response(JSON.stringify({ groups: [group("slow")] }))), + ); + await waitFor(() => + expect(result.current.data).toEqual([group("slow"), group("fast")]), + ); + expect(result.current.isFetching).toBe(false); + }); + it("includes the active trigger in the search request", async () => { const fetchMock = mockFetchJsonOnce({ ok: true, @@ -140,7 +202,17 @@ describe("usePluginMentionSearch", () => { projectId: "proj_1", threadId: null, }, - { enabled: true }, + { + enabled: true, + providers: [ + { + pluginId: "github", + id: "issue", + label: "GitHub issues", + triggers: ["#"], + }, + ], + }, ), { wrapper }, ); @@ -163,7 +235,7 @@ describe("usePluginMentionSearch", () => { ]); }); expect(fetchMock).toHaveBeenCalledWith( - "/api/v1/plugins/mentions/search?q=42&trigger=%23&projectId=proj_1", + "/api/v1/plugins/mentions/search?q=42&trigger=%23&pluginId=github&providerId=issue&projectId=proj_1", expect.objectContaining({ signal: expect.any(AbortSignal) }), ); }); diff --git a/apps/app/src/hooks/queries/plugin-contribution-queries.ts b/apps/app/src/hooks/queries/plugin-contribution-queries.ts index b80dab5a974..5ae47f79acd 100644 --- a/apps/app/src/hooks/queries/plugin-contribution-queries.ts +++ b/apps/app/src/hooks/queries/plugin-contribution-queries.ts @@ -1,4 +1,8 @@ -import { useQuery } from "@tanstack/react-query"; +import { + useQueries, + useQuery, + type UseQueryResult, +} from "@tanstack/react-query"; import { normalizePluginMentionTriggers, type PluginMentionTrigger, @@ -117,11 +121,14 @@ interface PluginMentionSearchArgs { async function fetchPluginMentionSearch( args: PluginMentionSearchArgs, + provider: PluginMentionProviderContribution, signal: AbortSignal, ): Promise { const params = new URLSearchParams({ q: args.query, trigger: args.trigger, + pluginId: provider.pluginId, + providerId: provider.id, }); if (args.projectId !== null) params.set("projectId", args.projectId); if (args.threadId !== null) params.set("threadId", args.threadId); @@ -138,20 +145,40 @@ async function fetchPluginMentionSearch( export function usePluginMentionSearch( args: PluginMentionSearchArgs, - options: { enabled: boolean }, + options: { + enabled: boolean; + providers: readonly PluginMentionProviderContribution[]; + }, ) { - return useQuery({ - queryKey: [ - "plugin-mention-search", - args.trigger, - args.query, - args.projectId, - args.threadId, - ], - queryFn: ({ signal }) => fetchPluginMentionSearch(args, signal), - enabled: options.enabled, - staleTime: 15_000, - placeholderData: (previous, previousQuery) => - previousQuery?.queryKey[1] === args.trigger ? previous : undefined, + return useQueries({ + queries: options.providers + .filter((provider) => provider.triggers.includes(args.trigger)) + .map((provider) => ({ + queryKey: [ + "plugin-mention-search", + args.trigger, + args.query, + args.projectId, + args.threadId, + provider.pluginId, + provider.id, + ], + queryFn: ({ signal }: { signal: AbortSignal }) => + fetchPluginMentionSearch(args, provider, signal), + enabled: options.enabled, + staleTime: 15_000, + })), + combine: combineMentionSearches, }); } + +function combineMentionSearches( + queries: UseQueryResult[], +) { + return { + data: queries.flatMap((query) => query.data ?? []), + isLoading: queries.some((query) => query.isLoading), + isFetching: queries.some((query) => query.isFetching), + isError: queries.some((query) => query.isError), + }; +} diff --git a/apps/app/src/hooks/usePromptMentions.ts b/apps/app/src/hooks/usePromptMentions.ts index e52ed4e49cf..6f10041b82f 100644 --- a/apps/app/src/hooks/usePromptMentions.ts +++ b/apps/app/src/hooks/usePromptMentions.ts @@ -169,6 +169,7 @@ export function usePromptMentions( threadId: options.currentThreadId ?? null, }, { + providers: pluginContributions.data?.mentionProviders ?? [], enabled: hasMentionProviders && pluginSearchMatchesInput && diff --git a/apps/host-daemon/src/command-handlers/file-list.ts b/apps/host-daemon/src/command-handlers/file-list.ts index 4d724098382..ea72285cd35 100644 --- a/apps/host-daemon/src/command-handlers/file-list.ts +++ b/apps/host-daemon/src/command-handlers/file-list.ts @@ -58,9 +58,45 @@ interface ListWorkspacePathsArgs extends PathListInclusion { includeHidden: boolean; excludeNames: readonly string[]; respectGitIgnore: boolean; + maxAgeMs: number; } -const pendingListings = new Map>(); +interface WorkspacePathListing { + root: string; + promise: Promise; + expiresAt: number | null; + pathCount: number; +} + +const workspaceListings = new Map(); +const MAX_CACHED_LISTINGS = 32; +const MAX_CACHED_PATHS = 100_000; + +export function invalidateWorkspacePathListings(root: string): void { + for (const [key, entry] of workspaceListings) { + if ( + entry.root === root || + entry.root.startsWith(`${root}${path.sep}`) || + root.startsWith(`${entry.root}${path.sep}`) + ) { + workspaceListings.delete(key); + } + } +} + +function trimWorkspaceListings(): void { + let pathCount = 0; + for (const entry of workspaceListings.values()) pathCount += entry.pathCount; + for (const [key, entry] of workspaceListings) { + if ( + workspaceListings.size <= MAX_CACHED_LISTINGS && + pathCount <= MAX_CACHED_PATHS + ) + break; + workspaceListings.delete(key); + pathCount -= entry.pathCount; + } +} export function listWorkspacePaths( args: ListWorkspacePathsArgs, @@ -73,13 +109,43 @@ export function listWorkspacePaths( args.includeFiles, args.includeDirectories, ]); - const existing = pendingListings.get(key); - if (existing) return existing; - const listing = discoverWorkspacePaths(args).finally(() => { - pendingListings.delete(key); - }); - pendingListings.set(key, listing); - return listing; + const existing = workspaceListings.get(key); + if ( + existing && + (existing.expiresAt === null || + (args.maxAgeMs > 0 && existing.expiresAt > Date.now())) + ) { + workspaceListings.delete(key); + workspaceListings.set(key, existing); + return existing.promise; + } + const entry: WorkspacePathListing = { + root: args.root, + expiresAt: null, + pathCount: 0, + promise: discoverWorkspacePaths(args).then( + (paths) => { + if (workspaceListings.get(key) === entry) { + if (args.maxAgeMs > 0 && paths.length <= MAX_CACHED_PATHS) { + entry.expiresAt = Date.now() + args.maxAgeMs; + entry.pathCount = paths.length; + trimWorkspaceListings(); + } else { + workspaceListings.delete(key); + } + } + return paths; + }, + (error: unknown) => { + if (workspaceListings.get(key) === entry) workspaceListings.delete(key); + throw error; + }, + ), + }; + workspaceListings.delete(key); + workspaceListings.set(key, entry); + trimWorkspaceListings(); + return entry.promise; } async function discoverWorkspacePaths( diff --git a/apps/host-daemon/src/command-handlers/host-files.ts b/apps/host-daemon/src/command-handlers/host-files.ts index 1d40039a81d..338911ef3be 100644 --- a/apps/host-daemon/src/command-handlers/host-files.ts +++ b/apps/host-daemon/src/command-handlers/host-files.ts @@ -79,6 +79,7 @@ export async function listHostFiles( respectGitIgnore: command.respectGitIgnore, includeFiles: true, includeDirectories: false, + maxAgeMs: command.query ? 2_000 : 0, }) ).map((entry) => entry.path), limit: command.limit, @@ -113,6 +114,7 @@ export async function listHostPaths( includeHidden: command.includeHidden, excludeNames: command.excludeNames, respectGitIgnore: command.respectGitIgnore, + maxAgeMs: command.query ? 2_000 : 0, }), limit: command.limit, includeFiles: command.includeFiles, diff --git a/apps/host-daemon/src/command-handlers/workspace-path-list.test.ts b/apps/host-daemon/src/command-handlers/workspace-path-list.test.ts index 7e3b4d29bbf..ce68d09f291 100644 --- a/apps/host-daemon/src/command-handlers/workspace-path-list.test.ts +++ b/apps/host-daemon/src/command-handlers/workspace-path-list.test.ts @@ -1,9 +1,12 @@ import fs from "node:fs/promises"; import os from "node:os"; import path from "node:path"; -import { afterEach, describe, expect, it } from "vitest"; +import { afterEach, describe, expect, it, vi } from "vitest"; import { runGit } from "@bb/host-workspace"; -import { listWorkspacePaths } from "./file-list.js"; +import { + invalidateWorkspacePathListings, + listWorkspacePaths, +} from "./file-list.js"; import { listHostPaths } from "./host-files.js"; const roots: string[] = []; @@ -30,6 +33,7 @@ function listingArgs(root: string) { includeDirectories: true, includeHidden: true, respectGitIgnore: true, + maxAgeMs: 0, excludeNames: [], }; } @@ -41,6 +45,7 @@ async function paths(root: string) { } afterEach(async () => { + vi.restoreAllMocks(); await Promise.all( roots .splice(0) @@ -49,6 +54,48 @@ afterEach(async () => { }); describe("workspace path discovery", () => { + it("refreshes cached suggestions after invalidation and expiry", async () => { + const root = await createRoot(); + await initRepo(root); + await write(root, "first.txt"); + const args = { ...listingArgs(root), maxAgeMs: 2_000 }; + await listWorkspacePaths(args); + await write(root, "second.txt"); + await write(root, ".gitignore", "first.txt\n"); + invalidateWorkspacePathListings(root); + expect((await listWorkspacePaths(args)).map((entry) => entry.path)).toEqual( + [".gitignore", "second.txt"], + ); + await write(root, "third.txt"); + vi.spyOn(Date, "now").mockReturnValue(Date.now() + 2_001); + expect((await listWorkspacePaths(args)).map((entry) => entry.path)).toEqual( + [".gitignore", "second.txt", "third.txt"], + ); + }); + + it("reuses discovery across successive file suggestion queries", async () => { + const root = await createRoot(); + await write(root, "alpha.md"); + await write(root, "beta.md"); + const readDirectory = vi.spyOn(fs, "readdir"); + const command = { + type: "host.list_paths" as const, + path: root, + includeFiles: true, + includeDirectories: false, + includeHidden: true, + respectGitIgnore: false, + excludeNames: [], + limit: 1, + }; + const alpha = await listHostPaths({ ...command, query: "alpha" }); + const reads = readDirectory.mock.calls.length; + const beta = await listHostPaths({ ...command, query: "beta" }); + expect(alpha.paths.map((entry) => entry.path)).toEqual(["alpha.md"]); + expect(beta.paths.map((entry) => entry.path)).toEqual(["beta.md"]); + expect(readDirectory.mock.calls.length).toBe(reads); + }); + it("includes tracked and non-ignored untracked files while pruning Git-ignored trees", async () => { const root = await createRoot(); await initRepo(root); diff --git a/apps/host-daemon/src/watch-manager.ts b/apps/host-daemon/src/watch-manager.ts index 775ec7ee840..55a27fbe252 100644 --- a/apps/host-daemon/src/watch-manager.ts +++ b/apps/host-daemon/src/watch-manager.ts @@ -17,6 +17,7 @@ import type { WorkspaceWatchError, } from "@bb/host-watcher"; import { userExecutableProcessOptions } from "./user-executable-env.js"; +import { invalidateWorkspacePathListings } from "./command-handlers/file-list.js"; type StopWatching = () => void | Promise; @@ -257,6 +258,7 @@ export class WatchManager { }); }, onWatchError: (error) => { + invalidateWorkspacePathListings(workspace.path); this.options.onWorkspaceStatusWatchError?.({ error }); }, }); @@ -306,6 +308,7 @@ export class WatchManager { changeKinds: readonly WorkspaceStatusWatchChangeKind[]; entry: WorkspaceWatchEntry; }): void { + invalidateWorkspacePathListings(args.entry.workspace.path); if ( this.workspaceEntries.get(args.entry.target.environmentId) !== args.entry ) { @@ -484,6 +487,7 @@ export class WatchManager { this.threadStorageTargets.get(threadId) ?? null, onChange: (event) => { if (event.kind === "thread-storage-changed") { + invalidateWorkspacePathListings(threadStorageRootPath); this.options.onThreadStorageChanged?.({ environmentId: event.environmentId, threadId: event.threadId, diff --git a/apps/server/package.json b/apps/server/package.json index 602e75d84b4..7efe4237f53 100644 --- a/apps/server/package.json +++ b/apps/server/package.json @@ -4,7 +4,7 @@ "type": "module", "private": true, "scripts": { - "build": "node ../../scripts/build-node-entry.mjs src/index.ts dist/index.js --clean-dist --external ./start-server.js --copy-dir ../../packages/db/drizzle dist/drizzle --copy-dir src/assets dist/assets && node --import tsx scripts/copy-builtin-skills.ts && node ../../scripts/build-node-entry.mjs src/start-server.ts dist/start-server.js && node ../../scripts/build-node-entry.mjs ../../packages/plugin-sdk/src/index.ts dist/plugin-sdk-runtime.js && node ../../scripts/build-node-entry.mjs src/services/plugins/zod-runtime.ts dist/zod-runtime.js", + "build": "node ../../scripts/build-node-entry.mjs src/index.ts dist/index.js --clean-dist --external ./start-server.js --copy-dir ../../packages/db/drizzle dist/drizzle --copy-dir src/assets dist/assets && node --import tsx scripts/copy-builtin-skills.ts && node ../../scripts/build-node-entry.mjs src/start-server.ts dist/start-server.js && node ../../scripts/build-node-entry.mjs ../../packages/plugin-sdk/src/index.ts dist/plugin-sdk-runtime.js && node ../../scripts/build-node-entry.mjs src/services/plugins/zod-runtime.ts dist/zod-runtime.js && node ../../scripts/build-node-entry.mjs src/thread-search-worker.ts dist/thread-search-worker.js", "start": "node dist/index.js", "start:prod": "cross-env NODE_ENV=production node dist/index.js", "dev": "node --conditions=source --import tsx scripts/dev-supervisor.mjs", diff --git a/apps/server/src/routes/plugins.ts b/apps/server/src/routes/plugins.ts index 6d708f454a5..15baa17caad 100644 --- a/apps/server/src/routes/plugins.ts +++ b/apps/server/src/routes/plugins.ts @@ -477,11 +477,23 @@ export function registerPluginRoutes( 400, ); } + const pluginId = context.req.query("pluginId")?.trim() ?? ""; + const providerId = context.req.query("providerId")?.trim() ?? ""; + if (Boolean(pluginId) !== Boolean(providerId)) { + return context.json( + { + ok: false, + error: "pluginId and providerId must be supplied together", + }, + 400, + ); + } const groups = await plugins.searchMentions({ trigger, query, projectId: projectId !== null && projectId.length > 0 ? projectId : null, threadId: threadId !== null && threadId.length > 0 ? threadId : null, + provider: pluginId && providerId ? { pluginId, providerId } : null, }); return context.json({ ok: true, groups }); }); diff --git a/apps/server/src/routes/threads/base.ts b/apps/server/src/routes/threads/base.ts index 8138f7f464a..081c122ee81 100644 --- a/apps/server/src/routes/threads/base.ts +++ b/apps/server/src/routes/threads/base.ts @@ -1,3 +1,4 @@ +import { searchThreads } from "../../services/threads/thread-search.js"; import { countUnarchivedThreadDescendants } from "../../services/threads/thread-archive.js"; import { cancelAbandonedProviderCreations } from "../../services/threads/thread-environment-providers.js"; import { @@ -13,7 +14,6 @@ import { listThreadsWithPendingInteractionState, markThreadDeleted, listLifecycleThreadTree, - searchThreadsWithPendingInteractionState, updateThread, type ThreadSearchResultGroup as DbThreadSearchResultGroup, type UpdateThreadInput, @@ -301,7 +301,7 @@ export function registerThreadBaseRoutes(app: Hono, deps: AppDeps): void { ); }); - get(routes.search, (context, query) => { + get(routes.search, async (context, query) => { const searchQuery = query.query.trim(); if (countNonWhitespaceChars(searchQuery) < 2) { throw new ApiError( @@ -313,10 +313,14 @@ export function registerThreadBaseRoutes(app: Hono, deps: AppDeps): void { const limitPerGroup = parseSearchLimitPerGroup(query.limitPerGroup); return context.json( buildThreadSearchResponse(deps, { - ...searchThreadsWithPendingInteractionState(deps.db, { - query: searchQuery, - limitPerGroup, - }), + ...(await searchThreads( + deps.db, + { + query: searchQuery, + limitPerGroup, + }, + context.req.raw.signal, + )), }) satisfies ThreadSearchResponse, ); }); diff --git a/apps/server/src/services/plugins/plugin-service.ts b/apps/server/src/services/plugins/plugin-service.ts index 24cb7f0bcd5..14d63d6906c 100644 --- a/apps/server/src/services/plugins/plugin-service.ts +++ b/apps/server/src/services/plugins/plugin-service.ts @@ -418,6 +418,7 @@ export interface PluginService { query: string; projectId: string | null; threadId: string | null; + provider: { pluginId: string; providerId: string } | null; }): Promise; resolveMention(args: { pluginId: string; @@ -2302,7 +2303,10 @@ export function createPluginService(deps: PluginServiceDeps): PluginService { if (entries.length === 0) return []; const tasks: Array> = []; for (const [id, plugin] of entries) { + if (args.provider !== null && args.provider.pluginId !== id) continue; for (const record of [...plugin.handle.mentionProviders]) { + if (args.provider !== null && args.provider.providerId !== record.id) + continue; if (!record.triggers.includes(args.trigger)) continue; tasks.push( (async () => { diff --git a/apps/server/src/services/threads/thread-search-protocol.ts b/apps/server/src/services/threads/thread-search-protocol.ts new file mode 100644 index 00000000000..a16cd8a7645 --- /dev/null +++ b/apps/server/src/services/threads/thread-search-protocol.ts @@ -0,0 +1,25 @@ +import { z } from "zod"; + +export const threadSearchRequestSchema = z.object({ + query: z.string(), + limitPerGroup: z.number().int().positive(), +}); + +export const threadSearchWorkerResponseSchema = z.discriminatedUnion("ok", [ + z.object({ + ok: z.literal(true), + rows: z.array( + z.object({ + archived: z.union([z.literal(0), z.literal(1)]), + segmentOrder: z.number(), + sourceKind: z.string(), + sourceSeq: z.number().nullable(), + text: z.string(), + threadId: z.string(), + threadOrder: z.number(), + total: z.number(), + }), + ), + }), + z.object({ ok: z.literal(false), error: z.string() }), +]); diff --git a/apps/server/src/services/threads/thread-search.ts b/apps/server/src/services/threads/thread-search.ts new file mode 100644 index 00000000000..948001338d4 --- /dev/null +++ b/apps/server/src/services/threads/thread-search.ts @@ -0,0 +1,116 @@ +import { Worker } from "node:worker_threads"; +import { + hydrateThreadSearchResults, + searchThreadsWithPendingInteractionState, + type DbConnection, + type ThreadSearchMatchRow, +} from "@bb/db"; +import { threadSearchWorkerResponseSchema } from "./thread-search-protocol.js"; + +interface SearchArgs { + query: string; + limitPerGroup: number; +} + +class ThreadSearchWorker { + private worker: Worker | null = null; + private tail: Promise = Promise.resolve(); + private closed = false; + + constructor(private readonly databasePath: string) {} + + search( + args: SearchArgs, + signal: AbortSignal, + ): Promise { + const result = this.tail.then(() => { + signal.throwIfAborted(); + if (this.closed) throw new Error("Thread search is closed"); + const worker = this.worker ?? this.start(); + return new Promise((resolve, reject) => { + const cleanup = () => { + worker.off("message", onMessage); + worker.off("error", onError); + worker.off("exit", onExit); + }; + const onMessage = (message: unknown) => { + cleanup(); + try { + signal.throwIfAborted(); + const response = threadSearchWorkerResponseSchema.parse(message); + if (!response.ok) throw new Error(response.error); + resolve(response.rows); + } catch (error) { + reject(error); + } + }; + const onError = (error: Error) => { + cleanup(); + this.worker = null; + reject(error); + }; + const onExit = (code: number) => + onError(new Error(`Thread search worker exited (${code})`)); + worker.once("message", onMessage); + worker.once("error", onError); + worker.once("exit", onExit); + worker.postMessage(args); + }); + }); + this.tail = result.catch(() => undefined); + return result; + } + + private start(): Worker { + const source = import.meta.url.endsWith(".ts"); + const entry = new URL( + source ? "../../thread-search-worker.ts" : "./thread-search-worker.js", + import.meta.url, + ); + const worker = source + ? new Worker( + `import("tsx/esm/api").then(({ register }) => { register(); return import(${JSON.stringify(entry.href)}); });`, + { eval: true, workerData: this.databasePath }, + ) + : new Worker(entry, { workerData: this.databasePath }); + worker.unref(); + worker.on("error", () => { + if (this.worker === worker) this.worker = null; + }); + worker.on("exit", () => { + if (this.worker === worker) this.worker = null; + }); + this.worker = worker; + return worker; + } + + async close(): Promise { + this.closed = true; + await this.worker?.terminate(); + this.worker = null; + } +} + +const workers = new WeakMap(); + +export async function searchThreads( + db: DbConnection, + args: SearchArgs, + signal: AbortSignal, +) { + signal.throwIfAborted(); + if (db.$client.memory) + return searchThreadsWithPendingInteractionState(db, args); + let worker = workers.get(db); + if (worker === undefined) { + worker = new ThreadSearchWorker(db.$client.name); + workers.set(db, worker); + } + const rows = await worker.search(args, signal); + return hydrateThreadSearchResults(db, { query: args.query, rows }); +} + +export async function closeThreadSearch(db: DbConnection): Promise { + await workers.get(db)?.close(); + workers.delete(db); +} diff --git a/apps/server/src/start-server.ts b/apps/server/src/start-server.ts index f6a84f53980..f101772f391 100644 --- a/apps/server/src/start-server.ts +++ b/apps/server/src/start-server.ts @@ -7,6 +7,7 @@ import { isLoopbackHostname } from "@bb/config/loopback"; import { toOptionalString } from "@bb/config/strings"; import { createLogger } from "@bb/logger"; import { getAppSettings, listRunningThreads } from "@bb/db"; +import { closeThreadSearch } from "./services/threads/thread-search.js"; import { initDb } from "./db.js"; import { createApp } from "./server.js"; import { PendingInteractionLifecycle } from "./services/interactions/pending-interactions.js"; @@ -394,6 +395,7 @@ export async function runServer(serverConfig: ServerConfig): Promise { return shutdownPromise; } shutdownPromise = (async () => { + await closeThreadSearch(db); serverMove.dispose(); appUpdate.dispose(); providerModelCatalogPrewarm?.stop(); diff --git a/apps/server/src/thread-search-worker.ts b/apps/server/src/thread-search-worker.ts new file mode 100644 index 00000000000..b964957389e --- /dev/null +++ b/apps/server/src/thread-search-worker.ts @@ -0,0 +1,22 @@ +import { parentPort, workerData } from "node:worker_threads"; +import { createConnection, searchThreadMatchRows } from "@bb/db"; +import { z } from "zod"; +import { threadSearchRequestSchema } from "./services/threads/thread-search-protocol.js"; + +const port = parentPort; +if (port === null) throw new Error("Thread search requires a worker port"); +const db = createConnection(z.string().min(1).parse(workerData), { + readonly: true, +}); +port.on("message", (message: unknown) => { + try { + const args = threadSearchRequestSchema.parse(message); + port.postMessage({ ok: true, rows: searchThreadMatchRows(db, args) }); + } catch (error) { + port.postMessage({ + ok: false, + error: error instanceof Error ? error.message : String(error), + }); + } +}); +port.on("close", () => db.$client.close()); diff --git a/apps/server/test/services/plugins/plugin-mention-providers.test.ts b/apps/server/test/services/plugins/plugin-mention-providers.test.ts index 389ccc08e53..6037b06fa62 100644 --- a/apps/server/test/services/plugins/plugin-mention-providers.test.ts +++ b/apps/server/test/services/plugins/plugin-mention-providers.test.ts @@ -261,6 +261,27 @@ describe("plugin mention providers (bb.ui.registerMentionProvider)", () => { expect(entry?.handlerStats.errorCount).toBe(1); }); + it("searches one provider without invoking unrelated providers", async () => { + const response = await harness.app.request( + `${BASE}/api/v1/plugins/mentions/search?q=fix&pluginId=mentions&providerId=issues`, + ); + expect(response.status).toBe(200); + const body = await response.json(); + expect(body.groups).toHaveLength(1); + expect(body.groups[0].providerId).toBe("issues"); + const entry = harness.pluginService + .list() + .find((plugin) => plugin.id === "mentions"); + expect(entry?.handlerStats.errorCount).toBe(0); + }); + + it("rejects incomplete provider selectors", async () => { + const response = await harness.app.request( + `${BASE}/api/v1/plugins/mentions/search?q=fix&providerId=issues`, + ); + expect(response.status).toBe(400); + }); + it("searches only providers registered for the requested trigger", async () => { const response = await harness.app.request( `${BASE}/api/v1/plugins/mentions/search?q=fix&trigger=%23&projectId=proj_1&threadId=thr_1`, @@ -721,6 +742,7 @@ describe("mention search time box", () => { query: "o", projectId: null, threadId: null, + provider: null, }); expect(groups).toEqual([ { diff --git a/apps/server/test/services/threads/thread-search.test.ts b/apps/server/test/services/threads/thread-search.test.ts new file mode 100644 index 00000000000..fd891bd467d --- /dev/null +++ b/apps/server/test/services/threads/thread-search.test.ts @@ -0,0 +1,65 @@ +import { mkdtemp, rm } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { + archiveThread, + createConnection, + createThread, + ensurePersonalProject, + migrate, + noopNotifier, + searchThreadsWithPendingInteractionState, + updateThread, +} from "@bb/db"; +import { expect, it } from "vitest"; +import { + closeThreadSearch, + searchThreads, +} from "../../../src/services/threads/thread-search.js"; + +it("searches persisted data without blocking timers and skips cancelled queued searches", async () => { + const dir = await mkdtemp(join(tmpdir(), "bb-thread-search-")); + const db = createConnection(join(dir, "bb.db")); + try { + migrate(db); + ensurePersonalProject(db); + const thread = createThread(db, noopNotifier, { + projectId: "proj_personal", + providerId: "codex", + title: "Search worker needle", + }); + const archived = createThread(db, noopNotifier, { + projectId: "proj_personal", + providerId: "codex", + title: "Archived needle", + }); + archiveThread(db, noopNotifier, archived.id); + const args = { query: "needle", limitPerGroup: 20 }; + const expected = searchThreadsWithPendingInteractionState(db, args); + let timerFired = false; + const timer = setTimeout(() => { + timerFired = true; + }, 0); + const first = searchThreads(db, args, new AbortController().signal); + const cancellation = new AbortController(); + const queued = searchThreads(db, args, cancellation.signal); + const cancelled = expect(queued).rejects.toMatchObject({ + name: "AbortError", + }); + cancellation.abort(); + expect(await first).toEqual(expected); + clearTimeout(timer); + expect(timerFired).toBe(true); + await cancelled; + updateThread(db, noopNotifier, thread.id, { title: "Renamed thread" }); + const updated = await searchThreads(db, args, new AbortController().signal); + expect(updated.active.total).toBe(0); + expect(updated.archived.results.map((result) => result.thread.id)).toEqual([ + archived.id, + ]); + } finally { + await closeThreadSearch(db); + db.$client.close(); + await rm(dir, { recursive: true, force: true }); + } +}); diff --git a/docs/search-performance.md b/docs/search-performance.md new file mode 100644 index 00000000000..bc465aba74f --- /dev/null +++ b/docs/search-performance.md @@ -0,0 +1,46 @@ +# Search performance comparison + +Run the manual CI workflow on a candidate branch with `search-profile-base` set +to its comparison commit. Builds and measurements run remotely: + +```sh +gh workflow run ci.yml --ref \ + -f search-profile-base= +``` + +The optional profiling job builds both revisions through Turbo, then runs their +packaged servers and enrolled host daemons sequentially on the same runner in +before/after/after/before order. It uses one SQLite fixture copied for each +process and one workspace containing 10,000 files. No production data or agents +are used. + +The artifact preserves the exact harness, revision and runtime archive hashes, +fixture identity, runtime versions, machine information, individual measurements, +result hashes, logs, CPU profiles, and a percentile summary. Download it with +`gh run download --name search-profile-`. CI retains +artifacts for 90 days; copy evidence that must outlive that retention window. + +| Measurement | Samples per revision | Meaning | +| ------------------------------------------ | -------------------- | ----------------------------------------------------------------------- | +| Server startup and first search | 10 | Fresh process and SQLite connection; OS page cache remains warm | +| Warm search and concurrent health response | 40 | Three query shapes; health request starts 15 ms after search | +| File suggestion bursts | 20 × 5 queries | Real HTTP, server-to-daemon transport, discovery, and ranking | +| Plugin mention first/all results | 20 | Real provider endpoints with deterministic 20 ms and 1,600 ms providers | + +File bursts wait 2.1 seconds between sequences and 80 ms between successive +queries. Each revision uses its own frontend request topology for plugin +mentions: one aggregate request before, independent provider requests after. +These API measurements exclude frontend debounce, DOM rendering, and remote +Connect latency. Use Browser Automation separately for keystroke-to-render +measurements; do not call API time an end-to-end UI latency. + +CPU profiles come from separate runs and are excluded from latency percentiles. +Open `.cpuprofile` files in a compatible profiler; `results/cpu-summary.json` +lists sampled self-time by function. Worker profiles must be considered with +the main-server profile: moving SQLite work to a worker can improve server +responsiveness without reducing total query CPU time. + +The harness verifies equivalent result hashes across revisions. Read individual +samples and query shapes alongside p50/p95 values. A slower packaged cold start, +stale results, missing profiles, or an unsuccessful process shutdown is evidence +to investigate, not a successful benchmark. diff --git a/packages/db/src/connection.ts b/packages/db/src/connection.ts index b3245220764..00df56d5cfe 100644 --- a/packages/db/src/connection.ts +++ b/packages/db/src/connection.ts @@ -16,6 +16,7 @@ export interface SlowDbQueryLogger { } export interface CreateConnectionOptions { + readonly?: boolean; slowQueryLogger?: SlowDbQueryLogger; slowQueryThresholdMs?: number; } @@ -156,10 +157,12 @@ export function createConnection( source: string | Buffer = "bb.db", options: CreateConnectionOptions = {}, ) { - const sqlite = new Database(source); + const sqlite = new Database(source, { readonly: options.readonly ?? false }); - sqlite.pragma("auto_vacuum = INCREMENTAL"); - sqlite.pragma("journal_mode = WAL"); + if (!options.readonly) { + sqlite.pragma("auto_vacuum = INCREMENTAL"); + sqlite.pragma("journal_mode = WAL"); + } sqlite.pragma("foreign_keys = ON"); sqlite.pragma("synchronous = NORMAL"); sqlite.pragma(`cache_size = -${SQLITE_CACHE_SIZE_KIB}`); diff --git a/packages/db/src/data/index.ts b/packages/db/src/data/index.ts index ae9a99f68f3..1d3b688410f 100644 --- a/packages/db/src/data/index.ts +++ b/packages/db/src/data/index.ts @@ -101,6 +101,8 @@ export { applyThreadLifecycleEventInTransaction, requireThreadLifecycleEventApplied, searchThreadsWithPendingInteractionState, + searchThreadMatchRows, + hydrateThreadSearchResults, THREAD_SEARCH_LIMIT_PER_GROUP_DEFAULT, THREAD_SEARCH_LIMIT_PER_GROUP_MAX, } from "./threads.js"; @@ -110,6 +112,7 @@ export type { ReorderPinnedThreadResult, RunningThreadRow, ThreadSearchResultGroup, + ThreadSearchMatchRow, ThreadWithPendingInteractionState, ThreadExecutionOverride, UpdateThreadInput, diff --git a/packages/db/src/data/threads.ts b/packages/db/src/data/threads.ts index 6f5ec83ef16..bfe0467e338 100644 --- a/packages/db/src/data/threads.ts +++ b/packages/db/src/data/threads.ts @@ -153,7 +153,7 @@ interface ListThreadSearchMatchRowsArgs { tokenMatchQueries: readonly string[]; } -interface ThreadSearchMatchRow { +export interface ThreadSearchMatchRow { archived: number; segmentOrder: number; sourceKind: string; @@ -1144,43 +1144,51 @@ function hydrateThreadSearchGroup( return { total: firstRow.total, results }; } -export function searchThreadsWithPendingInteractionState( +export function searchThreadMatchRows( db: DbConnection, args: SearchThreadsWithPendingInteractionStateArgs, -): ThreadSearchResults { +): ThreadSearchMatchRow[] { const tokens = listThreadSearchQueryTokens(args.query); const tokenMatchQueries = listThreadSearchTokenMatchQueries(tokens); - const anyTokenMatchQuery = - buildThreadSearchAnyTokenMatchQuery(tokenMatchQueries); - if (anyTokenMatchQuery === null) { - return { - active: { total: 0, results: [] }, - archived: { total: 0, results: [] }, - }; - } - const limitPerGroup = Math.min( - Math.max(args.limitPerGroup, 1), - THREAD_SEARCH_LIMIT_PER_GROUP_MAX, - ); - - const rows = listThreadSearchMatchRows(db, { + const anyTokenMatchQuery = buildThreadSearchAnyTokenMatchQuery(tokenMatchQueries); + if (anyTokenMatchQuery === null) return []; + return listThreadSearchMatchRows(db, { anyTokenMatchQuery, - limitPerGroup, tokenMatchQueries, + limitPerGroup: Math.min( + Math.max(args.limitPerGroup, 1), + THREAD_SEARCH_LIMIT_PER_GROUP_MAX, + ), }); +} +export function hydrateThreadSearchResults( + db: DbConnection, + args: { query: string; rows: readonly ThreadSearchMatchRow[] }, +): ThreadSearchResults { + const tokens = listThreadSearchQueryTokens(args.query); return { active: hydrateThreadSearchGroup(db, { tokens, - rows: rows.filter((row) => row.archived === 0), + rows: args.rows.filter((row) => row.archived === 0), }), archived: hydrateThreadSearchGroup(db, { tokens, - rows: rows.filter((row) => row.archived === 1), + rows: args.rows.filter((row) => row.archived === 1), }), }; } +export function searchThreadsWithPendingInteractionState( + db: DbConnection, + args: SearchThreadsWithPendingInteractionStateArgs, +): ThreadSearchResults { + return hydrateThreadSearchResults(db, { + query: args.query, + rows: searchThreadMatchRows(db, args), + }); +} + /** How `countThreads` buckets its result; omitted asks for the total only. */ export type CountThreadsGroupBy = "host" | "provider" | "project"; diff --git a/packages/db/test/data/thread-search.test.ts b/packages/db/test/data/thread-search.test.ts index 6226d59e357..85a54fd2f90 100644 --- a/packages/db/test/data/thread-search.test.ts +++ b/packages/db/test/data/thread-search.test.ts @@ -413,7 +413,7 @@ describe("thread search data", () => { } }); - it("limits and counts each group from one partitioned scan", () => { + it("keeps exact group totals and newest tied results at the limit", () => { const { db, project } = setup(); try { const activeIds: string[] = []; @@ -425,6 +425,7 @@ describe("thread search data", () => { title: `partitionneedle active ${index}`, }); activeIds.push(activeThread.id); + db.$client.prepare("UPDATE threads SET updated_at = ? WHERE id = ?").run(index, activeThread.id); const archivedThread = createThread(db, noopNotifier, { projectId: project.id, providerId: "codex", @@ -432,6 +433,7 @@ describe("thread search data", () => { }); archiveThread(db, noopNotifier, archivedThread.id); archivedIds.push(archivedThread.id); + db.$client.prepare("UPDATE threads SET updated_at = ? WHERE id = ?").run(index, archivedThread.id); } const results = searchThreadsWithPendingInteractionState(db, { @@ -443,6 +445,8 @@ describe("thread search data", () => { expect(results.archived.total).toBe(4); expect(results.active.results).toHaveLength(2); expect(results.archived.results).toHaveLength(2); + expect(results.active.results.map((result) => result.thread.id)).toEqual(activeIds.slice(-2).reverse()); + expect(results.archived.results.map((result) => result.thread.id)).toEqual(archivedIds.slice(-2).reverse()); for (const result of results.active.results) { expect(activeIds).toContain(result.thread.id); expect(result.thread.archivedAt).toBeNull(); diff --git a/packages/fuzzy-match/src/index.ts b/packages/fuzzy-match/src/index.ts index f19c5580038..5a800202844 100644 --- a/packages/fuzzy-match/src/index.ts +++ b/packages/fuzzy-match/src/index.ts @@ -553,30 +553,24 @@ function rankPlainQueryMatches( items: readonly NormalizedPathItem[], query: string, ): RankedPathMatch[] { - const tiebreakers: Tiebreaker>[] = [ - byPathStartAsc, - byPathLengthAsc, - ]; const matcher = new Fzf[]>(items, { selector: (item: NormalizedPathItem) => item.path, casing: "smart-case", forward: true, - tiebreakers, + sort: false, }); const matches: FzfResultItem>[] = matcher.find(query); - return matches - .map((match) => ({ - item: match.item.item, - path: match.item.path, - positions: [...match.positions].sort((left, right) => left - right), - score: - getRankedScore(PathIntentRank.PlainFzf, match.score) + - getPathRelevanceBonus(match.item.path, query), - start: match.start, - })) - .sort(compareRankedMatches); + return matches.map((match) => ({ + item: match.item.item, + path: match.item.path, + positions: [...match.positions].sort((left, right) => left - right), + score: + getRankedScore(PathIntentRank.PlainFzf, match.score) + + getPathRelevanceBonus(match.item.path, query), + start: match.start, + })); } function rankPathQueryMatches( diff --git a/scripts/search-profile.mjs b/scripts/search-profile.mjs new file mode 100644 index 00000000000..893917dd3d2 --- /dev/null +++ b/scripts/search-profile.mjs @@ -0,0 +1,609 @@ +import { spawn, execFileSync } from "node:child_process"; +import { createHash } from "node:crypto"; +import { createWriteStream } from "node:fs"; +import { + cp, + mkdir, + readFile, + readdir, + rm, + symlink, + writeFile, +} from "node:fs/promises"; +import { + availableParallelism, + cpus, + freemem, + loadavg, + platform, + totalmem, +} from "node:os"; +import { dirname, join, resolve } from "node:path"; +import { fileURLToPath, pathToFileURL } from "node:url"; + +const root = process.cwd(); +const output = resolve(".search-profile"); +const script = fileURLToPath(import.meta.url); +const sleep = (ms) => new Promise((done) => setTimeout(done, ms)); +const hash = (value) => createHash("sha256").update(value).digest("hex"); +const json = async (path, value) => { + await mkdir(dirname(path), { recursive: true }); + await writeFile(path, `${JSON.stringify(value, null, 2)}\n`); +}; +const command = (name, args) => + execFileSync(name, args, { cwd: root, stdio: "inherit" }); + +async function checkout(revision) { + for (let attempt = 0; ; attempt++) { + try { + execFileSync("git", ["checkout", "--detach", revision], { + cwd: root, + stdio: ["ignore", "inherit", "pipe"], + }); + return; + } catch (error) { + if (attempt >= 29 || !error.stderr?.toString().includes("index.lock")) + throw error; + await sleep(1000); + } + } +} + +async function seed() { + const { + createConnection, + migrate, + ensurePersonalProject, + createThread, + noopNotifier, + } = await import(pathToFileURL(resolve("packages/db/src/index.ts")).href); + const dir = join(output, "fixture"); + await mkdir(dir, { recursive: true }); + const db = createConnection(join(dir, "bb.db")); + migrate(db); + ensurePersonalProject(db); + const insert = db.$client.prepare( + "INSERT INTO thread_search_segments (id,thread_id,source_kind,source_key,source_seq,text,created_at,updated_at) VALUES (?,?,?,?,?,?,?,?)", + ); + const threadIds = []; + db.$client.transaction(() => { + for (let t = 0; t < 30; t++) { + const thread = createThread(db, noopNotifier, { + projectId: "proj_personal", + providerId: "codex", + title: `Search fixture ${t}`, + status: "idle", + }); + threadIds.push(thread.id); + db.$client + .prepare( + "UPDATE threads SET updated_at = ?, archived_at = ? WHERE id = ?", + ) + .run( + 1_700_000_000_000 + t, + t >= 15 ? 1_700_000_000_000 : null, + thread.id, + ); + for (let n = 0; n < 1800; n++) { + insert.run( + `segment-${t}-${n}`, + thread.id, + "user_message", + `event-${n}`, + n, + `search needle workspace component ${t} ${n} `.repeat(32), + 1_700_000_000_000, + 1_700_000_000_000, + ); + } + } + })(); + db.$client.pragma("wal_checkpoint(TRUNCATE)"); + db.$client.close(); + const workspace = join(dir, "workspace"); + await mkdir(workspace, { recursive: true }); + for (let d = 0; d < 100; d++) { + const path = join(workspace, `package-${String(d).padStart(3, "0")}`); + await mkdir(path, { recursive: true }); + await Promise.all( + Array.from({ length: 100 }, (_, f) => + writeFile( + join(path, `search-component-${String(f).padStart(3, "0")}.ts`), + "export {};\n", + ), + ), + ); + } + command("git", ["init", "--quiet", workspace]); + await json(join(output, "fixture-manifest.json"), { + threads: 30, + segments: 54000, + files: 10000, + directories: 100, + threadIds, + databaseSha256: hash(await readFile(join(dir, "bb.db"))), + }); +} + +function start(entry, env, label, profile) { + const log = createWriteStream(join(output, "logs", `${label}.log`)); + const args = profile + ? ["--cpu-prof", `--cpu-prof-dir=${join(output, "profiles", label)}`] + : []; + const child = spawn(process.execPath, [...args, entry], { + cwd: root, + env: { ...process.env, NODE_ENV: "production", ...env }, + stdio: ["ignore", "pipe", "pipe"], + }); + child.stdout.pipe(log); + child.stderr.pipe(log); + const exited = new Promise((done) => + child.once("exit", (code, signal) => done({ code, signal })), + ); + return { child, exited, log }; +} + +async function stop(process) { + if (process.child.exitCode === null) process.child.kill("SIGINT"); + const result = await Promise.race([ + process.exited, + sleep(15000).then(() => null), + ]); + if (result === null) { + process.child.kill("SIGKILL"); + await process.exited; + throw new Error("Profiled process did not stop cleanly"); + } + process.log.end(); +} + +async function request(url, body) { + const response = await fetch(url, { + signal: AbortSignal.timeout(60000), + ...(body === undefined + ? {} + : { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify(body), + }), + }); + const text = await response.text(); + if (!response.ok) + throw new Error(`${response.status} ${url}: ${text.slice(0, 1000)}`); + return JSON.parse(text); +} + +async function ready(server, url) { + for (let n = 0; n < 600; n++) { + if (server.child.exitCode !== null) + throw new Error(`Server exited: ${server.child.exitCode}`); + try { + await request(`${url}/health`); + return; + } catch {} + await sleep(100); + } + throw new Error("Server readiness timed out"); +} + +function resultProjection(value) { + return Object.fromEntries( + ["active", "archived"].map((group) => [ + group, + { + total: value[group].total, + results: value[group].results.map((result) => ({ + id: result.thread.id, + matches: result.matches, + })), + }, + ]), + ); +} + +async function searchSample(url, query) { + const start = performance.now(); + const search = request( + `${url}/api/v1/threads/search?${new URLSearchParams({ query, limitPerGroup: "20" })}`, + ).then((body) => ({ + ms: performance.now() - start, + body: resultProjection(body), + })); + await sleep(15); + const healthStart = performance.now(); + const health = request(`${url}/health`).then( + () => performance.now() - healthStart, + ); + const [result, healthMs] = await Promise.all([search, health]); + return { + query, + searchMs: result.ms, + healthMs, + resultSha256: hash(JSON.stringify(result.body)), + result: result.body, + }; +} + +async function measure(label, revision, round, profile = false) { + const name = `${label}-${round}${profile ? "-profile" : ""}`; + const build = join(output, "builds", label); + const samples = []; + const coldCount = profile ? 1 : 5; + for (let cold = 0; cold < coldCount; cold++) { + const run = `${name}-${cold}`; + const data = join(output, "runtime", run); + await mkdir(data, { recursive: true }); + await cp(join(output, "fixture", "bb.db"), join(data, "bb.db")); + await mkdir(join(output, "profiles", run), { recursive: true }); + const url = "http://127.0.0.1:29871"; + const env = { + BB_DATA_DIR: data, + BB_SERVER_PORT: "29871", + BB_SERVER_URL: url, + BB_APP_URL: url, + BB_SERVER_BIND_HOST: "127.0.0.1", + BB_TELEMETRY: "false", + BB_LOG_LEVEL: "warn", + BB_APP_UPDATE_MODE: "source", + }; + const startTime = performance.now(); + const server = start( + join(build, "server/dist/index.js"), + env, + run, + profile, + ); + let daemon; + try { + await ready(server, url); + const startupMs = performance.now() - startTime; + samples.push({ + kind: "cold-search", + startupMs, + ...(await searchSample(url, "search")), + }); + if (cold > 0) continue; + for (let i = 0; i < (profile ? 8 : 20); i++) { + samples.push({ + kind: "warm-search", + ...(await searchSample( + url, + ["search", "needle", "workspace component"][i % 3], + )), + }); + } + const enrollment = await request(`${url}/internal/hosts/enroll-key`, {}); + const daemonData = join(data, "daemon"); + await mkdir(join(output, "profiles", `${run}-daemon`), { + recursive: true, + }); + daemon = start( + join(build, "host-daemon/dist/daemon-bundle.mjs"), + { + ...env, + BB_DATA_DIR: daemonData, + BB_HOST_DAEMON_PORT: "29872", + BB_CLI_DIR: join(build, "host-daemon/dist"), + BB_BRIDGE_DIR: join(build, "host-daemon/dist"), + BB_HOST_ID: enrollment.hostId, + BB_HOST_ENROLL_KEY: enrollment.enrollKey, + BB_HOST_NAME: "Search profiling fixture", + BB_HOST_DAEMON_AUTO_UPDATE: "false", + }, + `${run}-daemon`, + profile, + ); + for (let i = 0; i < 300; i++) { + const hosts = await request(`${url}/api/v1/hosts`); + if ( + hosts.some( + (host) => + host.id === enrollment.hostId && host.status === "connected", + ) + ) + break; + if (daemon.child.exitCode !== null) + throw new Error("Fixture daemon exited before connecting"); + if (i === 299) throw new Error("Fixture daemon did not connect"); + await sleep(100); + } + for (let burst = 0; burst < (profile ? 5 : 10); burst++) { + await sleep(2100); + for (const [index, query] of [ + "se", + "sea", + "sear", + "searc", + "search", + ].entries()) { + const began = performance.now(); + const body = await request(`${url}/api/v1/files/paths`, { + hostId: enrollment.hostId, + path: join(output, "fixture/workspace"), + query, + limit: 16, + includeFiles: true, + includeDirectories: true, + includeHidden: false, + }); + samples.push({ + kind: index === 0 ? "cold-files" : "warm-files", + burst, + query, + ms: performance.now() - began, + resultSha256: hash(JSON.stringify(body)), + }); + await sleep(80); + } + } + const installed = await request(`${url}/api/v1/plugins/install`, { + source: join(output, "mention-plugin"), + }); + const pluginId = installed.plugin.id; + for (let i = 0; i < (profile ? 3 : 10); i++) { + const query = `arrival-${i}`; + const began = performance.now(); + const base = `${url}/api/v1/plugins/mentions/search?${new URLSearchParams({ q: query, trigger: "#" })}`; + if (label === "before") { + const body = await request(base); + samples.push({ + kind: "mentions", + query, + firstMs: performance.now() - began, + allMs: performance.now() - began, + resultSha256: hash(JSON.stringify(body.groups)), + }); + } else { + let firstMs; + const groups = await Promise.all( + ["fast", "slow"].map(async (providerId) => { + const body = await request( + `${base}&${new URLSearchParams({ pluginId, providerId })}`, + ); + firstMs ??= performance.now() - began; + return body.groups; + }), + ); + samples.push({ + kind: "mentions", + query, + firstMs, + allMs: performance.now() - began, + resultSha256: hash(JSON.stringify(groups.flat())), + }); + } + } + } finally { + await json(join(output, "results", `${name}.json`), { + label, + revision, + round, + profile, + samples, + }); + if (daemon) await stop(daemon); + await stop(server); + await rm(data, { recursive: true, force: true }); + } + } + await json(join(output, "results", `${name}.json`), { + label, + revision, + round, + profile, + samples, + }); +} + +async function compare(baseArg, headArg) { + const commit = (ref) => + execFileSync("git", ["rev-parse", "--verify", `${ref}^{commit}`], { + encoding: "utf8", + }).trim(); + const revisions = { before: commit(baseArg), after: commit(headArg) }; + for (const folder of [ + "logs", + "profiles", + "results", + "builds", + "runtime", + "mention-plugin", + ]) { + await mkdir(join(output, folder), { recursive: true }); + } + await json(join(output, "mention-plugin/package.json"), { + name: "search-profile-mentions", + version: "1.0.0", + bb: { + name: "Search profiling", + description: "Deterministic mention providers for search measurements", + branding: { icon: "Search" }, + server: "./server.js", + }, + }); + await writeFile( + join(output, "mention-plugin/server.js"), + `export default function(bb) { + for (const [id, delay] of [["fast", 20], ["slow", 1600]]) bb.ui.registerMentionProvider({ + id, label: id, triggers: ["#"], + async search(ctx) { await new Promise(r => setTimeout(r, delay)); return [{ id: ctx.query, title: ctx.query + " " + id }]; }, + async resolve(id) { return { context: id }; } + }); + }\n`, + ); + const manifest = { + revisions, + node: process.version, + platform: platform(), + cpus: cpus().map((cpu) => cpu.model), + parallelism: availableParallelism(), + totalMemory: totalmem(), + freeMemory: freemem(), + loadAverage: loadavg(), + startedAt: new Date().toISOString(), + order: ["before", "after", "after", "before"], + timing: + "Unprofiled; cold means fresh process and SQLite connection, OS page cache is not flushed", + mentions: + "Real endpoint requests using each revision's frontend request topology; excludes DOM rendering and debounce", + files: + "Real HTTP route and enrolled packaged daemon against 10,000 files; 2.1 s between bursts, 80 ms between queries", + artifacts: {}, + }; + for (const [label, revision] of Object.entries(revisions)) { + await checkout(revision); + command("pnpm", ["install", "--frozen-lockfile", "--prefer-offline"]); + if (label === "before") + command(process.execPath, [ + "--conditions=source", + "--import", + "tsx", + script, + "seed", + ]); + command("pnpm", [ + "exec", + "turbo", + "run", + "build", + "--filter=bb-app", + "--concurrency=4", + "--output-logs=errors-only", + ]); + const source = resolve("packages/bb-app"); + const target = join(output, "builds", label); + await mkdir(target, { recursive: true }); + for (const name of [ + "package.json", + "dist", + "server", + "host-daemon", + "app", + ]) { + await cp(join(source, name), join(target, name), { recursive: true }); + } + await symlink( + join(source, "node_modules"), + join(target, "node_modules"), + "dir", + ); + command("tar", [ + "-czf", + join(output, `${label}-runtime.tgz`), + "--exclude=node_modules", + "-C", + target, + ".", + ]); + manifest.artifacts[label] = { + lockfileSha256: hash(await readFile("pnpm-lock.yaml")), + runtimeSha256: hash(await readFile(join(output, `${label}-runtime.tgz`))), + }; + } + manifest.fixture = JSON.parse( + await readFile(join(output, "fixture-manifest.json"), "utf8"), + ); + command("tar", [ + "-czf", + join(output, "fixture.tgz"), + "-C", + output, + "fixture", + "mention-plugin", + ]); + manifest.fixture.archiveSha256 = hash( + await readFile(join(output, "fixture.tgz")), + ); + await json(join(output, "manifest.json"), manifest); + for (const [round, label] of manifest.order.entries()) + await measure(label, revisions[label], round); + for (const label of ["before", "after"]) + await measure(label, revisions[label], "cpu", true); + await summarize(); + const profiles = []; + for (const directory of await readdir(join(output, "profiles"))) { + for (const file of await readdir(join(output, "profiles", directory))) { + if (!file.endsWith(".cpuprofile")) continue; + const profile = JSON.parse( + await readFile(join(output, "profiles", directory, file), "utf8"), + ); + const nodes = new Map(profile.nodes.map((node) => [node.id, node])); + const costs = new Map(); + for (const [index, sample] of profile.samples.entries()) { + const frame = nodes.get(sample).callFrame; + const key = `${frame.functionName} ${frame.url}:${frame.lineNumber + 1}`; + costs.set( + key, + (costs.get(key) ?? 0) + profile.timeDeltas[index] / 1000, + ); + } + profiles.push({ + file: `${directory}/${file}`, + durationMs: (profile.endTime - profile.startTime) / 1000, + hottest: [...costs.entries()].sort((a, b) => b[1] - a[1]).slice(0, 30), + }); + } + } + await json(join(output, "results/cpu-summary.json"), profiles); + await checkout(revisions.after); +} + +async function summarize() { + const runs = await Promise.all( + (await readdir(join(output, "results"))) + .filter((name) => name.endsWith(".json")) + .map(async (name) => + JSON.parse(await readFile(join(output, "results", name), "utf8")), + ), + ); + const percentile = (values, p) => + [...values].sort((a, b) => a - b)[Math.ceil(values.length * p) - 1]; + const rows = [ + "| Measurement | Before p50 / p95 (ms) | After p50 / p95 (ms) | n per revision |", + "| --- | --- | --- | --- |", + ]; + for (const [kind, field] of [ + ["cold-search", "startupMs"], + ["cold-search", "searchMs"], + ["warm-search", "searchMs"], + ["warm-search", "healthMs"], + ["cold-files", "ms"], + ["warm-files", "ms"], + ["mentions", "firstMs"], + ["mentions", "allMs"], + ]) { + const values = ["before", "after"].map((label) => + runs + .filter((run) => run.label === label && !run.profile) + .flatMap((run) => run.samples) + .filter((sample) => sample.kind === kind) + .map((sample) => sample[field]), + ); + rows.push( + `| ${kind} ${field} | ${values.map((v) => `${percentile(v, 0.5).toFixed(1)} / ${percentile(v, 0.95).toFixed(1)}`).join(" | ")} | ${values[0].length} / ${values[1].length} |`, + ); + } + const hashes = new Map(); + for (const run of runs) + for (const sample of run.samples) { + const key = `${sample.kind.replace("cold-", "").replace("warm-", "")}:${sample.query}`; + const existing = hashes.get(key); + if (existing && existing !== sample.resultSha256) + throw new Error(`Result mismatch: ${key}`); + hashes.set(key, sample.resultSha256); + } + await writeFile( + join(output, "summary.md"), + `${rows.join("\n")}\n\nAll comparable result hashes matched. CPU-profile runs are excluded from timings.\n`, + ); +} + +const [mode, ...args] = process.argv.slice(2); +if (mode === "seed") await seed(); +else if (mode === "compare" && args.length === 2) await compare(...args); +else if (mode === "summarize") await summarize(); +else + throw new Error( + "Usage: search-profile.mjs compare ", + );