diff --git a/migrations/0002-meta.sql b/migrations/0002-meta.sql new file mode 100644 index 0000000..42df7cc --- /dev/null +++ b/migrations/0002-meta.sql @@ -0,0 +1,5 @@ +CREATE TABLE IF NOT EXISTS meta ( + key TEXT PRIMARY KEY, + value INTEGER NOT NULL, + updated_at INTEGER NOT NULL +); diff --git a/schema.sql b/schema.sql index 43bc147..adc1179 100644 --- a/schema.sql +++ b/schema.sql @@ -22,3 +22,9 @@ CREATE TABLE IF NOT EXISTS segment_collections ( CREATE INDEX IF NOT EXISTS segment_collections_nsid ON segment_collections (nsid, idx); CREATE INDEX IF NOT EXISTS segment_collections_ns ON segment_collections (ns, idx); + +CREATE TABLE IF NOT EXISTS meta ( + key TEXT PRIMARY KEY, + value INTEGER NOT NULL, + updated_at INTEGER NOT NULL +); diff --git a/src/index.ts b/src/index.ts index e20f876..a9ba669 100644 --- a/src/index.ts +++ b/src/index.ts @@ -126,6 +126,13 @@ async function ingest(request: Request, env: Env): Promise { for (const segment of parsed.value.segments) { await env.DB.batch(segmentStatements(env.DB, segment, ingestedAt)); } + if (parsed.value.archiveSegments !== undefined) { + await env.DB.prepare( + "INSERT INTO meta (key, value, updated_at) VALUES ('archive_segments', ?1, ?2) ON CONFLICT(key) DO UPDATE SET value = excluded.value, updated_at = excluded.updated_at", + ) + .bind(parsed.value.archiveSegments, ingestedAt) + .run(); + } return json({ ingested: parsed.value.segments.length }); } @@ -138,12 +145,15 @@ type ProgressRow = { }; async function progress(env: Env): Promise { - const [row, kinds] = await env.DB.batch([ + const [row, kinds, meta] = await env.DB.batch([ env.DB.prepare( "SELECT COUNT(*) AS segments, MAX(idx) AS max_idx, MAX(max_seq) AS max_seq, SUM(event_count) AS events, MAX(ingested_at) AS updated_at FROM segments", ), env.DB.prepare("SELECT COUNT(DISTINCT nsid) AS collections, COUNT(DISTINCT ns) AS namespaces FROM segment_collections"), + env.DB.prepare("SELECT value FROM meta WHERE key = 'archive_segments'"), ]); + // SAFETY: the third statement selects exactly this one column. + const archive = meta?.results[0] as { value: number } | undefined; // SAFETY: the first statement selects exactly ProgressRow's columns. const p = row?.results[0] as ProgressRow | undefined; // SAFETY: the second statement selects exactly these two counts. @@ -156,6 +166,7 @@ async function progress(env: Env): Promise { updatedAt: p?.updated_at ?? null, collections: k?.collections ?? 0, namespaces: k?.namespaces ?? 0, + archiveSegments: archive?.value ?? null, }); } diff --git a/src/ingest.ts b/src/ingest.ts index b156c48..b353ad7 100644 --- a/src/ingest.ts +++ b/src/ingest.ts @@ -9,6 +9,7 @@ import { integer, integerRange, object, + optional, string, stringLength, type InferOutput, @@ -37,6 +38,8 @@ export const SegmentRecordSchema = object({ export const IngestBatchSchema = object({ segments: constrain(array(SegmentRecordSchema), [arrayLength(1, 200)]), + /** how many sealed segments the archive reports in total, so the page can say how partial it is */ + archiveSegments: optional(count), }); export type CollectionCount = InferOutput; diff --git a/src/ui.js b/src/ui.js index 9c5f2b0..91e06c3 100644 --- a/src/ui.js +++ b/src/ui.js @@ -66,7 +66,10 @@ async function loadProgress() { el.lede.innerHTML = `strata — what's in the archive. every app on the network writes its own kinds of records; this is how much of each there is, and when it arrived. so far: ${compact.format(p.events)} records from ${fmt.format(p.namespaces)} apps, in ${fmt.format(p.collections)} record types.`; if (p.updatedAt) { const d = new Date(p.updatedAt); - el.live.innerHTML = `· ${fmt.format(p.segments)} chunks read`; + const partial = p.archiveSegments && p.archiveSegments > p.segments; + el.live.innerHTML = partial + ? `· showing the first ${fmt.format(p.segments)} of ${fmt.format(p.archiveSegments)} chunks — still reading` + : `· ${fmt.format(p.segments)} chunks read`; el.updated.textContent = `updated ${dayFmt.format(d)} ${timeFmt.format(d)} utc`; } view.step = STEPS.find((s) => Math.ceil(p.segments / s) <= 120) ?? STEPS[STEPS.length - 1];