feat(marketplace): track star history for trending and new plugin surfaces

This commit is contained in:
Ogulcan Celik 2026-08-06 20:07:44 +03:00
parent 3825c0c300
commit d16af2660a
3 changed files with 455 additions and 13 deletions

View File

@ -19,7 +19,11 @@ type TreeFixture = {
class MemoryR2 {
objects = new Map<string, { value: string; options: unknown }>();
putCount = 0;
failPuts = false;
async put(key: string, value: string, options?: unknown): Promise<void> {
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",

View File

@ -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<RepositoryListing, "defaultBranch"> & {
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<PluginManifestL
}
}
export function updateStarHistory(
history: StarHistory,
cards: Array<{ id: number; fullName: string; stars: number }>,
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<number, StarHistoryEntry>();
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<StarHistory> {
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<Console, "error">,

View File

@ -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"