Skip to content
Open
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
33 changes: 12 additions & 21 deletions apps/server/src/pullRequest/PullRequestService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,6 @@ import {
} from "@t3tools/shared/sourceControl";
import { normalizeGitRemoteUrl } from "@t3tools/shared/git";
import * as Cache from "effect/Cache";
import * as Cause from "effect/Cause";
import * as Clock from "effect/Clock";
import * as Context from "effect/Context";
import * as DateTime from "effect/DateTime";
Expand Down Expand Up @@ -83,6 +82,7 @@ import * as PullRequestFilesViewed from "../persistence/PullRequestFilesViewed.t
import * as RepositoryIdentityResolver from "../project/RepositoryIdentityResolver.ts";
import * as SourceControlProviderRegistry from "../sourceControl/SourceControlProviderRegistry.ts";
import * as SourceControlRateLimit from "../sourceControl/SourceControlRateLimit.ts";
import { getCached } from "../utils/getCached.ts";
import {
type ProviderChangeRequest,
type ProviderListCursor,
Expand Down Expand Up @@ -1077,13 +1077,13 @@ export const make = Effect.gen(function* () {
// roots are the same lookup, and putting them on separate flights would spawn two of
// this host's CLIs on a cold page load, which is the coalescing this exists for.
const key = JSON.stringify([host, api.kind, [...new Set(roots)].sort()]);
if (options?.allowPaused === true) return Cache.get(viewerFlights, key);
if (options?.allowPaused === true) return getCached(viewerFlights, key);
// The pause is checked here rather than inside the lookup, so that it holds back the
// callers nobody is waiting on without splitting the flight they share with a press.
// A failed lookup is held nowhere, so letting a background read through would spawn
// this host's CLI on every refresh for as long as the pause lasted, and re-extend it.
return rateLimits.check({ provider: api.kind, host }).pipe(
Effect.flatMap(() => Cache.get(viewerFlights, key)),
Effect.flatMap(() => getCached(viewerFlights, key)),
Effect.catch((error) =>
Effect.succeed<ResolvedViewer>({
host,
Expand Down Expand Up @@ -2948,12 +2948,7 @@ export const make = Effect.gen(function* () {
? null
: Object.entries(input.cursors).toSorted(([left], [right]) => left.localeCompare(right)),
]);
// A replacement reader can join a lookup still finishing its previous reader's cancellation.
return Cache.get(listCache, key).pipe(
Effect.catchCause((cause) =>
Cause.hasInterruptsOnly(cause) ? Cache.get(listCache, key) : Effect.failCause(cause),
),
);
return getCached(listCache, key);
};

const checksCache = yield* Cache.makeWith(
Expand Down Expand Up @@ -3070,7 +3065,7 @@ export const make = Effect.gen(function* () {
// `serveHeld` returns immediately. Skip the write when that read is older
// than a later strict summary — display reuse would otherwise keep the
// regression and never ask the host again.
const read = Cache.get(detailCache, key).pipe(
const read = getCached(detailCache, key).pipe(
Effect.tap((value) => {
const summary = summaryFromDetail(value, lastGoodSummary.peek(key));
return shouldReplaceHeldSummary(key, summary)
Expand Down Expand Up @@ -3110,15 +3105,15 @@ export const make = Effect.gen(function* () {
return Cache.getSuccess(detailCache, key).pipe(
Effect.flatMap(
Option.match({
onNone: () => Cache.get(previewCache, key),
onNone: () => getCached(previewCache, key),
onSome: (detail) => Effect.succeed(previewFields(detail)),
}),
),
);
};
const activity: PullRequestService["Service"]["activity"] = (input) => {
const key = refCacheKey(input);
return Cache.get(activityCache, key);
return getCached(activityCache, key);
};

const diffCache = yield* Cache.makeWith(
Expand Down Expand Up @@ -3148,7 +3143,7 @@ export const make = Effect.gen(function* () {
? (lastGoodSummary.peek(refCacheKey(input))?.updatedAt ?? null)
: null,
]);
const read = Cache.get(diffCache, key).pipe(
const read = getCached(diffCache, key).pipe(
Effect.tap((value) =>
canCacheDiff(value)
? Effect.void
Expand Down Expand Up @@ -3180,7 +3175,7 @@ export const make = Effect.gen(function* () {
const filesViewed: PullRequestService["Service"]["filesViewed"] = (input) =>
canonicalRef(input).pipe(
Effect.flatMap((ref) =>
Cache.get(filesViewedCache, JSON.stringify([refCacheKey(ref), filesViewedEpoch(ref)])),
getCached(filesViewedCache, JSON.stringify([refCacheKey(ref), filesViewedEpoch(ref)])),
),
);

Expand Down Expand Up @@ -3228,11 +3223,7 @@ export const make = Effect.gen(function* () {
}
if (missing.size === 0) return { stats: held };
const key = statsBatchKey(missing.values());
const { result, at } = yield* Cache.get(listStatsCache, key).pipe(
Effect.catchCause((cause) =>
Cause.hasInterruptsOnly(cause) ? Cache.get(listStatsCache, key) : Effect.failCause(cause),
),
);
const { result, at } = yield* getCached(listStatsCache, key);
for (const [key, ref] of missing) {
const stat = result.stats.find(
(stat) =>
Expand Down Expand Up @@ -3379,9 +3370,9 @@ export const make = Effect.gen(function* () {
),
refreshAfterTurn,
detail: credentialCached(detail),
checks: credentialCached((input) => Cache.get(checksCache, refCacheKey(input))),
checks: credentialCached((input) => getCached(checksCache, refCacheKey(input))),
watchFingerprint: credentialCached((input) =>
Cache.get(watchFingerprintCache, refCacheKey(input)),
getCached(watchFingerprintCache, refCacheKey(input)),
),
activity: credentialCached(activity),
preview: credentialCached(preview),
Expand Down
61 changes: 61 additions & 0 deletions apps/server/src/utils/getCached.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,61 @@
import { assert, describe, it } from "@effect/vitest";
import * as Cache from "effect/Cache";
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 { getCached } from "./getCached.ts";

describe("getCached", () => {
it.effect("a read that joins an abandoned lookup gets a fresh result", () =>
Effect.gen(function* () {
const started = yield* Deferred.make<void>();
const stopped = yield* Deferred.make<void>();
let lookups = 0;
const cache = yield* Cache.make<string, string>({
capacity: 10,
lookup: () =>
++lookups === 1
? // Like `gh`, the first lookup takes a moment to stop once interrupted.
Effect.yieldNow.pipe(
Effect.andThen(Deferred.succeed(started, undefined)),
Effect.andThen(Effect.never),
Effect.onInterrupt(() => Deferred.await(stopped)),
)
: Effect.succeed("fresh"),
});
const first = yield* getCached(cache, "detail").pipe(Effect.forkChild);
yield* Deferred.await(started);
yield* Fiber.interrupt(first).pipe(Effect.forkChild({ startImmediately: true }));
const second = yield* getCached(cache, "detail").pipe(
Effect.forkChild({ startImmediately: true }),
);
yield* Deferred.succeed(stopped, undefined);

assert.deepStrictEqual(yield* Fiber.await(second), Exit.succeed("fresh"));
assert.strictEqual(lookups, 2);
}),
);

it.effect("an interrupted read does not start another lookup", () =>
Effect.gen(function* () {
const started = yield* Deferred.make<void>();
let lookups = 0;
const cache = yield* Cache.make<string, string>({
capacity: 10,
lookup: () => {
lookups++;
return Deferred.succeed(started, undefined).pipe(Effect.andThen(Effect.never));
},
});
const read = yield* getCached(cache, "detail").pipe(Effect.forkChild);
yield* Deferred.await(started);

yield* Fiber.interrupt(read);

assert.isTrue(Exit.hasInterrupts(yield* Fiber.await(read)));
assert.strictEqual(lookups, 1);
}),
);
});
20 changes: 20 additions & 0 deletions apps/server/src/utils/getCached.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
import * as Cache from "effect/Cache";
import * as Cause from "effect/Cause";
import * as Effect from "effect/Effect";

/**
* `Cache.get` for a caller that must not inherit another caller's interrupt.
*
* When a lookup's last caller is interrupted, `Cache` interrupts the lookup but
* keeps its entry until the lookup has stopped. A caller that arrives in that
* window joins the dying lookup and fails with a bare interrupt. The entry is
* gone by the time that failure arrives, so one more `get` starts a fresh
* lookup. A caller that was itself interrupted never reaches the retry.
*/
export const getCached = <Key, A, E, R>(
cache: Cache.Cache<Key, A, E, R>,
key: Key,
): Effect.Effect<A, E, R> =>
Cache.get(cache, key).pipe(
Effect.catchCauseIf(Cause.hasInterruptsOnly, () => Cache.get(cache, key)),
);
Loading