diff --git a/lexicons-generated/rsvp/atmo/space/defs.json b/lexicons-generated/rsvp/atmo/space/defs.json index 01f50b6..a9f68e3 100644 --- a/lexicons-generated/rsvp/atmo/space/defs.json +++ b/lexicons-generated/rsvp/atmo/space/defs.json @@ -127,6 +127,35 @@ } } }, + "blobInfo": { + "type": "object", + "required": [ + "cid", + "mimeType", + "size", + "authorDid", + "createdAt" + ], + "properties": { + "cid": { + "type": "string", + "format": "cid" + }, + "mimeType": { + "type": "string" + }, + "size": { + "type": "integer" + }, + "authorDid": { + "type": "string", + "format": "did" + }, + "createdAt": { + "type": "integer" + } + } + }, "inviteView": { "type": "object", "required": [ diff --git a/lexicons-generated/rsvp/atmo/space/getBlob.json b/lexicons-generated/rsvp/atmo/space/getBlob.json new file mode 100644 index 0000000..a4b45c5 --- /dev/null +++ b/lexicons-generated/rsvp/atmo/space/getBlob.json @@ -0,0 +1,42 @@ +{ + "lexicon": 1, + "id": "rsvp.atmo.space.getBlob", + "defs": { + "main": { + "type": "query", + "description": "Read a blob from a space. Requires read access via a service-auth JWT or a read-grant invite token.", + "parameters": { + "type": "params", + "required": [ + "spaceUri", + "cid" + ], + "properties": { + "spaceUri": { + "type": "string", + "format": "at-uri" + }, + "cid": { + "type": "string", + "format": "cid" + }, + "inviteToken": { + "type": "string", + "description": "Read-grant invite token for anonymous bearer access." + } + } + }, + "output": { + "encoding": "*/*" + }, + "errors": [ + { + "name": "NotFound" + }, + { + "name": "Forbidden" + } + ] + } + } +} diff --git a/lexicons-generated/rsvp/atmo/space/listBlobs.json b/lexicons-generated/rsvp/atmo/space/listBlobs.json new file mode 100644 index 0000000..b19d743 --- /dev/null +++ b/lexicons-generated/rsvp/atmo/space/listBlobs.json @@ -0,0 +1,57 @@ +{ + "lexicon": 1, + "id": "rsvp.atmo.space.listBlobs", + "defs": { + "main": { + "type": "query", + "description": "List blob metadata for a space. Members only.", + "parameters": { + "type": "params", + "required": [ + "spaceUri" + ], + "properties": { + "spaceUri": { + "type": "string", + "format": "at-uri" + }, + "byUser": { + "type": "string", + "format": "did", + "description": "Only blobs uploaded by this DID." + }, + "limit": { + "type": "integer", + "minimum": 1, + "maximum": 200, + "default": 50 + }, + "cursor": { + "type": "string" + } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": [ + "blobs" + ], + "properties": { + "blobs": { + "type": "array", + "items": { + "type": "ref", + "ref": "rsvp.atmo.space.defs#blobInfo" + } + }, + "cursor": { + "type": "string" + } + } + } + } + } + } +} diff --git a/lexicons-generated/rsvp/atmo/space/uploadBlob.json b/lexicons-generated/rsvp/atmo/space/uploadBlob.json new file mode 100644 index 0000000..c7517dc --- /dev/null +++ b/lexicons-generated/rsvp/atmo/space/uploadBlob.json @@ -0,0 +1,50 @@ +{ + "lexicon": 1, + "id": "rsvp.atmo.space.uploadBlob", + "defs": { + "main": { + "type": "procedure", + "description": "Upload a blob into a space. Returns a standard atproto BlobRef that records in this space can then reference. The record will be rejected at putRecord time if it references a blob that was not uploaded to the same space.", + "parameters": { + "type": "params", + "required": [ + "spaceUri" + ], + "properties": { + "spaceUri": { + "type": "string", + "format": "at-uri" + } + } + }, + "input": { + "encoding": "*/*" + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": [ + "blob" + ], + "properties": { + "blob": { + "type": "blob" + } + } + } + }, + "errors": [ + { + "name": "Forbidden" + }, + { + "name": "BlobTooLarge" + }, + { + "name": "InvalidMimeType" + } + ] + } + } +} diff --git a/spaces-lexicon-templates/defs.json b/spaces-lexicon-templates/defs.json index 4eaee52..2d5c9cf 100644 --- a/spaces-lexicon-templates/defs.json +++ b/spaces-lexicon-templates/defs.json @@ -47,6 +47,17 @@ "apps": { "type": "array", "items": { "type": "string" } } } }, + "blobInfo": { + "type": "object", + "required": ["cid", "mimeType", "size", "authorDid", "createdAt"], + "properties": { + "cid": { "type": "string", "format": "cid" }, + "mimeType": { "type": "string" }, + "size": { "type": "integer" }, + "authorDid": { "type": "string", "format": "did" }, + "createdAt": { "type": "integer" } + } + }, "inviteView": { "type": "object", "required": ["tokenHash", "spaceUri", "kind", "usedCount", "createdBy", "createdAt"], diff --git a/spaces-lexicon-templates/getBlob.json b/spaces-lexicon-templates/getBlob.json new file mode 100644 index 0000000..b866979 --- /dev/null +++ b/spaces-lexicon-templates/getBlob.json @@ -0,0 +1,27 @@ +{ + "lexicon": 1, + "id": "tools.atmo.space.getBlob", + "defs": { + "main": { + "type": "query", + "description": "Read a blob from a space. Requires read access via a service-auth JWT or a read-grant invite token.", + "parameters": { + "type": "params", + "required": ["spaceUri", "cid"], + "properties": { + "spaceUri": { "type": "string", "format": "at-uri" }, + "cid": { "type": "string", "format": "cid" }, + "inviteToken": { + "type": "string", + "description": "Read-grant invite token for anonymous bearer access." + } + } + }, + "output": { "encoding": "*/*" }, + "errors": [ + { "name": "NotFound" }, + { "name": "Forbidden" } + ] + } + } +} diff --git a/spaces-lexicon-templates/listBlobs.json b/spaces-lexicon-templates/listBlobs.json new file mode 100644 index 0000000..ce7afd0 --- /dev/null +++ b/spaces-lexicon-templates/listBlobs.json @@ -0,0 +1,31 @@ +{ + "lexicon": 1, + "id": "tools.atmo.space.listBlobs", + "defs": { + "main": { + "type": "query", + "description": "List blob metadata for a space. Members only.", + "parameters": { + "type": "params", + "required": ["spaceUri"], + "properties": { + "spaceUri": { "type": "string", "format": "at-uri" }, + "byUser": { "type": "string", "format": "did", "description": "Only blobs uploaded by this DID." }, + "limit": { "type": "integer", "minimum": 1, "maximum": 200, "default": 50 }, + "cursor": { "type": "string" } + } + }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["blobs"], + "properties": { + "blobs": { "type": "array", "items": { "type": "ref", "ref": "tools.atmo.space.defs#blobInfo" } }, + "cursor": { "type": "string" } + } + } + } + } + } +} diff --git a/spaces-lexicon-templates/uploadBlob.json b/spaces-lexicon-templates/uploadBlob.json new file mode 100644 index 0000000..53a68f6 --- /dev/null +++ b/spaces-lexicon-templates/uploadBlob.json @@ -0,0 +1,33 @@ +{ + "lexicon": 1, + "id": "tools.atmo.space.uploadBlob", + "defs": { + "main": { + "type": "procedure", + "description": "Upload a blob into a space. Returns a standard atproto BlobRef that records in this space can then reference. The record will be rejected at putRecord time if it references a blob that was not uploaded to the same space.", + "parameters": { + "type": "params", + "required": ["spaceUri"], + "properties": { + "spaceUri": { "type": "string", "format": "at-uri" } + } + }, + "input": { "encoding": "*/*" }, + "output": { + "encoding": "application/json", + "schema": { + "type": "object", + "required": ["blob"], + "properties": { + "blob": { "type": "blob" } + } + } + }, + "errors": [ + { "name": "Forbidden" }, + { "name": "BlobTooLarge" }, + { "name": "InvalidMimeType" } + ] + } + } +} diff --git a/src/core/spaces/adapter.ts b/src/core/spaces/adapter.ts index 8c839d3..b5f6b3c 100644 --- a/src/core/spaces/adapter.ts +++ b/src/core/spaces/adapter.ts @@ -10,10 +10,13 @@ import { import { getDialect } from "../dialect"; import type { AppPolicy, + BlobMetaRow, CollectionCount, CreateInviteInput, InviteKind, InviteRow, + ListBlobsOptions, + ListBlobsResult, ListOptions, ListResult, ListSpacesOptions, @@ -63,6 +66,17 @@ function mapMemberRow(row: any): SpaceMemberRow { }; } +function mapBlobMetaRow(row: any): BlobMetaRow { + return { + spaceUri: row.space_uri, + cid: row.cid, + mimeType: row.mime_type, + size: Number(row.size), + authorDid: row.author_did, + createdAt: toNum(row.created_at), + }; +} + function mapInviteRow(row: any): InviteRow { return { tokenHash: row.token_hash, @@ -601,4 +615,114 @@ export class HostedAdapter implements StorageAdapter { } return results; } + + async putBlobMeta(row: BlobMetaRow): Promise { + const sql = `INSERT INTO spaces_blobs (space_uri, cid, mime_type, size, author_did, created_at) + VALUES (?, ?, ?, ?, ?, ?) + ON CONFLICT (space_uri, cid) DO NOTHING`; + await this.db + .prepare(sql) + .bind(row.spaceUri, row.cid, row.mimeType, row.size, row.authorDid, row.createdAt) + .run(); + } + + async getBlobMeta(spaceUri: string, cid: string): Promise { + const r = await this.db + .prepare(`SELECT * FROM spaces_blobs WHERE space_uri = ? AND cid = ?`) + .bind(spaceUri, cid) + .first(); + return r ? mapBlobMetaRow(r) : null; + } + + async listBlobMeta( + spaceUri: string, + options: ListBlobsOptions = {} + ): Promise { + const limit = Math.min(options.limit ?? 50, 200); + const clauses: string[] = ["space_uri = ?"]; + const params: any[] = [spaceUri]; + if (options.byUser) { + clauses.push("author_did = ?"); + params.push(options.byUser); + } + if (options.cursor) { + clauses.push("created_at < ?"); + params.push(Number(options.cursor)); + } + const sql = `SELECT * FROM spaces_blobs + WHERE ${clauses.join(" AND ")} + ORDER BY created_at DESC + LIMIT ?`; + params.push(limit + 1); + const { results } = await this.db.prepare(sql).bind(...params).all(); + const blobs = results.map(mapBlobMetaRow); + let cursor: string | undefined; + if (blobs.length > limit) { + const next = blobs.pop()!; + cursor = String(next.createdAt); + } + return { blobs, cursor }; + } + + async deleteBlobMeta(spaceUri: string, cid: string): Promise { + await this.db + .prepare(`DELETE FROM spaces_blobs WHERE space_uri = ? AND cid = ?`) + .bind(spaceUri, cid) + .run(); + } + + async findOrphanBlobs( + spaceUri: string, + cutoff: number, + limit: number + ): Promise { + if (!this.config) return []; + // Gather candidate blobs older than cutoff, then filter out any whose CID + // appears in any record JSON in this space. We use a cheap substring probe + // (LIKE) per collection — false positives are OK because an orphan that + // survives GC just gets collected next cycle; false negatives (deleting + // a referenced blob) would be a bug, and substring search over the full + // CID is safe enough for that. + const { results } = await this.db + .prepare( + `SELECT * FROM spaces_blobs + WHERE space_uri = ? AND created_at < ? + ORDER BY created_at ASC + LIMIT ?` + ) + .bind(spaceUri, cutoff, limit) + .all(); + const candidates = results.map(mapBlobMetaRow); + if (candidates.length === 0) return []; + + const tables: string[] = []; + for (const [short, colConfig] of Object.entries(this.config.collections)) { + if (colConfig.allowInSpaces === false) continue; + tables.push(spacesRecordsTableName(short)); + } + + const orphans: BlobMetaRow[] = []; + for (const blob of candidates) { + let referenced = false; + const pattern = `%${blob.cid}%`; + for (const table of tables) { + try { + const row = await this.db + .prepare( + `SELECT 1 FROM ${table} WHERE space_uri = ? AND record LIKE ? LIMIT 1` + ) + .bind(spaceUri, pattern) + .first(); + if (row) { + referenced = true; + break; + } + } catch { + // table missing — ignore + } + } + if (!referenced) orphans.push(blob); + } + return orphans; + } } diff --git a/src/core/spaces/blob-adapter.ts b/src/core/spaces/blob-adapter.ts new file mode 100644 index 0000000..632b566 --- /dev/null +++ b/src/core/spaces/blob-adapter.ts @@ -0,0 +1,94 @@ +/** + * Bytes-only storage adapter for space blobs. Metadata (CID, mime, size, + * author, space) lives in the `spaces_blobs` table on the main StorageAdapter; + * this interface only moves bytes in and out of a backend (R2, S3, fs, …). + * + * Keys are opaque strings formed by the router as `blobKey(spaceUri, cid)`. + */ + +export interface BlobUploadMeta { + mimeType: string; + size: number; +} + +export interface BlobAdapter { + put(key: string, bytes: Uint8Array, meta: BlobUploadMeta): Promise; + get(key: string): Promise; + /** Bulk delete. Adapters that don't support batch can implement serially. */ + delete(keys: string[]): Promise; +} + +/** In-memory adapter. Useful for tests and local development. */ +export class MemoryBlobAdapter implements BlobAdapter { + private readonly store = new Map(); + + async put(key: string, bytes: Uint8Array): Promise { + this.store.set(key, bytes.slice()); + } + + async get(key: string): Promise { + const v = this.store.get(key); + return v ? v.slice() : null; + } + + async delete(keys: string[]): Promise { + for (const k of keys) this.store.delete(k); + } + + /** Test helper. */ + size(): number { + return this.store.size; + } +} + +/** Minimal Cloudflare R2 bucket shape — matches @cloudflare/workers-types' R2Bucket + * without forcing a types dependency here. */ +export interface R2BucketLike { + put( + key: string, + value: ArrayBuffer | ArrayBufferView | ReadableStream | Blob, + options?: { httpMetadata?: { contentType?: string }; customMetadata?: Record } + ): Promise; + get(key: string): Promise<{ arrayBuffer(): Promise } | null>; + delete(keys: string | string[]): Promise; +} + +/** Cloudflare R2 adapter. Pass the `env.BLOBS` binding from your Worker. */ +export class R2BlobAdapter implements BlobAdapter { + constructor(private readonly bucket: R2BucketLike) {} + + async put(key: string, bytes: Uint8Array, meta: BlobUploadMeta): Promise { + await this.bucket.put(key, bytes, { + httpMetadata: { contentType: meta.mimeType }, + }); + } + + async get(key: string): Promise { + const obj = await this.bucket.get(key); + if (!obj) return null; + const buf = await obj.arrayBuffer(); + return new Uint8Array(buf); + } + + async delete(keys: string[]): Promise { + if (keys.length === 0) return; + await this.bucket.delete(keys); + } +} + +/** Hash a space URI to a short, filesystem/R2-safe key segment. + * Used as the first segment of a blob key so all blobs for one space + * share a common prefix (enables bulk delete on space deletion). */ +export async function spaceKeyPrefix(spaceUri: string): Promise { + const bytes = new TextEncoder().encode(spaceUri); + const digest = await crypto.subtle.digest("SHA-256", bytes); + const hex = Array.from(new Uint8Array(digest), (b) => b.toString(16).padStart(2, "0")).join(""); + return hex.slice(0, 16); +} + +/** Compose an adapter key from a space URI and CID. + * Shape: `<16-hex-chars-of-sha256(spaceUri)>/`. */ +export async function blobKey(spaceUri: string, cid: string): Promise { + const prefix = await spaceKeyPrefix(spaceUri); + return `${prefix}/${cid}`; +} diff --git a/src/core/spaces/blob-gc.ts b/src/core/spaces/blob-gc.ts new file mode 100644 index 0000000..d523fcd --- /dev/null +++ b/src/core/spaces/blob-gc.ts @@ -0,0 +1,38 @@ +import type { BlobAdapter } from "./blob-adapter"; +import { blobKey } from "./blob-adapter"; +import type { StorageAdapter } from "./types"; + +export interface BlobGcOptions { + /** Orphan rows created before this timestamp are eligible for deletion. */ + olderThan: number; + /** Maximum number of blobs to delete in this pass. Defaults to 500. */ + batchSize?: number; +} + +export interface BlobGcResult { + deleted: number; + cids: string[]; +} + +/** Delete blob bytes + metadata for any blob older than `olderThan` that + * is not referenced by any record in the space. Safe to run periodically. */ +export async function gcOrphanBlobs( + storage: StorageAdapter, + blobs: BlobAdapter, + spaceUri: string, + options: BlobGcOptions +): Promise { + const batchSize = options.batchSize ?? 500; + const orphans = await storage.findOrphanBlobs(spaceUri, options.olderThan, batchSize); + if (orphans.length === 0) return { deleted: 0, cids: [] }; + + const keys: string[] = []; + for (const row of orphans) { + keys.push(await blobKey(row.spaceUri, row.cid)); + } + await blobs.delete(keys); + for (const row of orphans) { + await storage.deleteBlobMeta(row.spaceUri, row.cid); + } + return { deleted: orphans.length, cids: orphans.map((o) => o.cid) }; +} diff --git a/src/core/spaces/blob-refs.ts b/src/core/spaces/blob-refs.ts new file mode 100644 index 0000000..62cf693 --- /dev/null +++ b/src/core/spaces/blob-refs.ts @@ -0,0 +1,27 @@ +/** + * Walk a record JSON and collect every atproto blob ref. + * + * Blob refs look like: + * { "$type": "blob", "ref": { "$link": "" }, "mimeType": "...", "size": N } + * + * We return the CID strings. + */ +export function collectBlobCids(value: unknown, out: Set = new Set()): Set { + if (value == null) return out; + if (Array.isArray(value)) { + for (const v of value) collectBlobCids(v, out); + return out; + } + if (typeof value !== "object") return out; + + const obj = value as Record; + if (obj["$type"] === "blob") { + const ref = obj["ref"] as { $link?: unknown } | undefined; + if (ref && typeof ref["$link"] === "string") out.add(ref["$link"]); + // Don't descend — a blob ref's own shape has no nested blobs. + return out; + } + + for (const v of Object.values(obj)) collectBlobCids(v, out); + return out; +} diff --git a/src/core/spaces/router.ts b/src/core/spaces/router.ts index c1c102c..286c81f 100644 --- a/src/core/spaces/router.ts +++ b/src/core/spaces/router.ts @@ -13,7 +13,17 @@ import { import { nextTid } from "./tid"; import { generateInviteToken, hashInviteToken } from "./invite-token"; import { buildSpaceUri } from "./uri"; -import type { InviteKind, InviteRow, SpaceRow, SpacesConfig, StorageAdapter } from "./types"; +import { + DEFAULT_BLOB_MAX_SIZE, + type InviteKind, + type InviteRow, + type SpaceRow, + type SpacesConfig, + type StorageAdapter, +} from "./types"; +import { blobKey } from "./blob-adapter"; +import { collectBlobCids } from "./blob-refs"; +import { create as createCid, toString as cidToString } from "@atcute/cid"; import type { Did } from "@atcute/lexicons"; export interface SpacesRoutesOptions { @@ -230,6 +240,27 @@ export function registerSpacesRoutes( }); if (!result.allow) return c.json({ error: "Forbidden", reason: result.reason }, 403); + // Validate that every blob referenced by this record has already been + // uploaded into this space. This mirrors how PDSes require uploadBlob + // before putRecord, and prevents forging refs to blobs the caller never + // actually claimed. + if (spacesConfig.blobs) { + const cids = collectBlobCids(body.record); + for (const cid of cids) { + const meta = await adapter.getBlobMeta(body.spaceUri, cid); + if (!meta) { + return c.json( + { + error: "InvalidRequest", + reason: "unknown-blob-ref", + message: `Record references blob ${cid} that has not been uploaded to this space.`, + }, + 400 + ); + } + } + } + const rkey = body.rkey ?? nextTid(); const now = Date.now(); await adapter.putRecord({ @@ -270,6 +301,154 @@ export function registerSpacesRoutes( return c.json({ ok: true }); }); + // Blobs (only registered when a blob adapter is configured) + if (spacesConfig.blobs) { + const blobsCfg = spacesConfig.blobs; + const blobAdapter = blobsCfg.adapter; + const maxSize = blobsCfg.maxSize ?? DEFAULT_BLOB_MAX_SIZE; + const accept = blobsCfg.accept; + + app.post(`/xrpc/${SPACE}.uploadBlob`, auth, async (c) => { + const sa = getAuth(c); + const spaceUri = c.req.query("spaceUri"); + if (!spaceUri) { + return c.json({ error: "InvalidRequest", message: "spaceUri required" }, 400); + } + const space = await adapter.getSpace(spaceUri); + if (!space) return c.json({ error: "NotFound" }, 404); + + const member = await adapter.getMember(spaceUri, sa.issuer); + const aclResult = checkAccess({ + op: "write", + space, + callerDid: sa.issuer, + member, + clientId: sa.clientId, + }); + if (!aclResult.allow) { + return c.json({ error: "Forbidden", reason: aclResult.reason }, 403); + } + + const mimeType = c.req.header("content-type") ?? "application/octet-stream"; + if (accept && !accept.includes(mimeType)) { + return c.json( + { error: "InvalidMimeType", message: `MIME type ${mimeType} is not accepted.` }, + 400 + ); + } + + const declaredLen = c.req.header("content-length"); + if (declaredLen && Number(declaredLen) > maxSize) { + return c.json( + { error: "BlobTooLarge", message: `Blob exceeds max size of ${maxSize} bytes.` }, + 413 + ); + } + + const buf = await c.req.arrayBuffer(); + const bytes = new Uint8Array(buf); + if (bytes.byteLength > maxSize) { + return c.json( + { error: "BlobTooLarge", message: `Blob exceeds max size of ${maxSize} bytes.` }, + 413 + ); + } + + const cid = await createCid(0x55, bytes); + const cidString = cidToString(cid); + const key = await blobKey(spaceUri, cidString); + + await blobAdapter.put(key, bytes, { mimeType, size: bytes.byteLength }); + await adapter.putBlobMeta({ + spaceUri, + cid: cidString, + mimeType, + size: bytes.byteLength, + authorDid: sa.issuer, + createdAt: Date.now(), + }); + + return c.json({ + blob: { + $type: "blob", + ref: { $link: cidString }, + mimeType, + size: bytes.byteLength, + }, + }); + }); + + app.get(`/xrpc/${SPACE}.getBlob`, readAuth, async (c) => { + const spaceUri = c.req.query("spaceUri"); + const cid = c.req.query("cid"); + if (!spaceUri || !cid) { + return c.json({ error: "InvalidRequest", message: "spaceUri and cid required" }, 400); + } + const space = await adapter.getSpace(spaceUri); + if (!space) return c.json({ error: "NotFound" }, 404); + + const authz = await authorizeRead(c, spaceUri); + if (authz instanceof Response) return authz; + + if (authz.via === "jwt") { + const sa = authz.sa; + const member = await adapter.getMember(spaceUri, sa.issuer); + const aclResult = checkAccess({ + op: "read", + space, + callerDid: sa.issuer, + member, + clientId: sa.clientId, + }); + if (!aclResult.allow) { + return c.json({ error: "Forbidden", reason: aclResult.reason }, 403); + } + } + + const meta = await adapter.getBlobMeta(spaceUri, cid); + if (!meta) return c.json({ error: "NotFound" }, 404); + const key = await blobKey(spaceUri, cid); + const bytes = await blobAdapter.get(key); + if (!bytes) return c.json({ error: "NotFound" }, 404); + + return new Response(bytes, { + headers: { + "content-type": meta.mimeType, + "content-length": String(meta.size), + }, + }); + }); + + app.get(`/xrpc/${SPACE}.listBlobs`, auth, async (c) => { + const sa = getAuth(c); + const spaceUri = c.req.query("spaceUri"); + if (!spaceUri) { + return c.json({ error: "InvalidRequest", message: "spaceUri required" }, 400); + } + const space = await adapter.getSpace(spaceUri); + if (!space) return c.json({ error: "NotFound" }, 404); + + const member = await adapter.getMember(spaceUri, sa.issuer); + const aclResult = checkAccess({ + op: "read", + space, + callerDid: sa.issuer, + member, + clientId: sa.clientId, + }); + if (!aclResult.allow) { + return c.json({ error: "Forbidden", reason: aclResult.reason }, 403); + } + + const result = await adapter.listBlobMeta(spaceUri, { + byUser: c.req.query("byUser") ?? undefined, + cursor: c.req.query("cursor") ?? undefined, + limit: c.req.query("limit") ? Number(c.req.query("limit")) : undefined, + }); + return c.json(result); + }); + } + // Space management (owner-gated) app.post(`/xrpc/${SPACE}.createSpace`, auth, async (c) => { const sa = getAuth(c); diff --git a/src/core/spaces/schema.ts b/src/core/spaces/schema.ts index d6cf4a1..a08ada6 100644 --- a/src/core/spaces/schema.ts +++ b/src/core/spaces/schema.ts @@ -34,6 +34,18 @@ export function buildSpacesBaseSchema(dialect: SqlDialect): string[] { )`, `CREATE INDEX IF NOT EXISTS idx_spaces_members_did ON spaces_members(did)`, + `CREATE TABLE IF NOT EXISTS spaces_blobs ( + space_uri TEXT NOT NULL, + cid TEXT NOT NULL, + mime_type TEXT NOT NULL, + size INTEGER NOT NULL, + author_did TEXT NOT NULL, + created_at ${dialect.bigintType} NOT NULL, + PRIMARY KEY (space_uri, cid) + )`, + `CREATE INDEX IF NOT EXISTS idx_spaces_blobs_author ON spaces_blobs(space_uri, author_did)`, + `CREATE INDEX IF NOT EXISTS idx_spaces_blobs_created ON spaces_blobs(space_uri, created_at)`, + `CREATE TABLE IF NOT EXISTS spaces_invites ( token_hash TEXT PRIMARY KEY, space_uri TEXT NOT NULL, diff --git a/src/core/spaces/types.ts b/src/core/spaces/types.ts index ff5584f..397e6a8 100644 --- a/src/core/spaces/types.ts +++ b/src/core/spaces/types.ts @@ -1,5 +1,6 @@ import type { Database } from "../types"; import type { DidDocumentResolver } from "@atcute/identity-resolver"; +import type { BlobAdapter } from "./blob-adapter"; export type AppPolicyMode = "allow" | "deny"; @@ -8,6 +9,22 @@ export interface AppPolicy { apps: string[]; } +export interface SpacesBlobsConfig { + /** Bytes backend (R2, S3, in-memory, …). */ + adapter: BlobAdapter; + /** Max blob size in bytes. Defaults to 2 MiB. */ + maxSize?: number; + /** MIME allowlist. If set, only these content types are accepted. */ + accept?: string[]; + /** Orphan blobs (those with no referencing record) are kept this long before + * GC can delete them, to allow upload-then-putRecord flows. + * Defaults to 24 hours. */ + gcOrphanAfterMs?: number; +} + +export const DEFAULT_BLOB_MAX_SIZE = 2 * 1024 * 1024; +export const DEFAULT_BLOB_GC_ORPHAN_AFTER_MS = 24 * 60 * 60 * 1000; + export interface SpacesConfig { /** NSID that identifies the kind of space this service hosts, e.g. "tools.atmo.event.space". */ type: string; @@ -18,6 +35,8 @@ export interface SpacesConfig { /** DID document resolver for service-auth JWT verification. * Defaults to a composite PLC + did:web resolver if omitted. */ resolver?: DidDocumentResolver; + /** Blob-upload backend. When omitted, blob XRPCs are not exposed. */ + blobs?: SpacesBlobsConfig; } export interface SpaceRow { @@ -106,6 +125,26 @@ export interface RedeemInviteResult { spaceUri: string; } +export interface BlobMetaRow { + spaceUri: string; + cid: string; + mimeType: string; + size: number; + authorDid: string; + createdAt: number; +} + +export interface ListBlobsOptions { + byUser?: string; + cursor?: string; + limit?: number; +} + +export interface ListBlobsResult { + blobs: BlobMetaRow[]; + cursor?: string; +} + export interface StorageAdapter { // Space lifecycle createSpace(space: Omit): Promise; @@ -144,6 +183,15 @@ export interface StorageAdapter { listRecords(spaceUri: string, collection: string, options?: ListOptions): Promise; deleteRecord(spaceUri: string, collection: string, authorDid: string, rkey: string): Promise; listCollections(spaceUri: string, options?: { byUser?: string }): Promise; + + // Blobs (metadata only; bytes live on BlobAdapter) + putBlobMeta(row: BlobMetaRow): Promise; + getBlobMeta(spaceUri: string, cid: string): Promise; + listBlobMeta(spaceUri: string, options?: ListBlobsOptions): Promise; + deleteBlobMeta(spaceUri: string, cid: string): Promise; + /** Find blob rows older than `cutoff` whose CIDs are not referenced in any + * record JSON in this space. Capped at `limit` to bound a single GC pass. */ + findOrphanBlobs(spaceUri: string, cutoff: number, limit: number): Promise; } export interface AdapterContext { diff --git a/tests/spaces-blobs.test.ts b/tests/spaces-blobs.test.ts new file mode 100644 index 0000000..63af7a5 --- /dev/null +++ b/tests/spaces-blobs.test.ts @@ -0,0 +1,325 @@ +import { describe, it, expect, beforeAll } from "vitest"; +import { Hono } from "hono"; +import type { MiddlewareHandler } from "hono"; +import { createSqliteDatabase } from "../src/adapters/sqlite"; +import { initSchema } from "../src/core/db/schema"; +import { createApp } from "../src/core/router"; +import { resolveConfig } from "../src/core/types"; +import type { ContrailConfig } from "../src/core/types"; +import { MemoryBlobAdapter } from "../src/core/spaces/blob-adapter"; +import { HostedAdapter } from "../src/core/spaces/adapter"; +import { gcOrphanBlobs } from "../src/core/spaces/blob-gc"; + +const ALICE = "did:plc:alice"; +const BOB = "did:plc:bob"; +const CHARLIE = "did:plc:charlie"; + +function makeConfig(blobs: MemoryBlobAdapter, maxSize = 2 * 1024 * 1024): ContrailConfig { + return { + namespace: "test.blobs", + collections: { + photo: { collection: "app.event.photo" }, + }, + spaces: { + type: "tools.atmo.event.space", + serviceDid: "did:web:test.example#svc", + blobs: { adapter: blobs, maxSize }, + }, + }; +} + +function fakeAuth(aud: string): MiddlewareHandler { + return async (c, next) => { + const did = c.req.header("X-Test-Did"); + if (!did) return c.json({ error: "AuthRequired" }, 401); + c.set("serviceAuth", { + issuer: did, + audience: aud, + lxm: undefined, + clientId: c.req.header("X-Test-App") ?? undefined, + }); + await next(); + }; +} + +async function makeApp( + blobs: MemoryBlobAdapter, + maxSize = 2 * 1024 * 1024 +): Promise<{ app: Hono; db: any; config: ReturnType }> { + const db = createSqliteDatabase(":memory:"); + const cfg = makeConfig(blobs, maxSize); + const resolved = resolveConfig(cfg); + await initSchema(db, resolved); + const app = createApp(db, resolved, { + spaces: { authMiddleware: fakeAuth(cfg.spaces!.serviceDid) }, + }); + return { app, db, config: resolved }; +} + +function call( + app: Hono, + method: string, + path: string, + did: string, + body?: BodyInit, + contentType?: string +): Promise { + const headers: Record = { "X-Test-Did": did }; + if (contentType) headers["Content-Type"] = contentType; + return app.fetch( + new Request(`http://localhost${path}`, { + method, + headers, + body, + }) + ); +} + +async function callJson( + app: Hono, + method: string, + path: string, + did: string, + body?: any +): Promise { + return call( + app, + method, + path, + did, + body === undefined ? undefined : JSON.stringify(body), + body === undefined ? undefined : "application/json" + ); +} + +async function createSpace(app: Hono, owner: string, key: string): Promise { + const res = await callJson(app, "POST", "/xrpc/test.blobs.space.createSpace", owner, { key }); + expect(res.status).toBe(200); + const body = (await res.json()) as any; + return body.space.uri; +} + +describe("spaces blobs", () => { + let app: Hono; + let blobs: MemoryBlobAdapter; + let spaceUri: string; + + beforeAll(async () => { + blobs = new MemoryBlobAdapter(); + const out = await makeApp(blobs); + app = out.app; + spaceUri = await createSpace(app, ALICE, "album"); + }); + + it("owner uploads a blob and gets back a valid BlobRef", async () => { + const bytes = new TextEncoder().encode("hello world"); + const res = await call( + app, + "POST", + `/xrpc/test.blobs.space.uploadBlob?spaceUri=${encodeURIComponent(spaceUri)}`, + ALICE, + bytes, + "image/png" + ); + expect(res.status).toBe(200); + const body = (await res.json()) as any; + expect(body.blob.$type).toBe("blob"); + expect(body.blob.mimeType).toBe("image/png"); + expect(body.blob.size).toBe(bytes.byteLength); + expect(typeof body.blob.ref.$link).toBe("string"); + expect(body.blob.ref.$link.startsWith("b")).toBe(true); + }); + + it("non-member cannot upload", async () => { + const bytes = new TextEncoder().encode("intruder"); + const res = await call( + app, + "POST", + `/xrpc/test.blobs.space.uploadBlob?spaceUri=${encodeURIComponent(spaceUri)}`, + BOB, + bytes, + "image/png" + ); + expect(res.status).toBe(403); + }); + + it("upload is rejected when bytes exceed maxSize", async () => { + const tiny = new MemoryBlobAdapter(); + const { app: smallApp } = await makeApp(tiny, 64); + const uri = await createSpace(smallApp, ALICE, "small"); + const bytes = new Uint8Array(128); + const res = await call( + smallApp, + "POST", + `/xrpc/test.blobs.space.uploadBlob?spaceUri=${encodeURIComponent(uri)}`, + ALICE, + bytes, + "application/octet-stream" + ); + expect(res.status).toBe(413); + }); + + it("getBlob returns bytes for a member, 403 for non-member, 404 for bogus cid", async () => { + const payload = new TextEncoder().encode("download me"); + const up = await call( + app, + "POST", + `/xrpc/test.blobs.space.uploadBlob?spaceUri=${encodeURIComponent(spaceUri)}`, + ALICE, + payload, + "text/plain" + ); + const { blob } = (await up.json()) as any; + const cid = blob.ref.$link; + + // Member (owner) — gets bytes back. + const okRes = await call( + app, + "GET", + `/xrpc/test.blobs.space.getBlob?spaceUri=${encodeURIComponent(spaceUri)}&cid=${cid}`, + ALICE + ); + expect(okRes.status).toBe(200); + expect(okRes.headers.get("content-type")).toBe("text/plain"); + const got = new Uint8Array(await okRes.arrayBuffer()); + expect(new TextDecoder().decode(got)).toBe("download me"); + + // Non-member — 403. + const deny = await call( + app, + "GET", + `/xrpc/test.blobs.space.getBlob?spaceUri=${encodeURIComponent(spaceUri)}&cid=${cid}`, + CHARLIE + ); + expect(deny.status).toBe(403); + + // Made-up CID — 404 even for the owner (never enumerate). + const miss = await call( + app, + "GET", + `/xrpc/test.blobs.space.getBlob?spaceUri=${encodeURIComponent(spaceUri)}&cid=bafynotareal`, + ALICE + ); + expect(miss.status).toBe(404); + }); + + it("putRecord rejects records referencing blobs that were not uploaded here", async () => { + const res = await callJson(app, "POST", "/xrpc/test.blobs.space.putRecord", ALICE, { + spaceUri, + collection: "app.event.photo", + record: { + caption: "bogus", + image: { + $type: "blob", + ref: { $link: "bafynotareal" }, + mimeType: "image/png", + size: 10, + }, + }, + }); + expect(res.status).toBe(400); + const body = (await res.json()) as any; + expect(body.reason).toBe("unknown-blob-ref"); + }); + + it("putRecord accepts records referencing a previously uploaded blob", async () => { + const bytes = new TextEncoder().encode("real image bytes"); + const up = await call( + app, + "POST", + `/xrpc/test.blobs.space.uploadBlob?spaceUri=${encodeURIComponent(spaceUri)}`, + ALICE, + bytes, + "image/png" + ); + const { blob } = (await up.json()) as any; + + const res = await callJson(app, "POST", "/xrpc/test.blobs.space.putRecord", ALICE, { + spaceUri, + collection: "app.event.photo", + record: { caption: "ok", image: blob }, + }); + expect(res.status).toBe(200); + }); + + it("listBlobs returns metadata for members", async () => { + const res = await call( + app, + "GET", + `/xrpc/test.blobs.space.listBlobs?spaceUri=${encodeURIComponent(spaceUri)}`, + ALICE + ); + expect(res.status).toBe(200); + const body = (await res.json()) as any; + expect(Array.isArray(body.blobs)).toBe(true); + expect(body.blobs.length).toBeGreaterThan(0); + const row = body.blobs[0]; + expect(row).toHaveProperty("cid"); + expect(row).toHaveProperty("mimeType"); + expect(row).toHaveProperty("size"); + expect(row).toHaveProperty("authorDid"); + }); + + it("gcOrphanBlobs deletes unreferenced blobs older than cutoff but keeps referenced ones", async () => { + const isolatedBlobs = new MemoryBlobAdapter(); + const { app: isoApp, db } = await makeApp(isolatedBlobs); + const uri = await createSpace(isoApp, ALICE, "gc-test"); + const storage = new HostedAdapter(db, makeConfig(isolatedBlobs)); + + // Upload two blobs; reference only one of them. + const keepBytes = new TextEncoder().encode("keep me"); + const dropBytes = new TextEncoder().encode("drop me"); + + const keepUp = await call( + isoApp, + "POST", + `/xrpc/test.blobs.space.uploadBlob?spaceUri=${encodeURIComponent(uri)}`, + ALICE, + keepBytes, + "image/png" + ); + const dropUp = await call( + isoApp, + "POST", + `/xrpc/test.blobs.space.uploadBlob?spaceUri=${encodeURIComponent(uri)}`, + ALICE, + dropBytes, + "image/png" + ); + const keep = (await keepUp.json()) as any; + const drop = (await dropUp.json()) as any; + + // Reference the "keep" blob from a record. + const put = await callJson(isoApp, "POST", "/xrpc/test.blobs.space.putRecord", ALICE, { + spaceUri: uri, + collection: "app.event.photo", + record: { caption: "keep", image: keep.blob }, + }); + expect(put.status).toBe(200); + + // Run GC with a cutoff far in the future so every blob is eligible by age. + const result = await gcOrphanBlobs(storage, isolatedBlobs, uri, { + olderThan: Date.now() + 60_000, + }); + + expect(result.deleted).toBe(1); + expect(result.cids).toContain(drop.blob.ref.$link); + expect(result.cids).not.toContain(keep.blob.ref.$link); + + // The "keep" blob is still retrievable; the dropped one is gone. + const okRes = await call( + isoApp, + "GET", + `/xrpc/test.blobs.space.getBlob?spaceUri=${encodeURIComponent(uri)}&cid=${keep.blob.ref.$link}`, + ALICE + ); + expect(okRes.status).toBe(200); + const missRes = await call( + isoApp, + "GET", + `/xrpc/test.blobs.space.getBlob?spaceUri=${encodeURIComponent(uri)}&cid=${drop.blob.ref.$link}`, + ALICE + ); + expect(missRes.status).toBe(404); + }); +});