diff --git a/apps/server/src/orchestration/ThreadPullRequestReactor.test.ts b/apps/server/src/orchestration/ThreadPullRequestReactor.test.ts index d6572083dad2..003bfb7bae0e 100644 --- a/apps/server/src/orchestration/ThreadPullRequestReactor.test.ts +++ b/apps/server/src/orchestration/ThreadPullRequestReactor.test.ts @@ -388,6 +388,93 @@ 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("keeps the snapshot identity when a turn-end refresh fails", () => + Effect.scoped( + Effect.gen(function* () { + const current = thread("refresh-failed"); + let refreshes = 0; + let pullRequestOpened = false; + const fixture = yield* makeHarness({ + threads: [current], + branchPullRequest: (_input, options) => + Effect.sync(() => { + if (options?.refresh) pullRequestOpened = true; + return pullRequestOpened ? branchPullRequest() : null; + }), + // Only the turn-end refresh fails; the pre-save recheck succeeds. + resolveRepositoryIdentity: (_cwd, options) => + Effect.sync(() => + options?.refresh && refreshes++ === 0 ? null : project.repositoryIdentity, + ), + }); + yield* Effect.gen(function* () { + const reactor = yield* fixture.start(); + yield* fixture.publish({ + type: "thread.unsettled", + sequence: 2, + eventId: EventId.make("unsettled"), + aggregateKind: "thread", + aggregateId: current.id, + occurredAt: NOW, + commandId: null, + causationEventId: null, + correlationId: null, + metadata: {}, + payload: { threadId: current.id, reason: "activity", updatedAt: 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 17efc74e4f40..27aca2d24686 100644 --- a/apps/server/src/orchestration/ThreadPullRequestReactor.ts +++ b/apps/server/src/orchestration/ThreadPullRequestReactor.ts @@ -131,8 +131,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..6a41dfaca142 100644 --- a/apps/server/src/project/RepositoryIdentityResolver.test.ts +++ b/apps/server/src/project/RepositoryIdentityResolver.test.ts @@ -170,6 +170,75 @@ it.layer(NodeServices.layer)("RepositoryIdentityResolverLive", (it) => { }).pipe(Effect.provide(resolverLayer)); }); + it.effect("retries a failed remote lookup instead of caching no repository", () => { + let remoteAttempts = 0; + const processRunner = Layer.succeed(ProcessRunner.ProcessRunner, { + run: (input) => + Effect.sync(() => { + const rootLookup = input.args.includes("rev-parse"); + const failed = !rootLookup && remoteAttempts++ === 0; + return { + stdout: rootLookup + ? "/repo\n" + : failed + ? "" + : "origin\tgit@github.com:T3Tools/t3code.git (fetch)\n", + stderr: failed ? "temporary Git failure" : "", + code: ChildProcessSpawner.ExitCode(failed ? 1 : 0), + timedOut: false, + stdoutTruncated: false, + stderrTruncated: false, + stdoutInvalidUtf8: false, + stderrInvalidUtf8: false, + }; + }), + }); + const resolverLayer = Layer.effect( + RepositoryIdentityResolver.RepositoryIdentityResolver, + RepositoryIdentityResolver.make(), + ).pipe(Layer.provide(processRunner)); + + return Effect.gen(function* () { + const resolver = yield* RepositoryIdentityResolver.RepositoryIdentityResolver; + expect(yield* resolver.resolve("/repo")).toBeNull(); + expect((yield* resolver.resolve("/repo"))?.canonicalKey).toBe("github.com/t3tools/t3code"); + }).pipe(Effect.provide(resolverLayer)); + }); + + it.effect("caches non-git roots until refreshed", () => { + const calls: Array> = []; + const processRunner = Layer.succeed(ProcessRunner.ProcessRunner, { + run: (input) => + Effect.sync(() => { + calls.push(input.args); + return { + stdout: "", + stderr: "fatal: not a git repository (or any of the parent directories): .git\n", + code: ChildProcessSpawner.ExitCode(128), + timedOut: false, + stdoutTruncated: false, + stderrTruncated: false, + stdoutInvalidUtf8: false, + stderrInvalidUtf8: false, + }; + }), + }); + const resolverLayer = Layer.effect( + RepositoryIdentityResolver.RepositoryIdentityResolver, + RepositoryIdentityResolver.make(), + ).pipe(Layer.provide(processRunner)); + + return Effect.gen(function* () { + const resolver = yield* RepositoryIdentityResolver.RepositoryIdentityResolver; + expect(yield* resolver.resolve("/scratch")).toBeNull(); + expect(yield* resolver.resolve("/scratch")).toBeNull(); + expect(calls).toHaveLength(1); + + expect(yield* resolver.resolve("/scratch", { refresh: true })).toBeNull(); + expect(calls).toHaveLength(2); + }).pipe(Effect.provide(resolverLayer)); + }); + it.effect("normalizes equivalent GitHub remotes into a stable repository identity", () => Effect.gen(function* () { const fileSystem = yield* FileSystem.FileSystem; diff --git a/apps/server/src/project/RepositoryIdentityResolver.ts b/apps/server/src/project/RepositoryIdentityResolver.ts index 2d7f5d02d02e..b758f83d2fa7 100644 --- a/apps/server/src/project/RepositoryIdentityResolver.ts +++ b/apps/server/src/project/RepositoryIdentityResolver.ts @@ -9,12 +9,32 @@ import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as Exit from "effect/Exit"; import * as Layer from "effect/Layer"; +import * as Schema from "effect/Schema"; import * as ProcessRunner from "../processRunner.ts"; const DEFAULT_REPOSITORY_IDENTITY_CACHE_CAPACITY = 512; -const DEFAULT_POSITIVE_CACHE_TTL = Duration.minutes(1); -const DEFAULT_NEGATIVE_CACHE_TTL = Duration.minutes(1); +// Settlement and pull request sweeps resolve every project each minute, so these +// must outlast the sweep interval or every sweep respawns git per project. +// Callers that know the repository changed pass `refresh: true`. +const DEFAULT_POSITIVE_CACHE_TTL = Duration.minutes(10); +const DEFAULT_NEGATIVE_CACHE_TTL = Duration.minutes(5); + +// Lookup failures stay uncached so the next resolve retries them. +class RepositoryIdentityLookupError extends Schema.TaggedError()( + "RepositoryIdentityLookupError", + { + stage: Schema.Literals(["root", "remote"]), + cwd: Schema.String, + exitCode: Schema.optional(Schema.NullOr(Schema.Number)), + timedOut: Schema.optional(Schema.Boolean), + cause: Schema.optional(Schema.Defect()), + }, +) { + override get message(): string { + return `Git ${this.stage} lookup failed in '${this.cwd}'`; + } +} export interface RepositoryIdentityResolverOptions { readonly cacheCapacity?: number; @@ -103,14 +123,29 @@ const resolveRepositoryIdentityCacheKey = Effect.fn("RepositoryIdentityResolver. .run({ command: "git", args: ["-C", cwd, "rev-parse", "--show-toplevel"], + // Keep git's "not a git repository" diagnostic untranslated. + env: { LC_ALL: "C" }, timeoutBehavior: "timedOutResult", }) - .pipe(Effect.option); - if (topLevelResult._tag === "None" || topLevelResult.value.code !== 0) { - return null; + .pipe( + Effect.mapError( + (cause) => new RepositoryIdentityLookupError({ stage: "root", cwd, cause }), + ), + ); + if (topLevelResult.code !== 0) { + // Remember git's "not a repository" verdict; retry timeouts and other failures. + if (!topLevelResult.timedOut && /not a git repository/i.test(topLevelResult.stderr)) { + return null; + } + return yield* new RepositoryIdentityLookupError({ + stage: "root", + cwd, + exitCode: topLevelResult.code, + timedOut: topLevelResult.timedOut, + }); } - const candidate = topLevelResult.value.stdout.trim(); + const candidate = topLevelResult.stdout.trim(); return candidate.length > 0 ? candidate : null; }, ); @@ -119,7 +154,11 @@ const resolveRepositoryIdentityFromCacheKey = Effect.fn( "RepositoryIdentityResolver.resolveFromCacheKey", )(function* ( cacheKey: string, -): Effect.fn.Return { +): Effect.fn.Return< + RepositoryIdentity | null, + RepositoryIdentityLookupError, + ProcessRunner.ProcessRunner +> { const processRunner = yield* ProcessRunner.ProcessRunner; const remoteResult = yield* processRunner .run({ @@ -127,12 +166,22 @@ const resolveRepositoryIdentityFromCacheKey = Effect.fn( args: ["-C", cacheKey, "remote", "-v"], timeoutBehavior: "timedOutResult", }) - .pipe(Effect.option); - if (remoteResult._tag === "None" || remoteResult.value.code !== 0) { - return null; + .pipe( + Effect.mapError( + (cause) => new RepositoryIdentityLookupError({ stage: "remote", cwd: cacheKey, cause }), + ), + ); + if (remoteResult.code !== 0) { + // A repository without remotes exits 0 with no output; failures retry. + return yield* new RepositoryIdentityLookupError({ + stage: "remote", + cwd: cacheKey, + exitCode: remoteResult.code, + timedOut: remoteResult.timedOut, + }); } - const remote = pickPrimaryRemote(parseRemoteFetchUrls(remoteResult.value.stdout)); + const remote = pickPrimaryRemote(parseRemoteFetchUrls(remoteResult.stdout)); return remote ? buildRepositoryIdentity({ ...remote, rootPath: cacheKey }) : null; }); @@ -143,7 +192,11 @@ export const make = Effect.fn("RepositoryIdentityResolver.make")(function* ( const cacheCapacity = options.cacheCapacity ?? DEFAULT_REPOSITORY_IDENTITY_CACHE_CAPACITY; const refine = options.refine ?? Effect.succeed; - const repositoryRootCache = yield* Cache.makeWith( + const repositoryRootCache = yield* Cache.makeWith< + string, + string | null, + RepositoryIdentityLookupError + >( (cwd) => resolveRepositoryIdentityCacheKey(cwd).pipe( Effect.provideService(ProcessRunner.ProcessRunner, processRunner), @@ -152,13 +205,19 @@ export const make = Effect.fn("RepositoryIdentityResolver.make")(function* ( capacity: cacheCapacity, timeToLive: Exit.match({ onSuccess: (value) => - value === null ? Duration.zero : (options.positiveCacheTtl ?? DEFAULT_POSITIVE_CACHE_TTL), + value === null + ? (options.negativeCacheTtl ?? DEFAULT_NEGATIVE_CACHE_TTL) + : (options.positiveCacheTtl ?? DEFAULT_POSITIVE_CACHE_TTL), onFailure: () => Duration.zero, }), }, ); - const repositoryIdentityCache = yield* Cache.makeWith( + const repositoryIdentityCache = yield* Cache.makeWith< + string, + RepositoryIdentity | null, + RepositoryIdentityLookupError + >( (cacheKey) => resolveRepositoryIdentityFromCacheKey(cacheKey).pipe( Effect.provideService(ProcessRunner.ProcessRunner, processRunner), @@ -183,10 +242,14 @@ export const make = Effect.fn("RepositoryIdentityResolver.make")(function* ( "RepositoryIdentityResolver.resolve", )(function* (cwd, options) { if (options?.refresh) yield* Cache.invalidate(repositoryRootCache, cwd); - const cacheKey = yield* Cache.get(repositoryRootCache, cwd); + const cacheKey = yield* Cache.get(repositoryRootCache, cwd).pipe( + Effect.orElseSucceed(() => null), + ); if (cacheKey === null) return null; if (options?.refresh) yield* Cache.invalidate(repositoryIdentityCache, cacheKey); - return yield* Cache.get(repositoryIdentityCache, cacheKey); + return yield* Cache.get(repositoryIdentityCache, cacheKey).pipe( + Effect.orElseSucceed(() => null), + ); }); return RepositoryIdentityResolver.of({ resolve }); diff --git a/apps/server/src/server.test.ts b/apps/server/src/server.test.ts index 12f37780701d..3d71c4e5d7b8 100644 --- a/apps/server/src/server.test.ts +++ b/apps/server/src/server.test.ts @@ -7649,6 +7649,128 @@ it.layer(NodeServices.layer)("server router seam", (it) => { }).pipe(Effect.provide(NodeHttpServer.layerTest)), ); + it.effect.each(["init", "publish"] as const)( + "re-emits the project after %s adds its repository", + (action) => + Effect.gen(function* () { + const dispatched: Array = []; + const refreshedRoots: Array = []; + const metaUpdateDispatched = yield* Deferred.make(); + const project = makeDefaultOrchestrationReadModel().projects[0]!; + + yield* buildAppUnderTest({ + layers: { + orchestrationEngine: { + dispatch: (command) => + Effect.sync(() => { + dispatched.push(command.type); + return { sequence: dispatched.length }; + }).pipe( + Effect.tap(() => + command.type === "project.meta.update" + ? Deferred.succeed(metaUpdateDispatched, undefined) + : Effect.void, + ), + ), + }, + projectionSnapshotQuery: { + getActiveProjectByWorkspaceRoot: (workspaceRoot) => + Effect.succeed( + workspaceRoot === project.workspaceRoot + ? Option.some({ ...project, repositoryIdentity: null }) + : Option.none(), + ), + }, + repositoryIdentityResolver: { + resolve: (cwd, options) => + Effect.sync(() => { + if (options?.refresh) refreshedRoots.push(cwd); + return { + canonicalKey: "github.com/t3tools/t3code", + locator: { + source: "git-remote", + remoteName: "origin", + remoteUrl: "git@github.com:t3tools/t3code.git", + }, + } as const; + }), + }, + sourceControlRepositoryService: { + publishRepository: () => + Effect.succeed({ + repository: { + provider: "github", + nameWithOwner: "t3tools/t3code", + url: "https://github.com/t3tools/t3code", + sshUrl: "git@github.com:t3tools/t3code.git", + }, + remoteName: "origin", + remoteUrl: "git@github.com:t3tools/t3code.git", + branch: "main", + status: "pushed", + }), + }, + }, + }); + + const wsUrl = yield* getWsServerUrl("/ws"); + yield* Effect.scoped( + withWsRpcClient(wsUrl, (client) => + Effect.gen(function* () { + if (action === "init") { + yield* client[WS_METHODS.vcsInit]({ cwd: project.workspaceRoot }); + } else { + yield* client[WS_METHODS.sourceControlPublishRepository]({ + cwd: project.workspaceRoot, + provider: "github", + repository: "t3tools/t3code", + visibility: "private", + }); + } + }), + ), + ); + yield* Deferred.await(metaUpdateDispatched); + assert.deepEqual(refreshedRoots, [project.workspaceRoot]); + assert.deepEqual(dispatched, ["project.meta.update"]); + }).pipe(Effect.provide(NodeHttpServer.layerTest)), + ); + + it.effect("does not re-emit the project when the identity refresh finds nothing", () => + Effect.gen(function* () { + const dispatched: Array = []; + const project = makeDefaultOrchestrationReadModel().projects[0]!; + + yield* buildAppUnderTest({ + layers: { + orchestrationEngine: { + dispatch: (command) => + Effect.sync(() => { + dispatched.push(command.type); + return { sequence: dispatched.length }; + }), + }, + projectionSnapshotQuery: { + getActiveProjectByWorkspaceRoot: () => + Effect.succeed(Option.some({ ...project, repositoryIdentity: null })), + }, + repositoryIdentityResolver: { + resolve: () => Effect.succeed(null), + }, + }, + }); + + const wsUrl = yield* getWsServerUrl("/ws"); + // The refresh runs before the RPC returns, so any dispatch has landed. + yield* Effect.scoped( + withWsRpcClient(wsUrl, (client) => + client[WS_METHODS.vcsInit]({ cwd: project.workspaceRoot }), + ), + ); + assert.deepEqual(dispatched, []); + }).pipe(Effect.provide(NodeHttpServer.layerTest)), + ); + it.effect("records thread analytics only after a client command succeeds", () => Effect.gen(function* () { const effects: string[] = []; diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index 830bb35b86db..49fa0a9fd26c 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -1871,6 +1871,32 @@ const makeWsRpcLayer = ( vcsStatusBroadcaster .refreshStatus(cwd) .pipe(Effect.ignoreCause({ log: true }), Effect.forkDetach, Effect.asVoid); + // Publish and init add a remote or a repository the cached identity has not + // seen. Re-emitting the project shell carries the new identity to clients. + // A null identity (a failed lookup, or init without a remote) changes + // nothing clients can use, so it is not re-emitted. + const refreshRepositoryIdentity = (cwd: string) => + repositoryIdentityResolver.resolve(cwd, { refresh: true }).pipe( + Effect.flatMap((identity) => + identity === null + ? Effect.succeedNone + : projectionSnapshotQuery.getActiveProjectByWorkspaceRoot(cwd), + ), + Effect.flatMap((project) => + Option.isNone(project) + ? Effect.void + : Effect.gen(function* () { + const command = yield* normalizeDispatchCommand({ + type: "project.meta.update", + commandId: yield* serverCommandId("repository-identity-refresh"), + projectId: project.value.id, + }); + yield* dispatchNormalizedCommand(command); + }), + ), + Effect.ignoreCause({ log: true }), + Effect.provideContext(normalizerContext), + ); return WsRpcGroup.of({ [ORCHESTRATION_WS_METHODS.dispatchCommand]: (command) => @@ -3038,9 +3064,10 @@ const makeWsRpcLayer = ( [WS_METHODS.sourceControlPublishRepository]: (input) => observeRpcEffect( WS_METHODS.sourceControlPublishRepository, - sourceControlRepositories - .publishRepository(input) - .pipe(Effect.tap(() => refreshGitStatus(input.cwd))), + sourceControlRepositories.publishRepository(input).pipe( + Effect.tap(() => refreshRepositoryIdentity(input.cwd)), + Effect.tap(() => refreshGitStatus(input.cwd)), + ), { "rpc.aggregate": "source-control", }, @@ -3396,9 +3423,10 @@ const makeWsRpcLayer = ( [WS_METHODS.vcsInit]: (input) => observeRpcEffect( WS_METHODS.vcsInit, - vcsProvisioning - .initRepository(input) - .pipe(Effect.tap(() => refreshGitStatus(input.cwd))), + vcsProvisioning.initRepository(input).pipe( + Effect.tap(() => refreshRepositoryIdentity(input.cwd)), + Effect.tap(() => refreshGitStatus(input.cwd)), + ), { "rpc.aggregate": "vcs" }, ), [WS_METHODS.reviewGetDiffPreview]: (input) =>