diff --git a/package.json b/package.json index ed689f3..e44f6b5 100644 --- a/package.json +++ b/package.json @@ -11,7 +11,8 @@ "dev:client": "cd packages/client && npm run dev", "deploy": "cd packages/server && npm run deploy", "deploy:client": "cd packages/client && npm run pages:deploy", - "db:create": "cd packages/server && npm run db:create", + "db:init": "cd packages/server && npm run db:init", + "db:init:prod": "cd packages/server && npm run db:init:prod", "db:migrate": "cd packages/server && npm run db:migrate", "db:migrate:prod": "cd packages/server && npm run db:migrate:prod", "secret:set": "cd packages/server && npm run secret:set" diff --git a/packages/client/package.json b/packages/client/package.json index 8ea92b8..ede6228 100644 --- a/packages/client/package.json +++ b/packages/client/package.json @@ -6,7 +6,7 @@ "dev": "vite", "build": "tsc && vite build", "preview": "vite preview", - "pages:deploy": "vite build && wrangler pages deploy dist --project-name=atfeeds" + "deploy": "vite build && wrangler pages deploy dist" }, "dependencies": { "react": "^18.2.0", diff --git a/packages/client/src/App.tsx b/packages/client/src/App.tsx index 5155640..1a2f494 100644 --- a/packages/client/src/App.tsx +++ b/packages/client/src/App.tsx @@ -3,20 +3,41 @@ import { useEffect, useState } from "react"; // API base URL - empty for same-origin (local dev), or set via env var for production const API_URL = "https://atfeeds-api.stevedsimkins.workers.dev"; +interface BskyPostRef { + uri: string; + cid: string; +} + +interface Publication { + url: string; + name: string; + description?: string; + iconCid?: string; + iconUrl?: string; +} + interface Document { uri: string; did: string; rkey: string; title: string; - path: string | null; - site: string | null; - content: { + description?: string; + path?: string; + site?: string; + content?: { $type: string; markdown?: string; - } | null; - textContent: string | null; - publishedAt: string | null; - viewUrl: string | null; + }; + textContent?: string; + coverImageCid?: string; + coverImageUrl?: string; + bskyPostRef?: BskyPostRef; + tags?: string[]; + publishedAt?: string; + updatedAt?: string; + publication?: Publication; + viewUrl?: string; + pdsEndpoint?: string; } interface FeedResponse { @@ -52,7 +73,7 @@ function App() { fetchFeed(); }, []); - const formatDate = (dateString: string | null) => { + const formatDate = (dateString?: string) => { if (!dateString) return "Unknown date"; return new Date(dateString).toLocaleDateString("en-US", { year: "numeric", @@ -61,12 +82,16 @@ function App() { }); }; - const truncateText = (text: string | null, maxLength: number = 200) => { + const truncateText = (text?: string, maxLength: number = 200) => { if (!text) return ""; if (text.length <= maxLength) return text; return text.slice(0, maxLength) + "..."; }; + const getDescription = (doc: Document) => { + return doc.description || doc.textContent || ""; + }; + return (
@@ -130,6 +155,51 @@ function App() { )}
+ {/* Publication info */} + {doc.publication && ( +
+ {doc.publication.iconUrl && ( + {doc.publication.name} + )} + + {doc.publication.name} + +
+ )} + + {/* Cover image */} + {doc.coverImageUrl && ( +
+ {doc.title} +
+ )} + + {/* Date */}
Published: {formatDate(doc.publishedAt)} + {doc.updatedAt && doc.updatedAt !== doc.publishedAt && ( + <> | Updated: {formatDate(doc.updatedAt)} + )}
- {doc.textContent && ( + + {/* Description */} + {getDescription(doc) && (

- {truncateText(doc.textContent)} + {truncateText(getDescription(doc))}

)} - {doc.viewUrl && ( -
+ + {/* Tags */} + {doc.tags && doc.tags.length > 0 && ( +
+ {doc.tags.map((tag) => ( + + {tag} + + ))} +
+ )} + + {/* Actions */} +
+ {doc.bskyPostRef && ( + + )} + {doc.viewUrl && ( -
- )} + )} +
))} diff --git a/packages/server/migrations/001_add_document_fields.sql b/packages/server/migrations/001_add_document_fields.sql new file mode 100644 index 0000000..96289f0 --- /dev/null +++ b/packages/server/migrations/001_add_document_fields.sql @@ -0,0 +1,23 @@ +-- Migration: Add full Document and Publication fields to resolved_documents +-- Run with: wrangler d1 execute atfeeds-db --file=migrations/001_add_document_fields.sql --remote + +-- Document fields +ALTER TABLE resolved_documents ADD COLUMN description TEXT; +ALTER TABLE resolved_documents ADD COLUMN cover_image_cid TEXT; +ALTER TABLE resolved_documents ADD COLUMN cover_image_url TEXT; +ALTER TABLE resolved_documents ADD COLUMN bsky_post_ref TEXT; +ALTER TABLE resolved_documents ADD COLUMN tags TEXT; +ALTER TABLE resolved_documents ADD COLUMN updated_at TEXT; + +-- Publication fields +ALTER TABLE resolved_documents ADD COLUMN pub_url TEXT; +ALTER TABLE resolved_documents ADD COLUMN pub_name TEXT; +ALTER TABLE resolved_documents ADD COLUMN pub_description TEXT; +ALTER TABLE resolved_documents ADD COLUMN pub_icon_cid TEXT; +ALTER TABLE resolved_documents ADD COLUMN pub_icon_url TEXT; + +-- Metadata +ALTER TABLE resolved_documents ADD COLUMN pds_endpoint TEXT; + +-- Index for publication queries +CREATE INDEX IF NOT EXISTS idx_resolved_documents_pub_url ON resolved_documents(pub_url); diff --git a/packages/server/package.json b/packages/server/package.json index 7298afe..45b391d 100644 --- a/packages/server/package.json +++ b/packages/server/package.json @@ -4,7 +4,11 @@ "private": true, "scripts": { "dev": "wrangler dev", - "deploy": "wrangler deploy" + "deploy": "wrangler deploy", + "db:init": "wrangler d1 execute atfeeds-db --file=schema.sql --local", + "db:init:prod": "wrangler d1 execute atfeeds-db --file=schema.sql --remote", + "db:migrate": "wrangler d1 execute atfeeds-db --file=migrations/001_add_document_fields.sql --local", + "db:migrate:prod": "wrangler d1 execute atfeeds-db --file=migrations/001_add_document_fields.sql --remote" }, "dependencies": { "hono": "^4.0.0" diff --git a/packages/server/schema.sql b/packages/server/schema.sql index 74b9f7a..b17aa64 100644 --- a/packages/server/schema.sql +++ b/packages/server/schema.sql @@ -32,16 +32,32 @@ CREATE TABLE IF NOT EXISTS resolved_documents ( uri TEXT PRIMARY KEY, did TEXT NOT NULL, rkey TEXT NOT NULL, + -- Document fields title TEXT, + description TEXT, path TEXT, site TEXT, - content TEXT, -- JSON blob + content TEXT, -- JSON blob for content union text_content TEXT, + cover_image_cid TEXT, -- CID for cover image blob + cover_image_url TEXT, -- Full URL: {pds}/xrpc/com.atproto.sync.getBlob?did={did}&cid={cid} + bsky_post_ref TEXT, -- JSON blob for strong reference {uri, cid} + tags TEXT, -- JSON array of strings published_at TEXT, - view_url TEXT, + updated_at TEXT, + -- Publication fields (resolved from site at:// URI) + pub_url TEXT, -- Publication base URL + pub_name TEXT, + pub_description TEXT, + pub_icon_cid TEXT, -- CID for publication icon blob + pub_icon_url TEXT, -- Full URL to publication icon + -- Metadata + view_url TEXT, -- Constructed canonical URL (pub_url + path) + pds_endpoint TEXT, -- Cached PDS endpoint for this DID resolved_at TEXT DEFAULT (datetime('now')), stale_at TEXT -- When this record should be re-resolved ); CREATE INDEX IF NOT EXISTS idx_resolved_documents_rkey ON resolved_documents(rkey DESC); CREATE INDEX IF NOT EXISTS idx_resolved_documents_stale ON resolved_documents(stale_at); +CREATE INDEX IF NOT EXISTS idx_resolved_documents_pub_url ON resolved_documents(pub_url); diff --git a/packages/server/src/index.ts b/packages/server/src/index.ts index 04f4463..134831c 100644 --- a/packages/server/src/index.ts +++ b/packages/server/src/index.ts @@ -1,7 +1,7 @@ import { Hono } from "hono"; import { cors } from "hono/cors"; import type { Bindings } from "./types"; -import { health, webhook, feed, stats, records } from "./routes"; +import { health, webhook, feed, stats, records, admin } from "./routes"; import { processDocument } from "./utils"; const app = new Hono<{ Bindings: Bindings }>(); @@ -15,6 +15,7 @@ app.route("/webhook", webhook); app.route("/feed", feed); app.route("/stats", stats); app.route("/records", records); +app.route("/admin", admin); // Legacy alias: /feed-raw -> /feed/raw app.get("/feed-raw", async (c) => { diff --git a/packages/server/src/routes/admin.ts b/packages/server/src/routes/admin.ts new file mode 100644 index 0000000..84990d0 --- /dev/null +++ b/packages/server/src/routes/admin.ts @@ -0,0 +1,77 @@ +import { Hono } from "hono"; +import type { Bindings } from "../types"; + +const admin = new Hono<{ Bindings: Bindings }>(); + +// Queue all documents for re-processing +admin.post("/resolve-all", async (c) => { + try { + const db = c.env.DB; + const queue = c.env.RESOLUTION_QUEUE; + + // Get all records from repo_records + const { results } = await db + .prepare( + `SELECT did, rkey FROM repo_records + WHERE collection = 'site.standard.document'` + ) + .all<{ did: string; rkey: string }>(); + + if (!results || results.length === 0) { + return c.json({ message: "No documents to process", queued: 0 }); + } + + // Queue in batches of 100 (Cloudflare Queue limit) + const batchSize = 100; + let queued = 0; + + for (let i = 0; i < results.length; i += batchSize) { + const batch = results.slice(i, i + batchSize); + const messages = batch.map((row) => ({ + body: { + did: row.did, + collection: "site.standard.document", + rkey: row.rkey, + }, + })); + + await queue.sendBatch(messages); + queued += messages.length; + } + + return c.json({ + message: "Documents queued for re-processing", + queued, + }); + } catch (error) { + return c.json( + { error: "Failed to queue documents", details: String(error) }, + 500 + ); + } +}); + +// Mark all documents as stale (alternative - lets cron handle it) +admin.post("/mark-stale", async (c) => { + try { + const db = c.env.DB; + + const result = await db + .prepare( + `UPDATE resolved_documents SET stale_at = datetime('now', '-1 hour')` + ) + .run(); + + return c.json({ + message: "All documents marked as stale", + affected: result.meta.changes, + }); + } catch (error) { + return c.json( + { error: "Failed to mark documents as stale", details: String(error) }, + 500 + ); + } +}); + +export default admin; diff --git a/packages/server/src/routes/feed.ts b/packages/server/src/routes/feed.ts index 67164a6..dd2000e 100644 --- a/packages/server/src/routes/feed.ts +++ b/packages/server/src/routes/feed.ts @@ -1,8 +1,76 @@ import { Hono } from "hono"; -import type { Bindings } from "../types"; +import type { Bindings, ResolvedDocumentRow, Document, Publication, BskyPostRef } from "../types"; const feed = new Hono<{ Bindings: Bindings }>(); +/** + * Transforms a database row into a Document object for the API response. + */ +function rowToDocument(row: ResolvedDocumentRow): Document { + // Build publication object if we have publication data + let publication: Publication | undefined; + if (row.pub_url && row.pub_name) { + publication = { + url: row.pub_url, + name: row.pub_name, + description: row.pub_description || undefined, + iconCid: row.pub_icon_cid || undefined, + iconUrl: row.pub_icon_url || undefined, + }; + } + + // Parse bskyPostRef if present + let bskyPostRef: BskyPostRef | undefined; + if (row.bsky_post_ref) { + try { + bskyPostRef = JSON.parse(row.bsky_post_ref); + } catch { + // Ignore parse errors + } + } + + // Parse tags if present + let tags: string[] | undefined; + if (row.tags) { + try { + tags = JSON.parse(row.tags); + } catch { + // Ignore parse errors + } + } + + // Parse content if present + let content: unknown | undefined; + if (row.content) { + try { + content = JSON.parse(row.content); + } catch { + // Ignore parse errors + } + } + + return { + uri: row.uri, + did: row.did, + rkey: row.rkey, + title: row.title || "Untitled", + description: row.description || undefined, + path: row.path || undefined, + site: row.site || undefined, + content, + textContent: row.text_content || undefined, + coverImageCid: row.cover_image_cid || undefined, + coverImageUrl: row.cover_image_url || undefined, + bskyPostRef, + tags, + publishedAt: row.published_at || undefined, + updatedAt: row.updated_at || undefined, + publication, + viewUrl: row.view_url || undefined, + pdsEndpoint: row.pds_endpoint || undefined, + }; +} + // Get raw feed data (for client-side fetching) // Accessible at both /feed/raw and /feed-raw (via alias in index.ts) feed.get("/raw", async (c) => { @@ -44,37 +112,19 @@ feed.get("/", async (c) => { const { results } = await db .prepare( - `SELECT uri, did, rkey, title, path, site, content, text_content, published_at, view_url + `SELECT uri, did, rkey, title, description, path, site, content, text_content, + cover_image_cid, cover_image_url, bsky_post_ref, tags, + published_at, updated_at, pub_url, pub_name, pub_description, + pub_icon_cid, pub_icon_url, view_url, pds_endpoint, + resolved_at, stale_at FROM resolved_documents ORDER BY rkey DESC LIMIT ? OFFSET ?` ) .bind(limit, offset) - .all<{ - uri: string; - did: string; - rkey: string; - title: string | null; - path: string | null; - site: string | null; - content: string | null; - text_content: string | null; - published_at: string | null; - view_url: string | null; - }>(); + .all(); - const documents = (results || []).map((doc) => ({ - uri: doc.uri, - did: doc.did, - rkey: doc.rkey, - title: doc.title || "Untitled", - path: doc.path, - site: doc.site, - content: doc.content ? JSON.parse(doc.content) : null, - textContent: doc.text_content, - publishedAt: doc.published_at, - viewUrl: doc.view_url, - })); + const documents = (results || []).map(rowToDocument); return c.json({ count: documents.length, diff --git a/packages/server/src/routes/index.ts b/packages/server/src/routes/index.ts index 65c517d..5fbf8c4 100644 --- a/packages/server/src/routes/index.ts +++ b/packages/server/src/routes/index.ts @@ -3,3 +3,4 @@ export { default as webhook } from "./webhook"; export { default as feed } from "./feed"; export { default as stats } from "./stats"; export { default as records } from "./records"; +export { default as admin } from "./admin"; diff --git a/packages/server/src/types/index.ts b/packages/server/src/types/index.ts index 9f21719..060bcd2 100644 --- a/packages/server/src/types/index.ts +++ b/packages/server/src/types/index.ts @@ -32,15 +32,70 @@ export interface TapIdentityEvent { export type TapEvent = TapRecordEvent | TapIdentityEvent; +// Strong reference to a Bluesky post +export interface BskyPostRef { + uri: string; + cid: string; +} + +// Publication record from site.standard.publication +export interface Publication { + url: string; // Base publication URL + name: string; + description?: string; + iconCid?: string; // CID for icon blob + iconUrl?: string; // Resolved full URL to icon +} + +// Document record from site.standard.document export interface Document { uri: string; did: string; rkey: string; + // Document fields title: string; + description?: string; + path?: string; + site?: string; // at:// URI to publication or https:// URL + content?: unknown; + textContent?: string; + coverImageCid?: string; // CID for cover image blob + coverImageUrl?: string; // Resolved full URL to cover image + bskyPostRef?: BskyPostRef; + tags?: string[]; + publishedAt?: string; + updatedAt?: string; + // Resolved publication data + publication?: Publication; + // Metadata + viewUrl?: string; // Canonical URL (publication.url + path) + pdsEndpoint?: string; // PDS endpoint used for blob URLs +} + +// Database row for resolved_documents table +export interface ResolvedDocumentRow { + uri: string; + did: string; + rkey: string; + title: string | null; + description: string | null; path: string | null; site: string | null; - content: unknown; - textContent: string | null; - publishedAt: string | null; - viewUrl: string | null; + content: string | null; + text_content: string | null; + cover_image_cid: string | null; + cover_image_url: string | null; + bsky_post_ref: string | null; + tags: string | null; + published_at: string | null; + updated_at: string | null; + pub_url: string | null; + pub_name: string | null; + pub_description: string | null; + pub_icon_cid: string | null; + pub_icon_url: string | null; + view_url: string | null; + pds_endpoint: string | null; + resolved_at: string | null; + stale_at: string | null; } diff --git a/packages/server/src/utils/blob.ts b/packages/server/src/utils/blob.ts new file mode 100644 index 0000000..d5fa766 --- /dev/null +++ b/packages/server/src/utils/blob.ts @@ -0,0 +1,35 @@ +/** + * Constructs a full URL to fetch a blob from a PDS. + * Format: {pds}/xrpc/com.atproto.sync.getBlob?did={did}&cid={cid} + */ +export function buildBlobUrl(pds: string, did: string, cid: string): string { + const baseUrl = pds.endsWith("/") ? pds.slice(0, -1) : pds; + return `${baseUrl}/xrpc/com.atproto.sync.getBlob?did=${encodeURIComponent(did)}&cid=${encodeURIComponent(cid)}`; +} + +/** + * Extracts the CID from a blob reference object. + * Blob refs can be in different formats: + * - { $link: "cid" } (legacy) + * - { ref: { $link: "cid" } } (current) + * - { cid: "cid" } (simple) + */ +export function extractBlobCid(blob: unknown): string | null { + if (!blob || typeof blob !== "object") return null; + + const b = blob as Record; + + // Current format: { ref: { $link: "cid" } } + if (b.ref && typeof b.ref === "object") { + const ref = b.ref as Record; + if (typeof ref.$link === "string") return ref.$link; + } + + // Legacy format: { $link: "cid" } + if (typeof b.$link === "string") return b.$link; + + // Simple format: { cid: "cid" } + if (typeof b.cid === "string") return b.cid; + + return null; +} diff --git a/packages/server/src/utils/document.ts b/packages/server/src/utils/document.ts index a0613eb..3c6ee8c 100644 --- a/packages/server/src/utils/document.ts +++ b/packages/server/src/utils/document.ts @@ -1,11 +1,46 @@ import { resolvePds } from "./resolver"; import { parseAtUri } from "./at-uri"; +import { buildBlobUrl, extractBlobCid } from "./blob"; -export async function resolveViewUrl( +// Raw document record from PDS +interface DocumentRecord { + site?: string; + path?: string; + title?: string; + description?: string; + coverImage?: unknown; + content?: unknown; + textContent?: string; + bskyPostRef?: { uri: string; cid: string }; + tags?: string[]; + publishedAt?: string; + updatedAt?: string; +} + +// Raw publication record from PDS +interface PublicationRecord { + url?: string; + name?: string; + description?: string; + icon?: unknown; +} + +// Resolved publication data +interface ResolvedPublication { + url: string; + name: string; + description: string | null; + iconCid: string | null; + iconUrl: string | null; +} + +/** + * Fetches a publication record from an at:// URI + */ +async function fetchPublication( db: D1Database, - siteUri: string, - path: string -): Promise { + siteUri: string +): Promise { const parsed = parseAtUri(siteUri); if (!parsed) return null; @@ -18,22 +53,56 @@ export async function resolveViewUrl( )}&collection=${encodeURIComponent(parsed.collection)}&rkey=${encodeURIComponent( parsed.rkey )}`; + const response = await fetch(url); if (!response.ok) return null; - const data = (await response.json()) as { - value?: { url?: string; domain?: string }; - }; - const siteUrl = data.value?.url || data.value?.domain; - if (!siteUrl) return null; + const data = (await response.json()) as { value?: PublicationRecord }; + const pub = data.value; + if (!pub?.url || !pub?.name) return null; + + const iconCid = extractBlobCid(pub.icon); + const iconUrl = iconCid ? buildBlobUrl(pds, parsed.did, iconCid) : null; - const baseUrl = siteUrl.startsWith("http") ? siteUrl : `https://${siteUrl}`; - return `${baseUrl}${path}`; + return { + url: pub.url, + name: pub.name, + description: pub.description || null, + iconCid, + iconUrl, + }; } catch { return null; } } +/** + * Resolves the view URL for a document. + * If site is an at:// URI, fetches the publication to get the base URL. + * If site is an https:// URL, uses it directly. + */ +export async function resolveViewUrl( + db: D1Database, + siteUri: string, + path: string +): Promise { + // Check if site is an at:// URI or direct URL + if (siteUri.startsWith("at://")) { + const pub = await fetchPublication(db, siteUri); + if (!pub?.url) return null; + const baseUrl = pub.url.startsWith("http") ? pub.url : `https://${pub.url}`; + return `${baseUrl.replace(/\/$/, "")}${path}`; + } + + // Direct URL + const baseUrl = siteUri.startsWith("http") ? siteUri : `https://${siteUri}`; + return `${baseUrl.replace(/\/$/, "")}${path}`; +} + +/** + * Processes a document record: fetches from PDS, resolves publication, + * and stores all fields in resolved_documents table. + */ export async function processDocument( db: D1Database, did: string, @@ -48,16 +117,15 @@ export async function processDocument( return; } - // 2. Fetch Record + // 2. Fetch Document Record const url = `${pds}/xrpc/com.atproto.repo.getRecord?repo=${encodeURIComponent( did )}&collection=${encodeURIComponent(collection)}&rkey=${encodeURIComponent(rkey)}`; - + const response = await fetch(url); if (!response.ok) { if (response.status === 404) { - // Record deleted? - console.warn(`Record not found: ${did}/${collection}/${rkey}`); + console.warn(`Record not found: ${did}/${collection}/${rkey}`); } return; } @@ -65,15 +133,7 @@ export async function processDocument( const data = (await response.json()) as { uri: string; cid?: string; - value: { - title?: string; - path?: string; - site?: string; - content?: unknown; - textContent?: string; - publishedAt?: string; - [key: string]: unknown; - }; + value: DocumentRecord; }; const { value, cid } = data; @@ -90,52 +150,89 @@ export async function processDocument( .bind(did, rkey, collection, cid || null, cid || null) .run(); - // 4. Resolve View URL and Update resolved_documents - const uri = `at://${did}/${collection}/${rkey}`; + // 4. Extract document fields + const title = value.title || null; + const description = value.description || null; + const path = value.path || null; + const site = value.site || null; + const content = value.content ? JSON.stringify(value.content) : null; + const textContent = value.textContent || null; + const coverImageCid = extractBlobCid(value.coverImage); + const coverImageUrl = coverImageCid ? buildBlobUrl(pds, did, coverImageCid) : null; + const bskyPostRef = value.bskyPostRef ? JSON.stringify(value.bskyPostRef) : null; + const tags = value.tags ? JSON.stringify(value.tags) : null; + const publishedAt = value.publishedAt || null; + const updatedAt = value.updatedAt || null; + + // 5. Resolve publication if site is at:// URI + let pubUrl: string | null = null; + let pubName: string | null = null; + let pubDescription: string | null = null; + let pubIconCid: string | null = null; + let pubIconUrl: string | null = null; let viewUrl: string | null = null; - if (value.site && value.path) { - viewUrl = await resolveViewUrl(db, value.site, value.path); + + if (site) { + if (site.startsWith("at://")) { + // Fetch publication record + const pub = await fetchPublication(db, site); + if (pub) { + pubUrl = pub.url; + pubName = pub.name; + pubDescription = pub.description; + pubIconCid = pub.iconCid; + pubIconUrl = pub.iconUrl; + // Construct view URL + if (pubUrl && path) { + const baseUrl = pubUrl.startsWith("http") ? pubUrl : `https://${pubUrl}`; + viewUrl = `${baseUrl.replace(/\/$/, "")}${path}`; + } + } + } else { + // Site is a direct URL (loose document) + pubUrl = site; + if (path) { + const baseUrl = site.startsWith("http") ? site : `https://${site}`; + viewUrl = `${baseUrl.replace(/\/$/, "")}${path}`; + } + } } - // Set stale_at to 12 hours from now + // 6. Insert/update resolved_documents + const uri = `at://${did}/${collection}/${rkey}`; const STALE_OFFSET_HOURS = 12; await db .prepare( - `INSERT INTO resolved_documents (uri, did, rkey, title, path, site, content, text_content, published_at, view_url, resolved_at, stale_at) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, datetime('now'), datetime('now', '+${STALE_OFFSET_HOURS} hours')) - ON CONFLICT(uri) DO UPDATE SET - title = ?, - path = ?, - site = ?, - content = ?, - text_content = ?, - published_at = ?, - view_url = ?, - resolved_at = datetime('now'), - stale_at = datetime('now', '+${STALE_OFFSET_HOURS} hours')` + `INSERT INTO resolved_documents ( + uri, did, rkey, title, description, path, site, content, text_content, + cover_image_cid, cover_image_url, bsky_post_ref, tags, + published_at, updated_at, pub_url, pub_name, pub_description, + pub_icon_cid, pub_icon_url, view_url, pds_endpoint, + resolved_at, stale_at + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, datetime('now'), datetime('now', '+${STALE_OFFSET_HOURS} hours')) + ON CONFLICT(uri) DO UPDATE SET + title = ?, description = ?, path = ?, site = ?, content = ?, text_content = ?, + cover_image_cid = ?, cover_image_url = ?, bsky_post_ref = ?, tags = ?, + published_at = ?, updated_at = ?, pub_url = ?, pub_name = ?, pub_description = ?, + pub_icon_cid = ?, pub_icon_url = ?, view_url = ?, pds_endpoint = ?, + resolved_at = datetime('now'), stale_at = datetime('now', '+${STALE_OFFSET_HOURS} hours')` ) .bind( - uri, - did, - rkey, - value.title || null, - value.path || null, - value.site || null, - value.content ? JSON.stringify(value.content) : null, - value.textContent || null, - value.publishedAt || null, - viewUrl, - // Update bindings - value.title || null, - value.path || null, - value.site || null, - value.content ? JSON.stringify(value.content) : null, - value.textContent || null, - value.publishedAt || null, - viewUrl + // INSERT values + uri, did, rkey, title, description, path, site, content, textContent, + coverImageCid, coverImageUrl, bskyPostRef, tags, + publishedAt, updatedAt, pubUrl, pubName, pubDescription, + pubIconCid, pubIconUrl, viewUrl, pds, + // UPDATE values + title, description, path, site, content, textContent, + coverImageCid, coverImageUrl, bskyPostRef, tags, + publishedAt, updatedAt, pubUrl, pubName, pubDescription, + pubIconCid, pubIconUrl, viewUrl, pds ) .run(); + + console.log(`Processed document: ${uri}`); } catch (error) { console.error(`Error processing document ${did}/${collection}/${rkey}:`, error); } diff --git a/packages/server/src/utils/index.ts b/packages/server/src/utils/index.ts index 5e48159..aa181bb 100644 --- a/packages/server/src/utils/index.ts +++ b/packages/server/src/utils/index.ts @@ -1,3 +1,4 @@ export { parseAtUri, buildAtUri, type AtUriComponents } from "./at-uri"; export { resolvePds } from "./resolver"; export { resolveViewUrl, processDocument } from "./document"; +export { buildBlobUrl, extractBlobCid } from "./blob";