From 46f41956a4873800e0cf4c6d4195637c4f48c46a Mon Sep 17 00:00:00 2001 From: Florian <45694132+flo-bit@users.noreply.github.com> Date: Sat, 25 Apr 2026 19:23:20 +0200 Subject: [PATCH] small fixes labels --- README.md | 27 +-- docs/08-labels.md | 6 +- packages/contrail/src/contrail.ts | 62 ++++-- packages/contrail/src/core/labels/select.ts | 19 +- packages/contrail/src/core/labels/types.ts | 4 - packages/contrail/src/worker/index.ts | 8 +- packages/contrail/tests/labels-router.test.ts | 2 +- packages/contrail/tests/labels.test.ts | 9 +- todo/communities-as-labelers.md | 187 ++++++++++++++++++ 9 files changed, 253 insertions(+), 71 deletions(-) create mode 100644 todo/communities-as-labelers.md diff --git a/README.md b/README.md index d9605f8..efd97d0 100644 --- a/README.md +++ b/README.md @@ -8,7 +8,8 @@ a library for easily creating (serverless) atproto backends/appviews. - get automatic jetstream backfill and ingestion, typed XRPC endpoints - optional: permissioned spaces and group-controlled communities -mostly tested on cloudflare workers with d1 but should run in any node env too (+ has adapters for node:sqlite and postgres for the db). +mostly tested on cloudflare workers with d1 but should run in any node env too +(+ has adapters for node:sqlite and postgres for the db). ## Install @@ -18,7 +19,7 @@ pnpm add @atmo-dev/contrail ## Minimal example -a complete cloudflare worker that indexes public calendar events from the atproto network and serves them over a typed XRPC endpoint. two files + config. a runnable version lives in [`apps/cloudflare-workers`](./apps/cloudflare-workers) — clone, deploy, `pnpm contrail backfill --remote`, done. +a complete cloudflare worker that indexes public calendar events from the atproto network and serves them over a typed XRPC endpoint. two files + config. a runnable version lives in [`apps/cloudflare-workers`](https://github.com/flo-bit/contrail/tree/main/apps/cloudflare-workers) — clone, deploy, `pnpm contrail backfill --remote`, done. **`src/contrail.config.ts`** — picked up automatically by the `contrail` CLI: @@ -73,21 +74,21 @@ the worker keeps itself fresh from now on via the cron. hit: GET https://.workers.dev/xrpc/com.example.event.listRecords?startsAtMin=2026-01-01&limit=10 ``` -returns every `community.lexicon.calendar.event` record published anywhere on atproto that matches, as JSON. that's it — no PDS setup, no lexicon publishing, no relay configuration. everything scales from there: add filters, add full-text search, add more collections, turn on [spaces](./docs/05-spaces.md) for private records, mount the handler in sveltekit instead, swap the adapter for postgres. +returns every `community.lexicon.calendar.event` record published anywhere on atproto that matches, as JSON. that's it — no PDS setup, no lexicon publishing, no relay configuration. everything scales from there: add filters, add full-text search, add more collections, turn on [spaces](https://github.com/flo-bit/contrail/blob/main/docs/05-spaces.md) for private records, mount the handler in sveltekit instead, swap the adapter for postgres. -**not using workers?** same library, different `db`. see [adapters](./docs/01-indexing.md#adapters) for node:sqlite and postgres. +**not using workers?** same library, different `db`. see [adapters](https://github.com/flo-bit/contrail/blob/main/docs/01-indexing.md#adapters) for node:sqlite and postgres. ## Docs -- [Indexing](./docs/01-indexing.md) — the core: collections, ingestion, adapters -- [Querying](./docs/02-querying.md) — filters, sorts, hydration, search, pagination -- [Lexicons](./docs/03-lexicons.md) — `contrail-lex` CLI, codegen, publishing -- [Auth](./docs/04-auth.md) — service-auth JWTs, invite tokens, watch tickets, OAuth permission sets -- [Spaces](./docs/05-spaces.md) — permissioned records stored by the appview -- [Communities](./docs/06-communities.md) — group-controlled atproto DIDs -- [Sync](./docs/07-sync.md) — reactive client-side store over `watchRecords` -- [Labels](./docs/08-labels.md) — atproto-native moderation hydration from external labelers -- Frameworks: [SvelteKit + Cloudflare](./docs/frameworks/sveltekit-cloudflare.md) +- [Indexing](https://github.com/flo-bit/contrail/blob/main/docs/01-indexing.md) — the core: collections, ingestion, adapters +- [Querying](https://github.com/flo-bit/contrail/blob/main/docs/02-querying.md) — filters, sorts, hydration, search, pagination +- [Lexicons](https://github.com/flo-bit/contrail/blob/main/docs/03-lexicons.md) — `contrail-lex` CLI, codegen, publishing +- [Auth](https://github.com/flo-bit/contrail/blob/main/docs/04-auth.md) — service-auth JWTs, invite tokens, watch tickets, OAuth permission sets +- [Spaces](https://github.com/flo-bit/contrail/blob/main/docs/05-spaces.md) — permissioned records stored by the appview +- [Communities](https://github.com/flo-bit/contrail/blob/main/docs/06-communities.md) — group-controlled atproto DIDs +- [Sync](https://github.com/flo-bit/contrail/blob/main/docs/07-sync.md) — reactive client-side store over `watchRecords` +- [Labels](https://github.com/flo-bit/contrail/blob/main/docs/08-labels.md) — atproto-native moderation hydration from external labelers +- Frameworks: [SvelteKit + Cloudflare](https://github.com/flo-bit/contrail/blob/main/docs/frameworks/sveltekit-cloudflare.md) ## Packages diff --git a/docs/08-labels.md b/docs/08-labels.md index f6c92dd..5dad220 100644 --- a/docs/08-labels.md +++ b/docs/08-labels.md @@ -39,7 +39,7 @@ Per request, contrail picks accepted labelers in this order: 3. `config.labels.defaults` — operator policy. 4. Every entry in `config.labels.sources`. -The list is intersected with what's actually configured (unknowns dropped — see [`allowUserSupplied`](#allowusersupplied) below) and capped at `maxPerRequest` (default 20). Contrail echoes the applied set back via `atproto-content-labelers`. +The list is intersected with what's actually configured (unknowns dropped — only labelers we've subscribed to have rows to hydrate from) and capped at `maxPerRequest` (default 20). Contrail echoes the applied set back via `atproto-content-labelers`. ``` GET /xrpc/com.example.event.listRecords @@ -70,10 +70,6 @@ GET /xrpc/com.example.event.listRecords `labels` matches `com.atproto.label.defs#label` field-for-field — pass it straight to atproto SDK moderation helpers. -### `allowUserSupplied` - -Default: `false` — caller-supplied DIDs that aren't in `sources` are silently dropped. Set `true` to honor them anyway. The current request still only returns labels for already-indexed sources; lazy registration of new labelers is future work. - ### `defaults: []` Set defaults to an empty array if you want strict opt-in: callers that send no header / param see no labels at all. diff --git a/packages/contrail/src/contrail.ts b/packages/contrail/src/contrail.ts index 630fbc9..7c52aa0 100644 --- a/packages/contrail/src/contrail.ts +++ b/packages/contrail/src/contrail.ts @@ -15,7 +15,7 @@ import { runPersistent as runPersistentIngestion } from "./core/persistent"; import type { PersistentIngestOptions } from "./core/persistent"; import { runLabelIngestCycle, - runPersistentLabels, + runPersistentLabels as runPersistentLabelsImpl, type PersistentLabelsOptions, } from "./core/labels/subscribe"; import type { PubSub } from "./core/realtime/types"; @@ -84,29 +84,49 @@ export class Contrail { return queryRecords(this.getDb(db), this.config, { collection, ...options }); } - /** Run one Jetstream ingestion cycle (catches up to present, then stops). */ + /** Run one ingestion cycle: catches up records from Jetstream and — when + * `config.labels` is set — labels from each configured labeler in parallel. + * Both share the same `timeoutMs` budget; they're independent network + * operations so concurrency is free. */ async ingest(options?: { timeoutMs?: number }, db?: Database): Promise { - await runIngestCycle( - this.getDb(db), - this.config, - options?.timeoutMs, - this._ingestState, - this._pubsub ?? undefined - ); + const d = this.getDb(db); + const tasks: Promise[] = [ + runIngestCycle(d, this.config, options?.timeoutMs, this._ingestState, this._pubsub ?? undefined), + ]; + if (this.config.labels) { + tasks.push(runLabelIngestCycle(d, this.config, options?.timeoutMs)); + } + await Promise.all(tasks); } - /** Run persistent Jetstream ingestion (long-lived, stays connected). */ + /** Long-lived ingestion: streams records via Jetstream and — when + * `config.labels` is set — labels via per-labeler `subscribeLabels` sockets. + * Both honor the supplied `signal` and shut down cleanly together. */ async runPersistent(options?: Omit, db?: Database): Promise { - await runPersistentIngestion(this.getDb(db), this.config, { - ...options, - logger: this.config.logger, - pubsub: this._pubsub ?? undefined, - }); + const d = this.getDb(db); + const tasks: Promise[] = [ + runPersistentIngestion(d, this.config, { + ...options, + logger: this.config.logger, + pubsub: this._pubsub ?? undefined, + }), + ]; + if (this.config.labels) { + tasks.push( + runPersistentLabelsImpl(d, this.config, { + signal: options?.signal, + batchSize: options?.batchSize, + flushIntervalMs: options?.flushIntervalMs, + logger: this.config.logger, + }), + ); + } + await Promise.all(tasks); } - /** Run one labeler ingestion cycle — for every labeler in `config.labels.sources`, - * drains pending `subscribeLabels` frames and persists them to the `labels` - * table. No-op when `config.labels` is unset. Mirrors `ingest()`. */ + /** Run *only* the labeler ingestion cycle. Escape hatch for callers who + * want to run record and label ingestion in separate processes / workers. + * `ingest()` already covers the typical case. */ async ingestLabels( options?: { timeoutMs?: number }, db?: Database, @@ -115,14 +135,14 @@ export class Contrail { await runLabelIngestCycle(this.getDb(db), this.config, options?.timeoutMs); } - /** Long-lived label ingestion — one socket per labeler, auto-reconnect on drop. - * No-op when `config.labels` is unset. Mirrors `runPersistent()`. */ + /** Run *only* the persistent labeler ingestion. Escape hatch counterpart + * to `ingestLabels()`. `runPersistent()` covers the typical case. */ async runPersistentLabels( options?: Omit, db?: Database, ): Promise { if (!this.config.labels) return; - await runPersistentLabels(this.getDb(db), this.config, { + await runPersistentLabelsImpl(this.getDb(db), this.config, { ...options, logger: this.config.logger, }); diff --git a/packages/contrail/src/core/labels/select.ts b/packages/contrail/src/core/labels/select.ts index 4293109..da89d9c 100644 --- a/packages/contrail/src/core/labels/select.ts +++ b/packages/contrail/src/core/labels/select.ts @@ -9,20 +9,14 @@ import { DEFAULT_LABELS_MAX_PER_REQUEST } from "./types"; * 3. `config.defaults` (operator policy) * 4. every entry in `config.sources` * - * Each candidate DID is checked against `config.sources`. Unknown DIDs are - * dropped unless `allowUserSupplied: true`, in which case they're returned - * in `lazyAdd` for the caller to schedule a registration. The active - * request gets results only from already-indexed sources. + * Each candidate DID is checked against `config.sources`. Unknowns are + * dropped — we only have rows for labelers we've subscribed to. * * Header values can carry `;param` modifiers (e.g. `did:plc:...;redact`); * v1 strips and ignores those — only the bare DID is honored. */ export interface SelectedLabelers { /** DIDs to use for hydration this request. */ accepted: string[]; - /** DIDs the caller asked for that aren't configured (only populated when - * `allowUserSupplied: true`). The hydrator ignores these for the current - * request — the caller should enqueue registration as a follow-up. */ - lazyAdd: string[]; } export function selectAcceptedLabelers( @@ -43,20 +37,15 @@ export function selectAcceptedLabelers( } const accepted: string[] = []; - const lazyAdd: string[] = []; const seen = new Set(); for (const did of candidates) { if (seen.has(did)) continue; seen.add(did); - if (known.has(did)) { - accepted.push(did); - } else if (cfg.allowUserSupplied) { - lazyAdd.push(did); - } + if (known.has(did)) accepted.push(did); if (accepted.length >= cap) break; } - return { accepted, lazyAdd }; + return { accepted }; } /** Parse a comma-separated DID list. Returns null when the input is empty diff --git a/packages/contrail/src/core/labels/types.ts b/packages/contrail/src/core/labels/types.ts index c5833cf..f8bc2e6 100644 --- a/packages/contrail/src/core/labels/types.ts +++ b/packages/contrail/src/core/labels/types.ts @@ -19,10 +19,6 @@ export interface LabelsConfig { * `?labelers=`. Defaults to every entry in `sources`. Set `[]` for * opt-in-only — clients see no labels unless they ask. */ defaults?: string[]; - /** Honor caller-supplied DIDs that aren't in `sources`. Default: false. - * When true, unknown DIDs in the request are accepted (and the labeler - * may be lazily registered for ingest in a later release). */ - allowUserSupplied?: boolean; /** Per-request cap. Default: 20 (matches Bluesky). */ maxPerRequest?: number; } diff --git a/packages/contrail/src/worker/index.ts b/packages/contrail/src/worker/index.ts index c9213d7..26a412f 100644 --- a/packages/contrail/src/worker/index.ts +++ b/packages/contrail/src/worker/index.ts @@ -63,13 +63,9 @@ export function createWorker( ): Promise { const db = env[binding] as Database; await ensureReady(env, db); + // ingest() drives both record and label ingestion in parallel when + // labels are configured — single waitUntil covers the whole cron tick. ctx.waitUntil(contrail.ingest({}, db)); - // Run label ingest alongside the jetstream catch-up so a single cron - // tick keeps both data streams fresh. Only scheduled when configured — - // skipping the no-op promise keeps the worker task list tight. - if (config.labels) { - ctx.waitUntil(contrail.ingestLabels({}, db)); - } }, }; } diff --git a/packages/contrail/tests/labels-router.test.ts b/packages/contrail/tests/labels-router.test.ts index ff1cea6..86c0dc1 100644 --- a/packages/contrail/tests/labels-router.test.ts +++ b/packages/contrail/tests/labels-router.test.ts @@ -127,7 +127,7 @@ describe("labels router integration", () => { expect(body.records[0]!.labels).toBeUndefined(); }); - it("unknown caller-supplied DIDs are dropped (allowUserSupplied off)", async () => { + it("unknown caller-supplied DIDs are dropped", async () => { const { contrail } = await setup(); const app = contrail.app(); diff --git a/packages/contrail/tests/labels.test.ts b/packages/contrail/tests/labels.test.ts index 317064c..899bb08 100644 --- a/packages/contrail/tests/labels.test.ts +++ b/packages/contrail/tests/labels.test.ts @@ -129,7 +129,6 @@ describe("labels: selectAcceptedLabelers", () => { it("falls back to defaults (= sources) when caller sends nothing", () => { const sel = selectAcceptedLabelers(null, null, cfg); expect(sel.accepted).toEqual([SRC_A, SRC_B]); - expect(sel.lazyAdd).toEqual([]); }); it("honors header before query param", () => { @@ -142,20 +141,18 @@ describe("labels: selectAcceptedLabelers", () => { expect(sel.accepted).toEqual([SRC_B]); }); - it("drops unknown DIDs by default", () => { + it("drops unknown DIDs", () => { const sel = selectAcceptedLabelers("did:plc:strangerlabeler", null, cfg); expect(sel.accepted).toEqual([]); - expect(sel.lazyAdd).toEqual([]); }); - it("collects unknowns as lazyAdd when allowUserSupplied is true", () => { + it("mixes known + unknown — keeps only the known", () => { const sel = selectAcceptedLabelers( `did:plc:strangerlabeler,${SRC_A}`, null, - { ...cfg, allowUserSupplied: true }, + cfg, ); expect(sel.accepted).toEqual([SRC_A]); - expect(sel.lazyAdd).toEqual(["did:plc:strangerlabeler"]); }); it("strips ;param modifiers from header values", () => { diff --git a/todo/communities-as-labelers.md b/todo/communities-as-labelers.md new file mode 100644 index 0000000..fa57ca4 --- /dev/null +++ b/todo/communities-as-labelers.md @@ -0,0 +1,187 @@ +# Communities as labelers + +## Context + +Layer 1 of the labels module ships in v0.4: contrail subscribes to external +labelers, indexes their labels, hydrates `record.labels` onto responses +(see `docs/08-labels.md`). The next interesting question is whether +contrail-managed *community* DIDs can themselves act as labelers. + +This is a feature that doesn't really exist in atproto today: most labelers +are individual operator accounts. Group-controlled, role-gated labeling +where the labeler DID is shared by a moderation team — and from the outside +looks like any other atproto labeler — falls out almost for free from +combining the existing community module with the labels machinery. + +## What a labeler needs + +Three independent contracts in atproto-native fashion: + +1. **Advertise** — `#atproto_labeler` service entry in the DID doc so + clients can find the WS endpoint. Optionally, an + `app.bsky.labeler.service` record with custom-value metadata + (display name, blur policy, severity). +2. **Author** — a way to create signed labels with monotonic per-`src` + `seq` numbers, retractable via `neg`. +3. **Serve** — `com.atproto.label.queryLabels` (paginated reads) and + `com.atproto.label.subscribeLabels` (CBOR-framed WS firehose). + +We have most of (2) already — the `labels` table is the storage. We'd +just be writing to it from a new XRPC instead of from the WS consumer. +(3) is layer 1 in reverse: same wire format, our side. (1) is the only +genuinely new thing, and even that mostly reuses existing community-signing +infrastructure. + +## Mint vs adopt + +- **Mint**: contrail holds a rotation key, so updating the DID doc to + add `#atproto_labeler` is a routine signed PLC operation — same path + the mint flow already uses. ✅ Supported in v1. +- **Adopt**: we hold an app password, not the rotation keys. We can + write `app.bsky.labeler.service` to the existing PDS, but cannot + update the DID doc to add the service entry. The owner has to do + that manually (one-shot PLC operation, documented). 🟡 Documented + manual step. + +## The PDS gotcha (resolved) + +Earlier sketch assumed we'd publish `app.bsky.labeler.service` to the +community's repo. But minted communities have no PDS — they're DIDs in +PLC with no repo. Spent a while exploring "contrail hosts a tiny repo +for minted communities" (option B) and concluded the simpler answer: +**most labelers don't need a service record**. + +In atproto, you don't browse a directory of labelers — you find them +out-of-band and pass the DID to your client. The DID doc service entry +is what makes the labeler reachable; the service record is just metadata +for custom label values. + +For standard atproto label values (`!hide`, `!warn`, +`!no-unauthenticated`, `porn`, `nudity`, `graphic-media`, `sexual`, +etc.) clients have hardcoded behavior. The service record only matters +for **custom** label values, where without it third-party clients +display the raw string with default styling. + +So we tier the implementation: + +| Tier | Mechanism | Custom values | Labeler-name in clients | +|---|---|---|---| +| 1 | DID doc service entry only | raw strings | "anonymous" | +| 2 | Tier 1 + hosted service-record endpoint | full metadata | yes | +| 3 | Tier 2 + general mini-PDS | n/a | (out of scope) | + +**Tier 1 is shippable as v1.** Tier 2 is small and additive — one read +endpoint, one write XRPC, one new table — when someone actually wants +custom values. Tier 3 (general mini-PDS for minted communities) is a +separate architecture question that should ride or die on its own +merits, not on labels: it brings in the firehose-participation problem, +MST + commits, blob storage, and operator-as-durability-promise. Out +of scope. + +## Sketch — Tier 1 surface + +**Setup (admin / owner):** +``` +community.becomeLabeler { communityDid } + → updates DID doc (PLC op) to add service[id="#atproto_labeler", + serviceEndpoint=] +community.unbecomeLabeler { communityDid } + → reverse +``` + +**Authoring (configurable level — default `moderator`):** +``` +community.label.create { + communityDid, + subject: { uri: "at://...", cid?: "..." } | { did: "did:plc:..." }, + val, + exp? // ISO-8601 +} +community.label.negate { communityDid, subject, val } +community.label.list { communityDid, cursor?, limit? } // moderation UI +``` + +A `community.label.create` call: + +1. Verifies caller's access level via the existing community ACL. +2. Optionally checks `val` against a per-community allowlist (Tier 2). +3. Builds the label object, signs with the community's signing key + (existing `CredentialCipher` + signer). +4. Inserts into `labels` with `src = communityDid` and a freshly + minted per-src `seq`. +5. Publishes to a `labels:` pubsub topic for live + subscribers. + +**Public, proxy-routed:** +``` +com.atproto.label.queryLabels ?uriPatterns=...&sources=did:plc:... +com.atproto.label.subscribeLabels WS, ?cursor=N +``` + +Both read `Atproto-Proxy: #atproto_labeler` to know which +community-as-labeler the caller is asking about, then filter +`labels WHERE src = ?` and serve. Live mode subscribes to +`labels:` and emits CBOR frames — structurally the +inverse of layer 1's WS consumer, with the same frame protocol. + +## Storage additions + +```sql +ALTER TABLE labels ADD COLUMN seq INTEGER; +-- nullable: NULL for ingested-from-remote rows, allocated for our own +-- outbound labels. seq IS NOT NULL is the marker for "we wrote this." + +CREATE TABLE labeler_outbound ( + src TEXT PRIMARY KEY, -- community DID acting as labeler + next_seq INTEGER NOT NULL DEFAULT 1 +); +-- atomic increment per label publish; serves as the seq source for +-- subscribeLabels. +``` + +Skip ingesting our own outbound stream — one extra check in the WS +consumer (`if known-community-DID, skip subscription`). + +## Suggested order + +Each step is independently shippable: + +1. **`labels.seq` column + `labeler_outbound` counter.** Schema-only; + no behavioural change. +2. **`community.label.create / .negate / .list`.** Operator-internal + use only — labels go into the table, hydration starts surfacing them + on caller responses, but no external discovery yet. Useful by itself + if the operator's own appview is the only consumer. +3. **`community.becomeLabeler`.** DID-doc update. Now external clients + can discover. +4. **`com.atproto.label.queryLabels` (proxy-routed).** External read + API. Most clients use this for one-shot lookups. +5. **`com.atproto.label.subscribeLabels` (WS).** External firehose. + Reuses the realtime pubsub. +6. **Realtime topic + live emission.** Wires `labels:` + events into `subscribeLabels` for live updates. + +Roughly 1–2 days of work per step. + +## Open questions + +1. **Which access level can label?** Add a config field: + `community.labelers: ["moderator", "admin", "owner"]` ranked. + Default to `manager` upward. +2. **Per-community label-value allowlist?** Bluesky enforces "labelers + can only emit values they declared in their service record." Without + a service record (Tier 1), there's nothing to enforce against. When + we layer Tier 2, plumb the allowlist check through. +3. **Permission set hook.** Auto-generated `.permissionSet` + should pick up `community.label.create` etc. so OAuth consent + renders correctly. Should fall out of the existing permission-set + generator since it already enumerates all XRPCs. + +## What this gets you that's actually new + +A moderation team gets actual access-level structure (junior moderators +can label `!warn` only, senior can label anything, owners can rotate +keys), and from the outside looks like a single labeler DID — Bluesky's +app, other appviews, third-party clients all consume it normally. That's +a reasonably interesting primitive that falls out almost for free from +layers 1 + community. -- 2.51.2