Skip to content
Closed
87 changes: 87 additions & 0 deletions apps/server/src/orchestration/ThreadPullRequestReactor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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* () {
Expand Down
15 changes: 13 additions & 2 deletions apps/server/src/orchestration/ThreadPullRequestReactor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Comment thread
t3dotgg marked this conversation as resolved.
}
: snapshotProject;
const repository = sourceControlRepositorySelector(project.repositoryIdentity);
if (first.branch !== null && repository === null) return finishBackfill(group);
const worktreeExists =
Expand Down
69 changes: 69 additions & 0 deletions apps/server/src/project/RepositoryIdentityResolver.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<ReadonlyArray<string>> = [];
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;
Expand Down
95 changes: 79 additions & 16 deletions apps/server/src/project/RepositoryIdentityResolver.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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>()(
"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;
Expand Down Expand Up @@ -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;
},
);
Expand All @@ -119,20 +154,34 @@ const resolveRepositoryIdentityFromCacheKey = Effect.fn(
"RepositoryIdentityResolver.resolveFromCacheKey",
)(function* (
cacheKey: string,
): Effect.fn.Return<RepositoryIdentity | null, never, ProcessRunner.ProcessRunner> {
): Effect.fn.Return<
RepositoryIdentity | null,
RepositoryIdentityLookupError,
ProcessRunner.ProcessRunner
> {
const processRunner = yield* ProcessRunner.ProcessRunner;
const remoteResult = yield* processRunner
.run({
command: "git",
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;
});

Expand All @@ -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<string, string | null>(
const repositoryRootCache = yield* Cache.makeWith<
string,
string | null,
RepositoryIdentityLookupError
>(
(cwd) =>
resolveRepositoryIdentityCacheKey(cwd).pipe(
Effect.provideService(ProcessRunner.ProcessRunner, processRunner),
Expand All @@ -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<string, RepositoryIdentity | null>(
const repositoryIdentityCache = yield* Cache.makeWith<
string,
RepositoryIdentity | null,
RepositoryIdentityLookupError
>(
(cacheKey) =>
resolveRepositoryIdentityFromCacheKey(cacheKey).pipe(
Effect.provideService(ProcessRunner.ProcessRunner, processRunner),
Expand All @@ -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 });
Expand Down
Loading
Loading