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
12 changes: 9 additions & 3 deletions src/app/api/cron/codex-stats/route.ts
Original file line number Diff line number Diff line change
@@ -1,25 +1,31 @@
import { revalidateTag } from "next/cache";
import { revalidatePath } from "next/cache";
import { after } from "next/server";

import { env } from "@/env";
import { getPublicCodexStats } from "@/lib/codex/public-stats";
import { syncCodexAccounts } from "@/lib/codex/sync";
import { reportOperationalError } from "@/lib/operational-error";
import { hasBearerSecret } from "@/lib/request-auth";
import { warmPublicPages } from "@/lib/warm-public-pages";

export const maxDuration = 60;

export async function POST(request: Request) {
if (!hasBearerSecret(request.headers.get("authorization"), env.CRON_SECRET)) {
return Response.json({ ok: false }, { status: 401 });
}
const warmingDeadlineAt = Date.now() + maxDuration * 1000 - 1000;

try {
const result = await syncCodexAccounts();
if (result.updated > 0) {
revalidateTag("public-codex-stats", "max");
revalidatePath("/");
revalidatePath("/tokens");

after(async () => {
await getPublicCodexStats();
if ((await getPublicCodexStats()) !== null) {
await warmPublicPages(["/", "/tokens"], warmingDeadlineAt);
}
});
}
return Response.json({ ok: true, result });
Expand Down
12 changes: 9 additions & 3 deletions src/app/api/cron/github-worker/route.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
import { revalidateTag } from "next/cache";
import { revalidatePath } from "next/cache";
import { after } from "next/server";

import { env } from "@/env";
Expand All @@ -9,6 +9,7 @@ import { GITHUB_WORKER_EXECUTION_DURATION_MS } from "@/lib/github-cron-config";
import type { GITHUB_WORKER_MAX_DURATION_SECONDS } from "@/lib/github-cron-config";
import { reportOperationalError } from "@/lib/operational-error";
import { hasBearerSecret } from "@/lib/request-auth";
import { warmPublicPages } from "@/lib/warm-public-pages";

export const dynamic = "force-dynamic";
export const maxDuration =
Expand All @@ -25,6 +26,7 @@ export async function POST(request: Request) {
if (batchSize === null) {
return Response.json({ ok: false }, { status: 400 });
}
const warmingDeadlineAt = Date.now() + maxDuration * 1000 - 1000;

try {
const activity = await runGitHubActivityWorker({
Expand All @@ -42,9 +44,13 @@ export async function POST(request: Request) {
}),
});
if (activity.projection?.feedRevisionChanged === true) {
revalidateTag("public-github-activity", { expire: 0 });
revalidatePath("/");
revalidatePath("/work");

after(async () => {
await getInitialGitHubActivity();
if ((await getInitialGitHubActivity()) !== null) {
await warmPublicPages(["/", "/work"], warmingDeadlineAt);
}
});
}
return Response.json({
Expand Down
133 changes: 84 additions & 49 deletions src/lib/codex/public-stats.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,4 @@
import "server-only";
import { unstable_cache } from "next/cache";
import { cache } from "react";

import { tokenPreferences } from "@/content/tokens";
Expand All @@ -12,69 +11,105 @@ import {
import type { ClosedCodexHistory } from "@/lib/codex/public-stats-store";
import { reportOperationalError } from "@/lib/operational-error";
import { readPublicSnapshot } from "@/lib/public-snapshot";
import { readRuntimeCache, writeRuntimeCache } from "@/lib/runtime-cache";

const readVersionedClosedDays = unstable_cache(
async (today: string, _historyRevision: string) =>
await readClosedCodexHistory(today),
["codex-closed-history-versioned-v1"],
{ revalidate: 86_400, tags: ["codex-history"] }
);

const readViews = async (
type PublicViews = Awaited<ReturnType<typeof readCodexPublicViews>>;
const preferencesKey = JSON.stringify(tokenPreferences);
const viewsKey = (
today: string,
closedHint?: ClosedCodexHistory,
revision?: Awaited<ReturnType<typeof readCodexPublicRevision>>
revision: Pick<PublicViews, "viewsRevision" | "historyRevision">
) =>
await unstable_cache(
// History is only a validated query-saving hint. The transaction stamps
// the result with the revision of the data it actually read.
async () => await readCodexPublicViews(today, closedHint),
[
revision === undefined
? "public-codex-fallback-v1"
: "public-codex-versioned-v1",
today,
JSON.stringify(tokenPreferences),
...(revision === undefined
? []
: [revision.viewsRevision, revision.historyRevision]),
],
{
revalidate: 900,
...(revision === undefined ? {} : { tags: ["public-codex-stats"] }),
}
)();
JSON.stringify([
"public-codex-v2",
today,
preferencesKey,
revision.viewsRevision,
revision.historyRevision,
]);
const historyKey = (today: string, historyRevision: string) =>
JSON.stringify(["codex-closed-history-v2", today, historyRevision]);

const getViews = cache(async () => {
if (!tokenPreferences.enabled || !isDatabaseConfigured()) {
return null;
}
const today = new Date().toISOString().slice(0, 10);
try {
return await readPublicSnapshot(async () => {
const revision = await readCodexPublicRevision(today);
// Next bypasses nested unstable_cache reads. Prefetch history outside the views cache.
const closed = await readVersionedClosedDays(
today,
revision.historyRevision
const fallbackKey = JSON.stringify([
"public-codex-fallback-v2",
today,
preferencesKey,
]);
const fallbackRead = readRuntimeCache<PublicViews>(fallbackKey);
const cacheWrites: Promise<void>[] = [];
let bodyRead: Promise<PublicViews> | undefined;
const readBody = async (closed?: ClosedCodexHistory) =>
await (bodyRead ??= readCodexPublicViews(today, closed));
const healthyRead = (async () => {
const revision = await readCodexPublicRevision(today);
const [versioned, fallback] = await Promise.all([
readRuntimeCache<PublicViews>(viewsKey(today, revision)),
fallbackRead,
]);
const covers = (value: PublicViews | null) =>
value !== null &&
BigInt(value.viewsRevision) >= BigInt(revision.viewsRevision) &&
BigInt(value.historyRevision) >= BigInt(revision.historyRevision);
const cached = covers(fallback) ? fallback : versioned;
let value = cached;
if (!covers(value)) {
let closed = await readRuntimeCache<ClosedCodexHistory>(
historyKey(today, revision.historyRevision)
);
const fallback = await readViews(today, closed);
if (
BigInt(fallback.viewsRevision) >= BigInt(revision.viewsRevision) &&
BigInt(fallback.historyRevision) >= BigInt(revision.historyRevision)
) {
return fallback.views;
if (closed === null) {
closed = await readClosedCodexHistory(today);
cacheWrites.push(
writeRuntimeCache(
historyKey(today, closed.historyRevision),
closed,
86_400
)
);
}
return (await readViews(today, closed, revision)).views;
});
value = await readBody(closed);
// Both store reads stamp their actual transaction revision. A publication
// between reads must never put newer data under an older revision key.
cacheWrites.push(writeRuntimeCache(viewsKey(today, value), value, 900));
}
if (value === null) {
throw new Error("Missing Codex public views");
}
if (
fallback?.viewsRevision !== value.viewsRevision ||
fallback.historyRevision !== value.historyRevision
) {
cacheWrites.push(writeRuntimeCache(fallbackKey, value));
}
return value.views;
})();
try {
const views = await readPublicSnapshot(async () => await healthyRead);
// Cache writes are optional, and must not trigger an older outage snapshot.
await Promise.all(cacheWrites);
return views;
} catch (error) {
reportOperationalError("public_codex_stats", error);
// Publication never invalidates this successful outage snapshot. Failed
// rebuilds throw rather than replacing it with null.
const fallback = await fallbackRead;
if (fallback !== null) {
return fallback.views;
}
try {
return (await readViews(today)).views;
// Keep a cold timeout on the same history/body read already in flight.
const views = await healthyRead;
await Promise.all(cacheWrites);
return views;
} catch {
return null;
try {
const value = await readBody();
await writeRuntimeCache(fallbackKey, value);
return value.views;
} catch {
return null;
}
}
}
});
Expand Down
113 changes: 66 additions & 47 deletions src/lib/github-activity-feed.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,4 @@
import "server-only";
import { unstable_cache } from "next/cache";

import type { GitHubActivityCursor } from "@/lib/github-activity-cursor";
import {
PUBLIC_GITHUB_ACTIVITY_DAY_PAGE_SIZE,
Expand All @@ -10,61 +8,82 @@ import {
import type { PublicGitHubActivityPage } from "@/lib/github-activity-types";
import { reportOperationalError } from "@/lib/operational-error";
import { readPublicSnapshot } from "@/lib/public-snapshot";
import { readRuntimeCache, writeRuntimeCache } from "@/lib/runtime-cache";

// Keep a successful snapshot available even when the current head cannot be read.
// Publication invalidation must not discard this outage fallback.
// Capture the successful page without putting it in the stable cache key.
const readCachedInitialGitHubActivity = async (
snapshot?: Promise<PublicGitHubActivityPage>
const FALLBACK_KEY = "public-github-activity-fallback-v3";
const contentKey = (
page: Pick<PublicGitHubActivityPage, "head" | "orderingRevision">
) =>
await unstable_cache(
async () =>
await (snapshot ??
readPublicGitHubActivityPage(
null,
PUBLIC_GITHUB_ACTIVITY_DAY_PAGE_SIZE
)),
["public-github-activity-fallback-v2"],
{ revalidate: 60 }
)();
JSON.stringify([
"public-github-activity-v3",
page.head.feedRevision,
page.orderingRevision,
]);

const readVersionedInitialGitHubActivity = unstable_cache(
async (_feedRevision: string, _orderingRevision: string) =>
await readPublicGitHubActivityPage(
export const getInitialGitHubActivity = async () => {
const fallbackRead = readRuntimeCache<PublicGitHubActivityPage>(FALLBACK_KEY);
let cacheWrites: Promise<unknown> | undefined;
let bodyRead: Promise<PublicGitHubActivityPage> | undefined;
const readBody = async () =>
await (bodyRead ??= readPublicGitHubActivityPage(
null,
PUBLIC_GITHUB_ACTIVITY_DAY_PAGE_SIZE
),
["public-github-activity-versioned-v2"],
{ revalidate: 3600, tags: ["public-github-activity"] }
);

export const getInitialGitHubActivity = async () => {
));
const healthyRead = (async () => {
const live = await readPublicGitHubActivityHead();
const [versioned, fallback] = await Promise.all([
readRuntimeCache<PublicGitHubActivityPage>(contentKey(live)),
fallbackRead,
]);
const cached =
fallback?.head.feedRevision === live.head.feedRevision &&
fallback.orderingRevision === live.orderingRevision
? fallback
: versioned;
const snapshot = cached ?? (await readBody());
const page =
snapshot.head.feedRevision === live.head.feedRevision &&
snapshot.orderingRevision === live.orderingRevision &&
BigInt(live.head.revision) >= BigInt(snapshot.head.revision)
? { ...snapshot, head: live.head }
: snapshot;
cacheWrites = Promise.all([
cached === null
? writeRuntimeCache(contentKey(snapshot), snapshot, 3600)
: undefined,
fallback?.head.feedRevision === page.head.feedRevision &&
fallback.head.revision === page.head.revision &&
fallback.orderingRevision === page.orderingRevision
? undefined
: writeRuntimeCache(FALLBACK_KEY, page),
]);
return page;
})();
try {
return await readPublicSnapshot(async () => {
const { head, orderingRevision } = await readPublicGitHubActivityHead();
const snapshot = readVersionedInitialGitHubActivity(
head.feedRevision,
orderingRevision
).then((page) =>
page.head.feedRevision === head.feedRevision &&
page.orderingRevision === orderingRevision &&
BigInt(head.revision) >= BigInt(page.head.revision)
? { ...page, head }
: page
);
// Check both caches concurrently; an expired fallback waits for the same successful snapshot.
const [page] = await Promise.all([
snapshot,
readCachedInitialGitHubActivity(snapshot),
]);
return page;
});
const page = await readPublicSnapshot(async () => await healthyRead);
// Optional persistence must not turn a fresh, masked page into an old fallback.
await cacheWrites;
return page;
} catch (error) {
reportOperationalError("github_activity_initial", error);
const fallback = await fallbackRead;
if (fallback !== null) {
return fallback;
}
// A timeout leaves the original read running. Reuse it instead of downloading
// the page again. If metadata failed, one direct page read can still succeed.
try {
return await readCachedInitialGitHubActivity();
const page = await healthyRead;
await cacheWrites;
return page;
} catch {
return null;
try {
const page = await readBody();
await writeRuntimeCache(FALLBACK_KEY, page);
return page;
} catch {
return null;
}
}
}
};
Expand Down
Loading
Loading