From ae663c8da4371b533fb3a2fadbe1bbf4c953f2ea Mon Sep 17 00:00:00 2001 From: zzstoatzz Date: Sat, 29 Aug 2026 01:12:47 -0500 Subject: [PATCH] drain: segment rows to namespace-partitioned parquet in R2 via PUT /archive Worker binds strata-archive (public) and strata-archive-private; app.bsky raw rows are the only private partition. drain.sh runs on the stream box, resumable, 16-17 s/segment measured. Also: ns column + namespace/collection heat modes, /api/namespaces, richer /api/progress. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01B927dNwNKNYsbQdUoNJwHS --- docs/how-it-works.md | 19 ++++ drain/drain.sh | 57 +++++++++++ drain/segment.sql | 15 +++ migrations/0001-ns.sql | 6 ++ schema.sql | 2 + src/index.ts | 212 ++++++++++++++++++++++++++++++++--------- src/ingest.ts | 6 ++ tests/ingest.test.ts | 8 +- wrangler.jsonc | 4 + 9 files changed, 285 insertions(+), 44 deletions(-) create mode 100755 drain/drain.sh create mode 100644 drain/segment.sql create mode 100644 migrations/0001-ns.sql diff --git a/docs/how-it-works.md b/docs/how-it-works.md index b752e1f..2e80c26 100644 --- a/docs/how-it-works.md +++ b/docs/how-it-works.md @@ -40,3 +40,22 @@ ingest-era marker. `STRATA_INGEST_TOKEN` and the archive key `stream-archive-key-strata` live in the sops store and reach the flow as prefect Secret blocks; the worker's copy of the ingest token is a wrangler secret set from the same store value. + +## the drain (decided 2026-08-29) + +the long tail needs rows, not footer counts. `drain/drain.sh` runs on the +stream box (`/opt/strata`, with the duckdb CLI and the jetstream SDK's +`segment-export` example): for every sealed segment it decodes rows from +local disk (`seq, time_us, did, kind, op, collection, rkey` — witnessed +time, no record bodies), writes parquet partitioned by namespace with +duckdb, and PUTs each partition to this worker's `/archive/` with the +ingest bearer. keys: `events/ns=/seg=.parquet`, +`agg/seg=.parquet` (per collection/kind/op: events, distinct dids, seq +and time bounds), `manifest/seg=.json`. resumable via `done.txt`; +the active segment is never read. measured: 16–17 s per segment at nice 19. + +placement policy: `events/ns=app.bsky/*` goes to the private bucket +(`strata-archive-private`); every other partition, all aggregates and +manifests are public at `pub-735e1688181b45e49c925eed81fe2da7.r2.dev` +(bucket `strata-archive`, CORS open for range reads). the tail is what the +page queries live with duckdb-wasm; `app.bsky` is served from aggregates. diff --git a/drain/drain.sh b/drain/drain.sh new file mode 100755 index 0000000..23ef9d6 --- /dev/null +++ b/drain/drain.sh @@ -0,0 +1,57 @@ +#!/usr/bin/env bash +# drain: every sealed segment on this box → parquet partitioned by namespace → strata via the Worker. +# +# state is done.txt (one segment index per line). resumable; idempotent per segment (PUTs overwrite). +# the active (highest-index) segment is never touched. everything runs at nice 19 / idle IO. +# +# env: STRATA_URL (default https://strata.zat.dev), STRATA_INGEST_TOKEN (required) +set -euo pipefail +cd "$(dirname "$0")" +: "${STRATA_INGEST_TOKEN:?set STRATA_INGEST_TOKEN}" +STRATA_URL="${STRATA_URL:-https://strata.zat.dev}" +SEGMENTS=/data/stream/segments +WORK=./work +LIMIT="${1:-0}" # max segments this run; 0 = all +touch done.txt + +idx_of() { python3 -c 'import sys;print(int(sys.argv[1][4:14],36))' "$(basename "$1" .jss)"; } + +put() { # key file + local code + code=$(curl -sS -m 600 -o /tmp/drain.put.json -w "%{http_code}" -X PUT "$STRATA_URL/archive/$1" \ + -H "Authorization: Bearer $STRATA_INGEST_TOKEN" -H "Content-Type: application/octet-stream" \ + -A "strata-drain/1" --data-binary "@$2") + [ "$code" = "200" ] || { echo "PUT $1 -> $code $(cat /tmp/drain.put.json)" >&2; return 1; } +} + +mapfile -t all < <(ls "$SEGMENTS"/seg_*.jss | sort) +unset 'all[-1]' # the active segment +n=0 +for seg in "${all[@]}"; do + idx=$(idx_of "$seg") + grep -qx "$idx" done.txt && continue + [ "$LIMIT" -gt 0 ] && [ "$n" -ge "$LIMIT" ] && break + rm -rf "$WORK"; mkdir -p "$WORK" + t0=$(date +%s) + nice -n 19 ionice -c 3 ./segment-export "$seg" > "$WORK/rows.tsv" 2> "$WORK/export.log" + nice -n 19 ionice -c 3 ./duckdb < <(sed "s#__WORK__#$WORK#g; s#__IDX__#$idx#g" segment.sql) > "$WORK/duckdb.log" 2>&1 + rm -f "$WORK/rows.tsv" + for part in "$WORK"/events/ns=*/; do + ns=$(basename "$part"); ns=${ns#ns=} + put "events/ns=$ns/seg=$idx.parquet" "$part/data_0.parquet" + done + put "agg/seg=$idx.parquet" "$WORK/agg.parquet" + python3 - "$idx" "$seg" "$WORK" "$t0" > "$WORK/manifest.json" <<'EOF' +import json, os, sys, time, glob +idx, seg, work, t0 = int(sys.argv[1]), sys.argv[2], sys.argv[3], int(sys.argv[4]) +parts = {os.path.basename(p)[3:]: os.path.getsize(os.path.join(p, "data_0.parquet")) for p in glob.glob(f"{work}/events/ns=*")} +print(json.dumps({"segment": idx, "name": os.path.basename(seg), "segment_bytes": os.path.getsize(seg), + "rows": int(next(l for l in open(f"{work}/export.log") if " rows from " in l).split()[0]), "partitions": parts, + "agg_bytes": os.path.getsize(f"{work}/agg.parquet"), "seconds": int(time.time()) - t0, "drained_at": int(time.time())})) +EOF + put "manifest/seg=$idx.json" "$WORK/manifest.json" + echo "$idx" >> done.txt + n=$((n + 1)) + echo "$(date -u +%FT%TZ) seg $idx done in $(( $(date +%s) - t0 ))s: $(python3 -c 'import json,sys;m=json.load(open(sys.argv[1]));print(m["rows"],"rows,",len(m["partitions"]),"partitions,",sum(m["partitions"].values())//1024,"KB events,",m["agg_bytes"]//1024,"KB agg")' "$WORK/manifest.json")" +done +echo "drained $n segment(s); $(wc -l < done.txt) total done; $(( ${#all[@]} - $(wc -l < done.txt) )) remaining" diff --git a/drain/segment.sql b/drain/segment.sql new file mode 100644 index 0000000..9887ce2 --- /dev/null +++ b/drain/segment.sql @@ -0,0 +1,15 @@ +CREATE TABLE ev AS + SELECT *, CASE WHEN kind = 'commit' THEN regexp_extract(collection, '^([^.]+\.[^.]+)', 1) ELSE '_' || kind END AS ns + FROM read_csv('__WORK__/rows.tsv', delim = '\t', header = false, + columns = {seq: 'BIGINT', time_us: 'BIGINT', did: 'VARCHAR', kind: 'VARCHAR', op: 'VARCHAR', collection: 'VARCHAR', rkey: 'VARCHAR'}, + quote = '', escape = '', null_padding = true); + +COPY (SELECT seq, time_us, did, kind, op, collection, rkey, ns FROM ev ORDER BY seq) + TO '__WORK__/events' (FORMAT parquet, PARTITION_BY (ns), COMPRESSION zstd); + +COPY ( + SELECT __IDX__ AS segment, ns, collection, kind, op, + count(*) AS events, count(DISTINCT did) AS dids, + min(seq) AS min_seq, max(seq) AS max_seq, min(time_us) AS min_time_us, max(time_us) AS max_time_us + FROM ev GROUP BY ALL ORDER BY ns, collection, kind, op +) TO '__WORK__/agg.parquet' (FORMAT parquet, COMPRESSION zstd); diff --git a/migrations/0001-ns.sql b/migrations/0001-ns.sql new file mode 100644 index 0000000..c099f70 --- /dev/null +++ b/migrations/0001-ns.sql @@ -0,0 +1,6 @@ +ALTER TABLE segment_collections ADD COLUMN ns TEXT; +UPDATE segment_collections SET ns = CASE + WHEN instr(substr(nsid, instr(nsid, '.') + 1), '.') = 0 THEN nsid + ELSE substr(nsid, 1, instr(nsid, '.') + instr(substr(nsid, instr(nsid, '.') + 1), '.') - 1) +END WHERE ns IS NULL; +CREATE INDEX IF NOT EXISTS segment_collections_ns ON segment_collections (ns, idx); diff --git a/schema.sql b/schema.sql index 49a11f4..43bc147 100644 --- a/schema.sql +++ b/schema.sql @@ -15,8 +15,10 @@ CREATE TABLE IF NOT EXISTS segments ( CREATE TABLE IF NOT EXISTS segment_collections ( idx INTEGER NOT NULL REFERENCES segments(idx), nsid TEXT NOT NULL, + ns TEXT NOT NULL, count INTEGER NOT NULL, PRIMARY KEY (idx, nsid) ); 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); diff --git a/src/index.ts b/src/index.ts index 76935f1..9902e0a 100644 --- a/src/index.ts +++ b/src/index.ts @@ -4,14 +4,38 @@ * collection index, and POSTs the result here. this worker only stores and * aggregates; it holds no schedule and never talks to the archive itself. */ -import { IngestParseError, parseIngestBatch, type SegmentRecord } from "./ingest"; +import { IngestParseError, namespaceOf, parseIngestBatch, type SegmentRecord } from "./ingest"; import page from "./page.html"; interface Env { readonly DB: D1Database; + readonly PUBLIC: R2Bucket; + readonly PRIVATE: R2Bucket; readonly STRATA_INGEST_TOKEN: string; } +const MAX_ARCHIVE_OBJECT_BYTES = 200 << 20; +const ARCHIVE_KEY = /^(events\/ns=[a-z0-9._-]+\/seg=\d+\.parquet|agg\/seg=\d+\.parquet|manifest\/seg=\d+\.json)$/; + +/** the drain uploads one object per (segment, partition). app.bsky raw rows are the only private thing: + * everything else — tail namespaces, marker rows, per-segment aggregates — is public by decision (2026-08-29). */ +function bucketFor(key: string, env: Env): R2Bucket { + return key.startsWith("events/ns=app.bsky/") ? env.PRIVATE : env.PUBLIC; +} + +async function putArchive(request: Request, key: string, env: Env): Promise { + if (!authorized(request, env)) return json({ error: "unauthorized" }, 401); + if (!ARCHIVE_KEY.test(key)) return json({ error: "key must be events/ns=/seg=.parquet, agg/seg=.parquet, or manifest/seg=.json" }, 400); + const length = Number(request.headers.get("content-length") ?? "0"); + if (!Number.isFinite(length) || length <= 0 || length > MAX_ARCHIVE_OBJECT_BYTES) return json({ error: "content-length required, ≤ 200 MiB" }, 413); + if (request.body === null) return json({ error: "empty body" }, 400); + const bucket = bucketFor(key, env); + const object = await bucket.put(key, request.body, { + httpMetadata: { contentType: key.endsWith(".json") ? "application/json" : "application/vnd.apache.parquet" }, + }); + return json({ key, size: object.size, etag: object.httpEtag, bucket: bucket === env.PRIVATE ? "private" : "public" }); +} + const MAX_BODY_BYTES = 8 << 20; const MAX_HEAT_BUCKETS = 2000; @@ -61,8 +85,12 @@ function segmentStatements(db: D1Database, segment: SegmentRecord, ingestedAt: n ingestedAt, ); const clear = db.prepare("DELETE FROM segment_collections WHERE idx = ?1").bind(segment.idx); - const upsertCollection = db.prepare("INSERT INTO segment_collections (idx, nsid, count) VALUES (?1, ?2, ?3)"); - return [upsertSegment, clear, ...segment.collections.map((c) => upsertCollection.bind(segment.idx, c.nsid, c.count))]; + const upsertCollection = db.prepare("INSERT INTO segment_collections (idx, nsid, ns, count) VALUES (?1, ?2, ?3, ?4)"); + return [ + upsertSegment, + clear, + ...segment.collections.map((c) => upsertCollection.bind(segment.idx, c.nsid, namespaceOf(c.nsid), c.count)), + ]; } async function ingest(request: Request, env: Env): Promise { @@ -87,23 +115,34 @@ interface ProgressRow { readonly segments: number; readonly max_idx: number | null; readonly max_seq: number | null; + readonly events: number | null; + readonly updated_at: number | null; } async function progress(env: Env): Promise { - const row = await env.DB.prepare( - "SELECT COUNT(*) AS segments, MAX(idx) AS max_idx, MAX(max_seq) AS max_seq FROM segments", - ).first(); - return json({ segments: row?.segments ?? 0, maxIdx: row?.max_idx ?? null, maxSeq: row?.max_seq ?? null }); + const [row, kinds] = 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"), + ]); + // SAFETY: the two statements select exactly these columns. + const p = row?.results[0] as ProgressRow | undefined; + const k = kinds?.results[0] as { collections: number; namespaces: number } | undefined; + return json({ + segments: p?.segments ?? 0, + maxIdx: p?.max_idx ?? null, + maxSeq: p?.max_seq ?? null, + events: p?.events ?? 0, + updatedAt: p?.updated_at ?? null, + collections: k?.collections ?? 0, + namespaces: k?.namespaces ?? 0, + }); } interface HeatRow { readonly bucket: number; - readonly nsid: string; - readonly count: number; -} - -interface RankRow { - readonly nsid: string; + readonly key: string; readonly count: number; } @@ -118,17 +157,101 @@ interface BucketRow { readonly event_count: number; } -const MAX_TOP = 40; +const MAX_TOP = 60; -async function heat(url: URL, env: Env): Promise { +interface HeatQuery { + readonly step: number; + readonly top: number; + readonly scope: { readonly kind: "namespaces" } | { readonly kind: "namespace"; readonly ns: string } | { readonly kind: "collections" }; +} + +function parseHeatQuery(url: URL): HeatQuery | string { const step = Number(url.searchParams.get("step") ?? "10"); const top = Number(url.searchParams.get("top") ?? "12"); - if (!Number.isSafeInteger(step) || step < 1) return json({ error: "step must be a positive integer" }, 400); - if (!Number.isSafeInteger(top) || top < 1 || top > MAX_TOP) return json({ error: `top must be 1..${MAX_TOP}` }, 400); + if (!Number.isSafeInteger(step) || step < 1) return "step must be a positive integer"; + if (!Number.isSafeInteger(top) || top < 1 || top > MAX_TOP) return `top must be 1..${MAX_TOP}`; + const ns = url.searchParams.get("ns"); + const group = url.searchParams.get("group") ?? "namespaces"; + if (ns !== null) { + if (!/^[a-z0-9-]+\.[a-z0-9-]+$/i.test(ns)) return "ns must be an authority like fm.teal"; + return { step, top, scope: { kind: "namespace", ns } }; + } + if (group === "collections") return { step, top, scope: { kind: "collections" } }; + if (group === "namespaces") return { step, top, scope: { kind: "namespaces" } }; + return "group must be namespaces or collections"; +} + +/** rows of the heat map: (bucket, key, count) where key is a namespace, or an nsid inside one namespace, or a + * top-N nsid; everything outside the top N folds into "other" so the response stays small at any archive size. */ +function heatRowsStatement(db: D1Database, q: HeatQuery): D1PreparedStatement { + switch (q.scope.kind) { + case "namespaces": + return db + .prepare( + `WITH ranked AS (SELECT ns FROM segment_collections GROUP BY ns ORDER BY SUM(count) DESC LIMIT CAST(?2 AS INTEGER)) + SELECT idx / CAST(?1 AS INTEGER) AS bucket, + CASE WHEN ns IN (SELECT ns FROM ranked) THEN ns ELSE 'other' END AS key, SUM(count) AS count + FROM segment_collections GROUP BY bucket, 2 ORDER BY bucket, 2`, + ) + .bind(q.step, q.top); + case "namespace": + return db + .prepare( + `SELECT idx / CAST(?1 AS INTEGER) AS bucket, nsid AS key, SUM(count) AS count + FROM segment_collections WHERE ns = ?2 GROUP BY bucket, nsid ORDER BY bucket, nsid`, + ) + .bind(q.step, q.scope.ns); + case "collections": + return db + .prepare( + `WITH ranked AS (SELECT nsid FROM segment_collections GROUP BY nsid ORDER BY SUM(count) DESC LIMIT CAST(?2 AS INTEGER)) + SELECT idx / CAST(?1 AS INTEGER) AS bucket, + CASE WHEN nsid IN (SELECT nsid FROM ranked) THEN nsid ELSE 'other' END AS key, SUM(count) AS count + FROM segment_collections GROUP BY bucket, 2 ORDER BY bucket, 2`, + ) + .bind(q.step, q.top); + } +} + +function heatKeysStatement(db: D1Database, q: HeatQuery): D1PreparedStatement { + switch (q.scope.kind) { + case "namespaces": + return db + .prepare( + `SELECT ns AS key, SUM(count) AS count, COUNT(DISTINCT nsid) AS members + FROM segment_collections GROUP BY ns ORDER BY count DESC LIMIT CAST(?1 AS INTEGER)`, + ) + .bind(q.top); + case "namespace": + return db + .prepare( + `SELECT nsid AS key, SUM(count) AS count, 1 AS members FROM segment_collections WHERE ns = ?1 + GROUP BY nsid ORDER BY count DESC`, + ) + .bind(q.scope.ns); + case "collections": + return db + .prepare( + `SELECT nsid AS key, SUM(count) AS count, 1 AS members FROM segment_collections + GROUP BY nsid ORDER BY count DESC LIMIT CAST(?1 AS INTEGER)`, + ) + .bind(q.top); + } +} + +interface KeyRow { + readonly key: string; + readonly count: number; + readonly members: number; +} + +async function heat(url: URL, env: Env): Promise { + const q = parseHeatQuery(url); + if (typeof q === "string") return json({ error: q }, 400); const progressRow = await env.DB.prepare("SELECT MAX(idx) AS max_idx FROM segments").first<{ max_idx: number | null }>(); const maxIdx = progressRow?.max_idx ?? -1; - if (Math.floor(maxIdx / step) + 1 > MAX_HEAT_BUCKETS) { - return json({ error: `step ${step} yields more than ${MAX_HEAT_BUCKETS} buckets` }, 400); + if (Math.floor(maxIdx / q.step) + 1 > MAX_HEAT_BUCKETS) { + return json({ error: `step ${q.step} yields more than ${MAX_HEAT_BUCKETS} buckets` }, 400); } const results = await env.DB.batch([ env.DB.prepare( @@ -136,46 +259,47 @@ async function heat(url: URL, env: Env): Promise { MAX(max_seq) AS max_seq, MIN(min_witnessed_at) AS min_witnessed_at, MAX(max_witnessed_at) AS max_witnessed_at, SUM(event_count) AS event_count FROM segments GROUP BY bucket ORDER BY bucket`, - ).bind(step), - env.DB.prepare( - `WITH ranked AS ( - SELECT nsid FROM segment_collections GROUP BY nsid ORDER BY SUM(count) DESC LIMIT CAST(?2 AS INTEGER) - ) - SELECT idx / CAST(?1 AS INTEGER) AS bucket, - CASE WHEN nsid IN (SELECT nsid FROM ranked) THEN nsid ELSE 'other' END AS nsid, - SUM(count) AS count - FROM segment_collections GROUP BY bucket, 2 ORDER BY bucket, 2`, - ).bind(step, top), - env.DB.prepare( - "SELECT nsid, SUM(count) AS count FROM segment_collections GROUP BY nsid ORDER BY count DESC LIMIT CAST(?1 AS INTEGER)", - ).bind(top), + ).bind(q.step), + heatRowsStatement(env.DB, q), + heatKeysStatement(env.DB, q), ]); - const [buckets, cells, ranked] = results; - if (buckets === undefined || cells === undefined || ranked === undefined) { + const [buckets, cells, keys] = results; + if (buckets === undefined || cells === undefined || keys === undefined) { return json({ error: "batch returned fewer results than statements" }, 500); } - // SAFETY: the statements above select exactly the columns of BucketRow, HeatRow, and RankRow. + // SAFETY: the statements above select exactly the columns of BucketRow, HeatRow, and KeyRow. const bucketRows = buckets.results as BucketRow[]; const heatRows = cells.results as HeatRow[]; - const rankRows = ranked.results as RankRow[]; + const keyRows = keys.results as KeyRow[]; return json( - { step, top, buckets: bucketRows, cells: heatRows, collections: rankRows }, + { step: q.step, top: q.top, scope: q.scope, buckets: bucketRows, cells: heatRows, keys: keyRows }, 200, { "cache-control": "public, max-age=300" }, ); } -async function collections(env: Env): Promise { - const { results } = await env.DB.prepare( - "SELECT nsid, SUM(count) AS count, COUNT(*) AS segments FROM segment_collections GROUP BY nsid ORDER BY count DESC", - ).all<{ nsid: string; count: number; segments: number }>(); +async function collections(url: URL, env: Env): Promise { + const ns = url.searchParams.get("ns"); + if (ns !== null && !/^[a-z0-9-]+\.[a-z0-9-]+$/i.test(ns)) return json({ error: "ns must be an authority like fm.teal" }, 400); + const stmt = ns === null + ? env.DB.prepare("SELECT nsid, ns, SUM(count) AS count, COUNT(*) AS segments FROM segment_collections GROUP BY nsid ORDER BY count DESC") + : env.DB.prepare("SELECT nsid, ns, SUM(count) AS count, COUNT(*) AS segments FROM segment_collections WHERE ns = ?1 GROUP BY nsid ORDER BY count DESC").bind(ns); + const { results } = await stmt.all<{ nsid: string; ns: string; count: number; segments: number }>(); return json({ collections: results }, 200, { "cache-control": "public, max-age=300" }); } +async function namespaces(env: Env): Promise { + const { results } = await env.DB.prepare( + "SELECT ns, SUM(count) AS count, COUNT(DISTINCT nsid) AS members FROM segment_collections GROUP BY ns ORDER BY count DESC", + ).all<{ ns: string; count: number; members: number }>(); + return json({ namespaces: results }, 200, { "cache-control": "public, max-age=300" }); +} + export default { async fetch(request: Request, env: Env): Promise { const url = new URL(request.url); if (request.method === "POST" && url.pathname === "/ingest") return ingest(request, env); + if (request.method === "PUT" && url.pathname.startsWith("/archive/")) return putArchive(request, url.pathname.slice("/archive/".length), env); if (request.method !== "GET") return json({ error: "method not allowed" }, 405); switch (url.pathname) { case "/api/progress": @@ -183,14 +307,16 @@ export default { case "/api/heat": return heat(url, env); case "/api/collections": - return collections(env); + return collections(url, env); + case "/api/namespaces": + return namespaces(env); case "/": return new Response(page, { headers: { "content-type": "text/html; charset=utf-8", "cache-control": "public, max-age=300" } }); case "/api": return json({ name: "strata", what: "per-segment, per-collection event counts of the stream.waow.tech archive", - endpoints: ["/api/progress", "/api/heat?step=N", "/api/collections"], + endpoints: ["/api/progress", "/api/heat?step=N&top=N&group=namespaces|collections", "/api/heat?step=N&ns=", "/api/namespaces", "/api/collections?ns="], source: "https://tangled.org/zat.dev/strata", }); default: diff --git a/src/ingest.ts b/src/ingest.ts index cd63937..43ae8e5 100644 --- a/src/ingest.ts +++ b/src/ingest.ts @@ -9,6 +9,12 @@ export interface CollectionCount { readonly count: number; } +/** the authority prefix of an nsid: its first two labels (`app.bsky.feed.like` → `app.bsky`). */ +export function namespaceOf(nsid: string): string { + const parts = nsid.split("."); + return parts.length <= 2 ? nsid : `${parts[0]}.${parts[1]}`; +} + export interface SegmentRecord { readonly idx: number; readonly name: string; diff --git a/tests/ingest.test.ts b/tests/ingest.test.ts index c17253d..6090fdb 100644 --- a/tests/ingest.test.ts +++ b/tests/ingest.test.ts @@ -1,7 +1,7 @@ import assert from "node:assert/strict"; import { test } from "node:test"; -import { IngestParseError, parseIngestBatch } from "../src/ingest"; +import { IngestParseError, namespaceOf, parseIngestBatch } from "../src/ingest"; const segment = { idx: 6000, @@ -43,3 +43,9 @@ test("rejects non-integer counts", () => { const bad = { ...segment, collections: [{ nsid: "app.bsky.feed.like", count: "many" }] }; assert.throws(() => parseIngestBatch({ segments: [bad] }), /count: expected a non-negative integer/); }); + +test("namespaceOf takes the first two labels", () => { + assert.equal(namespaceOf("app.bsky.feed.like"), "app.bsky"); + assert.equal(namespaceOf("net.anisota.beta.game.log"), "net.anisota"); + assert.equal(namespaceOf("sh.tangled"), "sh.tangled"); +}); diff --git a/wrangler.jsonc b/wrangler.jsonc index 6d8d04a..2ed24c5 100644 --- a/wrangler.jsonc +++ b/wrangler.jsonc @@ -7,6 +7,10 @@ "observability": { "enabled": true }, "workers_dev": false, "routes": [{ "pattern": "strata.zat.dev", "custom_domain": true }], + "r2_buckets": [ + { "binding": "PUBLIC", "bucket_name": "strata-archive" }, + { "binding": "PRIVATE", "bucket_name": "strata-archive-private" } + ], "d1_databases": [ { "binding": "DB", -- 2.51.2