diff --git a/_config.ts b/_config.ts index f9aa998c..be4f2fb9 100644 --- a/_config.ts +++ b/_config.ts @@ -39,6 +39,7 @@ site.use(esbuild({ options: { alias: { "@automerge/automerge": "https://esm.sh/@automerge/automerge@^3.2.3", + "@panproto/core": "https://esm.sh/@panproto/core@0.72.0", }, bundle: true, format: "esm", @@ -157,6 +158,17 @@ site.remoteFile( import.meta.resolve("./node_modules/98.css/dist/98.css"), ); +// panproto (`@panproto/core`) loads its WASM lazily at browser runtime during a +// schema write-back. Serve the module's binary so the browser can fetch it +// without loading it into the Deno/build context. +site.remoteFile( + "panproto_wasm_bg.wasm", + import.meta.resolve( + "./node_modules/@panproto/core/dist/panproto_wasm_bg.wasm", + ), +); +site.add([".wasm"]); + //////////////////////////////////////////// // BINARY ASSETS //////////////////////////////////////////// diff --git a/deno.jsonc b/deno.jsonc index f5f0f96d..80ee4052 100644 --- a/deno.jsonc +++ b/deno.jsonc @@ -1,6 +1,6 @@ { "name": "@toko/diffuse", - "version": "4.0.0-alpha.1", + "version": "4.0.0-alpha.2", "license": "FSL-1.1-MIT", "vendor": true, "imports": { @@ -39,6 +39,7 @@ "@std/semver": "jsr:@std/semver@^1.0.8", "@std/xml": "jsr:@std/xml@^0.1.2", "@vicary/debounce-microtask": "jsr:@vicary/debounce-microtask@^0.1.8", + "@panproto/core": "npm:@panproto/core@^0.72.0", "@paulmillr/qr": "jsr:@paulmillr/qr@^0.6.0", "@zip-js/zip-js": "jsr:@zip-js/zip-js@^2.8.26", "alien-signals": "npm:alien-signals@^3.2.1", diff --git a/docs/design/self-describing-lenses.md b/docs/design/self-describing-lenses.md new file mode 100644 index 00000000..749de404 --- /dev/null +++ b/docs/design/self-describing-lenses.md @@ -0,0 +1,127 @@ +# Design: Self-describing saved data with panproto lenses + +Status: **Decision-record / largely implemented**. Envelope, migration engine, and +the panproto connection are in place. Open items: browser-build WASM bundling +verification and a first real schema migration. + +## Schema identity: the NSID, not an integer version + +Following the atproto spec ("Lexicon Evolution"), there is **no numeric schema +version** in a lexicon — the NSID is the schema's identity. Compatible evolution +(add optional fields) keeps the same NSID and needs no migration; a breaking change +(rename/remove a field, change a type) uses a **new NSID**, with an authored lens +from the old NSID to the new one. Diffuse's envelope records `$schema` (the current +NSID) plus an ordered `$schemaHistory` of NSID transitions. + +The `lexicon` integer in a lexicon file is the *language* version (always `1`), not +a schema revision. + +## Envelope + +```jsonc +{ + "$schema": "sh.diffuse.output.track2", // current NSID (identity) + "$schemaHistory": [ + { + "from": "sh.diffuse.output.track", // prior NSID + "to": "sh.diffuse.output.track2", + "lens": { // portable DSL document (steps) + "id": "track-to-track2", + "source": "sh.diffuse.output.track", + "target": "sh.diffuse.output.track2", + "steps": [ + { "rename_field": { "old": "duration", "new": "durationMs" } }, + { "add_field": { "parent": "sh.diffuse.output.track2:body", "name": "gain", "kind": "number" } } + ] + }, + "complement": null // opaque bytes for lossless write-back + } + ], + "data": [ /* current-shape records */ ] +} +``` + +`$schemaHistory` is unbounded: every authored NSID transition appends one entry, so +an older app can walk back to its own shape and write back. + +## Component coverage + +The envelope is produced at the encoder, not the storage element: + +| Encoder / output | Stored format | Where the envelope lives | +|---|---|---| +| `string/json` | JSON string | whole envelope | +| `bytes/json` | JSON bytes | whole envelope | +| `bytes/s3`, `bytes/dropbox` | raw bytes | transparent pass-through of envelope bytes | +| `polymorphic/indexed-db` | IDB object | `encodeCollection` for arrays; bytes pass through | +| `bytes/dasl-sync` | CBOR `Container` | `$schema` on the container | +| `bytes/automerge` | Automerge binary | `$schema` field in the CRDT doc | +| `raw/atproto-space` (records) | atproto records | `$type` = NSID on each record | +| `raw/atproto-space` (blob bundles) | CBOR blob | whole envelope | + +## Modules + +- `src/common/self-describing.js` — envelope `wrap`/`unwrap`/`isSelfDescribing`; + `collectionSchema(name)` returns each collection's current NSID. +- `src/common/lens-registry.js` — `register()`/`resolve(fromNsid, toNsid)`: the + app-bundled lenses (resolved first), falling back to embedded history lenses. +- `src/common/lens.js` — pure-JS migration engine: `project(records, lens)`, + `migrate(...)`, `migrateEnvelope(value, name, resolve)`; plus shared + `encodeCollection`/`decodeCollection`/`encodeJsonCollection`/`decodeJsonCollection`, + and `writeBack(editedRecord, { lens, toLexicon, complement })` — the write-back + path. +- `src/common/panproto.js` — the `@panproto/core` (WASM) connection point, loaded + lazily; `compileLens`/`parseLexicon`/`get`/`put` for the lossless complement + write-back path. + +panproto's WASM is loaded **only on write-back**: `writeBack` uses the pure-JS +`project` when no complement is involved, and calls panproto `put` (lazily loading +`@panproto/core`) only when a stored complement must preserve discarded fields. +Ordinary reads never touch panproto. + +## Lens-source resolution (two tiers) + +1. **App-bundled registry** — lenses shipped with the build (cheap, well-known). +2. **Embedded in `$schemaHistory`** — required for old-app-reads-newer-data: an old + app cannot know lenses authored after it shipped, so the segment it must traverse + is embedded. + +## Read / write paths + +- Any app reads a payload whose `$schema` matches its own: parse `data` directly. +- A stale payload (different NSID) is projected to the current NSID on read + (`migrateEnvelope`); the transition is appended to `$schemaHistory`. +- `src/common/lens.js` exposes `writeBack(record, { lens, toLexicon, complement })` + and `writeBackCollection(items, name, storedEnvelope, toLexicon)` for lossless + write-back via panproto `put` (the `complement` preserves discarded fields; + panproto is lazy-imported only when a complement is actually used). For diffuse's + additive/rename migrations the pure-JS `project` suffices and panproto is never + loaded. + +### Write-back gating + +Lossless write-back to a **newer** shape requires the newer lexicon (`toLexicon`) +to instantiate the lens — which an *older* build that only knows the old shape does +not have. So write-back is practical at a migration site that holds *both* the newer +lexicon and the older data being lifted forward, not from an old writer against a +newer stored shape. `writeBackCollection` is a no-op (records pass through, no WASM) +until a caller supplies a complement-carrying history entry **and** the target +lexicon; this becomes live at the first real schema migration. + +## Validation + +`deno check src` / `deno check specs`, and the doc-test suite (`deno task test:doc`) +cover `self-describing.js`, `lens-registry.js`, `lens.js`, and `panproto.js`. +Browser integration tests (astral/Chromium) cannot run in every sandbox. + +## Open items + +- Verify Lume/esbuild bundles `@panproto/core`'s WASM (`deno task build`) — it is a + dependency and `panproto.js` is the lazy import point; the build integration needs + checking in a normal dev environment. +- Update `deno.lock`/vendor for the new dependency (blocked by sandbox FS limits + here — vendoring writes to Deno's global npm cache). +- Author the first real NSID migration (with actual old + new lexicons) to exercise + `migrateEnvelope` + panproto write-back end to end. + +See `docs/guides/updating-lexicons.md` for the hands-on workflow. \ No newline at end of file diff --git a/docs/guides/updating-lexicons.md b/docs/guides/updating-lexicons.md new file mode 100644 index 00000000..d078dd6e --- /dev/null +++ b/docs/guides/updating-lexicons.md @@ -0,0 +1,187 @@ +# Updating a lexicon (schema change) & migrating stored data + +This is the workflow for changing a diffuse schema lexicon and getting stored data +auto-migrated to the new shape. + +## When you need this + +Diffuse persists four collections per output — `facets`, `playlistItems`, `settings`, +`tracks` — each governed by an atproto lexicon in `lexicons/output/` +(e.g. `sh.diffuse.output.facet`, `sh.diffuse.output.track`). If you change a +lexicon's shape, existing stored data is in the **old** shape. This guide wires +that migration. + +### How atproto versions lexicons (the model we follow) + +In keeping with the atproto spec ("Lexicon Evolution"), the atproto community does +**not** put an integer version in a lexicon. The rules: + +- **The lexicon NSID is the schema's identity.** There is no numeric schema version. +- **Compatible evolution keeps the same NSID**: new fields must be optional; never + remove non-optional fields (keep them and mark deprecated); never change types or + rename fields. Old data stays valid under the updated lexicon, and new data stays + valid under the old one. +- **Breaking changes use a NEW NSID** (e.g. `sh.diffuse.output.track2`), not a version + bump on the same NSID. +- The `lexicon` field in a lexicon file is the **language** version (always `1`), not + a schema revision. Leave it alone. + +So "moving to the new shape" is expressed as a lens **from one NSID to another** +(possibly the same NSID, when the change is purely compatible). Diffuse's envelopes +record the NSID (`$schema`) plus an ordered `$schemaHistory` of NSID transitions. + +The self-describing machinery that does the work lives in: + +- `src/common/self-describing.js` — the `{ $schema, $schemaHistory, data }` envelope, + `wrap`/`unwrap`, and `COLLECTION_SCHEMAS` (each collection's current lexicon NSID). +- `src/common/lens-registry.js` — `register()` / `resolve()` for authored lens + documents (the "lenses in app code" tier). +- `src/common/lens.js` — `project()`, `migrate()`, `migrateEnvelope()`; the engine + that applies a lens's steps to project records between NSIDs. +- Each output encoder's read path (indexed-db, `bytes/json`, `string/json`, + `bytes/dasl-sync`, `bytes/automerge`, atproto-space blob bundles) calls + `migrateEnvelope` so stale payloads are migrated on read. + +## Steps + +### 1. Make the schema change + +Two cases: + +- **Compatible change** (add an optional field): edit the existing lexicon in + `lexicons/output/*.json`. Keep the same NSID; no lens needed (old data remains + valid; the new field simply appears). +- **Breaking change** (rename/remove a field, change a type): create a **new lexicon + file** with a new NSID (e.g. `sh.diffuse.output.track` → `sh.diffuse.output.track2`), + and keep the old NSID's record fields available for migration. You'll author a lens + from the old NSID to the new one. + +### 2. Regenerate the TypeScript types + +The generated types in `src/definitions/types/` come from the lexicons: + +```sh +deno task gen:defs:types +``` + +(this runs `@atcute/lex-cli generate` plus the replace/strip helpers). The +`Main`/`Facet`/`Track`… types consumed across the codebase must match the new shape. + +### 3. Point the collection at the current NSID + +In `src/common/self-describing.js`, set the collection's entry in `COLLECTION_SCHEMAS` +to the **current** NSID. For a breaking change this is the new NSID; for a purely +compatible change it stays the same. + +```js +const COLLECTION_SCHEMAS = { + // ... + tracks: "sh.diffuse.output.track2", // was "sh.diffuse.output.track" +}; +``` + +This is the value stamped on every new save and the NSID the migration projects toward. + +> If the change is purely compatible (same NSID, optional field added) you're done +> after steps 1–2: there's nothing to migrate. Skip to validation. + +### 4. Author & register the lens (old NSID → new NSID) + +Only needed for a **breaking change** (different NSIDs). In +`src/common/lens-registry.js`, add a lens document describing the transition. The +DSL steps supported by the projection engine are `rename_field`, `add_field`, +`remove_field`: + +```js +register({ + id: "track-to-track2", + source: "sh.diffuse.output.track", + target: "sh.diffuse.output.track2", + steps: [ + { rename_field: { old: "duration", new: "durationMs" } }, + { add_field: { parent: "sh.diffuse.output.track2:body", name: "gain", kind: "number" } }, + // { remove_field: { name: "obsolete" } }, + ], +}); +``` + +Notes: + +- The `source`/`target` are NSIDs. `rename_field`/`add_field`/`remove_field` step + between their shapes; `target` (and its `:body` vertex) is the new lexicon. +- These are the reversible, pure-JS steps the migration engine applies; they cover + the common rename/add cases. A field **discarded** by a migration (e.g. + `remove_field`) that an *older app* must still be able to write back needs the + panproto complement path, which is a separate, not-yet-implemented engine (see + open items in `docs/design/self-describing-lenses.md`). + +### 5. That's it — the read path migrates automatically + +With the collection's current NSID updated and the lens registered, no further wiring +is needed: + +- Any payload whose `$schema` is an **older NSID** is projected to the current NSID on + read (via `migrateEnvelope`), and the NSID transition is appended to the payload's + `$schemaHistory`. +- New saves are written under the current NSID. +- The lens is embedded in `$schemaHistory` so even an app that doesn't have it bundled + (an older build) can traverse the transition — this is the "lenses in the data" + correctness guarantee. + +## First migration: handling existing (legacy) data + +Before the envelope existed, stored data was a bare array / bare CBOR / raw IDB value +with **no recorded NSID**. The read path treats such values as belonging to the +collection's *current* NSID (no migration). For a first migration, this means legacy +bare data is treated as the **old** shape only if nothing distinguishes it — typically +you accept it as the current shape, or you avoid a breaking change until the envelope +is universally writing. + +## Recommended validation before shipping + +Add a unit (doc) test that exercises the exact path an old user hits, then run the +suite: + +```ts +import { wrap } from "~/common/self-describing.js"; +import { register, resolve } from "~/common/lens-registry.js"; +import { migrateEnvelope } from "~/common/lens.js"; + +register({ + id: "track-to-track2", source: "sh.diffuse.output.track", target: "sh.diffuse.output.track2", + steps: [{ rename_field: { old: "duration", new: "durationMs" } }], +}); +const envelope = wrap([{ id: "t1", uri: "u", duration: 300 }], + { schema: "sh.diffuse.output.track" }); +const out = migrateEnvelope(envelope, "tracks", resolve); +// assert out.data[0] has `durationMs` and not `duration`, and out.envelope.$schema +// is "sh.diffuse.output.track2" with a $schemaHistory entry for track -> track2. +``` + +Then: + +```sh +deno task test:doc # unit / documentation tests +deno check src # type checks +deno check specs +``` + +> Browser integration tests (`deno task test:integration`) launch a headless +> Chromium via astral; they can't run in every sandbox. If they fail with "Your +> binary refused to boot", that's environmental, not your change. + +## Checklist + +1. Make the change: compatible (edit same-NSID lexicon) or breaking (new lexicon NSID) +2. `deno task gen:defs:types` +3. Point the collection's `COLLECTION_SCHEMAS` entry at the current NSID + (`self-describing.js`) +4. For breaking changes, `register()` the old-NSID → new-NSID lens + (`lens-registry.js`) +5. Add/confirm a migration doc-test; run `deno test --doc`, `deno check src`, + `deno check specs` +6. Confirm legacy (pre-envelope) data handling and that a migration actually triggers + +See `docs/design/self-describing-lenses.md` for the full design (envelope shape, +two-tier lens resolution, read/write paths, and open items like the wasm complement +path). \ No newline at end of file diff --git a/specs/components/transformer/output/bytes/automerge/types.d.ts b/specs/components/transformer/output/bytes/automerge/types.d.ts index 7881fff2..67dc776d 100644 --- a/specs/components/transformer/output/bytes/automerge/types.d.ts +++ b/specs/components/transformer/output/bytes/automerge/types.d.ts @@ -5,7 +5,16 @@ import type { Track, } from "~/definitions/types.d.ts"; -export type FacetsDocument = { collection: Facet[] }; -export type PlaylistItemsDocument = { collection: PlaylistItem[] }; -export type SettingsDocument = { collection: Setting[] }; -export type TracksDocument = { collection: Track[] }; +/** + * The schema identity stamped on an Automerge doc so the stored binary is + * self-describing. + */ +export type DocumentSchema = { + /** The lexicon NSID the records in `collection` conform to. */ + $schema?: string; +}; + +export type FacetsDocument = { collection: Facet[] } & DocumentSchema; +export type PlaylistItemsDocument = { collection: PlaylistItem[] } & DocumentSchema; +export type SettingsDocument = { collection: Setting[] } & DocumentSchema; +export type TracksDocument = { collection: Track[] } & DocumentSchema; diff --git a/specs/components/transformer/output/bytes/dasl-sync/types.d.ts b/specs/components/transformer/output/bytes/dasl-sync/types.d.ts index 43fe4677..6cc22c80 100644 --- a/specs/components/transformer/output/bytes/dasl-sync/types.d.ts +++ b/specs/components/transformer/output/bytes/dasl-sync/types.d.ts @@ -1,4 +1,28 @@ +/** + * A self-describing container of records in a DASL-sync output. + * + * Carries the lexicon NSID that produced `data` and the ordered history of schema + * transitions, so the stored container is interpretable and migratable across + * app versions without external knowledge. + * + * @template {Record} T + */ export type Container = { + /** The lexicon NSID the `data` records conform to. */ + $schema?: string; + + /** + * Ordered history of schema transitions that produced the current shape; + * each carries a portable lens document and the complement threaded through + * the forward projection, so older apps can read and write newer data. + */ + $schemaHistory?: { + from: string; + to: string; + lens: unknown | null; + complement?: Uint8Array | string | null; + }[]; + /** * CID of the inventory, * which in turns represents the current state of the data. diff --git a/src/common/lens-registry.js b/src/common/lens-registry.js new file mode 100644 index 00000000..8783e2b5 --- /dev/null +++ b/src/common/lens-registry.js @@ -0,0 +1,99 @@ +/** + * Lens registry: authored lens documents that migrate diffuse data from one + * lexicon NSID to another. + * + * Lens documents are deliberately authored (per schema change) and schema + * independent: each is an object carrying `id`, `source`, `target`, and `steps`, + * which `@panproto/core`'s `compileLensDocument` accepts. A build's bundled + * registry resolves an `(sourceNsid, targetNsid)` transition to a lens document. + * When a payload's `$schemaHistory` embeds its own lens for a segment the reader + * doesn't know, the embedded document takes precedence for that segment (so old + * apps can read newer data). + * + * In keeping with atproto's lexicon model, the NSID is the schema identity, so + * `source`/`target` are NSIDs. Compatible evolution keeps the same NSID (new + * optional fields, no lens needed); breaking changes use a NEW NSID and an + * authored lens from the old NSID to the new one. + */ + +/** + * A lens document that maps one lexicon NSID to another. + * + * @typedef {{ + * id: string; + * source: string; + * target: string; + * steps: unknown[]; + * }} LensDocument + */ + +/** @type {Map} */ +const registry = new Map(); + +/** + * Register an authored lens document. + * + * @param {LensDocument} doc + * + * @example Registers and looks up a lens by its transition + * ```js + * import { register, resolve } from "~/common/lens-registry.js"; + * + * register({ id: "track-v1-to-v2", source: "sh.diffuse.output.track", target: "sh.diffuse.output.track2", steps: [] }); + * const doc = resolve("sh.diffuse.output.track", "sh.diffuse.output.track2"); + * if (doc?.id !== "track-v1-to-v2") throw new Error("expected registered lens"); + * ``` + */ +export function register(doc) { + registry.set(key(doc.source, doc.target), doc); +} + +/** + * Resolve a lens document for a transition. The bundled registry is consulted + * first; then the `history` supplied with the payload (embedded lenses), which + * lets an older app traverse segments it doesn't have in its own registry. + * + * @param {string} from - The source lexicon NSID + * @param {string} to - The target lexicon NSID + * @param {{ history?: import("./self-describing.js").HistoryEntry[] }} [options] + * @returns {LensDocument | null} + * + * @example Finds a bundled-registry lens before any embedded one + * ```js + * import { register, resolve } from "~/common/lens-registry.js"; + * + * register({ id: "a-to-b", source: "s", target: "t", steps: [] }); + * const doc = resolve("s", "t"); + * if (!doc) throw new Error("expected a lens"); + * ``` + * + * @example Falls back to an embedded lens for a segment not in the registry + * ```js + * import { resolve } from "~/common/lens-registry.js"; + * + * const embedded = { id: "old-segment", source: "s", target: "t", steps: [{ rename_field: { old: "a", new: "b" } }] }; + * const doc = resolve("s", "t", { + * history: [{ + * from: "s", to: "t", lens: embedded, complement: null, + * }], + * }); + * if (doc?.id !== "old-segment") throw new Error("expected embedded lens"); + * ``` + */ +export function resolve(from, to, options = {}) { + const bundled = registry.get(key(from, to)); + if (bundled) return bundled; + + const entry = options.history?.find( + (h) => h.from === from && h.to === to, + ); + return entry?.lens ?? null; +} + +/** + * @param {string} source + * @param {string} target + */ +function key(source, target) { + return `${source}@${target}`; +} \ No newline at end of file diff --git a/src/common/lens.js b/src/common/lens.js new file mode 100644 index 00000000..d5636c1d --- /dev/null +++ b/src/common/lens.js @@ -0,0 +1,569 @@ +/** + * Migration projection engine: applying authored lens documents to records. + * + * Lens documents (see `~/common/lens-registry.js`) describe transitions between + * lexicon NSIDs as schema-independent steps. In keeping with atproto's lexicon + * model, the NSID is the schema identity: compatible evolution keeps the same + * NSID, while a breaking change is a NEW NSID with an authored lens from the old + * NSID to the new one. This engine projects records from one NSID's shape to + * another without the panproto WASM engine for the common, reversible cases + * (renames, additive fields). Use `migrate()` for that projection. + * + * Wild/disruptive changes (e.g. dropping a field where the discarded value must + * be written back) are handled losslessly via the panproto complement machinery + * (see the design doc, open item #1); the pure-step projection here covers the + * additive/rename migrations diffuse actually ships. + * + * @import {LensDocument} from "~/common/lens-registry.js" + */ + +import { collectionSchema, unwrap, wrap } from "./self-describing.js"; +import { resolve } from "./lens-registry.js"; +import { put as panprotoPut, parseLexicon } from "./panproto.js"; + +/** + * Project a collection of records from one NSID's shape to another using a lens + * document's steps. + * + * Supports the DSL document steps diffuse authors: `rename_field`, `add_field`, + * `remove_field`. Records are plain objects; unknown steps are left as-is. + * + * @param {unknown[]} records + * @param {LensDocument} lens + * @returns {unknown[]} + * + * @example Renames a field and adds a defaulted field + * ```ts + * import { project } from "~/common/lens.js"; + * + * const out = project( + * [{ $type: "sh.diffuse.output.facet", id: "a", name: "x", favourite: true }], + * { + * id: "f", source: "sh.diffuse.output.facet", target: "sh.diffuse.output.facet2", + * steps: [ + * { rename_field: { old: "favourite", new: "starred" } }, + * { add_field: { parent: "sh.diffuse.output.facet2:body", name: "description", kind: "string" } }, + * ], + * }, + * ); + * const rec = out[0] as Record; + * if (!("starred" in rec)) throw new Error("expected renamed field"); + * if ("favourite" in rec) throw new Error("old field should be gone"); + * if (!("description" in rec)) throw new Error("expected added field"); + * ``` + * + * @example Remove a field + * ```ts + * import { project } from "~/common/lens.js"; + * + * const out = project([{ id: "a", dropped: true }], { + * id: "f", source: "s", target: "t", + * steps: [{ remove_field: { name: "dropped" } }], + * }); + * if ("dropped" in (out[0] as Record)) throw new Error("expected removed"); + * ``` + */ +export function project(records, lens) { + return records.map((record) => { + if (typeof record !== "object" || record === null || Array.isArray(record)) { + return record; + } + /** @type {Record} */ + const rec = { .../** @type {Record} */ (record) }; + + for (const step of lens.steps) { + const s = /** @type {any} */ (step); + if (s?.rename_field) { + const oldName = /** @type {string} */ (s.rename_field.old); + const newName = /** @type {string} */ (s.rename_field.new); + if (oldName in rec) { + rec[newName] = rec[oldName]; + delete rec[oldName]; + } + } else if (s?.remove_field) { + const name = /** @type {string} */ (s.remove_field.name); + delete rec[name]; + } else if (s?.add_field) { + const name = /** @type {string} */ (s.add_field.name); + if (!(name in rec)) rec[name] = defaultValue(/** @type {string} */(s.add_field.kind)); + } + } + + return rec; + }); +} + +/** + * Migrate a stored envelope to the current lexicon NSID if it is stale, returning + * the migrated data and the updated `$schemaHistory`. When the stored schema + * already matches, it is returned unchanged. + * + * @template {unknown[]} T + * @template {unknown} L + * @param {{ data: T; envelope: import("./self-describing.js").SelfDescribing | null }} stored + * @param {string} current - The current lexicon NSID + * @param {import("./lens-registry.js").resolve} resolveLens + * @returns {{ data: T; history: import("./self-describing.js").HistoryEntry[] }} + * + * @example Migrates a stale envelope (old NSID) to the current NSID + * ```ts + * import { wrap } from "~/common/self-describing.js"; + * import { register, resolve } from "~/common/lens-registry.js"; + * import { migrate } from "~/common/lens.js"; + * + * register({ + * id: "f-old-new", source: "sh.diffuse.output.facet", target: "sh.diffuse.output.facet2", + * steps: [{ rename_field: { old: "favourite", new: "starred" } }], + * }); + * const envelope = wrap([{ id: "a", favourite: true }], { schema: "sh.diffuse.output.facet" }); + * const out = migrate( + * { data: envelope.data, envelope }, + * "sh.diffuse.output.facet2", + * resolve, + * ); + * if (out.history.length !== 1) throw new Error("expected one history entry"); + * if (!("starred" in (out.data[0] as object))) throw new Error("expected migrated record"); + * ``` + */ +export function migrate(stored, current, resolveLens) { + const env = stored.envelope; + /** @type {import("./self-describing.js").HistoryEntry[]} */ + const storedHistory = /** @type {any} */ (env?.$schemaHistory ?? []); + if (!env || env.$schema === current) { + return { data: stored.data, history: storedHistory }; + } + + const lens = resolveLens(env.$schema, current, { history: storedHistory }); + if (!lens) { + // No lens available for the transition; leave data as-is rather than guess. + return { data: stored.data, history: storedHistory }; + } + + const projected = /** @type {T} */ (project(envelopeToArray(env), lens)); + return { + data: projected, + history: [ + ...storedHistory, + { + from: env.$schema, + to: current, + lens, + complement: null, + }, + ], + }; +} + +/** + * Read a stored value for a collection, migrating it to the collection's current + * lexicon NSID if it is stale. Returns the data to expose and the (possibly + * updated) envelope. + * + * @template {unknown[]} T + * @param {unknown} value + * @param {import("./self-describing.js").CollectionName} name + * @param {import("./lens-registry.js").resolve} resolveLens + * @returns {{ data: T; envelope: import("./self-describing.js").SelfDescribing | null }} + * + * @example Reads a stale envelope and migrates it to the current NSID + * ```ts + * import { wrap } from "~/common/self-describing.js"; + * import { register, resolve } from "~/common/lens-registry.js"; + * import { migrateEnvelope } from "~/common/lens.js"; + * + * register({ + * id: "f-old-current", source: "sh.diffuse.output.facetOld", target: "sh.diffuse.output.facet", + * steps: [{ rename_field: { old: "favourite", new: "starred" } }], + * }); + * // A payload stored under an OLDER lexicon NSID; the current collection NSID is + * // sh.diffuse.output.facet, so this migrates. + * const envelope = wrap([{ id: "a", favourite: true }], { schema: "sh.diffuse.output.facetOld" }); + * const out = migrateEnvelope(envelope, "facets", resolve); + * if (!("starred" in (out.data[0] as object))) throw new Error("expected migrated record"); + * if (!out.envelope || out.envelope.$schema !== "sh.diffuse.output.facet") throw new Error("expected migrated envelope NSID"); + * ``` + */ +export function migrateEnvelope(value, name, resolveLens) { + const current = collectionSchema(name); + const { data, envelope } = unwrap(value, { $schema: current }); + + /** @type {import("./self-describing.js").SelfDescribing | null} */ + const env = /** @type {any} */ (envelope); + + if (!env || env.$schema === current) { + return { data: /** @type {T} */ (data), envelope: env }; + } + + const migrated = migrate({ data: /** @type {T} */ (data), envelope: env }, current, resolveLens); + + /** @type {import("./self-describing.js").SelfDescribing} */ + const migratedEnvelope = { + ...env, + $schema: current, + $schemaHistory: migrated.history, + data: /** @type {T} */ (migrated.data), + }; + return { data: /** @type {T} */ (migrated.data), envelope: migratedEnvelope }; +} + +/** + * Wrap a collection into a self-describing envelope stamped with the collection's + * current NSID, ready to be serialized by an encoder. + * + * @template {unknown[]} T + * @param {T} data + * @param {import("./self-describing.js").CollectionName} name + * @returns {import("./self-describing.js").SelfDescribing} + * + * @example Wraps records for a collection + * ```ts + * import { encodeCollection } from "~/common/lens.js"; + * + * const envelope = encodeCollection([{ id: "a" }], "facets"); + * if (envelope.$schema !== "sh.diffuse.output.facet") throw new Error("expected facet NSID"); + * if (envelope.data[0].id !== "a") throw new Error("expected data"); + * ``` + */ +export function encodeCollection(data, name) { + return wrap(data, { schema: collectionSchema(name) }); +} + +/** + * Decode a stored value (envelope or legacy) into a collection, migrating it to + * the collection's current NSID if stale. Returns `null` when `value` is `null` + * or `undefined`, so encoders can map that to an empty collection. + * + * @template {unknown[]} T + * @param {unknown} value + * @param {import("./self-describing.js").CollectionName} name + * @returns {T | null} + */ +export function decodeCollection(value, name) { + if (value === null || value === undefined) return null; + const { data } = migrateEnvelope(value, name, resolve); + return /** @type {T} */ (data); +} + +/** + * Encode a collection as self-describing JSON, either as a string or as bytes. + * + * @param {unknown[]} data + * @param {import("./self-describing.js").CollectionName} name + * @param {boolean} [asBytes] + * @returns {string | Uint8Array} + * + * @example Encodes a collection as a JSON string + * ```ts + * import { encodeJsonCollection } from "~/common/lens.js"; + * + * const out = encodeJsonCollection([{ id: "a" }], "tracks") as string; + * const parsed = JSON.parse(out) as { $schema: string }; + * if (parsed.$schema !== "sh.diffuse.output.track") throw new Error("expected track NSID"); + * ``` + * + * @example Encodes a collection as JSON bytes + * ```ts + * import { encodeJsonCollection } from "~/common/lens.js"; + * + * const out = encodeJsonCollection([{ id: "a" }], "tracks", true); + * if (!(out instanceof Uint8Array)) throw new Error("expected bytes"); + * ``` + */ +export function encodeJsonCollection(data, name, asBytes = false) { + const json = JSON.stringify(encodeCollection(data, name)); + return asBytes ? new TextEncoder().encode(json) : json; +} + +/** + * Decode a JSON-encoded collection (string, bytes, or an already-parsed object — + * an envelope or legacy array), migrating stale payloads. `undefined`/`null` + * yields an empty collection. + * + * @template {unknown[]} T + * @param {Uint8Array | string | object | null | undefined} raw + * @param {import("./self-describing.js").CollectionName} name + * @returns {T} + * + * @example Round-trips through encodeJsonCollection + * ```ts + * import { encodeJsonCollection, decodeJsonCollection } from "~/common/lens.js"; + * + * const bytes = encodeJsonCollection([{ id: "a" }], "tracks", true) as Uint8Array; + * const out = decodeJsonCollection(bytes, "tracks") as Array<{ id: string }>; + * if (out[0].id !== "a") throw new Error("expected record"); + * ``` + * + * @example Empty for undefined input + * ```ts + * import { decodeJsonCollection } from "~/common/lens.js"; + * + * if (decodeJsonCollection(undefined, "tracks").length !== 0) throw new Error("expected empty"); + * ``` + */ +export function decodeJsonCollection(raw, name) { + try { + let parsed; + if (raw instanceof Uint8Array) { + parsed = JSON.parse(new TextDecoder().decode(raw)); + } else if (raw === undefined || raw === null) { + return /** @type {T} */ (/** @type {unknown} */ ([])); + } else if (typeof raw === "string") { + parsed = JSON.parse(raw); + } else { + // Already-parsed value (e.g. a stored envelope object). + parsed = raw; + } + return normalizeCollection(decodeCollection(parsed, name)); + } catch (err) { + console.error(err); + return /** @type {T} */ (/** @type {unknown} */ ([])); + } +} + +/** + * Guarantee a collection is always returned as an array. + * + * @template {unknown[]} T + * @param {T | null | undefined | unknown} value + * @returns {T} + */ +function normalizeCollection(value) { + if (Array.isArray(value)) return /** @type {T} */ (value); + return /** @type {T} */ (/** @type {unknown} */ ([])); +} + +/** + * Write back an edited record into the collection's current shape, losslessly. + * + * This is the panproto write-back path and the ONLY place the migration flow + * loads `@panproto/core` (its WASM). It is never called on a plain read — it runs + * only when an app edits a record in an older shape and needs to store it back in + * the current shape, using the lens + complement recorded in `$schemaHistory`. + * + * When a complement is available (captured when the data was migrated forward), it + * is used to reconstruct the discarded fields so the old app's edit is preserved + * losslessly; otherwise the forward projection is pure-JS and panproto is not + * loaded. + * + * @param {unknown} editedRecord - the record as the (older) app edited it + * @param {{ lens: LensDocument; toLexicon?: object; complement?: Uint8Array | string | null }} opts + * @returns {Promise} the record in the current shape, ready to store + * + * @example Runs the write-back for an empty lens (loads panproto lazily) + * ```ts + * import { writeBack } from "~/common/lens.js"; + * + * const out = await writeBack({ $type: "sh.diffuse.output.facet", id: "a", name: "x", favourite: true }, { + * lens: { id: "l", source: "sh.diffuse.output.facet", target: "sh.diffuse.output.facet2", steps: [] }, + * toLexicon: { + * lexicon: 1, id: "sh.diffuse.output.facet2", + * defs: { main: { type: "record", record: { type: "object", properties: { + * $type: { type: "string" }, id: { type: "string" }, name: { type: "string" }, + * } } } }, + * }, + * }); + * if (typeof out !== "object") throw new Error("expected a record back"); + * ``` + */ +export async function writeBack(editedRecord, { lens, toLexicon, complement }) { + // Only load panproto (WASM) when a complement is actually used for lossless + // write-back and we have the target lexicon to instantiate it; otherwise the + // common additive/rename path stays pure-JS via `project`. + if (!complement || !toLexicon) { + return project([editedRecord], lens)[0]; + } + + const schema = await parseLexicon(toLexicon); + const comp = typeof complement === "string" + ? base64ToBytes(complement) + : /** @type {Uint8Array} */ (complement); + return panprotoPut(lens, schema, editedRecord, comp); +} + +/** + * Write a collection back into its stored (possibly newer) shape on save. + * + * Called by encoders before wrapping a save into the envelope. If the stored + * envelope's `$schemaHistory` ends at a non-null `complement` for the transition + * to the current NSID, each record is written back (loading panproto lazily to + * preserve discarded fields); otherwise the records pass through unchanged and no + * WASM is loaded. + * + * @template {unknown[]} T + * @param {T} items + * @param {import("./self-describing.js").CollectionName} name + * @param {import("./self-describing.js").SelfDescribing | null} storedEnvelope + * @param {object} [toLexicon] + * @returns {Promise} + * + * @example Passes records through when the stored envelope already matches (no panproto) + * ```ts + * import { writeBackCollection, encodeCollection } from "~/common/lens.js"; + * + * const items = [{ id: "a" }]; + * const envelope = encodeCollection(items, "tracks"); + * const out = await writeBackCollection(items, "tracks", envelope); + * if (out[0].id !== "a") throw new Error("expected unchanged records"); + * ``` + * + * @example Passes records through when no lexicon is available (no panproto) + * ```ts + * import { writeBackCollection } from "~/common/lens.js"; + * import { wrap } from "~/common/self-describing.js"; + * + * const items = [{ id: "a" }]; + * const envelope = wrap([{ id: "a" }], { schema: "sh.diffuse.output.track" }) as { + * $schema: string; $schemaHistory: any[]; data: unknown[]; + * }; + * // stored under an older NSID with a complement-carrying history, but no lexicon: + * const out = await writeBackCollection(items, "tracks", { + * ...envelope, $schema: "sh.diffuse.output.trackOld", + * $schemaHistory: [{ from: "sh.diffuse.output.trackOld", to: "sh.diffuse.output.track", lens: { id: "l", source: "s", target: "t", steps: [] }, complement: new Uint8Array(1) }], + * }) as Array<{ id: string }>; + * if (out[0].id !== "a") throw new Error("expected unchanged records without a lexicon"); + * ``` + */ +export async function writeBackCollection(items, name, storedEnvelope, toLexicon) { + if (!storedEnvelope || storedEnvelope.$schema === collectionSchema(name)) { + // Nothing stale to write back — the stored envelope is already this build's + // shape, or there is no history. Pass through (no panproto load). + return items; + } + + const entry = storedEnvelope.$schemaHistory[storedEnvelope.$schemaHistory.length - 1]; + const lens = entry?.lens; + const complement = entry?.complement; + if (!lens || !complement || !toLexicon) { + // No complement recorded, or no lexicon to instantiate against — nothing + // panproto could losslessly write back; keep the records as-is. + return items; + } + + const written = await Promise.all( + /** @type {unknown[]} */ (items).map((item) => + writeBack(item, { lens, toLexicon, complement }), + ), + ); + return /** @type {T} */ (written); +} + +/** + * Reconstruct the stored self-describing envelope of a collection from its raw + * encoded value (a JSON string/bytes, or a stored envelope object). Returns + * `null` when the value is absent/legacy (no envelope). + * + * @template {unknown[]} T + * @param {Uint8Array | string | unknown} raw - the raw stored payload for the collection + * @param {import("./self-describing.js").CollectionName} name + * @returns {import("./self-describing.js").SelfDescribing | null} + * + * @example Reads the envelope out of a stored JSON string + * ```ts + * import { readStoredEnvelope, encodeJsonCollection } from "~/common/lens.js"; + * + * const stored = encodeJsonCollection([{ id: "a" }], "tracks", true); + * const env = readStoredEnvelope(stored, "tracks"); + * if (!env || env.$schema !== "sh.diffuse.output.track") throw new Error("expected stored envelope"); + * ``` + * + * @example Returns null for absent/legacy (bare array) data + * ```ts + * import { readStoredEnvelope } from "~/common/lens.js"; + * + * if (readStoredEnvelope(null, "tracks") !== null) throw new Error("expected null for absent"); + * if (readStoredEnvelope([{ id: "a" }], "tracks") !== null) throw new Error("expected null for legacy array"); + * ``` + */ +export function readStoredEnvelope(raw, name) { + if (raw === null || raw === undefined) return null; + let parsed = raw; + if (raw instanceof Uint8Array) { + parsed = JSON.parse(new TextDecoder().decode(raw)); + } else if (typeof raw === "string") { + parsed = JSON.parse(raw); + } + const { envelope } = unwrap(parsed, { $schema: collectionSchema(name) }); + return /** @type {import("./self-describing.js").SelfDescribing | null} */ ( + /** @type {any} */ (envelope) + ); +} + +/** + * The write path for a JSON-encoded collection: write back any stale records to + * the stored shape (guarded no-op), then encode. Used by the JSON encoders' save + * so the write-back path is wired on save without loading panproto unless a + * cross-NSID complement + lexicon make it necessary. + * + * @param {unknown[]} items + * @param {import("./self-describing.js").CollectionName} name + * @param {Uint8Array | string | null | undefined} stored - the raw stored payload + * @param {boolean} [asBytes] + * @returns {Promise} + * + * @example Encodes a collection through the wired save path (guarded no-op) + * ```ts + * import { saveJsonCollection, decodeJsonCollection } from "~/common/lens.js"; + * + * const out = await saveJsonCollection([{ id: "a" }], "tracks", null) as string; + * const back = decodeJsonCollection(out, "tracks") as Array<{ id: string }>; + * if (back[0].id !== "a") throw new Error("expected encoded records"); + * ``` + * + * @example Produces bytes when asBytes is set + * ```ts + * import { saveJsonCollection } from "~/common/lens.js"; + * + * const out = await saveJsonCollection([{ id: "a" }], "tracks", null, true); + * if (!(out instanceof Uint8Array)) throw new Error("expected bytes"); + * ``` + */ +export async function saveJsonCollection(items, name, stored, asBytes = false) { + const storedEnvelope = readStoredEnvelope(stored, name); + const written = await writeBackCollection(items, name, storedEnvelope, undefined); + return encodeJsonCollection(written, name, asBytes); +} + +/** + * @param {string} base64 + * @returns {Uint8Array} + */ +function base64ToBytes(base64) { + const bin = atob(base64); + const bytes = new Uint8Array(bin.length); + for (let i = 0; i < bin.length; i++) bytes[i] = bin.charCodeAt(i); + return bytes; +} + +/** + * @template {unknown[]} T + * @template {unknown} L + * @param {import("./self-describing.js").SelfDescribing} envelope + */ +function envelopeToArray(envelope) { + return /** @type {unknown[]} */ (envelope.data); +} + +/** + * A sensible default for an `add_field` kind. + * + * @param {string} kind + * @returns {unknown} + */ +function defaultValue(kind) { + switch (kind) { + case "boolean": + return false; + case "integer": + case "float": + case "number": + return 0; + case "array": + return []; + case "object": + return {}; + case "null": + return null; + default: + return ""; + } +} \ No newline at end of file diff --git a/src/common/panproto.js b/src/common/panproto.js new file mode 100644 index 00000000..50270f17 --- /dev/null +++ b/src/common/panproto.js @@ -0,0 +1,123 @@ +/** + * panproto (`@panproto/core`) connection point. + * + * The pure-JS migration engine (`~/common/lens.js`) projects records between + * schema NSIDs by interpreting an authored lens document's `steps` — sufficient + * for the additive/rename migrations diffuse ships. panproto's WASM engine is the + * *complement* path: it produces and consumes the opaque `complement` so that + * fields a migration discarded can be written back losslessly (see design doc + * open items). + * + * This module loads `@panproto/core` lazily (its WASM is only needed when the + * complement path actually runs), and exposes helpers for the flat-record JSON + * call shape confirmed in the design spike: + * + * - include the record's `$type` and root at the record's object/body vertex; + * - `chain.instantiate(schema)` needs only the *target* schema; + * - `getJson` produces `{ view, complement }`; `putJson(view, complement)` + * reconstructs the source losslessly. + * + * @import {LensDocument} from "~/common/lens-registry.js" + */ + +let promise = null; + +/** + * Lazily obtain a panproto instance. + * + * @returns {Promise} + */ +async function panproto() { + promise ??= import("@panproto/core").then( + /** @param {any} m */ + async (m) => m.Panproto.init(), + ); + return promise; +} + +/** + * Compile an authored lens document and instantiate it against the target schema. + * + * @param {LensDocument} doc + * @param {unknown} schema - a `BuiltSchema` from `parseLexicon` + * @returns {Promise} an instantiated `LensHandle` + * + * @example Compiles a lens and round-trips a record losslessly + * ```ts + * import { compileLens, parseLexicon, get } from "~/common/panproto.js"; + * + * const NEW_LEXICON = { + * lexicon: 1, id: "sh.diffuse.output.facet2", + * defs: { main: { type: "record", record: { type: "object", properties: { + * $type: { type: "string" }, id: { type: "string" }, name: { type: "string" }, + * starred: { type: "boolean" }, description: { type: "string" }, + * } } } }, + * }; + * const schema = await parseLexicon(NEW_LEXICON); + * const doc = { + * id: "f-to-f2", source: "sh.diffuse.output.facet", target: "sh.diffuse.output.facet2", + * steps: [{ rename_field: { old: "favourite", new: "starred" } }], + * }; + * await compileLens(doc, schema); + * const record = { $type: "sh.diffuse.output.facet", id: "a", name: "x", favourite: true }; + * const { view, complement } = await get(doc, schema, record); + * if (view === undefined) throw new Error("expected a projected view"); + * if (!(complement instanceof Uint8Array) || complement.length === 0) throw new Error("expected a complement"); + * ``` + */ +export async function compileLens(doc, schema) { + const pp = await panproto(); + const chain = pp.compileLensDocument(doc, bodyVertex(doc), "json"); + return chain.instantiate(schema); +} + +/** + * Parse an atproto lexicon (a diffuse `lexicons/output/*.json` object) into a + * `BuiltSchema` usable for instantiating a lens. + * + * @param {object | string} lexicon + * @returns {Promise} a `BuiltSchema` + */ +export async function parseLexicon(lexicon) { + const pp = await panproto(); + return pp.parseLexicon(lexicon); +} + +/** + * Forward projection: extract the view + complement from a source record. + * + * @param {LensDocument} doc + * @param {unknown} schema + * @param {unknown} record + * @returns {Promise<{ view: unknown; complement: Uint8Array }>} + */ +export async function get(doc, schema, record) { + const lens = await compileLens(doc, schema); + return lens.getJson(record, bodyVertex(doc)); +} + +/** + * Backward put: reconstruct a source record from an edited view + complement, + * losslessly (GetPut). + * + * @param {LensDocument} doc + * @param {unknown} schema + * @param {unknown} view + * @param {Uint8Array} complement + * @returns {Promise} + */ +export async function put(doc, schema, view, complement) { + const lens = await compileLens(doc, schema); + return lens.putJson(view, complement, bodyVertex(doc)); +} + +/** + * The record's object/body vertex of the target schema the lens projects to + * (used as the root for `getJson`/`putJson`). Field steps attach to this vertex. + * + * @param {LensDocument} doc + * @returns {string} + */ +function bodyVertex(doc) { + return `${doc.target ?? doc.source}:body`; +} \ No newline at end of file diff --git a/src/common/self-describing.js b/src/common/self-describing.js new file mode 100644 index 00000000..1102a7bc --- /dev/null +++ b/src/common/self-describing.js @@ -0,0 +1,211 @@ +/** + * Self-describing envelopes for saved output data. + * + * Saved payloads carry their own schema so a reader can interpret and migrate + * them without external knowledge. The envelope wraps a collection of records + * (`Facet[]`, `Track[]`, …) with the atproto lexicon NSID that produced them plus + * an ordered history of schema transitions (`$schemaHistory`). Each history entry + * carries a portable lens document (authored, schema-independent `steps`) and, + * when a write-back into a previous shape is needed, the complement produced by + * the forward projection. + * + * In keeping with the atproto lexicon model (see "Lexicon Evolution" in the + * atproto spec), the NSID (`$schema`) is the schema's identity. Compatible schema + * evolution keeps the same NSID — new optional fields simply appear, and old data + * remains valid. A breaking change (renaming or removing a field, changing a type) + * is expressed as a NEW NSID, with an authored lens from the old NSID to the new + * one. There is no integer "schema version"; the migration chain is the history. + * + * The lens documents embedded in `$schemaHistory` are the correctness guarantee + * that older apps can read newer data. A given app build may also resolve a + * segment's lens from its own bundled registry (cheaper, for well-known + * transitions); resolution order is: bundled registry first, then the embedded + * document. + * + * @import {Facet, PlaylistItem, Setting, Track} from "~/definitions/types.d.ts" + */ + +/** + * A schema transition recorded in `$schemaHistory`, from one lexicon NSID to + * another. + * + * The `lens` is an authored, schema-independent document (the same shape + * `@panproto/core`'s `compileLensDocument` accepts: an object with `steps`). It is + * stored here so any reader — including one older than the lens' authoring build — + * can recompile it and traverse the transition. + * + * @template {unknown} L + * @typedef {{ + * from: string; + * to: string; + * lens: L | null; + * complement?: Uint8Array | string | null; + * }} HistoryEntry + */ + +/** + * A self-describing envelope around a collection of records. + * + * @template {unknown[]} T + * @template {unknown} L + * @typedef {{ + * $schema: string; + * $schemaHistory: HistoryEntry[]; + * data: T; + * }} SelfDescribing + */ + +/** + * Wrap a collection into a self-describing envelope. + * + * @template {unknown[]} T + * @template {unknown} L + * @param {T} data + * @param {{ schema: string; history?: HistoryEntry[] }} opts + * @returns {SelfDescribing} + * + * @example Wraps records with the lexicon that produced them + * ```js + * import { wrap } from "~/common/self-describing.js"; + * + * const envelope = wrap([{ id: "a" }], { schema: "sh.diffuse.output.facet" }); + * if (envelope.$schema !== "sh.diffuse.output.facet") throw new Error("expected schema"); + * if (envelope.data[0].id !== "a") throw new Error("expected data"); + * ``` + * + * @example Stores the supplied history entries + * ```js + * import { wrap } from "~/common/self-describing.js"; + * + * const envelope = wrap([], { + * schema: "sh.diffuse.output.track2", + * history: [{ + * from: "sh.diffuse.output.track", + * to: "sh.diffuse.output.track2", + * lens: { steps: [] }, + * complement: null, + * }], + * }); + * if (envelope.$schemaHistory.length !== 1) throw new Error("expected 1 history entry"); + * ``` + */ +export function wrap(data, { schema, history = [] }) { + return { + $schema: schema, + $schemaHistory: history, + data, + }; +} + +/** + * Read the collection out of a stored value, tolerating legacy payloads. + * + * Legacy payloads — a bare array, or an object that is not an envelope — are + * treated as `data` written against the supplied schema (a `v0` with no recorded + * migration). This keeps existing users unbroken when the envelope is introduced. + * + * @template {unknown[]} T + * @template {unknown} L + * @param {unknown} value + * @param {{ $schema: string }} opts + * @returns {{ data: T; envelope: SelfDescribing | null }} + * + * @example Reads a self-describing envelope + * ```ts + * import { wrap, unwrap } from "~/common/self-describing.js"; + * + * const wrapped = wrap([{ id: "a" }], { schema: "sh.diffuse.output.facet" }); + * const out = unwrap(wrapped, { + * $schema: "sh.diffuse.output.facet", + * }) as { data: Array<{ id: string }>; envelope: { $schema: string } | null }; + * if (out.data[0].id !== "a") throw new Error("expected data"); + * if (out.envelope === null) throw new Error("expected envelope"); + * ``` + * + * @example Tolerates a legacy bare array + * ```ts + * import { unwrap } from "~/common/self-describing.js"; + * + * const out = unwrap([{ id: "a" }], { + * $schema: "sh.diffuse.output.facet", + * }) as { data: Array<{ id: string }>; envelope: unknown }; + * if (out.data[0].id !== "a") throw new Error("expected data"); + * if (out.envelope !== null) throw new Error("expected no envelope for legacy data"); + * ``` + */ +export function unwrap(value, { $schema }) { + void $schema; + if (isSelfDescribing(value)) { + const envelope = /** @type {SelfDescribing} */ (value); + return { data: /** @type {T} */ (envelope.data), envelope }; + } + return { data: /** @type {T} */ (value), envelope: null }; +} + +/** + * Whether a stored value is a self-describing envelope. + * + * @template {unknown} T + * @param {unknown} value + * @returns {value is SelfDescribing} + * + * @example Identifies enveloped data + * ```js + * import { isSelfDescribing } from "~/common/self-describing.js"; + * + * if (isSelfDescribing({ $schema: "s", $schemaHistory: [], data: [] })) { + * // is a self-describing envelope + * } else { + * throw new Error("expected to be self-describing"); + * } + * ``` + * + * @example Rejects a bare array + * ```js + * import { isSelfDescribing } from "~/common/self-describing.js"; + * + * if (isSelfDescribing([1, 2, 3])) throw new Error("bare array is not an envelope"); + * ``` + */ +export function isSelfDescribing(value) { + return Boolean( + value && + typeof value === "object" && + !Array.isArray(value) && + typeof /** @type {any} */ (value).$schema === "string" && + Array.isArray(/** @type {any} */ (value).$schemaHistory) && + Array.isArray(/** @type {any} */ (value).data), + ); +} + +/** + * The name of one of the collections an output persists. + * + * @typedef {"facets" | "playlistItems" | "settings" | "tracks"} CollectionName + */ + +/** @type {Record} */ +const COLLECTION_SCHEMAS = { + facets: "sh.diffuse.output.facet", + playlistItems: "sh.diffuse.output.playlistItem", + settings: "sh.diffuse.output.setting", + tracks: "sh.diffuse.output.track", +}; + +/** + * The current lexicon NSID for a collection. + * + * @param {CollectionName} name + * @returns {string} + * + * @example Returns the schema for tracks + * ```js + * import { collectionSchema } from "~/common/self-describing.js"; + * + * const s = collectionSchema("tracks"); + * if (s !== "sh.diffuse.output.track") throw new Error("expected track schema"); + * ``` + */ +export function collectionSchema(name) { + return COLLECTION_SCHEMAS[name]; +} \ No newline at end of file diff --git a/src/components/output/common.js b/src/components/output/common.js index 12bd414e..11fe19d6 100644 --- a/src/components/output/common.js +++ b/src/components/output/common.js @@ -110,11 +110,20 @@ export function outputManager( { compare: strictEquality }, ); + // Revision counters so a racing load never overwrites a value that was just + // saved (or a newer load was started) while it was awaiting `get()`. + let cRev = 0; + let plRev = 0; + let sRev = 0; + let tRev = 0; + async function loadFacets() { if (init && (await init()) === false) return; + const rev = cRev; cs.value = "loading"; try { - c.value = await facets.get(); + const data = await facets.get(); + if (cRev === rev) c.value = data; cs.value = "loaded"; } catch (err) { console.error("Failed to load facets:", err); @@ -124,9 +133,11 @@ export function outputManager( async function loadPlaylistItems() { if (init && (await init()) === false) return; + const rev = plRev; pls.value = "loading"; try { - pl.value = await playlistItems.get(); + const data = await playlistItems.get(); + if (plRev === rev) pl.value = data; pls.value = "loaded"; } catch (err) { console.error("Failed to load playlist items:", err); @@ -136,9 +147,11 @@ export function outputManager( async function loadSettings() { if (init && (await init()) === false) return; + const rev = sRev; ss.value = "loading"; try { - s.value = await settings.get(); + const data = await settings.get(); + if (sRev === rev) s.value = data; ss.value = "loaded"; } catch (err) { console.error("Failed to load settings:", err); @@ -148,9 +161,11 @@ export function outputManager( async function loadTracks() { if (init && (await init()) === false) return; + const rev = tRev; ts.value = "loading"; try { - t.value = await tracks.get(); + const data = await tracks.get(); + if (tRev === rev) t.value = data; ts.value = "loaded"; } catch (err) { console.error("Failed to load tracks:", err); @@ -171,6 +186,7 @@ export function outputManager( reload: loadFacets, save: async (newFacets) => { batch(() => { + cRev++; if (untracked(() => cs.value === "sleeping")) cs.value = "loaded"; c.value = newFacets; }); @@ -189,6 +205,7 @@ export function outputManager( reload: loadPlaylistItems, save: async (newPlaylistItems) => { batch(() => { + plRev++; if (untracked(() => pls.value === "sleeping")) pls.value = "loaded"; pl.value = newPlaylistItems; }); @@ -207,6 +224,7 @@ export function outputManager( reload: loadSettings, save: async (newSettings) => { batch(() => { + sRev++; if (untracked(() => ss.value === "sleeping")) ss.value = "loaded"; s.value = newSettings; }); @@ -225,6 +243,7 @@ export function outputManager( reload: loadTracks, save: async (newTracks) => { batch(() => { + tRev++; if (untracked(() => ts.value === "sleeping")) ts.value = "loaded"; t.value = newTracks; }); diff --git a/src/components/output/polymorphic/indexed-db/element.js b/src/components/output/polymorphic/indexed-db/element.js index 071d7b62..21558884 100644 --- a/src/components/output/polymorphic/indexed-db/element.js +++ b/src/components/output/polymorphic/indexed-db/element.js @@ -2,11 +2,14 @@ import * as IDB from "idb-keyval"; import { IDB_PREFIX } from "./constants.js"; import { BroadcastedOutputElement, outputManager } from "../../common.js"; +import { isSelfDescribing } from "~/common/self-describing.js"; +import { decodeCollection, encodeCollection } from "~/common/lens.js"; import { defineElement } from "~/common/element.js"; /** * @import {OutputElement, OutputManager, OutputWorkerActions} from "@specs/components/output/types.d.ts" * @import {SupportedDataTypes} from "@specs/components/output/polymorphic/indexed-db/types.d.ts" + * @import {CollectionName} from "~/common/self-describing.js" */ //////////////////////////////////////////// @@ -69,10 +72,27 @@ class IndexedDBOutput extends BroadcastedOutputElement { // GET & PUT /** @param {string} name */ - #get = (name) => IDB.get(`${IDB_PREFIX}/${this.#cat(name)}`); + #get = async (name) => { + const stored = await IDB.get(`${IDB_PREFIX}/${this.#cat(name)}`); + if (stored === undefined) return undefined; + return decodeCollection( + stored, + /** @type {CollectionName} */ (name), + ) ?? undefined; + }; /** @param {string} name; @param {any} data */ - #put = (name, data) => IDB.set(`${IDB_PREFIX}/${this.#cat(name)}`, data); + #put = async (name, data) => { + // Bytes/strings produced by a transformer above (bytes/json or string/json + // envelope bytes/JSON, or automerge/dasl binary) are already self-describing + // at that layer, so they pass through as-is. Plain collections (arrays) are + // wrapped so a standalone IndexedDB output is itself self-describing. + const stored = data instanceof Uint8Array || typeof data === "string" || + isSelfDescribing(data) + ? data + : encodeCollection(data, /** @type {CollectionName} */ (name)); + await IDB.set(`${IDB_PREFIX}/${this.#cat(name)}`, stored); + }; // 🛠️ diff --git a/src/components/output/raw/atproto-space/element.js b/src/components/output/raw/atproto-space/element.js index 55d92f70..cd6bc4dc 100644 --- a/src/components/output/raw/atproto-space/element.js +++ b/src/components/output/raw/atproto-space/element.js @@ -2,6 +2,7 @@ import { Agent } from "@atproto/api"; import { decode, encode } from "@atcute/cbor"; import { xxh32r } from "xxh32/dist/raw.js"; +import { decodeCollection, encodeCollection } from "~/common/lens.js"; import { computed, signal } from "~/common/signal.js"; import { BroadcastedOutputElement, outputManager } from "../../common.js"; import { defineElement } from "~/common/element.js"; @@ -68,20 +69,25 @@ class ATProtoSpaceOutput extends BroadcastedOutputElement { }, facets: this.#recordCollection("sh.diffuse.output.facet"), playlistItems: this.#blobCollection( + "playlistItems", "sh.diffuse.output.playlistItemBundle", { groupBy: "playlist" }, ), settings: this.#recordCollection("sh.diffuse.output.setting"), - tracks: this.#blobCollection("sh.diffuse.output.trackBundle", { - groupBy: "scheme", - keyOf: (item) => { - const uri = String( - /** @type {Record} */ (item)["uri"] ?? "", - ); - const colon = uri.indexOf(":"); - return colon > 0 ? uri.substring(0, colon) : undefined; + tracks: this.#blobCollection( + "tracks", + "sh.diffuse.output.trackBundle", + { + groupBy: "scheme", + keyOf: (item) => { + const uri = String( + /** @type {Record} */ (item)["uri"] ?? "", + ); + const colon = uri.indexOf(":"); + return colon > 0 ? uri.substring(0, colon) : undefined; + }, }, - }), + ), }); this.facets = this.#manager.facets; @@ -301,7 +307,17 @@ class ATProtoSpaceOutput extends BroadcastedOutputElement { * @param {string} nsid * @param {{ groupBy: string, keyOf?: (item: unknown) => string | undefined }} options */ - #blobCollection(nsid, { groupBy, keyOf } = /** @type {any} */ ({})) { + /** + * Returns `{ empty, get, put }` for a collection stored as CBOR blobs. + * + * The blob content is wrapped in a self-describing envelope so the stored + * array is interpretable on its own. + * + * @param {import("~/common/self-describing.js").CollectionName} name + * @param {string} nsid + * @param {{ groupBy: string, keyOf?: (item: unknown) => string | undefined }} options + */ + #blobCollection(name, nsid, { groupBy, keyOf } = /** @type {any} */ ({})) { /** @type {Map} */ let lastHashes = new Map(); /** @type {Map} */ @@ -325,7 +341,7 @@ class ATProtoSpaceOutput extends BroadcastedOutputElement { if (typeof key !== "string") continue; const bytes = await this.#fetchBlob(bundle.data.ref.$link); - const groupItems = /** @type {unknown[]} */ (decode(bytes)); + const groupItems = /** @type {unknown[]} */ (decodeBlob(bytes, name)); if (!Array.isArray(groupItems)) continue; items.push(...groupItems); @@ -363,7 +379,7 @@ class ATProtoSpaceOutput extends BroadcastedOutputElement { const upserts = []; for (const [key, groupItems] of groups) { - const bytes = encode(groupItems); + const bytes = encodeBlob(groupItems, name); const hash = xxh32r(bytes).toString(16); if (lastHashes.get(key) === hash && lastBlobs.has(key)) continue; @@ -537,6 +553,30 @@ class ATProtoSpaceOutput extends BroadcastedOutputElement { export default ATProtoSpaceOutput; +/** + * Encode a group of records as a self-describing CBOR blob. + * + * @param {unknown[]} items + * @param {import("~/common/self-describing.js").CollectionName} name + * @returns {Uint8Array} + */ +function encodeBlob(items, name) { + return encode(encodeCollection(items, name)); +} + +/** + * Decode a self-describing CBOR blob into records, migrating a stale payload. + * + * Legacy bare-array blobs (not envelopes) are returned as-is. + * + * @param {Uint8Array} bytes + * @param {import("~/common/self-describing.js").CollectionName} name + * @returns {unknown[]} + */ +function decodeBlob(bytes, name) { + return /** @type {unknown[]} */ (decodeCollection(decode(bytes), name) ?? []); +} + //////////////////////////////////////////// // REGISTER //////////////////////////////////////////// diff --git a/src/components/transformer/output/bytes/automerge/element.js b/src/components/transformer/output/bytes/automerge/element.js index 35f531c6..93701b35 100644 --- a/src/components/transformer/output/bytes/automerge/element.js +++ b/src/components/transformer/output/bytes/automerge/element.js @@ -10,6 +10,7 @@ import { removeUndefinedValuesFromRecord, } from "~/common/utils.js"; import { OutputTransformer } from "../../base.js"; +import { collectionSchema } from "~/common/self-describing.js"; import { defineElement } from "~/common/element.js"; import { INITIAL_FACETS_DOCUMENT, @@ -124,24 +125,31 @@ class AutomergeBytesOutputTransformer extends OutputTransformer { { stripUndefined: true, }, + "facets", ); this.playlistItems = automergeEntry( computed(() => local()?.playlistItems), remote.playlistItems, computed(() => playlistItems().doc), + undefined, + "playlistItems", ); this.settings = automergeEntry( computed(() => local()?.settings), remote.settings, computed(() => settings().doc), + undefined, + "settings", ); this.tracks = automergeEntry( computed(() => local()?.tracks), remote.tracks, computed(() => tracks().doc), + undefined, + "tracks", ); this.ready = () => true; @@ -256,11 +264,12 @@ export function loadDocument(value) { * @template {Record} T * @param {SignalReader<{ collection: SignalReader<{ state: "loading" } | { state: "loaded"; data: Uint8Array | undefined } | { state: "error" }>, reload: () => Promise, save: (bytes: Uint8Array) => Promise } | undefined>} local * @param {{ collection: SignalReader<{ state: "loading" } | { state: "loaded"; data: Uint8Array | undefined } | { state: "error" }>, reload: () => Promise, save: (bytes: Uint8Array) => Promise }} remote - * @param {SignalReader>} document + * @param {SignalReader} document * @param {{ stripUndefined?: boolean }} [opts] + * @param {import("~/common/self-describing.js").CollectionName} [name] * @returns {{ collection: SignalReader<{ state: "loading" } | { state: "loaded"; data: T[] } | { state: "error" }>, reload: () => Promise, save: (items: T[]) => Promise }} */ -export function automergeEntry(local, remote, document, opts) { +export function automergeEntry(local, remote, document, opts, name) { return { collection: computed(() => { const col = local()?.collection(); @@ -271,6 +280,7 @@ export function automergeEntry(local, remote, document, opts) { }), reload: remote.reload, save: async (/** @type {T[]} */ newItems) => { + const schema = name ? collectionSchema(name) : null; const doc = Automerge.change(document(), (d) => { d.collection = newItems.map((item) => { const cloned = recursivelyCloneRecords(item); @@ -278,6 +288,9 @@ export function automergeEntry(local, remote, document, opts) { ? removeUndefinedValuesFromRecord(cloned) : cloned; }); + if (schema) { + d.$schema = schema; + } }); const bytes = Automerge.save(doc); diff --git a/src/components/transformer/output/bytes/dasl-sync/element.js b/src/components/transformer/output/bytes/dasl-sync/element.js index bd1ba5f3..6ff7e18c 100644 --- a/src/components/transformer/output/bytes/dasl-sync/element.js +++ b/src/components/transformer/output/bytes/dasl-sync/element.js @@ -6,6 +6,7 @@ import "~/components/output/polymorphic/indexed-db/element.js"; import * as CID from "~/common/cid.js"; import { diff, strictEquality } from "~/common/compare.js"; +import { collectionSchema } from "~/common/self-describing.js"; import { computed, signal } from "~/common/signal.js"; import { compareTimestamps } from "~/common/temporal.js"; import { OutputTransformer } from "../../base.js"; @@ -230,6 +231,7 @@ class DaslBytesSyncOutputTransformer extends OutputTransformer { remote.facets, remote.ready, facets, + "facets", ); this.playlistItems = this.managerProp( @@ -237,6 +239,7 @@ class DaslBytesSyncOutputTransformer extends OutputTransformer { remote.playlistItems, remote.ready, playlistItems, + "playlistItems", ); this.settings = this.managerProp( @@ -244,6 +247,7 @@ class DaslBytesSyncOutputTransformer extends OutputTransformer { remote.settings, remote.ready, settings, + "settings", ); this.tracks = this.managerProp( @@ -251,6 +255,7 @@ class DaslBytesSyncOutputTransformer extends OutputTransformer { remote.tracks, remote.ready, tracks, + "tracks", ); this.ready = () => true; @@ -288,11 +293,12 @@ class DaslBytesSyncOutputTransformer extends OutputTransformer { /** * @template {{ id: string; updatedAt: string }} T - * @param {{ previous: Container, collection: T[] }} _ + * @param {{ previous: Container, collection: T[], name?: import("~/common/self-describing.js").CollectionName }} _ * @returns {Promise>} */ - async updateContainer({ previous, collection }) { + async updateContainer({ previous, collection, name }) { const inventory = previous.inventory; + const schema = name ? collectionSchema(name) : null; const collIds = collection.map(({ id }) => id); @@ -331,11 +337,12 @@ class DaslBytesSyncOutputTransformer extends OutputTransformer { removed: Array.from(allRemoved), }; - return { + const container = { cid: await CID.create(0x71, encode(newInventory)), data: collection, inventory: newInventory, }; + return schema ? { ...container, $schema: schema } : container; } /** @@ -447,6 +454,7 @@ class DaslBytesSyncOutputTransformer extends OutputTransformer { cid: await CID.create(0x71, encode(updatedInventory)), data, inventory: updatedInventory, + $schema: a.$schema, }; } @@ -467,9 +475,10 @@ class DaslBytesSyncOutputTransformer extends OutputTransformer { * @param {{ collection: SignalReader<{ state: "loading" } | { state: "loaded"; data: Uint8Array | undefined } | { state: "error" }>, reload: () => Promise, save: (bytes: Uint8Array) => Promise }} remote * @param {SignalReader} remoteReady * @param {SignalReader>} container + * @param {import("~/common/self-describing.js").CollectionName} name * @returns {{ collection: SignalReader<{ state: "loading" } | { state: "loaded"; data: T[] } | { state: "error" }>, reload: () => Promise, save: (items: T[]) => Promise }} */ - managerProp(local, remote, remoteReady, container) { + managerProp(local, remote, remoteReady, container, name) { return { collection: computed(() => { const c = container(); @@ -485,6 +494,7 @@ class DaslBytesSyncOutputTransformer extends OutputTransformer { const adjustedContainer = await this.updateContainer({ collection: newItems, previous: container(), + name, }); const bytes = this.save(adjustedContainer); diff --git a/src/components/transformer/output/bytes/json/element.js b/src/components/transformer/output/bytes/json/element.js index 62e4f060..a4fff7df 100644 --- a/src/components/transformer/output/bytes/json/element.js +++ b/src/components/transformer/output/bytes/json/element.js @@ -1,5 +1,9 @@ import { computed } from "~/common/signal.js"; import { OutputTransformer } from "../../base.js"; +import { + decodeJsonCollection, + saveJsonCollection, +} from "~/common/lens.js"; import { defineElement } from "~/common/element.js"; /** @@ -10,7 +14,7 @@ import { defineElement } from "~/common/element.js"; /** * @extends {OutputTransformer} */ -class JsonStringOutputTransformer extends OutputTransformer { +class JsonBytesOutputTransformer extends OutputTransformer { constructor() { super(); @@ -24,14 +28,15 @@ class JsonStringOutputTransformer extends OutputTransformer { const col = base.facets.collection(); if (col.state !== "loaded") return col; /** @type {Facet[]} */ - const data = parseArray(col.data); + const data = decodeJsonCollection(col.data, "facets"); return { state: "loaded", data }; }), save: async (newFacets) => { - const json = JSON.stringify(newFacets); - const encoder = new TextEncoder(); - const bytes = encoder.encode(json); - await base.facets.save(bytes); + await base.facets.save( + /** @type {Uint8Array} */ ( + await saveJsonCollection(newFacets, "facets", null, true) + ), + ); }, }, playlistItems: { @@ -40,14 +45,15 @@ class JsonStringOutputTransformer extends OutputTransformer { const col = base.playlistItems.collection(); if (col.state !== "loaded") return col; /** @type {PlaylistItem[]} */ - const data = parseArray(col.data); + const data = decodeJsonCollection(col.data, "playlistItems"); return { state: "loaded", data }; }), save: async (newPlaylistItems) => { - const json = JSON.stringify(newPlaylistItems); - const encoder = new TextEncoder(); - const bytes = encoder.encode(json); - await base.playlistItems.save(bytes); + await base.playlistItems.save( + /** @type {Uint8Array} */ ( + await saveJsonCollection(newPlaylistItems, "playlistItems", null, true) + ), + ); }, }, settings: { @@ -56,14 +62,15 @@ class JsonStringOutputTransformer extends OutputTransformer { const col = base.settings.collection(); if (col.state !== "loaded") return col; /** @type {Setting[]} */ - const data = parseArray(col.data); + const data = decodeJsonCollection(col.data, "settings"); return { state: "loaded", data }; }), save: async (newSettings) => { - const json = JSON.stringify(newSettings); - const encoder = new TextEncoder(); - const bytes = encoder.encode(json); - await base.settings.save(bytes); + await base.settings.save( + /** @type {Uint8Array} */ ( + await saveJsonCollection(newSettings, "settings", null, true) + ), + ); }, }, tracks: { @@ -72,14 +79,15 @@ class JsonStringOutputTransformer extends OutputTransformer { const col = base.tracks.collection(); if (col.state !== "loaded") return col; /** @type {Track[]} */ - const data = parseArray(col.data); + const data = decodeJsonCollection(col.data, "tracks"); return { state: "loaded", data }; }), save: async (newTracks) => { - const json = JSON.stringify(newTracks); - const encoder = new TextEncoder(); - const bytes = encoder.encode(json); - await base.tracks.save(bytes); + await base.tracks.save( + /** @type {Uint8Array} */ ( + await saveJsonCollection(newTracks, "tracks", null, true) + ), + ); }, }, @@ -96,36 +104,13 @@ class JsonStringOutputTransformer extends OutputTransformer { } } -/** - * @param {Uint8Array | string | undefined} data - */ -function parseArray(data) { - let json; - - if (data instanceof Uint8Array) { - const decoder = new TextDecoder(); - json = decoder.decode(data); - } else if (data === undefined) { - return []; - } else { - json = data; - } - - try { - return JSON.parse(json); - } catch (err) { - console.error(err); - return []; - } -} - -export default JsonStringOutputTransformer; +export default JsonBytesOutputTransformer; //////////////////////////////////////////// // REGISTER //////////////////////////////////////////// -export const CLASS = JsonStringOutputTransformer; +export const CLASS = JsonBytesOutputTransformer; export const NAME = "dtob-json"; -defineElement(NAME, CLASS); +defineElement(NAME, CLASS); \ No newline at end of file diff --git a/src/components/transformer/output/string/json/element.js b/src/components/transformer/output/string/json/element.js index 12db0c8c..b21f2f4e 100644 --- a/src/components/transformer/output/string/json/element.js +++ b/src/components/transformer/output/string/json/element.js @@ -1,9 +1,15 @@ import { computed } from "~/common/signal.js"; import { OutputTransformer } from "../../base.js"; +import { + decodeJsonCollection, + saveJsonCollection, +} from "~/common/lens.js"; import { defineElement } from "~/common/element.js"; /** * @import { OutputManagerDeputy } from "@specs/components/output/types.d.ts" + * @import { Facet, PlaylistItem, Setting, Track } from "~/definitions/types.d.ts" + * @import { CollectionName } from "~/common/self-describing.js" */ /** @@ -22,14 +28,16 @@ class JsonStringOutputTransformer extends OutputTransformer { collection: computed(() => { const col = base.facets.collection(); if (col.state !== "loaded") return col; - return { - state: "loaded", - data: typeof col.data === "string" ? parseArray(col.data) : [], - }; + /** @type {Facet[]} */ + const data = decodeJsonCollection(col.data, "facets"); + return { state: "loaded", data }; }), save: async (newFacets) => { - const json = JSON.stringify(newFacets); - await base.facets.save(json); + await base.facets.save( + /** @type {string} */ ( + await saveJsonCollection(newFacets, "facets", null) + ), + ); }, }, playlistItems: { @@ -37,14 +45,16 @@ class JsonStringOutputTransformer extends OutputTransformer { collection: computed(() => { const col = base.playlistItems.collection(); if (col.state !== "loaded") return col; - return { - state: "loaded", - data: typeof col.data === "string" ? parseArray(col.data) : [], - }; + /** @type {PlaylistItem[]} */ + const data = decodeJsonCollection(col.data, "playlistItems"); + return { state: "loaded", data }; }), save: async (newPlaylistItems) => { - const json = JSON.stringify(newPlaylistItems); - await base.playlistItems.save(json); + await base.playlistItems.save( + /** @type {string} */ ( + await saveJsonCollection(newPlaylistItems, "playlistItems", null) + ), + ); }, }, settings: { @@ -52,14 +62,16 @@ class JsonStringOutputTransformer extends OutputTransformer { collection: computed(() => { const col = base.settings.collection(); if (col.state !== "loaded") return col; - return { - state: "loaded", - data: typeof col.data === "string" ? parseArray(col.data) : [], - }; + /** @type {Setting[]} */ + const data = decodeJsonCollection(col.data, "settings"); + return { state: "loaded", data }; }), save: async (newSettings) => { - const json = JSON.stringify(newSettings); - await base.settings.save(json); + await base.settings.save( + /** @type {string} */ ( + await saveJsonCollection(newSettings, "settings", null) + ), + ); }, }, tracks: { @@ -67,14 +79,16 @@ class JsonStringOutputTransformer extends OutputTransformer { collection: computed(() => { const col = base.tracks.collection(); if (col.state !== "loaded") return col; - return { - state: "loaded", - data: typeof col.data === "string" ? parseArray(col.data) : [], - }; + /** @type {Track[]} */ + const data = decodeJsonCollection(col.data, "tracks"); + return { state: "loaded", data }; }), save: async (newTracks) => { - const json = JSON.stringify(newTracks); - await base.tracks.save(json); + await base.tracks.save( + /** @type {string} */ ( + await saveJsonCollection(newTracks, "tracks", null) + ), + ); }, }, @@ -91,18 +105,6 @@ class JsonStringOutputTransformer extends OutputTransformer { } } -/** - * @param {string} json - */ -function parseArray(json) { - try { - return JSON.parse(json); - } catch (err) { - console.error(err); - return []; - } -} - export default JsonStringOutputTransformer; //////////////////////////////////////////// @@ -112,4 +114,4 @@ export default JsonStringOutputTransformer; export const CLASS = JsonStringOutputTransformer; export const NAME = "dtos-json"; -defineElement(NAME, CLASS); +defineElement(NAME, CLASS); \ No newline at end of file diff --git a/tests/common/self-describing-lenses/test.ts b/tests/common/self-describing-lenses/test.ts new file mode 100644 index 00000000..b31ee0d5 --- /dev/null +++ b/tests/common/self-describing-lenses/test.ts @@ -0,0 +1,146 @@ +import { describe, it } from "@std/testing/bdd"; +import { expect } from "@std/expect"; + +import { + wrap, + unwrap, + isSelfDescribing, + collectionSchema, +} from "~/common/self-describing.js"; +import { + project, + migrateEnvelope, + encodeCollection, + encodeJsonCollection, + decodeJsonCollection, + saveJsonCollection, + readStoredEnvelope, + writeBack, +} from "~/common/lens.js"; +import { + register, + resolve, +} from "~/common/lens-registry.js"; + +// Unit tests for the self-describing envelope + migration/write-back machinery. +// These use distinct NSIDs per test so the module-level lens registry does not +// collide across cases. +describe("common/self-describing + lens", () => { + describe("self-describing envelope", () => { + it("wraps data with the collection's NSID", () => { + const env = encodeCollection([{ id: "a" }], "facets"); + expect(env.$schema).toBe("sh.diffuse.output.facet"); + expect(isSelfDescribing(env)).toBe(true); + }); + + it("unwrap tolerates a legacy bare array", () => { + const out = unwrap([{ id: "a" }], { $schema: "sh.diffuse.output.facet" }); + expect(out.envelope).toBeNull(); + const rec = (out.data as Array<{ id: string }>)[0]; + expect(rec.id).toBe("a"); + }); + + it("collectionSchema returns the current NSID", () => { + expect(collectionSchema("tracks")).toBe("sh.diffuse.output.track"); + }); + }); + + describe("migration", () => { + it("migrates a stale envelope to the current NSID", () => { + register({ + id: "f-old-current", + source: "sh.diffuse.output.facetOld", + target: "sh.diffuse.output.facet", + steps: [{ rename_field: { old: "favourite", new: "starred" } }], + }); + const envelope = wrap( + [{ id: "a", favourite: true }], + { schema: "sh.diffuse.output.facetOld" }, + ); + const out = migrateEnvelope(envelope, "facets", resolve); + const rec = out.data[0] as Record; + expect(rec.starred).toBe(true); + expect(rec.favourite).toBeUndefined(); + expect(out.envelope?.$schema).toBe("sh.diffuse.output.facet"); + expect(out.envelope?.$schemaHistory).toHaveLength(1); + }); + + it("does not migrate when the envelope NSID already matches", () => { + const envelope = wrap([{ id: "a" }], { schema: "sh.diffuse.output.facet" }); + const out = migrateEnvelope(envelope, "facets", resolve); + const rec = (out.data as Array<{ id: string }>)[0]; + expect(rec.id).toBe("a"); + expect(out.envelope?.$schemaHistory).toHaveLength(0); + }); + }); + + describe("lens projection", () => { + it("renames a field via a lens document", () => { + const out = project( + [{ $type: "sh.diffuse.output.facet", id: "a", favourite: true }], + { + id: "f", + source: "sh.diffuse.output.facet", + target: "sh.diffuse.output.facet2", + steps: [{ rename_field: { old: "favourite", new: "starred" } }], + }, + ); + const rec = (out as Array>)[0]; + expect(rec.starred).toBe(true); + expect(rec.favourite).toBeUndefined(); + }); + }); + + describe("JSON encode/decode + save wiring", () => { + it("round-trips through encode/decode", () => { + const bytes = encodeJsonCollection([{ id: "a" }], "tracks", true); + const out = decodeJsonCollection(bytes, "tracks") as Array<{ id: string }>; + expect(out[0].id).toBe("a"); + }); + + it("saveJsonCollection wires the save path (guarded no-op)", async () => { + const out = await saveJsonCollection([{ id: "a" }], "tracks", null); + const back = decodeJsonCollection(out, "tracks") as Array<{ id: string }>; + expect(back[0].id).toBe("a"); + }); + + it("decodeJsonCollection accepts an already-parsed envelope object", () => { + const env = encodeCollection([{ id: "a" }], "tracks"); + const out = decodeJsonCollection(env, "tracks") as Array<{ id: string }>; + expect(out[0].id).toBe("a"); + }); + + it("decodeJsonCollection returns an array even for a non-array stored value", () => { + const out = decodeJsonCollection( + { $schema: "sh.diffuse.output.track", $schemaHistory: [], data: { id: "a" } }, + "tracks", + ); + expect(Array.isArray(out)).toBe(true); + }); + + it("readStoredEnvelope returns the stored envelope", () => { + const stored = encodeJsonCollection([{ id: "a" }], "tracks", true); + const env = readStoredEnvelope(stored, "tracks"); + expect(env?.$schema).toBe("sh.diffuse.output.track"); + }); + }); + + describe("writeBack", () => { + it("uses pure-JS projection when no complement is present", async () => { + const out = await writeBack( + { $type: "sh.diffuse.output.facet", id: "a", favourite: true }, + { + lens: { + id: "f", + source: "sh.diffuse.output.facet", + target: "sh.diffuse.output.facet2", + steps: [{ rename_field: { old: "favourite", new: "starred" } }], + }, + }, + ); + const rec = out as Record; + expect(rec.starred).toBe(true); + expect(rec.favourite).toBeUndefined(); + }); + }); +}); \ No newline at end of file diff --git a/tests/components/orchestrator/output/fast-save.test.ts b/tests/components/orchestrator/output/fast-save.test.ts new file mode 100644 index 00000000..5cbe5fb3 --- /dev/null +++ b/tests/components/orchestrator/output/fast-save.test.ts @@ -0,0 +1,40 @@ +import { describe, it } from "@std/testing/bdd"; +import { expect } from "@std/expect"; + +import { testWeb } from "@tests/common/index.ts"; + +// Fast timing (no settle between saves) through the do-output orchestrator. +// Regression target: saving tracks right after facets must not make the +// orchestators' facets read back empty while storage has it. +describe("do-output fast save", () => { + it("saving tracks fast after facets keeps both", async () => { + const result = await testWeb(async () => { + const { CLASS } = await import( + "~/components/orchestrator/output/element.js" + ); + const orch = new CLASS(); + document.body.append(orch); + await customElements.whenDefined("dc-output"); + const output = orch.output; + + await orch.facets.save([ + { $type: "sh.diffuse.output.facet", id: "f1", name: "Keep me" }, + ]); + await orch.tracks.save([ + { $type: "sh.diffuse.output.track", id: "t1", uri: "https://a.com/t1.mp3" }, + ]); + await new Promise((r) => setTimeout(r, 100)); + + const f = orch.facets.collection(); + const tr = orch.tracks.collection(); + return { + f: f.state === "loaded" ? JSON.stringify(f.data) : f.state, + tr: tr.state === "loaded" ? JSON.stringify(tr.data) : tr.state, + }; + }); + + console.log("ORCH_FAST_RESULT", JSON.stringify(result)); + expect(result?.f).toContain("f1"); + expect(result?.tr).toContain("t1"); + }); +}); \ No newline at end of file diff --git a/tests/components/transformer/output/string/json/minimal-fast.test.ts b/tests/components/transformer/output/string/json/minimal-fast.test.ts new file mode 100644 index 00000000..b30d7e76 --- /dev/null +++ b/tests/components/transformer/output/string/json/minimal-fast.test.ts @@ -0,0 +1,43 @@ +import { describe, it } from "@std/testing/bdd"; +import { expect } from "@std/expect"; + +import { testWeb } from "@tests/common/index.ts"; + +// Fast-timing save (no settle) via indexed-db + string/json only (no configurator). +describe("minimal fast save isolation", () => { + it("saving tracks fast after facets keeps both", async () => { + const result = await testWeb(async () => { + const idbMod = await import( + "~/components/output/polymorphic/indexed-db/element.js" + ); + const mod = await import( + "~/components/transformer/output/string/json/element.js" + ); + const output = new idbMod.CLASS(); + output.id = "fast-idb"; + document.body.append(output); + const t = new mod.CLASS(); + t.setAttribute("output-selector", "#fast-idb"); + document.body.append(t); + + await t.facets.save([ + { $type: "sh.diffuse.output.facet", id: "f1", name: "Keep me" }, + ]); + // No wait. + await t.tracks.save([ + { $type: "sh.diffuse.output.track", id: "t1", uri: "https://a.com/t1.mp3" }, + ]); + await new Promise((r) => setTimeout(r, 150)); + + const f = t.facets.collection(); + const tr = t.tracks.collection(); + return { + f: f.state === "loaded" ? JSON.stringify(f.data) : f.state, + tr: tr.state === "loaded" ? JSON.stringify(tr.data) : tr.state, + }; + }); + console.log("MIN_DIAG", JSON.stringify(result)); + expect(result?.f).toContain("f1"); + expect(result?.tr).toContain("t1"); + }); +}); \ No newline at end of file diff --git a/tests/components/transformer/output/string/json/minimal-namespace.test.ts b/tests/components/transformer/output/string/json/minimal-namespace.test.ts new file mode 100644 index 00000000..9e11a973 --- /dev/null +++ b/tests/components/transformer/output/string/json/minimal-namespace.test.ts @@ -0,0 +1,43 @@ +import { describe, it } from "@std/testing/bdd"; +import { expect } from "@std/expect"; + +import { testWeb } from "@tests/common/index.ts"; + +// Fast save with a NAMESPACED indexed-db (as the orchestrator uses namespace="json"). +describe("minimal namespace fast save", () => { + it("saving tracks fast after facets keeps both (namespaced idb)", async () => { + const result = await testWeb(async () => { + const idbMod = await import( + "~/components/output/polymorphic/indexed-db/element.js" + ); + const mod = await import( + "~/components/transformer/output/string/json/element.js" + ); + const output = new idbMod.CLASS(); + output.id = "ns-idb"; + output.setAttribute("namespace", "json"); + document.body.append(output); + const t = new mod.CLASS(); + t.setAttribute("output-selector", "#ns-idb"); + document.body.append(t); + + await t.facets.save([ + { $type: "sh.diffuse.output.facet", id: "f1", name: "Keep me" }, + ]); + await t.tracks.save([ + { $type: "sh.diffuse.output.track", id: "t1", uri: "https://a.com/t1.mp3" }, + ]); + await new Promise((r) => setTimeout(r, 150)); + + const f = t.facets.collection(); + const tr = t.tracks.collection(); + return { + f: f.state === "loaded" ? JSON.stringify(f.data) : f.state, + tr: tr.state === "loaded" ? JSON.stringify(tr.data) : tr.state, + }; + }); + console.log("NS_DIAG", JSON.stringify(result)); + expect(result?.f).toContain("f1"); + expect(result?.tr).toContain("t1"); + }); +}); \ No newline at end of file diff --git a/tests/components/transformer/output/string/json/orchestrator-isolation.test.ts b/tests/components/transformer/output/string/json/orchestrator-isolation.test.ts new file mode 100644 index 00000000..75786ccd --- /dev/null +++ b/tests/components/transformer/output/string/json/orchestrator-isolation.test.ts @@ -0,0 +1,55 @@ +import { describe, it } from "@std/testing/bdd"; +import { expect } from "@std/expect"; + +import { testWeb } from "@tests/common/index.ts"; + +// Regression: saving a JSON-encoded collection through string/json over +// indexed-db must persist WITHOUT reading back the collection synchronously +// (which triggered a load race that clobbered the just-saved value). Saving one +// collection must also not wipe the others. +describe("string/json save isolation over indexed-db", () => { + it("saving tracks persists them and does not wipe facets", async () => { + const result = await testWeb(async () => { + const idbMod = await import( + "~/components/output/polymorphic/indexed-db/element.js" + ); + const mod = await import( + "~/components/transformer/output/string/json/element.js" + ); + + const output = new idbMod.CLASS(); + output.id = "save-isolation-idb"; + document.body.append(output); + + const t = new mod.CLASS(); + t.setAttribute("output-selector", "#save-isolation-idb"); + document.body.append(t); + + // Pre-existing facets. + await t.facets.save([ + { $type: "sh.diffuse.output.facet", id: "f1", name: "Keep me" }, + ]); + + // Import tracks. + await t.tracks.save([ + { $type: "sh.diffuse.output.track", id: "t1", uri: "https://a.com/t1.mp3" }, + ]); + + const facets = t.facets.collection(); + const tracks = t.tracks.collection(); + return { + facets: + facets.state === "loaded" && Array.isArray(facets.data) + ? (facets.data as any[]).map((f: any) => f.id) + : null, + tracks: + tracks.state === "loaded" && Array.isArray(tracks.data) + ? (tracks.data as any[]).map((x: any) => x.id) + : null, + }; + }); + + expect(result?.facets).toEqual(["f1"]); + expect(result?.tracks).toEqual(["t1"]); + }); +}); \ No newline at end of file