diff --git a/apps/server/src/orchestration/ThreadPullRequestReactor.test.ts b/apps/server/src/orchestration/ThreadPullRequestReactor.test.ts index 6e8b930a1ae6..cefd5bcb42b4 100644 --- a/apps/server/src/orchestration/ThreadPullRequestReactor.test.ts +++ b/apps/server/src/orchestration/ThreadPullRequestReactor.test.ts @@ -406,6 +406,51 @@ describe("ThreadPullRequestReactor", () => { ), ); + it.effect("refreshes the project identity when a turn adds the remote", () => + Effect.scoped( + Effect.gen(function* () { + const current = thread("new-remote"); + const fixture = yield* makeHarness({ + threads: [current], + project: { ...project, repositoryIdentity: null }, + branchPullRequest: () => Effect.succeed(branchPullRequest()), + resolveRepositoryIdentity: (_cwd, options) => + Effect.succeed(options?.refresh ? project.repositoryIdentity : null), + }); + yield* Effect.gen(function* () { + const reactor = yield* fixture.start(); + expect(yield* Ref.get(fixture.commands)).toHaveLength(0); + + yield* fixture.publish({ + type: "thread.turn-diff-completed", + sequence: 2, + eventId: EventId.make("checkpoint-finished"), + aggregateKind: "thread", + aggregateId: current.id, + occurredAt: NOW, + commandId: null, + causationEventId: null, + correlationId: null, + metadata: {}, + payload: { + threadId: current.id, + turnId: TurnId.make("turn"), + checkpointTurnCount: 1, + checkpointRef: CheckpointRef.make("checkpoint"), + status: "ready", + files: [], + assistantMessageId: null, + completedAt: NOW, + }, + }); + yield* Queue.take(fixture.reads); + yield* reactor.drain; + expect((yield* Ref.get(fixture.commands))[0]?.branchPullRequest).toEqual(reference(42)); + }).pipe(Effect.provide(fixture.layer)); + }), + ), + ); + it.effect("uses live worktrees and falls back to the project for removed worktrees", () => Effect.scoped( Effect.gen(function* () { diff --git a/apps/server/src/orchestration/ThreadPullRequestReactor.ts b/apps/server/src/orchestration/ThreadPullRequestReactor.ts index 0d709a867e07..2d9d92d64364 100644 --- a/apps/server/src/orchestration/ThreadPullRequestReactor.ts +++ b/apps/server/src/orchestration/ThreadPullRequestReactor.ts @@ -164,8 +164,19 @@ export const make = Effect.gen(function* () { (group) => Effect.gen(function* () { const first = group[0]!; - const project = projects.get(first.projectId); - if (project === undefined) return finishBackfill(group); + const snapshotProject = projects.get(first.projectId); + if (snapshotProject === undefined) return finishBackfill(group); + // A finished turn may have added the remote this PR lives on. A failed + // refresh resolves to null, so keep the snapshot's identity then. + const project = request.refresh + ? { + ...snapshotProject, + repositoryIdentity: + (yield* repositoryIdentities.resolve(snapshotProject.workspaceRoot, { + refresh: true, + })) ?? snapshotProject.repositoryIdentity, + } + : snapshotProject; const repository = sourceControlRepositorySelector(project.repositoryIdentity); if (first.branch !== null && repository === null) return finishBackfill(group); const worktreeExists = diff --git a/apps/server/src/project/RepositoryIdentityResolver.test.ts b/apps/server/src/project/RepositoryIdentityResolver.test.ts index d6ddb0b9263f..58f199b834e2 100644 --- a/apps/server/src/project/RepositoryIdentityResolver.test.ts +++ b/apps/server/src/project/RepositoryIdentityResolver.test.ts @@ -94,6 +94,8 @@ it.layer(NodeServices.layer)("RepositoryIdentityResolverLive", (it) => { const resolver = yield* RepositoryIdentityResolver.RepositoryIdentityResolver; const first = yield* resolver.resolve("/repo/packages/web"); rootPath = "/repo/packages/web"; + // Longer than the one-minute cadence of the background sweeps. + yield* TestClock.adjust(Duration.minutes(10)); const second = yield* resolver.resolve("/repo/packages/web"); expect(first?.canonicalKey).toBe("github.com/t3tools/t3code"); @@ -123,10 +125,10 @@ it.layer(NodeServices.layer)("RepositoryIdentityResolverLive", (it) => { const unavailable = yield* resolver.resolve(rootPath, { refresh: true }); expect(unavailable?.webUrl).toBeUndefined(); expect(unavailable?.canonicalKey).toBe("ssh.forge.test/team/repo"); - }).pipe(Effect.provide(resolverLayer)); + }).pipe(Effect.provide(Layer.merge(TestClock.layer(), resolverLayer))); }); - it.effect("retries Git root discovery after a failed lookup", () => { + it.effect("retries Git root discovery after the negative TTL", () => { const calls: Array> = []; let rootAttempts = 0; const processRunner = Layer.succeed(ProcessRunner.ProcessRunner, { @@ -159,7 +161,9 @@ it.layer(NodeServices.layer)("RepositoryIdentityResolverLive", (it) => { return Effect.gen(function* () { const resolver = yield* RepositoryIdentityResolver.RepositoryIdentityResolver; expect(yield* resolver.resolve("/repo/packages/web")).toBeNull(); + expect(yield* resolver.resolve("/repo/packages/web")).toBeNull(); + yield* TestClock.adjust(Duration.minutes(1)); const recovered = yield* resolver.resolve("/repo/packages/web"); expect(recovered?.rootPath).toBe("/repo"); expect(calls).toEqual([ @@ -167,7 +171,7 @@ it.layer(NodeServices.layer)("RepositoryIdentityResolverLive", (it) => { ["-C", "/repo/packages/web", "rev-parse", "--show-toplevel"], ["-C", "/repo", "remote", "-v"], ]); - }).pipe(Effect.provide(resolverLayer)); + }).pipe(Effect.provide(Layer.merge(TestClock.layer(), resolverLayer))); }); it.effect("normalizes equivalent GitHub remotes into a stable repository identity", () => diff --git a/apps/server/src/project/RepositoryIdentityResolver.ts b/apps/server/src/project/RepositoryIdentityResolver.ts index 2d7f5d02d02e..5acafa47e2e2 100644 --- a/apps/server/src/project/RepositoryIdentityResolver.ts +++ b/apps/server/src/project/RepositoryIdentityResolver.ts @@ -13,7 +13,11 @@ import * as Layer from "effect/Layer"; import * as ProcessRunner from "../processRunner.ts"; const DEFAULT_REPOSITORY_IDENTITY_CACHE_CAPACITY = 512; -const DEFAULT_POSITIVE_CACHE_TTL = Duration.minutes(1); +// Background sweeps resolve every project each minute. A long TTL keeps them +// from spawning git each time. Clone, publish, and PR discovery (after a turn +// and before it saves links) resolve with `refresh: true`. +const DEFAULT_POSITIVE_CACHE_TTL = Duration.minutes(15); +// Short, so a folder that gains a repository or a remote shows up quickly. const DEFAULT_NEGATIVE_CACHE_TTL = Duration.minutes(1); export interface RepositoryIdentityResolverOptions { @@ -142,20 +146,23 @@ export const make = Effect.fn("RepositoryIdentityResolver.make")(function* ( const processRunner = yield* ProcessRunner.ProcessRunner; const cacheCapacity = options.cacheCapacity ?? DEFAULT_REPOSITORY_IDENTITY_CACHE_CAPACITY; const refine = options.refine ?? Effect.succeed; + // Git errors and timeouts resolve to null, so they use the negative TTL like + // "no repository" or "no remote". Only interrupts and defects skip the cache. + const timeToLive = (exit: Exit.Exit) => + Exit.match(exit, { + onSuccess: (value) => + value === null + ? (options.negativeCacheTtl ?? DEFAULT_NEGATIVE_CACHE_TTL) + : (options.positiveCacheTtl ?? DEFAULT_POSITIVE_CACHE_TTL), + onFailure: () => Duration.zero, + }); const repositoryRootCache = yield* Cache.makeWith( (cwd) => resolveRepositoryIdentityCacheKey(cwd).pipe( Effect.provideService(ProcessRunner.ProcessRunner, processRunner), ), - { - capacity: cacheCapacity, - timeToLive: Exit.match({ - onSuccess: (value) => - value === null ? Duration.zero : (options.positiveCacheTtl ?? DEFAULT_POSITIVE_CACHE_TTL), - onFailure: () => Duration.zero, - }), - }, + { capacity: cacheCapacity, timeToLive }, ); const repositoryIdentityCache = yield* Cache.makeWith( @@ -167,27 +174,20 @@ export const make = Effect.fn("RepositoryIdentityResolver.make")(function* ( (identity) => refine(identity).pipe(Effect.orElseSucceed(() => identity)), ), ), - { - capacity: cacheCapacity, - timeToLive: Exit.match({ - onSuccess: (value) => - value === null - ? (options.negativeCacheTtl ?? DEFAULT_NEGATIVE_CACHE_TTL) - : (options.positiveCacheTtl ?? DEFAULT_POSITIVE_CACHE_TTL), - onFailure: () => Duration.zero, - }), - }, + { capacity: cacheCapacity, timeToLive }, ); - const resolve: RepositoryIdentityResolver["Service"]["resolve"] = Effect.fn( - "RepositoryIdentityResolver.resolve", - )(function* (cwd, options) { - if (options?.refresh) yield* Cache.invalidate(repositoryRootCache, cwd); - const cacheKey = yield* Cache.get(repositoryRootCache, cwd); - if (cacheKey === null) return null; - if (options?.refresh) yield* Cache.invalidate(repositoryIdentityCache, cacheKey); - return yield* Cache.get(repositoryIdentityCache, cacheKey); - }); + // Untraced because almost every call is a cache hit. The lookups that spawn + // git keep their own spans. + const resolve: RepositoryIdentityResolver["Service"]["resolve"] = Effect.fnUntraced( + function* (cwd, options) { + if (options?.refresh) yield* Cache.invalidate(repositoryRootCache, cwd); + const cacheKey = yield* Cache.get(repositoryRootCache, cwd); + if (cacheKey === null) return null; + if (options?.refresh) yield* Cache.invalidate(repositoryIdentityCache, cacheKey); + return yield* Cache.get(repositoryIdentityCache, cacheKey); + }, + ); return RepositoryIdentityResolver.of({ resolve }); }); diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index 830bb35b86db..077087a7d84e 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -3038,9 +3038,13 @@ const makeWsRpcLayer = ( [WS_METHODS.sourceControlPublishRepository]: (input) => observeRpcEffect( WS_METHODS.sourceControlPublishRepository, - sourceControlRepositories - .publishRepository(input) - .pipe(Effect.tap(() => refreshGitStatus(input.cwd))), + sourceControlRepositories.publishRepository(input).pipe( + // A new remote can change the cached identity. Only the `cwd` entry + // refreshes, so after a publish from a linked worktree the project + // root entry waits for its TTL. + Effect.tap(() => repositoryIdentityResolver.resolve(input.cwd, { refresh: true })), + Effect.tap(() => refreshGitStatus(input.cwd)), + ), { "rpc.aggregate": "source-control", },