diff --git a/docs/plans/2026-01-15-hydration-helpers.md b/docs/plans/2026-01-15-hydration-helpers.md new file mode 100644 index 0000000..f2c7b7b --- /dev/null +++ b/docs/plans/2026-01-15-hydration-helpers.md @@ -0,0 +1,472 @@ +# Hydration Helpers Implementation Plan + +> **For Claude:** REQUIRED SUB-SKILL: Use superpowers:executing-plans to implement this plan task-by-task. + +**Goal:** Add hydration helpers to lex-gql that simplify building query adapters, following the hexagonal architecture port/adapter pattern. + +**Architecture:** lex-gql defines a query port interface that adapters implement. Hydration helpers (`hydrateBlobs`, `hydrateRecord`) are shared utilities for adapter authors that transform raw database rows into the format lex-gql expects. The standard records schema is documented so adapters have a clear contract. + +**Tech Stack:** JavaScript, graphql-js, vitest + +--- + +### Task 1: Add hydrateBlobs helper function + +**Files:** +- Modify: `packages/lex-gql/lex-gql.js` (add after utility functions ~line 212) +- Test: `packages/lex-gql/lex-gql.test.js` + +**Step 1: Write the failing test** + +Add to `lex-gql.test.js` in a new describe block: + +```javascript +describe('Hydration Helpers', () => { + describe('hydrateBlobs', () => { + it('injects did into blob objects', () => { + const { hydrateBlobs } = require('./lex-gql.js'); + + const record = { + text: 'hello', + avatar: { + $type: 'blob', + ref: { $link: 'bafyreiabc123' }, + mimeType: 'image/jpeg', + size: 12345, + }, + }; + + const result = hydrateBlobs(record, 'did:plc:user123'); + + expect(result.text).toBe('hello'); + expect(result.avatar.did).toBe('did:plc:user123'); + expect(result.avatar.ref).toBe('bafyreiabc123'); + }); + + it('handles nested blob objects', () => { + const { hydrateBlobs } = require('./lex-gql.js'); + + const record = { + embed: { + images: [ + { image: { $type: 'blob', ref: 'bafyrei1', mimeType: 'image/jpeg', size: 100 } }, + { image: { $type: 'blob', ref: 'bafyrei2', mimeType: 'image/png', size: 200 } }, + ], + }, + }; + + const result = hydrateBlobs(record, 'did:plc:user123'); + + expect(result.embed.images[0].image.did).toBe('did:plc:user123'); + expect(result.embed.images[1].image.did).toBe('did:plc:user123'); + }); + + it('returns primitives unchanged', () => { + const { hydrateBlobs } = require('./lex-gql.js'); + + expect(hydrateBlobs(null, 'did:plc:x')).toBe(null); + expect(hydrateBlobs('string', 'did:plc:x')).toBe('string'); + expect(hydrateBlobs(123, 'did:plc:x')).toBe(123); + }); + + it('handles blob without $type but with ref/mimeType/size', () => { + const { hydrateBlobs } = require('./lex-gql.js'); + + const record = { + avatar: { ref: 'bafyreiabc', mimeType: 'image/jpeg', size: 100 }, + }; + + const result = hydrateBlobs(record, 'did:plc:user'); + expect(result.avatar.did).toBe('did:plc:user'); + }); + }); +}); +``` + +**Step 2: Run test to verify it fails** + +Run: `cd packages/lex-gql && pnpm test -- -t "hydrateBlobs"` +Expected: FAIL - hydrateBlobs is not exported + +**Step 3: Implement hydrateBlobs** + +Add to `lex-gql.js` after the utility functions section (~line 212): + +```javascript +// ============================================================================ +// HYDRATION HELPERS +// ============================================================================ + +/** + * Inject DID into blob objects for URL resolution. + * Blobs need the parent record's DID to generate CDN URLs. + * + * @param {*} obj - Record object or value to hydrate + * @param {string} did - DID to inject into blob objects + * @returns {*} - Hydrated object with did added to blobs + * + * @example + * const record = JSON.parse(row.record); + * const hydrated = hydrateBlobs(record, row.did); + */ +export function hydrateBlobs(obj, did) { + if (!obj || typeof obj !== 'object') return obj; + + // Check if this is a blob (has $type: 'blob' or has ref + mimeType + size) + if (obj.$type === 'blob' || (obj.ref && obj.mimeType && obj.size)) { + return { + ...obj, + ref: obj.ref?.$link || obj.ref, // Normalize { $link: "..." } format + did, + }; + } + + // Recurse into arrays + if (Array.isArray(obj)) { + return obj.map((item) => hydrateBlobs(item, did)); + } + + // Recurse into object properties + const result = {}; + for (const [key, value] of Object.entries(obj)) { + result[key] = hydrateBlobs(value, did); + } + return result; +} +``` + +**Step 4: Run test to verify it passes** + +Run: `cd packages/lex-gql && pnpm test -- -t "hydrateBlobs"` +Expected: PASS + +**Step 5: Commit** + +```bash +git add packages/lex-gql/lex-gql.js packages/lex-gql/lex-gql.test.js +git commit -m "feat: add hydrateBlobs helper for blob DID injection" +``` + +--- + +### Task 2: Add hydrateRecord helper function + +**Files:** +- Modify: `packages/lex-gql/lex-gql.js` +- Test: `packages/lex-gql/lex-gql.test.js` + +**Step 1: Write the failing test** + +Add to the `Hydration Helpers` describe block: + +```javascript +describe('hydrateRecord', () => { + it('transforms a database row to lex-gql record format', () => { + const { hydrateRecord } = require('./lex-gql.js'); + + const row = { + uri: 'at://did:plc:user123/app.bsky.feed.post/abc', + did: 'did:plc:user123', + collection: 'app.bsky.feed.post', + rkey: 'abc', + cid: 'bafyreicid', + record: JSON.stringify({ text: 'hello', createdAt: '2024-01-01T00:00:00Z' }), + indexed_at: '2024-01-01T00:00:00Z', + handle: 'user.bsky.social', + }; + + const result = hydrateRecord(row); + + expect(result.uri).toBe('at://did:plc:user123/app.bsky.feed.post/abc'); + expect(result.did).toBe('did:plc:user123'); + expect(result.collection).toBe('app.bsky.feed.post'); + expect(result.cid).toBe('bafyreicid'); + expect(result.indexedAt).toBe('2024-01-01T00:00:00Z'); + expect(result.actorHandle).toBe('user.bsky.social'); + expect(result.text).toBe('hello'); + expect(result.createdAt).toBe('2024-01-01T00:00:00Z'); + }); + + it('hydrates blob fields with did', () => { + const { hydrateRecord } = require('./lex-gql.js'); + + const row = { + uri: 'at://did:plc:user/app.bsky.actor.profile/self', + did: 'did:plc:user', + collection: 'app.bsky.actor.profile', + rkey: 'self', + cid: 'bafyreicid', + record: JSON.stringify({ + displayName: 'Test', + avatar: { $type: 'blob', ref: { $link: 'bafyrei123' }, mimeType: 'image/jpeg', size: 100 }, + }), + indexed_at: '2024-01-01T00:00:00Z', + }; + + const result = hydrateRecord(row); + + expect(result.avatar.did).toBe('did:plc:user'); + expect(result.avatar.ref).toBe('bafyrei123'); + }); + + it('handles missing optional fields', () => { + const { hydrateRecord } = require('./lex-gql.js'); + + const row = { + uri: 'at://did:plc:user/col/rkey', + did: 'did:plc:user', + collection: 'col', + rkey: 'rkey', + record: '{}', + indexed_at: '2024-01-01T00:00:00Z', + // cid and handle are missing + }; + + const result = hydrateRecord(row); + + expect(result.cid).toBeUndefined(); + expect(result.actorHandle).toBeNull(); + }); + + it('accepts record as object instead of JSON string', () => { + const { hydrateRecord } = require('./lex-gql.js'); + + const row = { + uri: 'at://did:plc:user/col/rkey', + did: 'did:plc:user', + collection: 'col', + rkey: 'rkey', + record: { text: 'already parsed' }, + indexed_at: '2024-01-01T00:00:00Z', + }; + + const result = hydrateRecord(row); + expect(result.text).toBe('already parsed'); + }); +}); +``` + +**Step 2: Run test to verify it fails** + +Run: `cd packages/lex-gql && pnpm test -- -t "hydrateRecord"` +Expected: FAIL - hydrateRecord is not exported + +**Step 3: Implement hydrateRecord** + +Add to `lex-gql.js` after `hydrateBlobs`: + +```javascript +/** + * Transform a database row into lex-gql record format. + * Expects the standard records table schema. + * + * Standard schema: + * - uri: TEXT (record AT URI) + * - did: TEXT (author DID) + * - collection: TEXT (lexicon NSID) + * - rkey: TEXT (record key) + * - cid: TEXT (optional, content ID) + * - record: TEXT (JSON) or Object + * - indexed_at: TEXT (ISO timestamp) + * - handle: TEXT (optional, actor handle from actors table join) + * + * @param {Object} row - Database row + * @returns {Object} - Hydrated record for lex-gql + * + * @example + * const rows = db.query('SELECT r.*, a.handle FROM records r LEFT JOIN actors a ON r.did = a.did'); + * const records = rows.map(hydrateRecord); + */ +export function hydrateRecord(row) { + const record = typeof row.record === 'string' ? JSON.parse(row.record) : row.record; + const hydrated = hydrateBlobs(record, row.did); + + return { + uri: row.uri, + cid: row.cid, + did: row.did, + collection: row.collection, + indexedAt: row.indexed_at, + actorHandle: row.handle || null, + ...hydrated, + }; +} +``` + +**Step 4: Run test to verify it passes** + +Run: `cd packages/lex-gql && pnpm test -- -t "hydrateRecord"` +Expected: PASS + +**Step 5: Commit** + +```bash +git add packages/lex-gql/lex-gql.js packages/lex-gql/lex-gql.test.js +git commit -m "feat: add hydrateRecord helper for standard row transformation" +``` + +--- + +### Task 3: Document the port interface in README + +**Files:** +- Modify: `packages/lex-gql/README.md` + +**Step 1: Add Port Interface section** + +Add after the "API" section in README.md: + +```markdown +## Query Port Interface + +lex-gql follows the hexagonal architecture pattern. Your data layer implements the **query port** interface: + +### Operation Types + +```typescript +type Operation = + | { type: 'findMany'; collection: string; where: WhereClause[]; pagination: Pagination; sort?: SortClause[] } + | { type: 'aggregate'; collection: string; where: WhereClause[]; groupBy?: string[] } + | { type: 'create'; collection: string; rkey?: string; record: object } + | { type: 'update'; collection: string; rkey: string; record: object } + | { type: 'delete'; collection: string; rkey: string } + +type WhereClause = { field: string; op: 'eq' | 'in' | 'contains' | 'gt' | 'gte' | 'lt' | 'lte'; value: any } +type SortClause = { field: string; dir: 'asc' | 'desc' } +type Pagination = { first?: number; after?: string; last?: number; before?: string } +``` + +### Response Format + +```typescript +// For findMany +{ rows: Record[]; hasNext: boolean; hasPrev: boolean } + +// For aggregate +{ count: number; groups: { [field]: value; count: number }[] } + +// For mutations +Record | { uri: string } +``` + +### Standard Records Schema + +For SQL-based adapters, we recommend this schema: + +```sql +CREATE TABLE records ( + uri TEXT PRIMARY KEY, + did TEXT NOT NULL, + collection TEXT NOT NULL, + rkey TEXT NOT NULL, + cid TEXT, + record TEXT NOT NULL, -- JSON blob + indexed_at TEXT NOT NULL +); + +CREATE INDEX idx_records_collection ON records(collection); +CREATE INDEX idx_records_did ON records(did); + +CREATE TABLE actors ( + did TEXT PRIMARY KEY, + handle TEXT NOT NULL +); +``` + +### Hydration Helpers + +Use these helpers to transform database rows into lex-gql format: + +```javascript +import { hydrateBlobs, hydrateRecord } from 'lex-gql'; + +// hydrateBlobs - inject DID into blob fields for URL resolution +const record = JSON.parse(row.record); +const hydrated = hydrateBlobs(record, row.did); + +// hydrateRecord - full transformation from standard schema +const rows = db.query('SELECT r.*, a.handle FROM records r LEFT JOIN actors a ON r.did = a.did'); +const records = rows.map(hydrateRecord); +``` +``` + +**Step 2: Commit** + +```bash +git add packages/lex-gql/README.md +git commit -m "docs: add port interface and hydration helpers documentation" +``` + +--- + +### Task 4: Update tap example to use hydration helpers + +**Files:** +- Modify: `examples/tap/index.js` + +**Step 1: Update imports** + +Change import to include helpers: + +```javascript +import { createAdapter, parseLexicon, hydrateRecord } from 'lex-gql'; +``` + +**Step 2: Remove local injectDidIntoBlobs helper** + +Delete the `injectDidIntoBlobs` function (lines ~198-223). + +**Step 3: Simplify findMany transform** + +Replace the transform section: + +```javascript +// Transform rows to lex-gql format +const transformed = rows.map((row) => ({ + ...hydrateRecord({ + uri: `at://${row.did}/${row.collection}/${row.rkey}`, + did: row.did, + collection: row.collection, + rkey: row.rkey, + cid: row.cid, + record: row.record, + indexed_at: row.indexed_at, + handle: row.handle, + }), + _id: row.id, // For cursor +})); +``` + +**Step 4: Verify tap example still works** + +Run: `cd examples/tap && node index.js` +Test the GraphQL query in browser. + +**Step 5: Commit** + +```bash +git add examples/tap/index.js +git commit -m "refactor(tap): use lex-gql hydration helpers" +``` + +--- + +### Task 5: Run full test suite and verify + +**Step 1: Run all tests** + +Run: `cd packages/lex-gql && pnpm test` +Expected: All tests pass + +**Step 2: Verify exports in type definitions** + +Check that `lex-gql.d.ts` exports the new functions (if it exists), or note that types need updating. + +**Step 3: Final commit if any cleanup needed** + +```bash +git add -A +git commit -m "chore: cleanup after hydration helpers" +``` + +--- diff --git a/examples/tap/docker-compose.yml b/examples/tap/docker-compose.yml index 30511e2..7a84092 100644 --- a/examples/tap/docker-compose.yml +++ b/examples/tap/docker-compose.yml @@ -10,4 +10,4 @@ services: TAP_DATABASE_URL: sqlite:///data/tap.db TAP_RELAY_URL: https://relay1.us-east.bsky.network TAP_SIGNAL_COLLECTION: xyz.statusphere.status - TAP_COLLECTION_FILTERS: xyz.statusphere.status + TAP_COLLECTION_FILTERS: xyz.statusphere.status,app.bsky.actor.profile diff --git a/examples/tap/index.js b/examples/tap/index.js index 9014527..9ce5c09 100644 --- a/examples/tap/index.js +++ b/examples/tap/index.js @@ -9,10 +9,18 @@ * 2. This server connects to tap's websocket and stores records locally * 3. lex-gql provides GraphQL queries over the stored records * + * Collections synced: + * - xyz.statusphere.status (emoji statuses) + * - app.bsky.actor.profile (user profiles with avatar/banner blobs) + * * Prerequisites: * docker compose up -d # Start tap server * pnpm install * node index.js + * + * Tap configuration for this example: + * TAP_SIGNAL_COLLECTION=xyz.statusphere.status + * TAP_COLLECTION_FILTERS=xyz.statusphere.status,app.bsky.actor.profile */ import { createServer } from 'node:http'; @@ -21,7 +29,7 @@ import { createHandler } from 'graphql-http/lib/use/node'; import { createAdapter, parseLexicon } from 'lex-gql'; import WebSocket from 'ws'; -// 1. Define lexicon for xyz.statusphere.status +// 1. Define lexicons const lexicons = [ parseLexicon({ lexicon: 1, @@ -46,8 +54,31 @@ const lexicons = [ }, }, }), + parseLexicon({ + lexicon: 1, + id: 'app.bsky.actor.profile', + defs: { + main: { + type: 'record', + key: 'literal:self', + record: { + type: 'object', + properties: { + displayName: { type: 'string', maxGraphemes: 64, maxLength: 640 }, + description: { type: 'string', maxGraphemes: 256, maxLength: 2560 }, + avatar: { type: 'blob', accept: ['image/png', 'image/jpeg'], maxSize: 1000000 }, + banner: { type: 'blob', accept: ['image/png', 'image/jpeg'], maxSize: 1000000 }, + createdAt: { type: 'string', format: 'datetime' }, + }, + }, + }, + }, + }), ]; +// Collections we want to sync from tap +const SYNC_COLLECTIONS = ['xyz.statusphere.status', 'app.bsky.actor.profile']; + // 2. Create local SQLite database for storing records from tap const dbPath = './data/records.db'; const db = new Database(dbPath); @@ -118,8 +149,8 @@ function connectToTap() { if (event.type === 'record' && event.record) { const { did, collection, rkey, cid, action, record } = event.record; - // Only store records for our collection - if (collection !== 'xyz.statusphere.status') { + // Only store records for our collections + if (!SYNC_COLLECTIONS.includes(collection)) { return; } @@ -164,6 +195,33 @@ function scheduleReconnect() { // Start websocket connection connectToTap(); +// Helper to inject did into blob fields for URL resolution +function injectDidIntoBlobs(obj, did) { + if (!obj || typeof obj !== 'object') return obj; + + // Check if this is a blob (has ref and mimeType) + if (obj.$type === 'blob' || (obj.ref && obj.mimeType && obj.size)) { + return { + ...obj, + ref: obj.ref?.$link || obj.ref, // Handle { $link: "..." } format + did, + }; + } + + // Recurse into object properties + const result = {}; + for (const [key, value] of Object.entries(obj)) { + if (Array.isArray(value)) { + result[key] = value.map((item) => injectDidIntoBlobs(item, did)); + } else if (typeof value === 'object') { + result[key] = injectDidIntoBlobs(value, did); + } else { + result[key] = value; + } + } + return result; +} + // 4. Query adapter: translates lex-gql operations to SQLite queries async function query(op) { if (op.type === 'findMany') { @@ -183,40 +241,49 @@ function findMany(op) { const conditions = ['r.collection = ?']; const params = [collection]; + // System fields that are columns, not in JSON + const systemFields = { + did: 'r.did', + uri: 'r.uri', + collection: 'r.collection', + cid: 'r.cid', + indexedAt: 'r.indexed_at', + }; + for (const clause of where) { const { field, op: operator, value } = clause; - // Map field to JSON path in record blob - const jsonPath = `json_extract(r.record, '$.${field}')`; + // Map field to column or JSON path + const fieldPath = systemFields[field] || `json_extract(r.record, '$.${field}')`; switch (operator) { case 'eq': - conditions.push(`${jsonPath} = ?`); + conditions.push(`${fieldPath} = ?`); params.push(value); break; case 'contains': - conditions.push(`${jsonPath} LIKE ?`); + conditions.push(`${fieldPath} LIKE ?`); params.push(`%${value}%`); break; case 'gt': - conditions.push(`${jsonPath} > ?`); + conditions.push(`${fieldPath} > ?`); params.push(value); break; case 'gte': - conditions.push(`${jsonPath} >= ?`); + conditions.push(`${fieldPath} >= ?`); params.push(value); break; case 'lt': - conditions.push(`${jsonPath} < ?`); + conditions.push(`${fieldPath} < ?`); params.push(value); break; case 'lte': - conditions.push(`${jsonPath} <= ?`); + conditions.push(`${fieldPath} <= ?`); params.push(value); break; case 'in': if (Array.isArray(value) && value.length > 0) { const placeholders = value.map(() => '?').join(', '); - conditions.push(`${jsonPath} IN (${placeholders})`); + conditions.push(`${fieldPath} IN (${placeholders})`); params.push(...value); } break; @@ -256,6 +323,7 @@ function findMany(op) { // Transform rows to lex-gql format const transformed = rows.map((row) => { const record = JSON.parse(row.record); + const withBlobs = injectDidIntoBlobs(record, row.did); return { uri: `at://${row.did}/${row.collection}/${row.rkey}`, cid: row.cid, @@ -263,7 +331,7 @@ function findMany(op) { collection: row.collection, indexedAt: row.indexed_at, actorHandle: row.handle || null, - ...record, + ...withBlobs, _id: row.id, // For cursor }; }); @@ -359,22 +427,24 @@ const graphiqlHtml = ` const defaultQuery = \`# Welcome to Tap GraphQL # # Query AT Protocol records synced by tap. -# This example tracks xyz.statusphere.status records. +# This example tracks xyz.statusphere.status and app.bsky.actor.profile records. -query { +# Query statuses with author profiles (DID join) +query StatusesWithProfiles { xyzStatusphereStatus(first: 10) { edges { node { - uri - did status createdAt - indexedAt + actorHandle + appBskyActorProfileByDid { + displayName + avatar { + url(preset: "avatar") + } + } } } - pageInfo { - hasNextPage - } } } \`; diff --git a/packages/lex-gql/lex-gql.js b/packages/lex-gql/lex-gql.js index 5f1f708..3677ef1 100644 --- a/packages/lex-gql/lex-gql.js +++ b/packages/lex-gql/lex-gql.js @@ -1280,6 +1280,8 @@ function createRecordType( * @param {Record} recordTypes * @param {GraphQLObjectType} resolvedRecordType * @param {JoinCollector} joinCollector + * @param {GraphQLObjectType} blobType + * @param {Function} queryFn * @returns {GraphQLObjectType} */ function createRecordTypeWithResolvers( @@ -1290,6 +1292,8 @@ function createRecordTypeWithResolvers( recordTypes, resolvedRecordType, joinCollector, + blobType, + queryFn, ) { return new GraphQLObjectType({ name: typeName, @@ -1315,7 +1319,7 @@ function createRecordTypeWithResolvers( // Add lexicon properties for (const prop of recordDef.properties) { fields[prop.name] = { - type: getGraphQLType(prop), + type: getGraphQLType(prop, blobType), description: `Field from lexicon`, }; @@ -1342,17 +1346,41 @@ function createRecordTypeWithResolvers( const fieldName = `${nsidToFieldName(otherLexicon.id)}ByDid`; const isUnique = otherLexicon.defs.main.key === 'literal:self'; + const otherCollection = otherLexicon.id; // Use full NSID as collection name + if (isUnique) { // Return single object for literal:self collections fields[fieldName] = { type: recordTypes[otherLexicon.id], description: `${otherTypeName} for this DID`, + resolve: async (parent) => { + const did = parent.did; + if (!did) return null; + const result = await queryFn({ + type: 'findMany', + collection: otherCollection, + where: [{ field: 'did', op: 'eq', value: did }], + pagination: { first: 1 }, + }); + return result.rows?.[0] || null; + }, }; } else { // Return list for multi-record collections fields[fieldName] = { type: new GraphQLList(recordTypes[otherLexicon.id]), description: `${otherTypeName} records for this DID`, + resolve: async (parent) => { + const did = parent.did; + if (!did) return []; + const result = await queryFn({ + type: 'findMany', + collection: otherCollection, + where: [{ field: 'did', op: 'eq', value: did }], + pagination: { first: 100 }, + }); + return result.rows || []; + }, }; } } @@ -1571,6 +1599,7 @@ function buildSchemaWithResolvers(lexicons, queryFn, subscribeFn) { const sortDirectionEnum = createSortDirectionEnum(); const deleteResultType = createDeleteResultType(); const resolvedRecordType = createResolvedRecordType(); + const blobType = createBlobType(); // Create join collector for batching const joinCollector = new JoinCollector(queryFn); @@ -1590,6 +1619,8 @@ function buildSchemaWithResolvers(lexicons, queryFn, subscribeFn) { recordTypes, resolvedRecordType, joinCollector, + blobType, + queryFn, ); // Create where input type (keyed by typeName for self-reference lookup) whereInputTypes[typeName] = createWhereInputType(