diff --git a/dev/relay-broker.mjs b/dev/relay-broker.mjs index 0bf76662a..317c681d6 100644 --- a/dev/relay-broker.mjs +++ b/dev/relay-broker.mjs @@ -625,6 +625,7 @@ export function relayBrokerPlugin({ ...((await getAuthority(relay)).channelCreation ? [9007] : []), ], workflowReads: true, + attachmentUploads: true, sidebarPreferences: true, readState: true, agentLibrary: true, diff --git a/src/features/relay/attachments.test.ts b/src/features/relay/attachments.test.ts new file mode 100644 index 000000000..73ffe2135 --- /dev/null +++ b/src/features/relay/attachments.test.ts @@ -0,0 +1,297 @@ +import { assert, afterEach, expect, it, vi } from "vitest"; +import { + attachmentMessage, + brokerUpload, + UPLOAD_MAX_BYTES, + validateUploadResult, +} from "./attachments"; +import { foldMessages, parseAttachments } from "./fold"; + +const origin = "https://relay.test"; +const hash = "a".repeat(64); +const descriptor = { + url: `${origin}/media/${hash}.pdf`, + type: "application/pdf", + size: 3, + sha256: hash, +}; +const file = new File([new Uint8Array([0, 128, 255])], "report.pdf", { + type: "application/pdf", +}); +afterEach(() => vi.unstubAllGlobals()); + +it("uploads exact bytes through the selected broker and returns frozen named metadata", async () => { + const fetcher = vi.fn(async (_url, options) => { + expect(new Uint8Array(await options.body.arrayBuffer())).toEqual( + new Uint8Array([0, 128, 255]), + ); + return Response.json(descriptor); + }); + vi.stubGlobal("fetch", fetcher); + const result = await brokerUpload("/api/relay/selected", origin)( + file, + new AbortController().signal, + ); + expect(fetcher).toHaveBeenCalledWith( + "/api/relay/selected/upload", + expect.objectContaining({ + method: "POST", + credentials: "same-origin", + body: file, + }), + ); + expect(result).toEqual({ ...descriptor, name: "report.pdf" }); + expect(Object.isFrozen(result)).toBe(true); +}); + +it.each([ + { url: `https://other.test/media/${hash}.pdf` }, + { url: `${origin}/media/${hash}.pdf?token=x` }, + { url: `${origin}/media/${hash}.pdf#fragment` }, + { url: `https://user@relay.test/media/${hash}.pdf` }, + { url: `${origin}/media/${"b".repeat(64)}.pdf` }, + { url: `${origin}/other/${hash}.pdf` }, + { type: "text/html; extra" }, + { size: 4 }, + { sha256: "invalid" }, +])("rejects an invalid descriptor independently: %j", (patch) => { + expect(() => + validateUploadResult({ ...descriptor, ...patch }, origin, 3, "file"), + ).toThrow(/invalid upload/); +}); + +it.each([0, UPLOAD_MAX_BYTES + 1])( + "rejects size %s before network work", + async (size) => { + const fetcher = vi.fn(); + vi.stubGlobal("fetch", fetcher); + await expect( + brokerUpload("/api/relay", origin)( + new File([new Uint8Array(size)], "file"), + new AbortController().signal, + ), + ).rejects.toMatchObject({ code: "size" }); + expect(fetcher).not.toHaveBeenCalled(); + }, +); + +it.each(["metadata", "denied", "capacity", "rejected", "size", "failed"])( + "preserves broker failure %s", + async (code) => { + vi.stubGlobal("fetch", async () => + Response.json({ code }, { status: 400 }), + ); + await expect( + brokerUpload("/api/relay", origin)(file, new AbortController().signal), + ).rejects.toMatchObject({ code }); + }, +); + +it("rejects oversized/malformed bodies and cancels response consumption", async () => { + const cancel = vi.fn(); + vi.stubGlobal( + "fetch", + async () => + new Response( + new ReadableStream({ + start(controller) { + controller.enqueue(new Uint8Array(8193)); + }, + cancel, + }), + ), + ); + await expect( + brokerUpload("/api/relay", origin)(file, new AbortController().signal), + ).rejects.toMatchObject({ code: "invalid" }); + expect(cancel).toHaveBeenCalledTimes(1); + vi.stubGlobal("fetch", async () => new Response("not json")); + await expect( + brokerUpload("/api/relay", origin)(file, new AbortController().signal), + ).rejects.toMatchObject({ code: "invalid" }); +}); + +it("passes cancellation to fetch and rejects late completion", async () => { + const controller = new AbortController(); + vi.stubGlobal( + "fetch", + async (_url: RequestInfo | URL, init?: RequestInit) => { + assert.exists(init?.signal); + controller.abort(); + expect(init.signal.aborted).toBe(true); + return Response.json(descriptor); + }, + ); + await expect( + brokerUpload("/api/relay", origin)(file, controller.signal), + ).rejects.toMatchObject({ name: "AbortError" }); + const fetcher = vi.fn(); + vi.stubGlobal("fetch", fetcher); + await expect( + brokerUpload("/api/relay", origin)(file, controller.signal), + ).rejects.toMatchObject({ name: "AbortError" }); + expect(fetcher).not.toHaveBeenCalled(); +}); + +it("builds escaped message links and metadata understood by the existing receiver", () => { + const uploaded = { ...descriptor, name: "[report](x)\n!.pdf" }; + const result = attachmentMessage(" hello ", [uploaded], origin); + expect(result.content).toBe( + `hello\n\n[\\[report\\]\\(x\\) \\!.pdf](<${descriptor.url}>)`, + ); + expect(result.tags).toEqual([ + [ + "imeta", + `url ${descriptor.url}`, + "m application/pdf", + "size 3", + `x ${hash}`, + "filename [report](x) !.pdf", + ], + ]); + const received = parseAttachments( + { id: "", pubkey: "", kind: 9, created_at: 1, ...result }, + [], + ); + expect(received).toMatchObject([ + { + kind: "file", + name: "[report](x) !.pdf", + size: 3, + mime: "application/pdf", + }, + ]); + expect(() => + attachmentMessage("", [uploaded], "https://different.test"), + ).toThrow(); +}); + +it.each(["report ©.pdf", "report A.pdf", "report A.pdf"])( + "preserves the literal filename %s through message projection", + (name) => { + const message = attachmentMessage( + "hello", + [{ ...descriptor, name }], + origin, + ); + const [received] = foldMessages("c", "d".repeat(64), [ + { + id: "b".repeat(64), + pubkey: "c".repeat(64), + created_at: 1, + kind: 9, + content: message.content, + tags: [["h", "c"], ...message.tags], + }, + ]); + expect(received?.content).toBe("hello"); + expect(received?.attachments).toEqual([ + { + url: descriptor.url, + kind: "file", + mime: descriptor.type, + size: 3, + name, + }, + ]); + }, +); + +it("keeps ordinary text unchanged and marks prepared image/video links as media", () => { + expect(attachmentMessage(" hello ", [])).toEqual({ + content: "hello", + tags: [], + }); + for (const type of ["image/png", "video/mp4"]) { + const result = attachmentMessage( + "", + [{ ...descriptor, name: "media", type }], + origin, + ); + expect(result.content).toBe(`![media](<${descriptor.url}>)`); + } +}); + +it.each([ + [401, "denied"], + [403, "denied"], + [413, "size"], + [429, "capacity"], +])( + "preserves HTTP %s without requiring an error body", + async (status, code) => { + vi.stubGlobal( + "fetch", + async () => new Response(null, { status: Number(status) }), + ); + await expect( + brokerUpload("/api/relay", origin)(file, new AbortController().signal), + ).rejects.toMatchObject({ code }); + }, +); + +it("bounds a stalled request and rejects its late result after the upload deadline", async () => { + vi.useFakeTimers(); + const deadline = new AbortController(); + const timeout = vi + .spyOn(AbortSignal, "timeout") + .mockReturnValue(deadline.signal); + let release!: (response: Response) => void; + let signal!: AbortSignal; + vi.stubGlobal("fetch", (_url: RequestInfo | URL, init?: RequestInit) => { + assert.exists(init?.signal); + signal = init.signal; + return new Promise((resolve) => { + release = resolve; + }); + }); + try { + const pending = brokerUpload("/api/relay", origin)( + file, + new AbortController().signal, + ); + const rejected = expect(pending).rejects.toMatchObject({ + name: "TimeoutError", + }); + expect(timeout).toHaveBeenCalledWith(120_000); + deadline.abort(new DOMException("Upload deadline", "TimeoutError")); + expect(signal.aborted).toBe(true); + release(Response.json(descriptor)); + await rejected; + } finally { + timeout.mockRestore(); + vi.useRealTimers(); + } +}); + +it("keeps filename metadata within the relay basename and UTF-8 byte rules", () => { + for (const name of [ + "folder/name\\other\u0085.pdf", + "📎".repeat(100), + " /\\\u0085 ", + ]) { + const result = attachmentMessage("", [{ ...descriptor, name }], origin); + const field = result.tags[0]?.find((item) => item.startsWith("filename ")); + assert.exists(field); + const value = field.slice("filename ".length); + expect(new TextEncoder().encode(value).length).toBeLessThanOrEqual(255); + expect(value.length).toBeGreaterThan(0); + expect(value).not.toContain("/"); + expect(value).not.toContain("\\"); + expect(value).not.toContain("\u0085"); + } +}); + +it("removes deceptive bidi controls from both metadata and the sent Markdown label", () => { + const controls = + "\u200e\u200f\u202a\u202b\u202c\u202d\u202e\u2066\u2067\u2068\u2069"; + const result = attachmentMessage( + "", + [{ ...descriptor, name: `report${controls}.pdf` }], + origin, + ); + for (const char of controls) { + expect(result.content).not.toContain(char); + expect(result.tags.flat().join(" ")).not.toContain(char); + } +}); diff --git a/src/features/relay/attachments.ts b/src/features/relay/attachments.ts new file mode 100644 index 000000000..e38ab65e9 --- /dev/null +++ b/src/features/relay/attachments.ts @@ -0,0 +1,196 @@ +import { safeAttachmentName } from "./message-content"; + +export const UPLOAD_MAX_BYTES = 20 * 1024 * 1024; +export const UPLOAD_TIMEOUT_MS = 120_000; +export type UploadedAttachment = Readonly<{ + name: string; + url: string; + type: string; + size: number; + sha256: string; +}>; +export type AttachmentUpload = ( + file: File, + signal: AbortSignal, +) => Promise; +export const UPLOAD_FAILURES = { + unavailable: "Uploads are unavailable on this connection.", + size: "Choose a file between 1 byte and 20 MB.", + capacity: "Uploads are busy. Retry this file shortly.", + metadata: + "This file needs metadata cleanup before it can be uploaded. Choose an exported copy without metadata, or remove it for now.", + rejected: + "The server could not accept this file. Its format or metadata may not be supported.", + denied: "You don’t have permission to upload here.", + failed: "Upload did not finish. Retry or remove this file.", + invalid: "The server returned an invalid upload result. Retry this file.", + cancelled: "Upload cancelled. Retry to upload this file.", +} as const; +export type UploadCode = keyof typeof UPLOAD_FAILURES; +export class UploadError extends Error { + constructor(readonly code: UploadCode) { + super(UPLOAD_FAILURES[code]); + } +} +export function validateUploadResult( + value: unknown, + origin: string, + size: number, + name: string, +): UploadedAttachment { + const v = value as Partial | null; + if ( + !v || + typeof v.url !== "string" || + typeof v.type !== "string" || + !/^[a-z0-9.+-]+\/[a-z0-9.+-]+$/i.test(v.type) || + v.size !== size || + typeof v.sha256 !== "string" || + !/^[0-9a-f]{64}$/.test(v.sha256) || + !Number.isSafeInteger(size) || + size < 1 || + size > UPLOAD_MAX_BYTES + ) + throw new UploadError("invalid"); + let url: URL; + try { + url = new URL(v.url, origin); + } catch { + throw new UploadError("invalid"); + } + if ( + url.origin !== new URL(origin).origin || + url.username || + url.password || + url.search || + url.hash || + !new RegExp(`^/media/${v.sha256}(?:\\.[a-z0-9]{1,8})?$`).test(url.pathname) + ) + throw new UploadError("invalid"); + return Object.freeze({ + name, + url: url.href, + type: v.type.toLowerCase(), + size, + sha256: v.sha256, + }); +} +function attachmentMarkdown(name: string, result: UploadedAttachment) { + const label = Array.from(name, (char) => + char.charCodeAt(0) < 32 ? " " : char, + ) + .join("") + .replace(/[\\`*_{}[\]()!<>&]/g, "\\$&"); + const image = /^(image\/(png|jpeg|gif|webp)|video\/mp4)$/.test(result.type); + return `${image ? "!" : ""}[${label}](<${result.url}>)`; +} +export function brokerUpload( + endpoint: string, + origin: string, +): AttachmentUpload { + return async (file, signal) => { + signal.throwIfAborted(); + if (!file.size || file.size > UPLOAD_MAX_BYTES) + throw new UploadError("size"); + const bounded = AbortSignal.any([ + signal, + AbortSignal.timeout(UPLOAD_TIMEOUT_MS), + ]); + const response = await fetch(`${endpoint}/upload`, { + method: "POST", + credentials: "same-origin", + headers: { "Content-Type": file.type || "application/octet-stream" }, + body: file, + signal: bounded, + }); + if ([401, 403, 413, 429].includes(response.status)) { + try { + await response.body?.cancel(); + } catch { + /* Keep the known failure. */ + } + throw new UploadError( + response.status === 413 + ? "size" + : response.status === 429 + ? "capacity" + : "denied", + ); + } + let text = ""; + const reader = response.body?.getReader(); + if (!reader) throw new UploadError("invalid"); + try { + const decoder = new TextDecoder(); + let size = 0; + while (true) { + const { done, value } = await reader.read(); + if (done) break; + size += value.byteLength; + if (size > 8192) throw new UploadError("invalid"); + text += decoder.decode(value, { stream: true }); + } + text += decoder.decode(); + } finally { + try { + await reader.cancel(); + } catch { + /* Do not replace a parsing failure. */ + } + } + let body: unknown; + try { + body = JSON.parse(text); + } catch { + throw new UploadError("invalid"); + } + if (!response.ok) { + const code = (body as { code?: string })?.code; + throw new UploadError( + code && Object.hasOwn(UPLOAD_FAILURES, code) + ? (code as UploadCode) + : "failed", + ); + } + bounded.throwIfAborted(); + return validateUploadResult(body, origin, file.size, file.name); + }; +} + +/** Only completed uploads enter the existing message/outbox contract. */ +export function attachmentMessage( + content: string, + attachments: readonly UploadedAttachment[], + origin?: string, +) { + const tags: string[][] = []; + const links = attachments.map((item) => { + if (!origin || typeof item.name !== "string") + throw new UploadError("invalid"); + const result = validateUploadResult(item, origin, item.size, item.name); + // Relay filename metadata is a basename capped at 255 UTF-8 bytes. + let name = ""; + let bytes = 0; + for (const char of item.name) { + const clean = + !safeAttachmentName(char) || char === "/" || char === "\\" ? " " : char; + bytes += new TextEncoder().encode(clean).length; + if (bytes > 255) break; + name += clean; + } + name = name.trim() || "File"; + tags.push([ + "imeta", + `url ${result.url}`, + `m ${result.type}`, + `size ${result.size}`, + `x ${result.sha256}`, + `filename ${name}`, + ]); + return attachmentMarkdown(name, result); + }); + return { + content: [content.trim(), ...links].filter(Boolean).join("\n\n"), + tags, + }; +} diff --git a/src/features/relay/broker.integration.test.ts b/src/features/relay/broker.integration.test.ts index 230a5c764..1f11bfb60 100644 --- a/src/features/relay/broker.integration.test.ts +++ b/src/features/relay/broker.integration.test.ts @@ -1,3 +1,4 @@ +import { createHash } from "node:crypto"; import { brokerSocket, openBrokerSocket, @@ -26,6 +27,9 @@ it("profiles a first slow publish through real local IPC, signing, authenticated }); const events = new Map(); let published = 0; + let uploaded = 0; + const attachmentBytes = new Uint8Array([0, 128, 255]); + const hash = createHash("sha256").update(attachmentBytes).digest("hex"); const socket = brokerSocket(async (event: { id: string }) => { published++; if (published === 1) await delayed; @@ -33,6 +37,17 @@ it("profiles a first slow publish through real local IPC, signing, authenticated return ""; }); const upstream: typeof fetch = async (_input, init) => { + if (String(_input).endsWith("/upload")) { + uploaded++; + expect(new Uint8Array(init?.body as Buffer)).toEqual(attachmentBytes); + expect(init?.method).toBe("PUT"); + return Response.json({ + url: `${fixtureRelayUrl}/media/${hash}.pdf`, + type: "application/pdf", + size: 3, + sha256: hash, + }); + } const body = JSON.parse(init?.body as string); return Response.json( (body[0]?.ids ?? []).flatMap((id: string) => events.get(id) ?? []), @@ -75,12 +90,20 @@ it("profiles a first slow publish through real local IPC, signing, authenticated try { const outbox = owner.session.outbox; assert.exists(outbox); - const first = owner.session.messages.send("c", "first"); + assert.exists(owner.session.attachments); + const attachment = await owner.session.attachments.upload( + new File([attachmentBytes], "report.pdf", { type: "application/pdf" }), + "c", + new AbortController().signal, + ); + expect(attachment.name).toBe("report.pdf"); + const first = owner.session.messages.send("c", "first", [], [attachment]); await vi.waitFor(() => expect(published).toBe(1), { timeout: 3000 }); const pending = owner.session.profiling .snapshot() .find((sample) => sample.stage === "send.publish" && sample.id === first); expect(pending).toMatchObject({ outcome: "pending" }); + expect(uploaded).toBe(1); release(); await vi.waitFor(() => expect(outbox.snapshot()).toHaveLength(0), { timeout: 3000, @@ -101,6 +124,19 @@ it("profiles a first slow publish through real local IPC, signing, authenticated assert.exists(secondUpstream); expect(firstUpstream.duration).toBeGreaterThan(0); expect(secondUpstream.duration).toBeLessThan(firstUpstream.duration); + expect(events.get(first)).toMatchObject({ + content: `first\n\n[report.pdf](<${fixtureRelayUrl}/media/${hash}.pdf>)`, + tags: expect.arrayContaining([ + [ + "imeta", + `url ${fixtureRelayUrl}/media/${hash}.pdf`, + "m application/pdf", + "size 3", + `x ${hash}`, + "filename report.pdf", + ], + ]), + }); expect( timings.some( (sample) => sample.id === first && sample.stage === "send.sign", diff --git a/src/features/relay/messages.ts b/src/features/relay/messages.ts index 38f91aacd..54c1b4794 100644 --- a/src/features/relay/messages.ts +++ b/src/features/relay/messages.ts @@ -1,3 +1,4 @@ +import { attachmentMessage, type UploadedAttachment } from "./attachments"; import { validReactionContent } from "./emoji"; import type { EventData } from "./events"; import type { Outbox } from "./outbox"; @@ -10,6 +11,7 @@ export function createMessages( emojiTags: (content: string) => string[][], validateMentions: (channelId: string, pubkeys: readonly string[]) => void, canParticipate: (channelId: string) => boolean = () => true, + relayOrigin?: string, ) { const writer = (kind: number, channelId: string) => { if (!canParticipate(channelId)) @@ -34,15 +36,22 @@ export function createMessages( return unique.map((key) => ["p", key]); }; return Object.freeze({ - send(channelId: string, content: string, mentions: readonly string[] = []) { + send( + channelId: string, + content: string, + mentions: readonly string[] = [], + attachments: readonly UploadedAttachment[] = [], + ) { if (!channelId) throw new Error("A channel is required"); + const message = attachmentMessage(content, attachments, relayOrigin); return writer(9, channelId).send({ kind: 9, - content: text(content), + content: text(message.content), tags: [ ["h", channelId], ...mentionTags(channelId, mentions), ...emojiTags(content), + ...message.tags, ], }); }, @@ -51,18 +60,21 @@ export function createMessages( rootId: string, content: string, mentions: readonly string[] = [], + attachments: readonly UploadedAttachment[] = [], ) { if (!channelId) throw new Error("A channel is required"); if (!/^[0-9a-f]{64}$/.test(rootId)) throw new Error("A valid thread root is required"); + const message = attachmentMessage(content, attachments, relayOrigin); return writer(9, channelId).send({ kind: 9, - content: text(content), + content: text(message.content), tags: [ ["h", channelId], ["e", rootId, "", "reply"], ...mentionTags(channelId, mentions), ...emojiTags(content), + ...message.tags, ], }); }, diff --git a/src/features/relay/session-attachments.test.ts b/src/features/relay/session-attachments.test.ts new file mode 100644 index 000000000..da9517bdf --- /dev/null +++ b/src/features/relay/session-attachments.test.ts @@ -0,0 +1,179 @@ +import { assert, afterEach, expect, it, vi } from "vitest"; +import type { EventTemplate } from "nostr-tools"; +import type { RelayEvent } from "./events"; +import type { AttachmentUpload, UploadedAttachment } from "./attachments"; +import { createRelaySession } from "./session"; +import { PublishRejected } from "./outbox"; +import { keypair, metadata, roster, signed } from "./testing"; + +const viewer = keypair(), + relay = keypair(), + other = keypair(); +const origin = "https://relay.test"; +const attachment: UploadedAttachment = { + name: "notes.pdf", + url: `${origin}/media/${"a".repeat(64)}.pdf`, + type: "application/pdf", + size: 3, + sha256: "a".repeat(64), +}; +const file = new File(["pdf"], attachment.name); +const owners: ReturnType[] = []; +afterEach(() => { + for (const owner of owners.splice(0)) owner.dispose(); +}); +function setup(uploadAttachment?: AttachmentUpload, writable = true) { + let members = [viewer.pubkey, other.pubkey]; + let time = 1700000000; + const sign = vi.fn(async (template: EventTemplate) => + signed(viewer, template), + ); + const publish = vi.fn(async (_event: RelayEvent) => {}); + const owner = createRelaySession( + { + viewer: viewer.pubkey, + relayAuthor: relay.pubkey, + scope: origin, + media: (url) => url, + async query(filters) { + if (filters[0]?.limit === 1 && filters[0]?.kinds?.[0] === 39002) + return [roster(relay, "c", members, time)]; + return filters.some((f) => f.kinds?.includes(39002)) + ? [ + roster(relay, "c", members, time), + metadata(relay, "c", "General", time), + ] + : []; + }, + ...(uploadAttachment ? { uploadAttachment } : {}), + ...(writable ? { writer: { sign, publish } } : {}), + }, + { outboxStorage: { load: () => [], save() {} } }, + ); + owners.push(owner); + async function membership(next = members) { + members = next; + time++; + await owner.session.read([ + { kinds: [39002, 39000], "#d": ["c"], limit: 10 }, + ]); + } + return { ...owner, sign, publish, membership }; +} + +it("does not advertise upload support on read-only or unsupported connections", () => { + expect(setup().session.attachments).toBeUndefined(); + expect( + setup(async () => attachment, false).session.attachments, + ).toBeUndefined(); +}); + +it.each([false, true])( + "uploads once, builds attachment-only root/reply (%s), and retries the same signed event", + async (reply) => { + const upload = vi.fn(async () => attachment); + const h = setup(upload); + assert.exists(h.session.attachments); + await h.membership(); + h.publish.mockRejectedValueOnce(new PublishRejected("try again")); + const result = await h.session.attachments.upload( + file, + "c", + new AbortController().signal, + ); + expect(h.publish).not.toHaveBeenCalled(); + const root = "b".repeat(64); + const id = reply + ? h.session.messages.reply("c", root, "", [other.pubkey], [result]) + : h.session.messages.send("c", "", [other.pubkey], [result]); + await vi.waitFor(() => + expect(h.session.outbox?.snapshot()[0]?.delivery).toBe("failed"), + ); + const first = h.publish.mock.calls[0]?.[0]; + assert.exists(first); + expect(first.content).toBe(`[notes.pdf](<${attachment.url}>)`); + expect(first.tags).toContainEqual(["p", other.pubkey]); + expect(first.tags).toContainEqual([ + "imeta", + `url ${attachment.url}`, + "m application/pdf", + "size 3", + `x ${attachment.sha256}`, + "filename notes.pdf", + ]); + if (reply) expect(first.tags).toContainEqual(["e", root, "", "reply"]); + h.session.messages.retry(id); + await vi.waitFor(() => expect(h.publish).toHaveBeenCalledTimes(2)); + expect(h.publish.mock.calls[1]?.[0]).toEqual(first); + expect(h.sign).toHaveBeenCalledTimes(1); + expect(upload).toHaveBeenCalledTimes(1); + }, +); + +it.each(["caller", "dispose", "clear", "access"])( + "cancels upload on %s and fences an uncooperative late result", + async (action) => { + let release!: (value: UploadedAttachment) => void; + let signal!: AbortSignal; + const h = setup(async (_file, supplied) => { + signal = supplied; + return new Promise((resolve) => { + release = resolve; + }); + }); + await h.membership(); + assert.exists(h.session.attachments); + const caller = new AbortController(); + const pending = h.session.attachments.upload(file, "c", caller.signal); + const rejected = expect(pending).rejects.toMatchObject({ + name: "AbortError", + }); + if (action === "caller") caller.abort(); + if (action === "dispose") h.dispose(); + if (action === "clear") await h.clearCache(); + if (action === "access") await h.membership([other.pubkey]); + expect(signal.aborted).toBe(true); + release(attachment); + await rejected; + expect(h.publish).not.toHaveBeenCalled(); + }, +); + +it("blocks upload and attachment sending after membership loss or session disposal", async () => { + const upload = vi.fn(async () => attachment); + const h = setup(upload); + assert.exists(h.session.attachments); + await h.membership([other.pubkey]); + await expect( + h.session.attachments.upload(file, "c", new AbortController().signal), + ).rejects.toMatchObject({ code: "denied" }); + expect(() => h.session.messages.send("c", "", [], [attachment])).toThrow( + /Join/, + ); + expect(upload).not.toHaveBeenCalled(); + h.dispose(); + expect(() => h.session.messages.send("c", "hello")).toThrow(/Join/); +}); + +it("rejects other-community attachment URLs, invalid roots and genuinely empty messages before enqueue", async () => { + const h = setup(async () => attachment); + await h.membership(); + expect(() => + h.session.messages.send( + "c", + "", + [], + [ + { + ...attachment, + url: attachment.url.replace("relay.test", "other.test"), + }, + ], + ), + ).toThrow(/invalid/); + expect(() => + h.session.messages.reply("c", "bad", "", [], [attachment]), + ).toThrow(/thread root/); + expect(() => h.session.messages.send("c", " ")).toThrow(/empty/); + expect(h.sign).not.toHaveBeenCalled(); +}); diff --git a/src/features/relay/session.ts b/src/features/relay/session.ts index a645bdc28..2666038c4 100644 --- a/src/features/relay/session.ts +++ b/src/features/relay/session.ts @@ -33,6 +33,7 @@ import { createSidebarPreferencesStore } from "./sidebar-preferences-store"; import { createEmojiDirectory } from "./emoji-directory"; import { createProfileDirectory } from "./profile-directory"; import { createChannelStore, type ChannelStoreOptions } from "./store"; +import { UploadError } from "./attachments"; import type { ReadTransport } from "./transport"; import type { LiveSnapshot, LiveSubscription } from "./live"; import { @@ -166,6 +167,11 @@ export function createRelaySession( ) { let closed = false; const lifetime = new AbortController(); + let uploadLifetime = new AbortController(); + function cancelUploads() { + uploadLifetime.abort(); + uploadLifetime = new AbortController(); + } const profiling = options.profiling ?? transport?.profiling ?? createRelayProfiler(); const requests = createRelayReader(transport, { profiling }); @@ -201,6 +207,7 @@ export function createRelaySession( const views = new Map<() => void, (clear?: boolean) => void>(); const threads = new Set>(); const writer = transport?.writer; + const uploadAttachment = transport?.uploadAttachment; const writes = transport && writer ? createOutbox( @@ -286,6 +293,7 @@ export function createRelaySession( revoking++; try { accessEpoch++; + cancelUploads(); typing.clear(); // Filters cannot tell us ownership of broad/ID/reference reads. Infrequent // authoritative access loss cancels them all, not merely explicit #h reads. @@ -904,6 +912,26 @@ export function createRelaySession( sidebarPreferences: sidebarPreferences.queries, live, profiling, + attachments: + uploadAttachment && writes?.outbox.supports(9) + ? Object.freeze({ + async upload(file: File, channelId: string, signal: AbortSignal) { + const combined = AbortSignal.any([ + signal, + lifetime.signal, + uploadLifetime.signal, + ]); + combined.throwIfAborted(); + if (!channelId || closed || !channels.canParticipate(channelId)) + throw new UploadError("denied"); + const result = await uploadAttachment(file, combined); + combined.throwIfAborted(); + if (!channels.canParticipate(channelId)) + throw new UploadError("denied"); + return result; + }, + }) + : undefined, messages: createMessages( writes?.outbox, transport?.viewer, @@ -913,7 +941,8 @@ export function createRelaySession( retainedEvent(id), emoji.tags, validateMentions, - (id) => channels.canParticipate(id), + (id) => !closed && channels.canParticipate(id), + transport?.scope, ), /** An owned bounded thread reader. Dispose on close; the session retains access/lifetime authority. */ thread( @@ -1470,6 +1499,7 @@ export function createRelaySession( session, async clearCache() { accessEpoch++; + cancelUploads(); cacheClearEpoch++; activity.clear(); presence.clear(); diff --git a/src/features/relay/transport.ts b/src/features/relay/transport.ts index 4a3286f4b..d03706ef1 100644 --- a/src/features/relay/transport.ts +++ b/src/features/relay/transport.ts @@ -1,3 +1,4 @@ +import { brokerUpload, type AttachmentUpload } from "./attachments"; import { workflowHost } from "../workflows/http"; import type { WorkflowHost } from "../workflows/host"; import { readReceiptText } from "./receipt"; @@ -48,6 +49,7 @@ export interface RelayWriter { ): Promise | Promise; } export interface ReadTransport { + readonly uploadAttachment?: AttachmentUpload; readonly workflows?: WorkflowHost; /** Purpose-bound observer decoding on the shared host live stream. */ readonly agentActivity?: boolean; @@ -222,6 +224,7 @@ export async function connectBrokerTransport( archiveAuthority?: unknown; writeKinds?: number[]; workflowReads?: boolean; + attachmentUploads?: boolean; relayUrl?: string; live?: boolean; presence?: boolean; @@ -251,6 +254,9 @@ export async function connectBrokerTransport( }); return { profiling, + ...(session.attachmentUploads === true && session.relayUrl + ? { uploadAttachment: brokerUpload(endpoint, session.relayUrl) } + : {}), ...(session.presence && session.live ? { async presenceSnapshot(