diff --git a/workers/plugin-marketplace/src/index.test.ts b/workers/plugin-marketplace/src/index.test.ts index 9e3df19b..6dacbc91 100644 --- a/workers/plugin-marketplace/src/index.test.ts +++ b/workers/plugin-marketplace/src/index.test.ts @@ -19,7 +19,11 @@ type TreeFixture = { class MemoryR2 { objects = new Map(); + putCount = 0; + failPuts = false; async put(key: string, value: string, options?: unknown): Promise { + if (this.failPuts) throw new Error("simulated R2 put failure"); + this.putCount += 1; this.objects.set(key, { value, options }); } @@ -84,9 +88,10 @@ platforms = ["linux", "macos"] ${overrides}`; } -function env(bucket = new MemoryR2(), blacklist?: MemoryKV): Env { +function env(bucket = new MemoryR2(), blacklist?: MemoryKV, backupBucket = new MemoryR2()): Env { return { PLUGIN_MARKETPLACE_BUCKET: bucket, + PLUGIN_MARKETPLACE_BACKUP_BUCKET: backupBucket, PLUGIN_MARKETPLACE_BLACKLIST: blacklist, GITHUB_TOKEN: "token", }; @@ -313,10 +318,24 @@ describe("refreshPlugins", () => { ], }); expect(snapshot.plugins[0]).not.toHaveProperty("defaultBranch"); + expect(snapshot.plugins[0]).toMatchObject({ + firstSeenAt: "2026-06-20T12:00:00.000Z", + starsDelta7d: null, + starsDelta30d: null, + }); expect(bucket.objects.has("plugins/scan-cache.json")).toBe(true); - expect( - [...bucket.objects.keys()].some((key) => key.includes("history")), - ).toBe(false); + const history = JSON.parse(bucket.objects.get("plugins/star-history.json")?.value ?? ""); + expect(history).toEqual({ + schemaVersion: 1, + entries: [ + { + repositoryId: 1, + fullName, + firstSeenAt: "2026-06-20T12:00:00.000Z", + samples: [{ date: "2026-06-20", stars: 5 }], + }, + ], + }); }); test("reuses cached manifests without resolving or rescanning an unchanged repository", async () => { @@ -653,6 +672,227 @@ describe("refreshPlugins", () => { expect(recovered.snapshot.repositoryCount).toBe(50); }); + test("keeps one star sample per UTC day with the first observation winning", async () => { + const bucket = new MemoryR2(); + const backupBucket = new MemoryR2(); + const repository = repo({ stargazers_count: 5 }); + const logger = { error() {} }; + const run = (stars: number, now: string) => { + repository.stargazers_count = stars; + return refreshPlugins(env(bucket, undefined, backupBucket), { + fetch: repositoryFetch({ repositories: [repository] }), + now: new Date(now), + logger, + }); + }; + + expect((await run(5, "2026-06-20T00:10:00.000Z")).ok).toBe(true); + expect((await run(9, "2026-06-20T23:50:00.000Z")).ok).toBe(true); + const sameDay = JSON.parse(bucket.objects.get("plugins/star-history.json")?.value ?? ""); + expect(sameDay.entries[0].samples).toEqual([{ date: "2026-06-20", stars: 5 }]); + + const nextDay = await run(12, "2026-06-21T00:10:00.000Z"); + expect(nextDay.ok).toBe(true); + if (!nextDay.ok) return; + expect(nextDay.snapshot.plugins[0].firstSeenAt).toBe("2026-06-20T00:10:00.000Z"); + const history = JSON.parse(bucket.objects.get("plugins/star-history.json")?.value ?? ""); + expect(history.entries[0].samples).toEqual([ + { date: "2026-06-20", stars: 5 }, + { date: "2026-06-21", stars: 12 }, + ]); + + // One dated backup per UTC day, written on the first run of that day and + // never touched by later intraday runs. + expect(backupBucket.putCount).toBe(2); + expect([...backupBucket.objects.keys()]).toEqual([ + "backups/star-history/2026-06-20.json", + "backups/star-history/2026-06-21.json", + ]); + const firstBackup = JSON.parse( + backupBucket.objects.get("backups/star-history/2026-06-20.json")?.value ?? "", + ); + expect(firstBackup.entries[0].samples).toEqual([{ date: "2026-06-20", stars: 5 }]); + }); + + test("fails the refresh without touching the primary bucket when the backup put fails", async () => { + const bucket = new MemoryR2(); + await bucket.put("plugins/index.json", '{"schemaVersion":1,"plugins":[{"id":1}]}'); + const backupBucket = new MemoryR2(); + backupBucket.failPuts = true; + const errors: string[] = []; + const options = { + fetch: repositoryFetch({ repositories: [repo()] }), + now: new Date("2026-06-20T12:00:00.000Z"), + logger: { error(message: unknown) { errors.push(String(message)); } }, + }; + + const failed = await refreshPlugins(env(bucket, undefined, backupBucket), options); + expect(failed.ok).toBe(false); + expect(errors[0]).toContain("simulated R2 put failure"); + expect(bucket.objects.get("plugins/index.json")?.value).toBe( + '{"schemaVersion":1,"plugins":[{"id":1}]}', + ); + expect(bucket.objects.has("plugins/star-history.json")).toBe(false); + expect(bucket.objects.has("plugins/scan-cache.json")).toBe(false); + + backupBucket.failPuts = false; + const retried = await refreshPlugins(env(bucket, undefined, backupBucket), options); + expect(retried.ok).toBe(true); + expect(backupBucket.objects.has("backups/star-history/2026-06-20.json")).toBe(true); + expect(bucket.objects.has("plugins/star-history.json")).toBe(true); + }); + + test("computes star deltas from near-boundary baselines and rejects stale ones", async () => { + const bucket = new MemoryR2(); + await bucket.put( + "plugins/star-history.json", + JSON.stringify({ + schemaVersion: 1, + entries: [ + { + repositoryId: 1, + fullName: "ogulcancelik/herdr-plugin-example", + firstSeenAt: "2026-05-01T00:00:00.000Z", + samples: [ + { date: "2026-05-01", stars: 10 }, + { date: "2026-05-20", stars: 30 }, + { date: "2026-06-13", stars: 40 }, + { date: "2026-06-19", stars: 90 }, + ], + }, + { + repositoryId: 2, + fullName: "example/rejoined", + firstSeenAt: "2026-05-01T00:00:00.000Z", + samples: [{ date: "2026-05-01", stars: 3 }], + }, + ], + }), + ); + + const result = await refreshPlugins(env(bucket), { + fetch: repositoryFetch({ + repositories: [ + repo({ stargazers_count: 100 }), + repo({ + id: 2, + full_name: "example/rejoined", + owner: { login: "example" }, + name: "rejoined", + html_url: "https://github.com/example/rejoined", + stargazers_count: 50, + }), + ], + }), + now: new Date("2026-06-20T12:00:00.000Z"), + logger: { error() {} }, + }); + + expect(result.ok).toBe(true); + if (!result.ok) return; + expect(result.snapshot.plugins[0]).toMatchObject({ + firstSeenAt: "2026-05-01T00:00:00.000Z", + starsDelta7d: 60, + starsDelta30d: 70, + }); + // The rejoined repository only has a six-week-old sample: far past both + // window boundaries, so no delta is reported instead of an inflated one. + expect(result.snapshot.plugins[1]).toMatchObject({ + fullName: "example/rejoined", + starsDelta7d: null, + starsDelta30d: null, + }); + }); + + test("prunes aged samples and keeps history for delisted repositories", async () => { + const bucket = new MemoryR2(); + await bucket.put( + "plugins/star-history.json", + JSON.stringify({ + schemaVersion: 1, + entries: [ + { + repositoryId: 1, + fullName: "ogulcancelik/herdr-plugin-example", + firstSeenAt: "2026-01-01T00:00:00.000Z", + samples: [ + { date: "2026-01-01", stars: 1 }, + { date: "2026-06-19", stars: 4 }, + ], + }, + { + repositoryId: 2, + fullName: "example/delisted-recently", + firstSeenAt: "2026-01-01T00:00:00.000Z", + samples: [ + { date: "2026-01-05", stars: 5 }, + { date: "2026-06-01", stars: 7 }, + ], + }, + { + repositoryId: 3, + fullName: "example/delisted-long-ago", + firstSeenAt: "2026-01-01T00:00:00.000Z", + samples: [{ date: "2026-01-05", stars: 3 }], + }, + ], + }), + ); + + const result = await refreshPlugins(env(bucket), { + fetch: repositoryFetch({ repositories: [repo()] }), + now: new Date("2026-06-20T12:00:00.000Z"), + logger: { error() {} }, + }); + + expect(result.ok).toBe(true); + const history = JSON.parse(bucket.objects.get("plugins/star-history.json")?.value ?? ""); + expect(history.entries.map((entry: { repositoryId: number }) => entry.repositoryId)).toEqual([ + 1, 2, + ]); + expect(history.entries[0].samples.map((sample: { date: string }) => sample.date)).toEqual([ + "2026-06-19", + "2026-06-20", + ]); + expect(history.entries[1].samples).toEqual([{ date: "2026-06-01", stars: 7 }]); + }); + + for (const { name, record } of [ + { name: "malformed JSON", record: "broken" }, + { name: "an unsupported shape", record: '{"schemaVersion":2,"entries":[]}' }, + { + name: "an invalid entry", + record: JSON.stringify({ + schemaVersion: 1, + entries: [{ repositoryId: -5, samples: "broken" }], + }), + }, + ]) { + test(`fails the refresh without overwriting anything when star history has ${name}`, async () => { + const bucket = new MemoryR2(); + await bucket.put("plugins/index.json", '{"schemaVersion":1,"plugins":[{"id":1}]}'); + await bucket.put("plugins/star-history.json", record); + const backupBucket = new MemoryR2(); + const errors: string[] = []; + + const result = await refreshPlugins(env(bucket, undefined, backupBucket), { + fetch: repositoryFetch({ repositories: [repo()] }), + now: new Date("2026-06-20T12:00:00.000Z"), + logger: { error(message: unknown) { errors.push(String(message)); } }, + }); + + expect(result.ok).toBe(false); + expect(errors[0]).toContain("invalid plugin marketplace star history"); + expect(bucket.objects.get("plugins/star-history.json")?.value).toBe(record); + expect(bucket.objects.get("plugins/index.json")?.value).toBe( + '{"schemaVersion":1,"plugins":[{"id":1}]}', + ); + expect(bucket.objects.has("plugins/scan-cache.json")).toBe(false); + // Broken history must never be preserved as a "backup". + expect(backupBucket.objects.size).toBe(0); + }); + } + for (const { name, fetch } of [ { name: "GitHub failure", diff --git a/workers/plugin-marketplace/src/index.ts b/workers/plugin-marketplace/src/index.ts index a504e1af..2e5b9dfb 100644 --- a/workers/plugin-marketplace/src/index.ts +++ b/workers/plugin-marketplace/src/index.ts @@ -24,6 +24,12 @@ const PLUGIN_VERSION_MAX_CHARS = 64; const PLUGIN_DESCRIPTION_MAX_CHARS = 500; const REQUEST_TIMEOUT_MS = 10_000; const SCAN_CACHE_SCHEMA_VERSION = 1; +const STAR_HISTORY_KEY = "plugins/star-history.json"; +const STAR_HISTORY_BACKUP_KEY_PREFIX = "backups/star-history/"; +const STAR_HISTORY_SCHEMA_VERSION = 1; +const STAR_HISTORY_MAX_AGE_DAYS = 60; +const STAR_BASELINE_TOLERANCE_DAYS = 2; +const DAY_MS = 86_400_000; const PLUGIN_PLATFORMS = new Set(["linux", "macos", "windows"]); const REGULAR_BLOB_MODES = new Set(["100644", "100755"]); @@ -57,6 +63,7 @@ type ScheduledController = unknown; export type Env = { PLUGIN_MARKETPLACE_BUCKET: R2Bucket; + PLUGIN_MARKETPLACE_BACKUP_BUCKET: R2Bucket; PLUGIN_MARKETPLACE_BLACKLIST?: KVNamespace; GITHUB_TOKEN?: string; }; @@ -101,9 +108,29 @@ export type PluginManifestListing = { export type PluginListing = Omit & { headCommit: string; + firstSeenAt: string; + starsDelta7d: number | null; + starsDelta30d: number | null; manifests: PluginManifestListing[]; }; +export type StarHistorySample = { + date: string; + stars: number; +}; + +export type StarHistoryEntry = { + repositoryId: number; + fullName: string; + firstSeenAt: string; + samples: StarHistorySample[]; +}; + +export type StarHistory = { + schemaVersion: 1; + entries: StarHistoryEntry[]; +}; + export type PluginSnapshot = { schemaVersion: 1; generatedAt: string; @@ -252,14 +279,30 @@ export async function refreshPlugins( }; }); const entriesById = new Map(nextEntries.map((entry) => [entry.repositoryId, entry])); - const plugins = resolved - .map(({ repository }) => { - const entry = entriesById.get(repository.id); - if (!entry || entry.manifests.length === 0) return null; - const { defaultBranch: _defaultBranch, ...card } = repository; - return { ...card, headCommit: entry.headCommit, manifests: entry.manifests }; - }) - .filter((plugin): plugin is PluginListing => plugin !== null); + const cards = resolved.flatMap(({ repository }) => { + const entry = entriesById.get(repository.id); + if (!entry || entry.manifests.length === 0) return []; + const { defaultBranch: _defaultBranch, ...card } = repository; + return [{ ...card, headCommit: entry.headCommit, manifests: entry.manifests }]; + }); + + const now = options.now ?? new Date(); + // History updates are read-modify-write without a conditional put. + // Overlapping runs are rare (30-minute cron, sub-minute refreshes) and the + // worst case is a daily sample recording a slightly later observation, so + // the race is accepted instead of building compare-and-swap on R2. + const history = await readStarHistory(env.PLUGIN_MARKETPLACE_BUCKET); + const nextHistory = updateStarHistory(history, cards, now); + const historyById = new Map(nextHistory.entries.map((entry) => [entry.repositoryId, entry])); + const plugins: PluginListing[] = cards.map((card) => { + const entry = historyById.get(card.id); + return { + ...card, + firstSeenAt: entry?.firstSeenAt ?? now.toISOString(), + starsDelta7d: entry ? starsDeltaSince(entry, card.stars, now, 7) : null, + starsDelta30d: entry ? starsDeltaSince(entry, card.stars, now, 30) : null, + }; + }); const pluginCount = plugins.reduce((total, plugin) => total + plugin.manifests.length, 0); const warnings = nextEntries.flatMap((entry) => (entry.warning ? [entry.warning] : [])); @@ -268,7 +311,7 @@ export async function refreshPlugins( `GitHub returned ${search.totalCount} results; only the first ${search.repositories.length} were collected.`, ); } - const generatedAt = (options.now ?? new Date()).toISOString(); + const generatedAt = now.toISOString(); const snapshot: PluginSnapshot = { schemaVersion: 1, generatedAt, @@ -300,6 +343,30 @@ export async function refreshPlugins( }; const nextCache: ScanCache = { schemaVersion: SCAN_CACHE_SCHEMA_VERSION, entries: nextEntries }; + // Star history is the only object that cannot be rebuilt, so the first run + // of each UTC day copies it to a dated key in the private backup bucket. + // The dated key's existence marks the day as backed up, which stays correct + // when history is empty or a repository joins mid-day. Backups are kept + // forever; the backup put runs before the primary puts so a failure leaves + // the primary untouched and the next run retries the same dated key. The + // get-then-put pair is not atomic; like the history read-modify-write + // above, overlapping runs are accepted because the worst case is the dated + // backup holding an observation from seconds later. + const backupKey = `${STAR_HISTORY_BACKUP_KEY_PREFIX}${isoDate(now)}.json`; + if ((await env.PLUGIN_MARKETPLACE_BACKUP_BUCKET.get(backupKey)) === null) { + await env.PLUGIN_MARKETPLACE_BACKUP_BUCKET.put(backupKey, JSON.stringify(nextHistory), { + httpMetadata: { + contentType: "application/json; charset=utf-8", + cacheControl: "no-store", + }, + }); + } + await env.PLUGIN_MARKETPLACE_BUCKET.put(STAR_HISTORY_KEY, JSON.stringify(nextHistory), { + httpMetadata: { + contentType: "application/json; charset=utf-8", + cacheControl: "no-store", + }, + }); await env.PLUGIN_MARKETPLACE_BUCKET.put(SCAN_CACHE_KEY, JSON.stringify(nextCache), { httpMetadata: { contentType: "application/json; charset=utf-8", @@ -695,6 +762,137 @@ export function parseManifestSummary(manifestText: string): Omit, + now: Date, +): StarHistory { + const today = isoDate(now); + const oldestKeptMs = now.getTime() - STAR_HISTORY_MAX_AGE_DAYS * DAY_MS; + const entriesById = new Map(history.entries.map((entry) => [entry.repositoryId, entry])); + const nextEntries = new Map(); + + for (const card of cards) { + const existing = entriesById.get(card.id); + const samples = (existing?.samples ?? []).filter( + (sample) => dateMs(sample.date) >= oldestKeptMs, + ); + // One sample per UTC day, first observation wins, so intraday cron runs + // keep a stable baseline for delta computation. + if (samples[samples.length - 1]?.date !== today) { + samples.push({ date: today, stars: card.stars }); + } + nextEntries.set(card.id, { + repositoryId: card.id, + fullName: card.fullName, + firstSeenAt: existing?.firstSeenAt ?? now.toISOString(), + samples, + }); + } + + // Keep history for repositories that dropped out of the snapshot (renamed, + // temporarily invalid manifest, blacklist experiments) so they rejoin with + // their record intact, until their samples age out. + for (const entry of history.entries) { + if (nextEntries.has(entry.repositoryId)) continue; + const samples = entry.samples.filter((sample) => dateMs(sample.date) >= oldestKeptMs); + if (samples.length > 0) { + nextEntries.set(entry.repositoryId, { ...entry, samples }); + } + } + + return { + schemaVersion: STAR_HISTORY_SCHEMA_VERSION, + entries: [...nextEntries.values()].sort((a, b) => a.repositoryId - b.repositoryId), + }; +} + +export function starsDeltaSince( + entry: StarHistoryEntry, + currentStars: number, + now: Date, + windowDays: number, +): number | null { + // The baseline must sit close to the window boundary. A repository that was + // delisted for weeks would otherwise report months of growth as a short + // window and dominate trending rankings. + const cutoffMs = now.getTime() - windowDays * DAY_MS; + const oldestAcceptedMs = cutoffMs - STAR_BASELINE_TOLERANCE_DAYS * DAY_MS; + let baseline: StarHistorySample | null = null; + for (const sample of entry.samples) { + const sampleMs = dateMs(sample.date); + if (sampleMs > cutoffMs) break; + if (sampleMs >= oldestAcceptedMs) baseline = sample; + } + return baseline ? currentStars - baseline.stars : null; +} + +// Unlike the scan cache, star history cannot be rebuilt from GitHub, so any +// invalid state fails the refresh instead of being discarded and overwritten. +async function readStarHistory(bucket: R2Bucket): Promise { + const object = await bucket.get(STAR_HISTORY_KEY); + if (!object) return { schemaVersion: STAR_HISTORY_SCHEMA_VERSION, entries: [] }; + const invalid = (reason: string): Error => + new Error( + `invalid plugin marketplace star history (${reason}); remove ${STAR_HISTORY_KEY} to accept losing recorded history`, + ); + let value: unknown; + try { + value = JSON.parse(await object.text()); + } catch { + throw invalid("malformed JSON"); + } + if ( + !isObject(value) || + value.schemaVersion !== STAR_HISTORY_SCHEMA_VERSION || + !Array.isArray(value.entries) + ) { + throw invalid("unsupported shape"); + } + const entries = value.entries.map(readStarHistoryEntry); + if (entries.some((entry) => entry === null)) { + throw invalid("invalid entry"); + } + return { schemaVersion: STAR_HISTORY_SCHEMA_VERSION, entries: entries as StarHistoryEntry[] }; +} + +function readStarHistoryEntry(value: unknown): StarHistoryEntry | null { + if (!isObject(value) || !Array.isArray(value.samples)) return null; + const repositoryId = readInteger(value.repositoryId); + const fullName = readString(value.fullName); + const firstSeenAt = readIsoString(value.firstSeenAt); + const samples = value.samples.map(readStarHistorySample); + if ( + repositoryId === null || + repositoryId <= 0 || + !fullName || + !firstSeenAt || + samples.some((sample) => sample === null) + ) { + return null; + } + return { + repositoryId, + fullName, + firstSeenAt, + samples: (samples as StarHistorySample[]).sort((a, b) => a.date.localeCompare(b.date)), + }; +} + +function readStarHistorySample(value: unknown): StarHistorySample | null { + if (!isObject(value)) return null; + const date = readString(value.date); + const stars = readNonNegativeIntegerOrNull(value.stars); + if (!date || !/^\d{4}-\d{2}-\d{2}$/.test(date) || Number.isNaN(Date.parse(date)) || stars === null) { + return null; + } + return { date, stars }; +} + +function isoDate(now: Date): string { + return now.toISOString().slice(0, 10); +} + async function readScanCache( bucket: R2Bucket, logger: Pick, diff --git a/workers/plugin-marketplace/wrangler.toml b/workers/plugin-marketplace/wrangler.toml index 76f118d7..2367055b 100644 --- a/workers/plugin-marketplace/wrangler.toml +++ b/workers/plugin-marketplace/wrangler.toml @@ -10,6 +10,10 @@ subrequests = 3500 binding = "PLUGIN_MARKETPLACE_BUCKET" bucket_name = "herdr-plugin-marketplace" +[[r2_buckets]] +binding = "PLUGIN_MARKETPLACE_BACKUP_BUCKET" +bucket_name = "herdr-plugin-marketplace-backup" + [[kv_namespaces]] binding = "PLUGIN_MARKETPLACE_BLACKLIST" id = "6504b3b84171492db56e805f6aad1686"