diff --git a/actions/publishToPublication.ts b/actions/publishToPublication.ts index 1755090b..1bf07483 100644 --- a/actions/publishToPublication.ts +++ b/actions/publishToPublication.ts @@ -271,6 +271,9 @@ export async function publishToPublication({ // Create record based on the document type let record: PubLeafletDocument.Record | SiteStandardDocument.Record; + // The record we persist locally always holds the fully-inflated content; when + // we offload pages to a blob, recordForPDS gets a shrunk copy below. + let recordForPDS: PubLeafletDocument.Record | SiteStandardDocument.Record; if (documentType === "site.standard.document") { // site.standard.document format @@ -305,6 +308,34 @@ export async function publishToPublication({ pages: pagesArray, }, } satisfies SiteStandardDocument.Record; + + // If the inline pages would push the record past the PDS's per-record size + // limits, offload them to a JSON blob and reference it via blobPages. We + // also lift every BlobRef found inside the pages onto a top-level `blobs` + // array — the PDS only scans the record itself for blob references when + // deciding what to garbage-collect, so any image/etc. blob that now lives + // inside the opaque JSON blob would otherwise look orphaned. + const CONTENT_BLOB_THRESHOLD = 100 * 1024; + const inlinePagesJson = JSON.stringify(pagesArray); + const inlinePagesBytes = Buffer.byteLength(inlinePagesJson, "utf8"); + if (inlinePagesBytes > CONTENT_BLOB_THRESHOLD) { + const pagesBlob = await agent.com.atproto.repo.uploadBlob( + new Blob([inlinePagesJson], { type: "application/json" }), + { headers: { "Content-Type": "application/json" } }, + ); + const referencedBlobs = collectBlobRefs(pagesArray); + recordForPDS = { + ...record, + content: { + $type: "pub.leaflet.content" as const, + pages: [], + blobPages: pagesBlob.data.blob, + ...(referencedBlobs.length > 0 && { blobs: referencedBlobs }), + }, + }; + } else { + recordForPDS = record; + } } else { // pub.leaflet.document format (legacy) record = { @@ -328,13 +359,14 @@ export async function publishToPublication({ pages: pagesArray, publishedAt: resolvedPublishedAt, } satisfies PubLeafletDocument.Record; + recordForPDS = record; } let { data: result } = await agent.com.atproto.repo.putRecord({ rkey, repo: credentialSession.did!, - collection: record.$type, - record, + collection: recordForPDS.$type, + record: recordForPDS, validate: false, //TODO publish the lexicon so we can validate! }); @@ -429,6 +461,27 @@ export async function publishToPublication({ return { success: true, rkey, record: JSON.parse(JSON.stringify(record)) }; } +// Walks an arbitrary value and returns every BlobRef instance reachable from +// it. Used to hoist image/etc. blob refs out of pages content so they remain +// referenced by the record after pages are offloaded to a JSON blob. +function collectBlobRefs(value: unknown): BlobRef[] { + const out: BlobRef[] = []; + const visit = (v: unknown) => { + if (v instanceof BlobRef) { + out.push(v); + return; + } + if (Array.isArray(v)) { + for (const item of v) visit(item); + return; + } + if (v && typeof v === "object") { + for (const item of Object.values(v)) visit(item); + } + }; + visit(value); + return out; +} async function extractThemeFromFacts( facts: Fact[], diff --git a/app/api/inngest/functions/sync_document_metadata.ts b/app/api/inngest/functions/sync_document_metadata.ts index 1ab7965f..7765faa9 100644 --- a/app/api/inngest/functions/sync_document_metadata.ts +++ b/app/api/inngest/functions/sync_document_metadata.ts @@ -2,6 +2,7 @@ import { inngest, events } from "../client"; import { supabaseServerClient } from "supabase/serverClient"; import { AtpAgent, AtUri } from "@atproto/api"; import { idResolver } from "app/(home-pages)/reader/idResolver"; +import type { Json } from "supabase/database.types"; // 1m, 2m, 4m, 8m, 16m, 32m, 1h, 2h, 4h, 8h, 8h, 8h (~37h total) const SLEEP_INTERVALS = [ @@ -49,6 +50,69 @@ export const sync_document_metadata = inngest.createFunction( return { handle: handleResult.handle, deleted: true }; } + await step.run("inflate-blob-pages", async () => { + // If the publisher offloaded pages to a JSON blob (see + // publishToPublication.ts), fetch the blob, splice it back into + // content.pages, and drop blobPages so downstream readers can treat + // documents.data as fully inflated. + const { data: doc } = await supabaseServerClient + .from("documents") + .select("data") + .eq("uri", document_uri) + .single(); + if (!doc?.data || typeof doc.data !== "object") return { skipped: true }; + + const record = doc.data as Record; + if (record.$type !== "site.standard.document") return { skipped: true }; + const content = record.content as Record | undefined; + const blobPages = content?.blobPages as + | { ref?: { $link?: string }; mimeType?: string } + | undefined; + if (!content || !blobPages) return { skipped: true }; + + const cid = blobPages.ref?.$link; + if (!cid) return { skipped: true }; + + const service = handleResult.doc?.service?.find( + (f) => f.id === "#atproto_pds", + ); + if (!service || typeof service.serviceEndpoint !== "string") { + return { skipped: "no_pds" as const }; + } + + const pdsAgent = new AtpAgent({ service: service.serviceEndpoint }); + const blobResponse = await pdsAgent.com.atproto.sync.getBlob({ + did, + cid, + }); + const inflatedPages = JSON.parse( + new TextDecoder().decode(blobResponse.data), + ); + if (!Array.isArray(inflatedPages)) { + return { skipped: "blob_not_array" as const }; + } + + // Strip both blobPages and the top-level `blobs` mirror — after inflation + // the blob refs live back inside pages, so the mirror would just be + // duplicate data. + const { blobPages: _blobPages, blobs: _blobs, ...contentWithoutBlob } = + content; + const inflatedRecord = { + ...record, + content: { + ...contentWithoutBlob, + pages: inflatedPages, + }, + }; + + await supabaseServerClient + .from("documents") + .update({ data: inflatedRecord as Json }) + .eq("uri", document_uri); + + return { inflated: true, pageCount: inflatedPages.length }; + }); + await step.run("set-indexed", async () => { return await supabaseServerClient .from("documents") diff --git a/appview/index.ts b/appview/index.ts index 5b9b0fc5..992b63b7 100644 --- a/appview/index.ts +++ b/appview/index.ts @@ -5,6 +5,7 @@ const idResolver = new IdResolver(); import { Firehose, MemoryRunner, Event } from "@atproto/sync"; import { ids } from "lexicons/api/lexicons"; import { + PubLeafletContent, PubLeafletDocument, PubLeafletGraphSubscription, PubLeafletPublication, @@ -276,11 +277,20 @@ async function handleEvent(evt: Event) { console.log(record.error); return; } - let docResult = await supabase.from("documents").upsert({ - uri: evt.uri.toString(), - data: record.value as Json, - }); - if (docResult.error) console.log(docResult.error); + // When the record offloads pages to a blob, skip the documents upsert — + // the firehose record has `pages: []` and would clobber the fully + // inflated copy that publishToPublication wrote optimistically. The + // inngest sync_document_metadata function is the writer in that case. + const hasBlobPages = + PubLeafletContent.isMain(record.value.content) && + !!record.value.content.blobPages; + if (!hasBlobPages) { + let docResult = await supabase.from("documents").upsert({ + uri: evt.uri.toString(), + data: record.value as Json, + }); + if (docResult.error) console.log(docResult.error); + } await inngest.send({ name: "appview/sync-document-metadata", data: { diff --git a/lexicons/api/lexicons.ts b/lexicons/api/lexicons.ts index 38a10025..6a9da079 100644 --- a/lexicons/api/lexicons.ts +++ b/lexicons/api/lexicons.ts @@ -1767,6 +1767,23 @@ export const schemaDict = { ], }, }, + blobPages: { + type: 'blob', + accept: ['application/json'], + maxSize: 5000000, + description: + 'JSON-encoded array of pages. When the inline pages array would be too large to store on the PDS, the pages are uploaded as a blob and referenced here. When set, consumers MUST ignore `pages` and use the decoded blob contents as the page array; the inline `pages` field will be empty or a stub.', + }, + blobs: { + type: 'array', + description: + 'Blobs referenced inside `blobPages`. Load-bearing when `blobPages` is set: the PDS only scans the top level of a record for blob references when deciding what to garbage-collect, so any image/etc. blob now living inside the opaque JSON blob must be mirrored here to remain referenced.', + items: { + type: 'blob', + accept: ['*/*'], + maxSize: 10000000, + }, + }, }, }, }, diff --git a/lexicons/api/types/pub/leaflet/content.ts b/lexicons/api/types/pub/leaflet/content.ts index 1dba8ce6..3f8784a5 100644 --- a/lexicons/api/types/pub/leaflet/content.ts +++ b/lexicons/api/types/pub/leaflet/content.ts @@ -20,6 +20,10 @@ export interface Main { | $Typed | { $type: string } )[] + /** JSON-encoded array of pages. When the inline pages array would be too large to store on the PDS, the pages are uploaded as a blob and referenced here. When set, consumers MUST ignore `pages` and use the decoded blob contents as the page array; the inline `pages` field will be empty or a stub. */ + blobPages?: BlobRef + /** Blobs referenced inside `blobPages`. Load-bearing when `blobPages` is set: the PDS only scans the top level of a record for blob references when deciding what to garbage-collect, so any image/etc. blob now living inside the opaque JSON blob must be mirrored here to remain referenced. */ + blobs?: BlobRef[] } const hashMain = 'main' diff --git a/lexicons/pub/leaflet/content.json b/lexicons/pub/leaflet/content.json index 7978c068..03596ffb 100644 --- a/lexicons/pub/leaflet/content.json +++ b/lexicons/pub/leaflet/content.json @@ -20,6 +20,25 @@ "pub.leaflet.pages.canvas" ] } + }, + "blobPages": { + "type": "blob", + "accept": [ + "application/json" + ], + "maxSize": 5000000, + "description": "JSON-encoded array of pages. When the inline pages array would be too large to store on the PDS, the pages are uploaded as a blob and referenced here. When set, consumers MUST ignore `pages` and use the decoded blob contents as the page array; the inline `pages` field will be empty or a stub." + }, + "blobs": { + "type": "array", + "description": "Blobs referenced inside `blobPages`. Load-bearing when `blobPages` is set: the PDS only scans the top level of a record for blob references when deciding what to garbage-collect, so any image/etc. blob now living inside the opaque JSON blob must be mirrored here to remain referenced.", + "items": { + "type": "blob", + "accept": [ + "*/*" + ], + "maxSize": 10000000 + } } } } diff --git a/lexicons/src/content.ts b/lexicons/src/content.ts index 4c8c7895..3f49d5dc 100644 --- a/lexicons/src/content.ts +++ b/lexicons/src/content.ts @@ -23,6 +23,23 @@ export const PubLeafletContent: LexiconDoc = { ], }, }, + blobPages: { + type: "blob", + accept: ["application/json"], + maxSize: 5000000, + description: + "JSON-encoded array of pages. When the inline pages array would be too large to store on the PDS, the pages are uploaded as a blob and referenced here. When set, consumers MUST ignore `pages` and use the decoded blob contents as the page array; the inline `pages` field will be empty or a stub.", + }, + blobs: { + type: "array", + description: + "Blobs referenced inside `blobPages`. Load-bearing when `blobPages` is set: the PDS only scans the top level of a record for blob references when deciding what to garbage-collect, so any image/etc. blob now living inside the opaque JSON blob must be mirrored here to remain referenced.", + items: { + type: "blob", + accept: ["*/*"], + maxSize: 10000000, + }, + }, }, }, },