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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
27 changes: 14 additions & 13 deletions apps/server/src/git/GitManager.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import * as LookupResultCache from "./LookupResultCache.ts";
import * as Arr from "effect/Array";
import * as Cache from "effect/Cache";
import * as Context from "effect/Context";
Expand Down Expand Up @@ -1092,7 +1093,7 @@ export const make = Effect.gen(function* () {
prLookupFailureStreakByKey.set(key, streak);
return prLookupFailureTtl(streak);
};
const prLookupCache = yield* Cache.makeWith(
const prLookupCache = yield* LookupResultCache.make(
(key: string) => {
const [
cwd = "",
Expand Down Expand Up @@ -1213,14 +1214,14 @@ export const make = Effect.gen(function* () {
const branchKey = `${cwd}\u0000${details.branch}`;
const cacheKey = prLookupCacheKey(cwd, details);
if (refreshMissingPullRequest) {
const cached = yield* Cache.getOption(prLookupCache, cacheKey).pipe(
Effect.orElseSucceed(() => Option.none()),
);
const cached = yield* prLookupCache
.getOption(cacheKey)
.pipe(Effect.orElseSucceed(() => Option.none()));
if (Option.isSome(cached) && cached.value.latest === null) {
yield* Cache.invalidate(prLookupCache, cacheKey);
yield* prLookupCache.invalidate(cacheKey);
}
}
return yield* Cache.get(prLookupCache, cacheKey).pipe(
return yield* prLookupCache.get(cacheKey).pipe(
Effect.map(({ latest, headContext }) => {
if (!latest) return { pr: null, headContext };
// On the default branch, only surface open PRs.
Expand Down Expand Up @@ -2225,12 +2226,12 @@ export const make = Effect.gen(function* () {
if (options?.refresh) {
// A completed turn can create a PR or reuse a merged PR's branch.
// Refresh successful answers, but keep failed lookups' retry backoff.
const cached = yield* Cache.getOption(prLookupCache, cacheKey).pipe(
Effect.orElseSucceed(() => Option.none()),
);
if (Option.isSome(cached)) yield* Cache.invalidate(prLookupCache, cacheKey);
const cached = yield* prLookupCache
.getOption(cacheKey)
.pipe(Effect.orElseSucceed(() => Option.none()));
if (Option.isSome(cached)) yield* prLookupCache.invalidate(cacheKey);
}
let cached = yield* Cache.get(prLookupCache, cacheKey);
let cached = yield* prLookupCache.get(cacheKey);
// The cached head context may have resolved on a different remote than
// the saved upstream: a branch tracking origin/main but pushed to a fork
// is looked up on the fork. Verify against the remote the lookup used.
Expand All @@ -2257,8 +2258,8 @@ export const make = Effect.gen(function* () {
});
}
if (!hasSameIdentity(cached.headContext, currentIdentity)) {
yield* Cache.invalidate(prLookupCache, cacheKey);
cached = yield* Cache.get(prLookupCache, cacheKey);
yield* prLookupCache.invalidate(cacheKey);
cached = yield* prLookupCache.get(cacheKey);
const refreshedIdentity = yield* resolvePrLookupRepositoryIdentity(
cacheCwd,
branch,
Expand Down
23 changes: 23 additions & 0 deletions apps/server/src/git/LookupResultCache.memory.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
// @effect-diagnostics nodeBuiltinImport:off - retention assertions need an isolated process with explicit GC.
import * as NodeChildProcess from "node:child_process";
import * as NodeURL from "node:url";
import * as NodeUtil from "node:util";
import { expect, it } from "vite-plus/test";
const execFile = NodeUtil.promisify(NodeChildProcess.execFile);
it.each([false, true])(
"releases the traced caller snapshot while caching a result (failure=%s)",
async (failure) => {
const fixture = NodeURL.fileURLToPath(
new URL("./testing/LookupRetention.fixture.mjs", import.meta.url),
);
const { stdout } = await execFile(process.execPath, [
"--expose-gc",
fixture,
...(failure ? ["failure"] : []),
]);
expect(JSON.parse(stdout)).toEqual({
retained: false,
cachedResult: failure ? "unavailable" : 1,
});
},
);
112 changes: 112 additions & 0 deletions apps/server/src/git/LookupResultCache.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,112 @@
import { expect, it } from "@effect/vitest";
import * as Deferred from "effect/Deferred";
import * as Effect from "effect/Effect";
import * as Exit from "effect/Exit";
import * as Fiber from "effect/Fiber";
import * as Option from "effect/Option";
import { TestClock } from "effect/testing";
import * as LookupResultCache from "./LookupResultCache.ts";

it.effect("shares pending lookups and expires results from completion", () =>
Effect.gen(function* () {
const started = yield* Deferred.make<void>();
const finish = yield* Deferred.make<void>();
let calls = 0;
const cache = yield* LookupResultCache.make(
(_key: string) =>
Effect.gen(function* () {
calls++;
yield* Deferred.succeed(started, undefined);
yield* Deferred.await(finish);
return calls;
}),
{ capacity: 2, timeToLive: () => "1 minute" },
);
const first = yield* Effect.forkChild(cache.get("a"));
yield* Deferred.await(started);
const second = yield* Effect.forkChild(cache.get("a"));
yield* TestClock.adjust("2 minutes");
yield* Deferred.succeed(finish, undefined);
expect(yield* Fiber.join(first)).toBe(1);
expect(yield* Fiber.join(second)).toBe(1);
expect(yield* cache.get("a")).toBe(1);
yield* TestClock.adjust("1 minute");
expect(yield* cache.getOption("a")).toEqual(Option.none());
expect(yield* cache.get("a")).toBe(2);
}).pipe(Effect.scoped),
);

it.effect("preserves failure TTL and invalidates failed entries", () =>
Effect.gen(function* () {
let calls = 0;
const cache = yield* LookupResultCache.make(
(_key: string) =>
Effect.suspend(() => {
calls++;
return Effect.fail("unavailable");
}),
{ capacity: 2, timeToLive: (exit) => (Exit.isFailure(exit) ? "10 seconds" : "1 minute") },
);
expect(yield* Effect.flip(cache.get("a"))).toBe("unavailable");
expect(yield* Effect.flip(cache.getOption("a"))).toBe("unavailable");
expect(calls).toBe(1);
yield* TestClock.adjust("10 seconds");
yield* Effect.exit(cache.get("a"));
expect(calls).toBe(2);
yield* cache.invalidate("a");
yield* Effect.exit(cache.get("a"));
expect(calls).toBe(3);
}).pipe(Effect.scoped),
);

it.effect("evicts the least recently read result and cannot repopulate an invalidated lookup", () =>
Effect.gen(function* () {
const started = yield* Deferred.make<void>();
const finish = yield* Deferred.make<void>();
let calls = 0;
const cache = yield* LookupResultCache.make(
(key: string) =>
Effect.gen(function* () {
const result = ++calls;
if (key === "pending") {
yield* Deferred.succeed(started, undefined);
yield* Deferred.await(finish);
}
return result;
}),
{ capacity: 2, timeToLive: () => "1 minute" },
);
yield* cache.get("a");
yield* cache.get("b");
yield* cache.get("a");
yield* cache.get("c");
expect(yield* cache.getOption("b")).toEqual(Option.none());
const pending = yield* Effect.forkChild(cache.get("pending"));
yield* Deferred.await(started);
yield* cache.invalidate("pending");
yield* Deferred.succeed(finish, undefined);
yield* Fiber.join(pending);
expect(yield* cache.getOption("pending")).toEqual(Option.none());
}).pipe(Effect.scoped),
);

it.effect("a cancelled caller does not cancel another waiter", () =>
Effect.gen(function* () {
const started = yield* Deferred.make<void>();
const finish = yield* Deferred.make<number>();
const cache = yield* LookupResultCache.make(
(_key: string) =>
Effect.gen(function* () {
yield* Deferred.succeed(started, undefined);
return yield* Deferred.await(finish);
}),
{ capacity: 2, timeToLive: () => "1 minute" },
);
const first = yield* Effect.forkChild(cache.get("a"));
yield* Deferred.await(started);
yield* Fiber.interrupt(first);
const second = yield* Effect.forkChild(cache.get("a"));
yield* Deferred.succeed(finish, 42);
expect(yield* Fiber.join(second)).toBe(42);
}).pipe(Effect.scoped),
);
90 changes: 90 additions & 0 deletions apps/server/src/git/LookupResultCache.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,90 @@
import * as Clock from "effect/Clock";
import * as Deferred from "effect/Deferred";
import * as Duration from "effect/Duration";
import * as Effect from "effect/Effect";
import * as Exit from "effect/Exit";
import * as Option from "effect/Option";
import * as References from "effect/References";

const snapshotStack = (
frame: References.StackFrame | undefined,
): References.StackFrame | undefined => {
if (frame === undefined) return undefined;
const stack = frame.stack();
return { name: frame.name, stack: () => stack, parent: snapshotStack(frame.parent) };
};

// Completed Effect cache fibers can retain a caller's lazy trace stack and its
// request snapshot. Store only the result; running lookups belong to the service scope.
export const make = Effect.fnUntraced(function* <Key, A, E>(
lookup: (key: Key) => Effect.Effect<A, E>,
options: {
readonly capacity: number;
readonly timeToLive: (exit: Exit.Exit<A, E>, key: Key) => Duration.Input;
},
) {
const scope = yield* Effect.scope;
const clock = yield* Clock.Clock;
const entries = new Map<
Key,
{
readonly result: Deferred.Deferred<A, E>;
expiresAt: number;
}
>();
const existing = (key: Key) => {
const entry = entries.get(key);
if (entry === undefined) return undefined;
if (entry.expiresAt <= clock.currentTimeMillisUnsafe()) {
entries.delete(key);
return undefined;
}
entries.delete(key);
entries.set(key, entry);
return entry;
};
const get = Effect.fnUntraced(function* (key: Key) {
return yield* Effect.uninterruptibleMask((restore) =>
Effect.gen(function* () {
const cached = existing(key);
if (cached !== undefined) return yield* restore(Deferred.await(cached.result));
const entry = { result: yield* Deferred.make<A, E>(), expiresAt: Infinity };
entries.set(key, entry);
if (entries.size > options.capacity) {
const oldest = entries.keys().next();
if (!oldest.done) entries.delete(oldest.value);
}
const stack = snapshotStack(yield* References.CurrentStackFrame);
yield* Effect.forkIn(
Effect.exit(Effect.interruptible(Effect.suspend(() => lookup(key)))).pipe(
Effect.provideService(References.CurrentStackFrame, stack),
Effect.flatMap((exit) =>
Effect.gen(function* () {
entry.expiresAt =
clock.currentTimeMillisUnsafe() +
Duration.toMillis(Duration.fromInputUnsafe(options.timeToLive(exit, key)));
yield* Deferred.done(entry.result, exit);
}),
),
),
scope,
);
return yield* restore(Deferred.await(entry.result));
}),
);
});
return {
get,
getOption: (key: Key) =>
Effect.suspend(() => {
const entry = existing(key);
return entry === undefined
? Effect.succeed(Option.none<A>())
: Effect.map(Deferred.await(entry.result), Option.some);
}),
invalidate: (key: Key) =>
Effect.sync(() => {
entries.delete(key);
}),
};
});
41 changes: 41 additions & 0 deletions apps/server/src/git/testing/LookupRetention.fixture.mjs
Original file line number Diff line number Diff line change
@@ -0,0 +1,41 @@
import { Cache, Effect, Exit, Scope } from "effect";
import * as LookupResultCache from "../LookupResultCache.ts";
const scope = Effect.runSync(Scope.make());
const original = process.argv[2] === "original";
const failure = process.argv.includes("failure");
const lookup = () => (failure ? Effect.fail("unavailable") : Effect.succeed(1));
const cache = await Effect.runPromise(
original
? Cache.makeWith(lookup, { capacity: 8, timeToLive: () => "1 hour" })
: LookupResultCache.make(lookup, {
capacity: 8,
timeToLive: () => "1 hour",
}).pipe(Effect.provideService(Scope.Scope, scope)),
);
let reference;
async function populate() {
await Effect.runPromise(
Effect.gen(function* () {
const snapshot = { values: Array.from({ length: 250000 }, (_, i) => i) };
reference = new WeakRef(snapshot);
const request = Effect.fn("LookupRetention.request")(function* () {
yield* Effect.exit(original ? Cache.get(cache, "key") : cache.get("key"));
return snapshot.values.length;
});
yield* request();
}),
);
}
await populate();
for (let i = 0; i < 12; i++) {
await new Promise((resolve) => setImmediate(resolve));
global.gc();
}
const retained = reference.deref() !== undefined;
const cachedResult = await Effect.runPromise(
(original ? Cache.get(cache, "key") : cache.get("key")).pipe(
Effect.catch(() => Effect.succeed("unavailable")),
),
);
await Effect.runPromise(Scope.close(scope, Exit.void));
process.stdout.write(JSON.stringify({ retained, cachedResult }));
34 changes: 34 additions & 0 deletions apps/server/src/orchestration-v2/ContextHandoffBudget.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -365,6 +365,40 @@ describe("handoff delivery", () => {
}),
);
}
it.effect("loads the history budget only when a handoff needs delivery", () =>
Effect.gen(function* () {
let reads = 0;
const input = {
providerThread,
budget: Effect.sync(() => {
reads++;
return 16_000;
}),
alreadyDeliveredItemIds: new Set<string>(),
persist: () => Effect.void,
};
yield* deliverContextHandoffs({ ...input, handoffs: [] });
yield* deliverContextHandoffs({ ...input, handoffs: [handoff], deferInline: true });
yield* deliverContextHandoffs({
...input,
handoffs: [
{
...handoff,
delivery: {
nativeThreadId: providerThread.nativeThreadRef!.nativeId!,
status: "injected",
itemIds: [],
},
},
],
});
assert.equal(reads, 0);
const result = yield* deliverContextHandoffs({ ...input, handoffs: [handoff] });
assert.equal(reads, 1);
assert.include(result.context, messages[0]!.text);
}),
);

it.effect("persists successful injection before turn start and skips it on retry", () =>
Effect.gen(function* () {
let durable = handoff;
Expand Down
Loading
Loading