source: Klonkt/src/services/ActivityPubService.js@ 7cc58bb

main
Last change on this file since 7cc58bb was 3dd99d3, checked in by Robin Genis <roboburr@…>, 3 months ago

harden(fediverse): SSRF guard, remote-URL XSS scheme-guard, scoped Delete, gate-by-default + queue fixes

From a 3-agent hardening review of this session's fediverse code:

  • SSRF: all outbound fetches (deliver/fetchActor/webfingerResolve) now go through safeFetch — http(s)-only, rejects hosts resolving to private/loopback/link-local ranges on the initial host AND every redirect hop (redirect:manual), + actor-doc size cap. Blocks inbox-driven SSRF to cloud-metadata/internal services.
  • Stored XSS: remote actor url/icon, timeline media + author urls, and remote-note images/object_uri are now run through an http(s) scheme-guard before storage, so a malicious actor can't smuggle javascript:/data: into owner-only-rendered href/src.
  • Cross-actor Delete: inbound Delete is now scoped to the signing actor (can't wipe another actor's replies/timeline rows).
  • Gate-by-default: Add/Remove/Update added to the signature-enforced activity list.
  • Delivery queue: re-entrancy guard (30 rows x 8s can exceed the 60s tick -> no double-delivery) + backoff off-by-one fix (1-min first retry no longer skipped).
  • Scheduler: delete-before-insert on FTS so a re-flipped post has no duplicate row.
  • /meldingen: don't mark-seen for a viewer (GET-side mutation the global guard misses).
  • Activity ids get a random suffix to avoid same-millisecond collisions.

Co-Authored-By: Claude <noreply@…>

  • Property mode set to 100644
File size: 60.3 KB
Line 
1/**
2 * ActivityPubService — Klonkt as a real ActivityPub actor (fediverse bridge).
3 *
4 * Phase 1 (this file): the PUBLISH/discoverable side.
5 * - per-site RSA keypair (Mastodon-compatible HTTP Signatures; separate from
6 * the Ed25519 keys used by the lighter Cirkels v1)
7 * - builders for the Actor document, Note objects and the Outbox collection
8 * - apWants(): HTTP content-negotiation helper (activity+json vs HTML)
9 *
10 * The interactive side (inbox: Follow/Accept, signature verify, delivery to
11 * followers) lands in the next step and is tested live against Mastodon.
12 *
13 * AP actor URLs live under /ap/* so they never clash with the human pages:
14 * actor = <base>/ap/users/<slug>
15 * inbox = <actor>/inbox outbox = <actor>/outbox
16 * note = <base>/ap/notes/<postId>
17 */
18import crypto from 'crypto';
19import dns from 'dns';
20import net from 'net';
21import db from '../config/database.js';
22import HtmlSanitizerService from './HtmlSanitizerService.js';
23
24const PUBLIC = 'https://www.w3.org/ns/activitystreams#Public';
25
26// Short random suffix so two activity ids minted in the same millisecond (e.g.
27// parallel saves) don't collide and get deduped by a receiver.
28const rid = () => crypto.randomBytes(4).toString('hex');
29
30// Keep only http(s) URLs — drops javascript:/data:/etc so a remote actor can't
31// smuggle a dangerous scheme into a stored href/src (rendered in owner-only views).
32const safeUrl = (u) => { const s = String(u == null ? '' : u).trim(); return /^https?:\/\//i.test(s) ? s : ''; };
33
34// ── SSRF guard for outbound fetches ───────────────────────────────
35// Remote URLs (actor/keyId/webfinger/inbox/inReplyTo) are attacker-controlled, so
36// every outbound fetch must refuse hosts that resolve to private/loopback ranges
37// (cloud metadata, internal services) — on the initial host AND each redirect hop.
38function isBlockedIp(ip) {
39 if (!ip) return true;
40 const v = net.isIP(ip);
41 if (v === 4) {
42 const o = ip.split('.').map(Number);
43 return o[0] === 127 || o[0] === 10 || o[0] === 0
44 || (o[0] === 172 && o[1] >= 16 && o[1] <= 31)
45 || (o[0] === 192 && o[1] === 168)
46 || (o[0] === 169 && o[1] === 254)
47 || (o[0] === 100 && o[1] >= 64 && o[1] <= 127); // CGNAT
48 }
49 if (v === 6) {
50 const s = ip.toLowerCase().replace(/^\[|\]$/g, '');
51 return s === '::1' || s === '::' || s.startsWith('fc') || s.startsWith('fd') || s.startsWith('fe80')
52 || s.startsWith('::ffff:127.') || s.startsWith('::ffff:10.') || s.startsWith('::ffff:192.168.')
53 || s.startsWith('::ffff:169.254.') || s.startsWith('::ffff:172.');
54 }
55 return true; // not an IP literal we recognise → refuse
56}
57async function assertPublicHost(hostname) {
58 if (net.isIP(hostname)) { if (isBlockedIp(hostname)) throw new Error('ssrf-blocked-ip'); return; }
59 const addrs = await dns.promises.lookup(hostname, { all: true });
60 if (!addrs.length || addrs.some((a) => isBlockedIp(a.address))) throw new Error('ssrf-blocked-host');
61}
62async function safeFetch(url, opts = {}, maxRedirects = 3) {
63 let target = url;
64 for (let hop = 0; ; hop++) {
65 const u = new URL(target); // throws on malformed → caller's catch
66 if (u.protocol !== 'https:' && u.protocol !== 'http:') throw new Error('ssrf-bad-scheme');
67 await assertPublicHost(u.hostname);
68 const r = await fetch(target, { ...opts, redirect: 'manual', signal: AbortSignal.timeout(8000) });
69 const loc = (r.status >= 300 && r.status < 400) ? r.headers.get('location') : null;
70 if (loc && hop < maxRedirects) { target = new URL(loc, target).toString(); continue; }
71 return r;
72 }
73}
74const MAX_OUTBOX = 20;
75// Cache-buster for the music listen-link → forces Mastodon to re-crawl a FRESH
76// (square) player card. Bump this whenever the twitter:player card dimensions change.
77const FEDI_CARD_VER = '2';
78
79// ── RSA keys per actor (lazy, cached in DB) ───────────────────────
80// Prepared lazily (NOT at module load) — the ap_keys table is created in
81// initializeDatabase(), which runs after this module is imported.
82let _sel, _ins;
83function keyStmts() {
84 if (!_sel) {
85 _sel = db.prepare('SELECT public_pem, private_pem FROM ap_keys WHERE slug = ?');
86 _ins = db.prepare('INSERT OR IGNORE INTO ap_keys (slug, public_pem, private_pem, created_at) VALUES (?,?,?,CURRENT_TIMESTAMP)');
87 }
88 return { sel: _sel, ins: _ins };
89}
90
91export function getOrCreateKeys(slug) {
92 const { sel, ins } = keyStmts();
93 const row = sel.get(slug);
94 if (row) return row;
95 const { publicKey, privateKey } = crypto.generateKeyPairSync('rsa', {
96 modulusLength: 2048,
97 publicKeyEncoding: { type: 'spki', format: 'pem' },
98 privateKeyEncoding: { type: 'pkcs8', format: 'pem' },
99 });
100 ins.run(slug, publicKey, privateKey);
101 return sel.get(slug) || { public_pem: publicKey, private_pem: privateKey };
102}
103
104// ── content negotiation ───────────────────────────────────────────
105// True when the caller wants ActivityPub JSON rather than the HTML page.
106export function apWants(req) {
107 const a = String(req.headers.accept || '').toLowerCase();
108 return a.includes('application/activity+json') ||
109 (a.includes('application/ld+json') && a.includes('activitystreams'));
110}
111
112const AP_CONTENT_TYPE = 'application/activity+json; charset=utf-8';
113export function sendAP(res, obj) {
114 res.type(AP_CONTENT_TYPE);
115 res.set('Cache-Control', 'public, max-age=120');
116 res.send(JSON.stringify(obj));
117}
118
119// ── document builders ─────────────────────────────────────────────
120export function actorId(base, slug) { return `${base}/ap/users/${encodeURIComponent(slug)}`; }
121export function noteId(base, postId) { return `${base}/ap/notes/${encodeURIComponent(postId)}`; }
122
123export function buildActor(base, site) {
124 const id = actorId(base, site.slug);
125 const keys = getOrCreateKeys(site.slug);
126 const actor = {
127 '@context': ['https://www.w3.org/ns/activitystreams', 'https://w3id.org/security/v1'],
128 id,
129 type: 'Person',
130 preferredUsername: site.slug,
131 name: site.title || site.slug,
132 summary: site.tagline || site.description || '',
133 url: `${base}/${site.slug === site.primary_slug ? '' : 'user/' + encodeURIComponent(site.slug)}`,
134 manuallyApprovesFollowers: false,
135 discoverable: true,
136 inbox: `${id}/inbox`,
137 outbox: `${id}/outbox`,
138 followers: `${id}/followers`,
139 featured: `${id}/featured`,
140 endpoints: { sharedInbox: `${base}/ap/inbox` },
141 publicKey: {
142 id: `${id}#main-key`,
143 owner: id,
144 publicKeyPem: keys.public_pem,
145 },
146 };
147 if (site.profile_photo) {
148 const u = /^https?:/.test(site.profile_photo) ? site.profile_photo : `${base}${site.profile_photo.startsWith('/') ? '' : '/'}${site.profile_photo}`;
149 actor.icon = { type: 'Image', url: u };
150 }
151 return actor;
152}
153
154// Does a post's audio shortcodes reference at least one PLAYABLE (file-backed)
155// track? Link-only tracks (external Spotify/YouTube, media_id NULL) don't count —
156// they have no Klonkt-hosted audio to embed, so no player card / cover-suppression.
157export function hasPlayableAudio(content, siteId) {
158 if (!content || !/\[\[(track|album|playlist):/i.test(content)) return false;
159 try {
160 for (const m of content.matchAll(/\[\[track:([A-Za-z0-9_-]+)\]\]/g)) { const r = db.prepare('SELECT media_id FROM audio_tracks WHERE id = ?').get(m[1]); if (r && r.media_id) return true; }
161 for (const m of content.matchAll(/\[\[album:([^\]]+)\]\]/g)) { if (db.prepare('SELECT 1 FROM audio_tracks WHERE site_id = ? AND album = ? AND media_id IS NOT NULL LIMIT 1').get(siteId, m[1].trim())) return true; }
162 for (const m of content.matchAll(/\[\[playlist:([A-Za-z0-9_-]+)\]\]/g)) { if (db.prepare('SELECT 1 FROM playlist_tracks pt JOIN audio_tracks t ON t.id = pt.track_id WHERE pt.playlist_id = ? AND t.media_id IS NOT NULL LIMIT 1').get(m[1])) return true; }
163 } catch { /* non-fatal */ }
164 return false;
165}
166
167// A single post as an AS2 Note (the object), and as a Create activity (for outbox/delivery).
168export function buildNote(base, site, post) {
169 const id = noteId(base, post.id);
170 const aId = actorId(base, site.slug);
171 const human = `${base}/${encodeURIComponent(post.slug)}`;
172 // Mastodon ignores a Note's `name`, so put the title INTO the content (bold
173 // first line) — the standard blog→fediverse convention. post.content is
174 // already sanitized HTML; the title is plain text, so escape it.
175 const escTitle = String(post.title || '').replace(/[<>&]/g, (c) => ({ '<': '&lt;', '>': '&gt;', '&': '&amp;' }[c]));
176 const titleHtml = post.title ? `<p><strong>${escTitle}</strong></p>` : '';
177
178 // Images travel as AP `attachment` (Mastodon strips <img> from content). Collect
179 // the cover + any inline <img>, make absolute, then strip <img> from the content
180 // to avoid duplicate rendering on clients that DO keep them.
181 const abs = (u) => !u ? null : (/^https?:/i.test(u) ? u : `${base}${u.startsWith('/') ? '' : '/'}${u}`);
182 const mediaType = (u) => {
183 const e = ((u || '').split('?')[0].match(/\.(\w+)$/) || [])[1];
184 return ({ jpg: 'image/jpeg', jpeg: 'image/jpeg', png: 'image/png', gif: 'image/gif', webp: 'image/webp', avif: 'image/avif' })[(e || '').toLowerCase()] || 'image/jpeg';
185 };
186 const hadAudio = /\[\[(track|album|playlist):/i.test(post.content || '');
187 const playable = hasPlayableAudio(post.content || '', site && site.id);
188 const urls = [];
189 // Posts with PLAYABLE hosted audio suppress image attachments so Mastodon renders
190 // the player CARD (twitter:player) instead of the cover — media attachment and
191 // link/player card are mutually exclusive on Mastodon. Link-only audio (external)
192 // keeps its cover (no player card to show).
193 if (post.cover_image_url && !playable) urls.push(abs(post.cover_image_url));
194 let body = post.content || '';
195 if (!playable) for (const m of body.matchAll(/<img\b[^>]*\bsrc="([^"]+)"[^>]*>/gi)) urls.push(abs(m[1]));
196 body = body.replace(/<img\b[^>]*>/gi, '');
197 // Audio shortcodes: do NOT federate the raw audio file — Klonkt deliberately
198 // gates audio (the /audio/stream URL has friction), and shipping it as an AP
199 // audio attachment would hand Mastodon a plain, downloadable mp3 URL. Instead,
200 // replace the shortcodes with a "🎵 listen on the site" link so the post invites
201 // a click-through to the protected player (discovery without leaking the file).
202 const esc = (s) => String(s == null ? '' : s).replace(/[<>&]/g, (c) => ({ '<': '&lt;', '>': '&gt;', '&': '&amp;' }[c]));
203 const audioLabels = [];
204 try {
205 for (const m of body.matchAll(/\[\[track:([A-Za-z0-9_-]+)\]\]/g)) { const r = db.prepare('SELECT title FROM audio_tracks WHERE id = ?').get(m[1]); if (r && r.title) audioLabels.push(r.title); }
206 for (const m of body.matchAll(/\[\[album:([^\]]+)\]\]/g)) audioLabels.push(m[1].trim());
207 } catch { /* non-fatal */ }
208 body = body.replace(/\[\[(track|album|playlist):[^\]]+\]\]/gi, '');
209 // External embeds ([[embed:url]]) → emit the bare URL as a link so Mastodon
210 // renders its OWN preview/player card (YouTube/Spotify/SoundCloud/etc) instead
211 // of federating the raw shortcode text.
212 body = body.replace(/\[\[embed:([^\]]+)\]\]/gi, (mm, raw) => {
213 const u = esc(raw.trim().replace(/&amp;/g, '&'));
214 return `<p><a href="${u}">${u}</a></p>`;
215 });
216 if (hadAudio) {
217 const lbl = audioLabels.length ? esc(audioLabels.slice(0, 4).join(', ')) : '';
218 // For playable posts, append a version param to the listen-link so Mastodon
219 // sees a NEW card URL and re-crawls it (fresh SQUARE player card) instead of
220 // reusing the cached landscape one. Invisible: the link TEXT stays clean, the
221 // page ignores the param. Bump FEDI_CARD_VER when the card dimensions change.
222 const listenHref = playable ? `${human}?fc=${FEDI_CARD_VER}` : human;
223 body += `<p>🎵 ${lbl ? `<strong>${lbl}</strong> — ` : ''}<a href="${listenHref}">listen on ${esc(site.title || 'the site')}</a></p>`;
224 }
225 const seen = new Set();
226 const attachment = urls.filter(Boolean)
227 .filter((u) => { if (seen.has(u)) return false; seen.add(u); return true; })
228 .map((u) => ({ type: 'Document', mediaType: mediaType(u), url: u }));
229
230 const note = {
231 id,
232 type: 'Note',
233 attributedTo: aId,
234 content: titleHtml + body,
235 url: human,
236 published: new Date(post.published_at || post.created_at || Date.now()).toISOString(),
237 to: [PUBLIC],
238 cc: [`${aId}/followers`],
239 tag: Array.isArray(post.tags) ? post.tags.map((t) => ({ type: 'Hashtag', name: '#' + String(t).replace(/\s+/g, '') })) : [],
240 replies: `${id}/replies`,
241 };
242 if (attachment.length) note.attachment = attachment;
243 return note;
244}
245
246// All reply note URIs on a local post (inbound fediverse replies + our own
247// outbound replies) — backs the Note's `replies` Collection so remote servers
248// can fetch the whole thread.
249export function getReplyUris(base, postId) {
250 const out = [];
251 try {
252 for (const r of db.prepare("SELECT object_uri FROM ap_interactions WHERE kind = 'reply' AND post_id = ? AND object_uri != '' ORDER BY created_at").all(postId)) out.push(r.object_uri);
253 for (const r of db.prepare('SELECT id FROM ap_outbox WHERE post_id = ? ORDER BY rowid').all(postId)) out.push(`${base}/ap/notes/${r.id}`);
254 } catch { /* non-fatal */ }
255 return out;
256}
257
258// Notifications "seen" tracking → a real bell badge. Stored per site in app_settings.
259export function markNotificationsSeen(slug) {
260 try {
261 db.prepare("INSERT INTO app_settings (key, value, updated_at) VALUES (?, ?, CURRENT_TIMESTAMP) ON CONFLICT(key) DO UPDATE SET value = excluded.value, updated_at = CURRENT_TIMESTAMP")
262 .run(`fedi_notif_seen:${slug}`, new Date().toISOString());
263 } catch { /* non-fatal */ }
264}
265export function countUnseenNotifications(slug) {
266 try {
267 const row = db.prepare('SELECT value FROM app_settings WHERE key = ?').get(`fedi_notif_seen:${slug}`);
268 const seen = row ? Date.parse(row.value) : 0;
269 let n = 0;
270 for (const it of getNotifications(slug, 50)) { if (Date.parse(it.created_at) > seen) n++; }
271 return n;
272 } catch { return 0; }
273}
274
275export function buildCreate(base, site, post) {
276 const note = buildNote(base, site, post);
277 return {
278 '@context': 'https://www.w3.org/ns/activitystreams',
279 id: note.id + '#create',
280 type: 'Create',
281 actor: actorId(base, site.slug),
282 published: note.published,
283 to: note.to,
284 cc: note.cc,
285 object: note,
286 };
287}
288
289export function buildOutbox(base, site, posts) {
290 const id = `${actorId(base, site.slug)}/outbox`;
291 const items = (posts || []).slice(0, MAX_OUTBOX).map((p) => buildCreate(base, site, p));
292 return {
293 '@context': 'https://www.w3.org/ns/activitystreams',
294 id,
295 type: 'OrderedCollection',
296 totalItems: items.length,
297 orderedItems: items,
298 };
299}
300
301export function buildFollowers(base, site, count) {
302 const id = `${actorId(base, site.slug)}/followers`;
303 return {
304 '@context': 'https://www.w3.org/ns/activitystreams',
305 id,
306 type: 'OrderedCollection',
307 totalItems: count || 0,
308 orderedItems: [], // hidden for privacy; count only
309 };
310}
311
312// Pinned posts → the actor's `featured` collection. Mastodon reads this and shows
313// these as the "Featured" tab (pinned to the profile). Posts come ordered by pin
314// rank; embedded as full Notes so a remote server doesn't need extra fetches.
315export function buildFeatured(base, site, posts) {
316 const id = `${actorId(base, site.slug)}/featured`;
317 const items = (posts || []).map((p) => buildNote(base, site, p));
318 return {
319 '@context': 'https://www.w3.org/ns/activitystreams',
320 id,
321 type: 'OrderedCollection',
322 totalItems: items.length,
323 orderedItems: items,
324 };
325}
326
327// ── followers store (lazy stmts) ──────────────────────────────────
328let _insF, _delF, _listF, _cntF;
329function fStmts() {
330 if (!_insF) {
331 _insF = db.prepare('INSERT OR IGNORE INTO ap_followers (slug, actor_uri, inbox, shared_inbox, created_at) VALUES (?,?,?,?,CURRENT_TIMESTAMP)');
332 _delF = db.prepare('DELETE FROM ap_followers WHERE slug = ? AND actor_uri = ?');
333 _listF = db.prepare('SELECT inbox, shared_inbox FROM ap_followers WHERE slug = ?');
334 _cntF = db.prepare('SELECT COUNT(*) n FROM ap_followers WHERE slug = ?');
335 }
336 return { ins: _insF, del: _delF, list: _listF, cnt: _cntF };
337}
338export function followerCount(slug) { return fStmts().cnt.get(slug).n; }
339
340// ── inbound interactions store (replies / likes / boosts) + our outbound replies ──
341let _insI, _delLA, _delReply, _listI, _getI, _insO, _listO, _getO;
342function iStmts() {
343 if (!_insI) {
344 _insI = db.prepare('INSERT OR IGNORE INTO ap_interactions (kind, post_id, object_uri, actor_uri, actor_name, actor_handle, actor_url, actor_icon, content, published, parent_uri, created_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,CURRENT_TIMESTAMP)');
345 _delLA = db.prepare('DELETE FROM ap_interactions WHERE kind = ? AND post_id = ? AND actor_uri = ?');
346 _delReply = db.prepare("DELETE FROM ap_interactions WHERE kind = 'reply' AND object_uri = ?");
347 _listI = db.prepare('SELECT id, kind, object_uri, parent_uri, actor_uri, actor_name, actor_handle, actor_url, actor_icon, content, published, created_at FROM ap_interactions WHERE post_id = ? ORDER BY created_at ASC');
348 _getI = db.prepare('SELECT * FROM ap_interactions WHERE id = ?');
349 _insO = db.prepare('INSERT INTO ap_outbox (id, site_slug, post_id, post_slug, in_reply_to, to_actor, to_handle, content, created_at) VALUES (?,?,?,?,?,?,?,?,CURRENT_TIMESTAMP)');
350 _listO = db.prepare('SELECT * FROM ap_outbox WHERE post_id = ? ORDER BY created_at ASC');
351 _getO = db.prepare('SELECT * FROM ap_outbox WHERE id = ?');
352 }
353 return { ins: _insI, delLA: _delLA, delReply: _delReply, list: _listI, getI: _getI, insO: _insO, listO: _listO, getO: _getO };
354}
355
356export function getInteractionById(id) { return iStmts().getI.get(id); }
357
358const localPostExists = (id) => { try { return !!db.prepare('SELECT 1 FROM posts WHERE id = ?').get(id); } catch { return false; } };
359// Extract our local post id from a note URL, but only if it's ours (base match).
360function postIdFromNoteUrl(url, base) {
361 const s = String(url || '');
362 if (base && !s.startsWith(base)) return null;
363 const m = s.match(/\/ap\/notes\/([^/?#]+)/);
364 return m ? decodeURIComponent(m[1]) : null;
365}
366function deriveHandle(actorUri) {
367 try { const u = new URL(actorUri); const seg = u.pathname.split('/').filter(Boolean).pop() || ''; return `@${seg}@${u.host}`; } catch { return String(actorUri || ''); }
368}
369function actorInfo(doc, actorUri) {
370 let host = ''; try { host = new URL(actorUri).host; } catch { /* keep empty */ }
371 const handle = doc && doc.preferredUsername ? `@${doc.preferredUsername}@${host}` : deriveHandle(actorUri);
372 const icon = doc && doc.icon ? (doc.icon.url || (Array.isArray(doc.icon) && doc.icon[0] && doc.icon[0].url)) : null;
373 return {
374 name: (doc && (doc.name || doc.preferredUsername)) || handle,
375 handle,
376 url: safeUrl((doc && (doc.url || doc.id)) || actorUri) || null,
377 icon: safeUrl(icon) || null,
378 };
379}
380
381// Given an inReplyTo note URL, find which local post the thread belongs to + the
382// note being replied to (parent), so a reply-to-a-comment can be nested.
383function findThreadTarget(inReplyTo, base) {
384 if (!inReplyTo) return null;
385 const seg = postIdFromNoteUrl(inReplyTo, base); // our /ap/notes/<id> segment (if ours)
386 if (seg && localPostExists(seg)) return { post_id: seg, parent_uri: inReplyTo };
387 if (seg) {
388 try { const o = db.prepare('SELECT post_id FROM ap_outbox WHERE id = ?').get(seg); if (o && o.post_id) return { post_id: o.post_id, parent_uri: inReplyTo }; } catch { /* ignore */ }
389 }
390 try { const row = db.prepare("SELECT post_id FROM ap_interactions WHERE object_uri = ? AND kind = 'reply' LIMIT 1").get(inReplyTo); if (row && row.post_id) return { post_id: row.post_id, parent_uri: inReplyTo }; } catch { /* ignore */ }
391 return null;
392}
393
394// View-ready threaded view of a post's fediverse activity (inbound replies +
395// our outbound replies, nested), plus like/boost counts.
396export function getInteractions(postId, base, site) {
397 const s = iStmts();
398 const rows = s.list.all(postId);
399 const baseClean = (base || process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
400 const postNoteId = baseClean ? `${baseClean}/ap/notes/${postId}` : null;
401 // Our own (outbound) replies show the SITE identity for everyone (not "You").
402 let host = ''; try { host = new URL(baseClean).host; } catch { /* ignore */ }
403 const siteName = (site && (site.title || site.slug)) || '';
404 const siteHandle = (site && site.slug && host) ? `@${site.slug}@${host}` : '';
405 const siteUrl = baseClean ? `${baseClean}/` : '';
406 const siteIcon = (site && site.profile_photo) || null;
407
408 const nodes = [];
409 for (const r of rows) {
410 if (r.kind !== 'reply') continue;
411 nodes.push({
412 noteId: r.object_uri, parent: r.parent_uri || null, mine: false,
413 actor_name: r.actor_name, actor_handle: r.actor_handle, actor_url: r.actor_url,
414 actor_icon: r.actor_icon, content: r.content, created_at: r.published || r.created_at,
415 children: [],
416 });
417 }
418 for (const o of s.listO.all(postId)) {
419 nodes.push({
420 noteId: baseClean ? `${baseClean}/ap/notes/${o.id}` : o.id, parent: o.in_reply_to || null,
421 mine: true, outboxId: o.id, content: o.content, created_at: o.created_at,
422 actor_name: siteName, actor_handle: siteHandle, actor_url: siteUrl, actor_icon: siteIcon,
423 children: [],
424 });
425 }
426
427 const byId = new Map(nodes.map((n) => [n.noteId, n]));
428 const isTop = (n) => !n.parent || n.parent === postNoteId || !byId.has(n.parent);
429 const tops = [];
430 for (const n of nodes) {
431 if (isTop(n)) { tops.push(n); continue; }
432 let anc = n, guard = 0;
433 while (!isTop(anc) && guard++ < 12) anc = byId.get(anc.parent);
434 anc.children.push(n);
435 }
436 const byTime = (a, b) => new Date(a.created_at) - new Date(b.created_at);
437 tops.sort(byTime).forEach((t) => t.children.sort(byTime));
438
439 return {
440 thread: tops,
441 likeCount: rows.filter((r) => r.kind === 'like').length,
442 announceCount: rows.filter((r) => r.kind === 'announce').length,
443 total: nodes.length,
444 };
445}
446
447// ── HTTP Signatures + delivery ────────────────────────────────────
448const slugFromActorUrl = (url) => { const m = String(url || '').match(/\/ap\/users\/([^/?#]+)/); return m ? decodeURIComponent(m[1]) : null; };
449
450// Sign + POST an activity to a remote inbox (draft-cavage HTTP Signatures, RSA-SHA256).
451export async function deliver(inboxUrl, bodyObj, keyId, privatePem) {
452 const body = JSON.stringify(bodyObj);
453 const u = new URL(inboxUrl);
454 const date = new Date().toUTCString();
455 const digest = 'SHA-256=' + crypto.createHash('sha256').update(body).digest('base64');
456 const signingString = `(request-target): post ${u.pathname}\nhost: ${u.host}\ndate: ${date}\ndigest: ${digest}`;
457 const signature = crypto.sign('sha256', Buffer.from(signingString), privatePem).toString('base64');
458 const sig = `keyId="${keyId}",algorithm="rsa-sha256",headers="(request-target) host date digest",signature="${signature}"`;
459 const r = await safeFetch(inboxUrl, {
460 method: 'POST',
461 headers: { 'Content-Type': 'application/activity+json', Accept: 'application/activity+json', Date: date, Digest: digest, Signature: sig },
462 body,
463 });
464 return r.status;
465}
466
467export async function fetchActor(url) {
468 try {
469 const r = await safeFetch(url, { headers: { Accept: 'application/activity+json' } });
470 if (!r.ok) return null;
471 const len = Number(r.headers.get('content-length') || 0);
472 if (len > 2_000_000) return null; // refuse oversized actor docs
473 return await r.json();
474 } catch { return null; }
475}
476
477// ── Delivery queue with retries ───────────────────────────────────
478// Outbound deliveries are tried immediately; on failure (down server, timeout,
479// non-2xx) they're queued and retried with backoff so a briefly-offline follower
480// doesn't silently miss the post. The signing key is NOT stored — the worker
481// re-derives it from the actor slug at send time.
482const DELIVERY_MAX_ATTEMPTS = 6;
483const DELIVERY_BACKOFF_MIN = [1, 5, 15, 60, 180, 360];
484let _insDeliv, _dueDeliv, _delDeliv, _bumpDeliv;
485function deliveryStmts() {
486 if (!_insDeliv) {
487 _insDeliv = db.prepare('INSERT INTO ap_delivery (slug, inbox, body, attempts, next_at) VALUES (?,?,?,0,CURRENT_TIMESTAMP)');
488 _dueDeliv = db.prepare("SELECT * FROM ap_delivery WHERE datetime(next_at) <= datetime('now') ORDER BY next_at LIMIT 30");
489 _delDeliv = db.prepare('DELETE FROM ap_delivery WHERE id = ?');
490 _bumpDeliv = db.prepare('UPDATE ap_delivery SET attempts = ?, next_at = ? WHERE id = ?');
491 }
492 return { ins: _insDeliv, due: _dueDeliv, del: _delDeliv, bump: _bumpDeliv };
493}
494export function enqueueDelivery(slug, inbox, activity) {
495 if (!slug || !inbox || !activity) return;
496 try { deliveryStmts().ins.run(slug, inbox, JSON.stringify(activity)); } catch { /* ignore */ }
497}
498// Deliver now; queue for retry if it fails.
499export async function deliverWithRetry(slug, inbox, activity, keyId, privPem) {
500 if (!inbox) return;
501 try { const st = await deliver(inbox, activity, keyId, privPem); if (st >= 200 && st < 300) return; } catch { /* queue below */ }
502 enqueueDelivery(slug, inbox, activity);
503}
504let _processingDeliv = false;
505export async function processDeliveryQueue() {
506 if (_processingDeliv) return; // re-entrancy guard: 30 rows × 8s can exceed the 60s tick → no double-delivery
507 _processingDeliv = true;
508 try {
509 let rows;
510 try { rows = deliveryStmts().due.all(); } catch { return; }
511 if (!rows || !rows.length) return;
512 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
513 for (const row of rows) {
514 let ok = false;
515 try {
516 const keys = getOrCreateKeys(row.slug);
517 const st = await deliver(row.inbox, JSON.parse(row.body), `${actorId(base, row.slug)}#main-key`, keys.private_pem);
518 ok = st >= 200 && st < 300;
519 } catch { ok = false; }
520 if (ok) { deliveryStmts().del.run(row.id); continue; }
521 const attempts = row.attempts + 1;
522 if (attempts >= DELIVERY_MAX_ATTEMPTS) { deliveryStmts().del.run(row.id); console.warn('[AP] delivery gave up after', attempts, 'tries →', row.inbox); continue; }
523 // Index the backoff on the CURRENT attempt count (row.attempts) so the first
524 // retry uses the 1-min tier instead of skipping it.
525 const mins = DELIVERY_BACKOFF_MIN[Math.min(row.attempts, DELIVERY_BACKOFF_MIN.length - 1)];
526 deliveryStmts().bump.run(attempts, new Date(Date.now() + mins * 60000).toISOString(), row.id);
527 }
528 } finally { _processingDeliv = false; }
529}
530let _delivTimer = null;
531export function startDeliveryWorker() {
532 if (_delivTimer) return;
533 _delivTimer = setInterval(() => { processDeliveryQueue().catch(() => {}); }, 60 * 1000);
534 if (_delivTimer.unref) _delivTimer.unref();
535}
536
537// Best-effort verification of an incoming signed request. Returns the sender's
538// actor doc if the signature checks out, else null. (Not gating yet — MVP.)
539export async function verifyRequest(req) {
540 const sigH = req.headers['signature'];
541 if (!sigH) return null;
542 const p = Object.fromEntries([...sigH.matchAll(/([a-zA-Z]+)="([^"]*)"/g)].map((m) => [m[1], m[2]]));
543 if (!p.keyId || !p.signature) return null;
544 const actor = await fetchActor(p.keyId.split('#')[0]);
545 const pem = actor && actor.publicKey && actor.publicKey.publicKeyPem;
546 if (!pem) return null;
547 const hs = (p.headers || '(request-target) host date').split(/\s+/);
548 const line = hs.map((h) => h === '(request-target)'
549 ? `(request-target): ${req.method.toLowerCase()} ${req.originalUrl}`
550 : `${h}: ${req.headers[h] || ''}`).join('\n');
551 let ok = false;
552 try { ok = crypto.verify('sha256', Buffer.from(line), pem, Buffer.from(p.signature, 'base64')); } catch { ok = false; }
553 if (ok && hs.includes('digest') && req.rawBody) {
554 const exp = 'SHA-256=' + crypto.createHash('sha256').update(req.rawBody).digest('base64');
555 if (req.headers['digest'] !== exp) ok = false;
556 }
557 return ok ? actor : null;
558}
559
560// Handle an incoming inbox POST. slugParam = null for the shared /ap/inbox.
561export async function handleInbox(req, slugParam) {
562 const act = req.body || {};
563 const type = act.type;
564 const base = (process.env.PUBLIC_BASE_URL || `${req.protocol}://${req.get('host')}`).replace(/\/+$/, '');
565 const verified = await verifyRequest(req).catch(() => null);
566
567 // ENFORCE HTTP signatures: a data-affecting activity must be signed by the very
568 // actor it claims to be. No valid signature, or signer ≠ actor → reject (no
569 // forged replies/likes/follows/timeline posts). GET/discovery stays open.
570 const claimedActor = typeof act.actor === 'string' ? act.actor : (act.actor && act.actor.id);
571 // Blocked actor/domain → silently drop (202, don't reveal the block).
572 if (claimedActor && isBlockedAny(claimedActor)) { console.log('[AP] inbox dropped (blocked)', claimedActor); return 202; }
573 const GATED = ['Create', 'Like', 'Announce', 'Follow', 'Delete', 'Undo', 'Accept', 'Reject', 'Add', 'Remove', 'Update'];
574 if (GATED.includes(type)) {
575 if (!verified || !claimedActor || verified.id !== claimedActor) {
576 console.warn('[AP] inbox REJECTED (signature)', type, claimedActor || '?', verified ? '(signer mismatch)' : '(unsigned/invalid)');
577 return 401;
578 }
579 }
580
581 if (type === 'Follow') {
582 const who = typeof act.actor === 'string' ? act.actor : (act.actor && act.actor.id);
583 const slug = slugParam || slugFromActorUrl(typeof act.object === 'string' ? act.object : (act.object && act.object.id));
584 if (!who || !slug) return 400;
585 const remote = await fetchActor(who);
586 if (!remote || !remote.inbox) return 202; // can't reach them → drop quietly
587 fStmts().ins.run(slug, who, remote.inbox, (remote.endpoints && remote.endpoints.sharedInbox) || null);
588 const me = actorId(base, slug);
589 const keys = getOrCreateKeys(slug);
590 const accept = { '@context': 'https://www.w3.org/ns/activitystreams', id: `${me}#accept-${Date.now()}-${rid()}`, type: 'Accept', actor: me, object: act };
591 deliver(remote.inbox, accept, `${me}#main-key`, keys.private_pem).catch((e) => console.warn('[AP] Accept delivery failed:', e.message));
592 // Auto-backfill: send the new follower our recent posts as Create so their
593 // timeline isn't empty (Mastodon doesn't fetch history on follow). Fire-and-forget.
594 backfillNewFollower(base, slug, remote.inbox).catch(() => { /* best-effort */ });
595 console.log('[AP] Follow', who, '→', slug, verified ? '(sig ok)' : '(sig unverified)');
596 return 202;
597 }
598 if (type === 'Undo' && act.object) {
599 const who = typeof act.actor === 'string' ? act.actor : (act.actor && act.actor.id);
600 const ot = act.object.type;
601 if (ot === 'Follow') {
602 const obj = act.object.object;
603 const slug = slugParam || slugFromActorUrl(typeof obj === 'string' ? obj : (obj && obj.id));
604 if (who && slug) { fStmts().del.run(slug, who); console.log('[AP] Unfollow', who, '→', slug); }
605 return 202;
606 }
607 if (ot === 'Like' || ot === 'Announce') {
608 const tgt = act.object.object;
609 const pid = postIdFromNoteUrl(typeof tgt === 'string' ? tgt : (tgt && tgt.id), base);
610 if (who && pid) { iStmts().delLA.run(ot.toLowerCase(), pid, who); console.log('[AP] Undo', ot, who, '→', pid); }
611 return 202;
612 }
613 return 202;
614 }
615
616 const actorUri = typeof act.actor === 'string' ? act.actor : (act.actor && act.actor.id);
617 const resolveActor = async (uri) => ((verified && verified.id === uri) ? verified : await fetchActor(uri).catch(() => null));
618 // Activities from our OWN actors are already stored via ap_outbox — don't re-store.
619 const isLocalActor = !!(base && actorUri && actorUri.startsWith(`${base}/ap/users/`));
620
621 // Inbound reply: a Create whose object replies to one of our notes (post OR comment).
622 if (type === 'Create' && act.object && (act.object.type === 'Note' || act.object.type === 'Article')) {
623 const o = act.object;
624 const tgt = findThreadTarget(o.inReplyTo, base);
625 if (tgt && actorUri && !isLocalActor) {
626 const ai = actorInfo(await resolveActor(actorUri), actorUri);
627 const html = HtmlSanitizerService.sanitize(o.content || '');
628 iStmts().ins.run('reply', tgt.post_id, o.id || '', actorUri, ai.name, ai.handle, ai.url, ai.icon, html, o.published || null, tgt.parent_uri);
629 console.log('[AP] reply', actorUri, '→', tgt.post_id);
630 return 202;
631 }
632 // Home timeline (client): a top-level post from an account we follow.
633 if (actorUri && !isLocalActor && !o.inReplyTo && o.id) {
634 let subs = []; try { subs = db.prepare('SELECT slug FROM ap_following WHERE actor_uri = ?').all(actorUri); } catch { /* table may not exist yet */ }
635 if (subs.length) {
636 const ai = actorInfo(await resolveActor(actorUri), actorUri);
637 const html = HtmlSanitizerService.sanitize(o.content || '');
638 const media = JSON.stringify((Array.isArray(o.attachment) ? o.attachment : []).map((a) => ({ url: safeUrl(a && a.url), type: (a && a.mediaType) || '' })).filter((m) => m.url));
639 for (const s of subs) tlStmts().ins.run(o.id, s.slug, actorUri, ai.name, ai.handle, ai.icon, ai.url, html, o.url || null, o.published || null, media);
640 console.log('[AP] timeline +', actorUri, 'x' + subs.length);
641 }
642 }
643 return 202;
644 }
645 if (type === 'Like' || type === 'Announce') {
646 const tgt = act.object;
647 const pid = postIdFromNoteUrl(typeof tgt === 'string' ? tgt : (tgt && tgt.id), base);
648 if (pid && actorUri && !isLocalActor && localPostExists(pid)) {
649 const ai = actorInfo(await resolveActor(actorUri), actorUri);
650 iStmts().ins.run(type.toLowerCase(), pid, '', actorUri, ai.name, ai.handle, ai.url, ai.icon, null, null, null);
651 console.log('[AP]', type === 'Like' ? 'like' : 'boost', actorUri, '→', pid);
652 }
653 return 202;
654 }
655 if (type === 'Delete') {
656 // A remote note was deleted upstream → drop it from replies AND the timeline.
657 // Scope to the SIGNING actor so actor B can't delete actor A's content (the
658 // signature gate guarantees claimedActor == the verified signer here).
659 const oid = typeof act.object === 'string' ? act.object : (act.object && act.object.id);
660 if (oid && claimedActor) {
661 try { db.prepare('DELETE FROM ap_interactions WHERE object_uri = ? AND actor_uri = ?').run(oid, claimedActor); } catch { /* ignore */ }
662 try { db.prepare('DELETE FROM ap_timeline WHERE id = ? AND author_uri = ?').run(oid, claimedActor); } catch { /* ignore */ }
663 }
664 return 202;
665 }
666 // Accept/Reject of a Follow WE sent (client side).
667 if (type === 'Accept' && act.object) {
668 const fid = typeof act.object === 'string' ? act.object : (act.object && act.object.id);
669 if (fid) { try { fwStmts().acc.run(fid); } catch { /* ignore */ } }
670 console.log('[AP] follow accepted', actorUri);
671 return 202;
672 }
673 if (type === 'Reject' && act.object) {
674 const who = actorUri;
675 if (who && slugParam) { try { fwStmts().del.run(slugParam, who); } catch { /* ignore */ } }
676 return 202;
677 }
678
679 console.log('[AP] inbox', type || 'unknown', '→', slugParam || 'shared', '(ignored)');
680 return 202;
681}
682
683// Deliver a new post as Create(Note) to all followers' inboxes (fire-and-forget).
684// Needs PUBLIC_BASE_URL (absolute URLs); no-op without followers or base.
685export async function deliverCreate(site, post) {
686 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
687 if (!base || !site || !site.slug) return;
688 const followers = fStmts().list.all(site.slug);
689 if (!followers.length) return;
690 const inboxes = [...new Set(followers.map((f) => f.shared_inbox || f.inbox).filter(Boolean))];
691 const keys = getOrCreateKeys(site.slug);
692 const keyId = `${actorId(base, site.slug)}#main-key`;
693 const create = buildCreate(base, site, post);
694 for (const inbox of inboxes) deliverWithRetry(site.slug, inbox, create, keyId, keys.private_pem);
695}
696
697// On a new Follow, send that follower our most recent posts as Create so their
698// timeline shows our history (Mastodon does not backfill on follow). Oldest-first
699// so they sort into the follower's timeline at their original dates.
700async function backfillNewFollower(base, slug, inbox) {
701 if (!base || !slug || !inbox) return;
702 const site = db.prepare('SELECT * FROM sites WHERE slug = ?').get(slug);
703 if (!site) return;
704 const recent = db.prepare(
705 `SELECT id, slug, title, content, cover_image_url, published_at, created_at
706 FROM posts WHERE site_id = ? AND status = 'published' AND (fan_only IS NULL OR fan_only = 0)
707 ORDER BY COALESCE(published_at, created_at) DESC LIMIT 20`
708 ).all(site.id).reverse();
709 if (!recent.length) return;
710 const keys = getOrCreateKeys(slug);
711 const keyId = `${actorId(base, slug)}#main-key`;
712 for (const p of recent) {
713 try { await deliver(inbox, buildCreate(base, site, p), keyId, keys.private_pem); } catch { /* best-effort */ }
714 await new Promise((r) => setTimeout(r, 150));
715 }
716 console.log('[AP] backfilled', recent.length, 'posts to new follower of', slug);
717}
718
719// Tell followers a post is gone (Delete + Tombstone) so it's removed from their feeds.
720export async function deliverDelete(site, post) {
721 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
722 if (!base || !site || !site.slug || !post || !post.id) return;
723 const followers = fStmts().list.all(site.slug);
724 if (!followers.length) return;
725 const inboxes = [...new Set(followers.map((f) => f.shared_inbox || f.inbox).filter(Boolean))];
726 const keys = getOrCreateKeys(site.slug);
727 const me = actorId(base, site.slug);
728 const nid = noteId(base, post.id);
729 const del = {
730 '@context': 'https://www.w3.org/ns/activitystreams',
731 id: `${nid}#delete-${Date.now()}-${rid()}`,
732 type: 'Delete',
733 actor: me,
734 to: [PUBLIC],
735 object: { id: nid, type: 'Tombstone' },
736 };
737 for (const inbox of inboxes) deliverWithRetry(site.slug, inbox, del, `${me}#main-key`, keys.private_pem);
738}
739
740// Tell followers an already-published post changed (Update + edited Note) so
741// Mastodon refreshes the cached copy (e.g. after fixing content).
742export async function deliverUpdate(site, post) {
743 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
744 if (!base || !site || !site.slug || !post || !post.id) return;
745 const followers = fStmts().list.all(site.slug);
746 if (!followers.length) return;
747 const inboxes = [...new Set(followers.map((f) => f.shared_inbox || f.inbox).filter(Boolean))];
748 const keys = getOrCreateKeys(site.slug);
749 const me = actorId(base, site.slug);
750 const note = buildNote(base, site, post);
751 note.updated = new Date().toISOString();
752 const update = {
753 '@context': 'https://www.w3.org/ns/activitystreams',
754 id: `${noteId(base, post.id)}#update-${Date.now()}-${rid()}`,
755 type: 'Update', actor: me, to: [PUBLIC], cc: [`${me}/followers`],
756 object: note,
757 };
758 for (const inbox of inboxes) deliverWithRetry(site.slug, inbox, update, `${me}#main-key`, keys.private_pem);
759}
760
761// Tell followers the ACTOR changed (Update + Person) so Mastodon re-processes the
762// account AND re-fetches the featured (pinned) collection — there is no standard
763// "featured changed" activity, so this is how a pin/unpin propagates promptly.
764export async function deliverActorUpdate(site) {
765 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
766 if (!base || !site || !site.slug) return;
767 const followers = fStmts().list.all(site.slug);
768 if (!followers.length) return;
769 const inboxes = [...new Set(followers.map((f) => f.shared_inbox || f.inbox).filter(Boolean))];
770 const keys = getOrCreateKeys(site.slug);
771 const me = actorId(base, site.slug);
772 const update = {
773 '@context': ['https://www.w3.org/ns/activitystreams', 'https://w3id.org/security/v1'],
774 id: `${me}#update-${Date.now()}-${rid()}`,
775 type: 'Update', actor: me, to: [PUBLIC], cc: [`${me}/followers`],
776 object: buildActor(base, site),
777 };
778 for (const inbox of inboxes) deliverWithRetry(site.slug, inbox, update, `${me}#main-key`, keys.private_pem);
779}
780
781// Reliably set the pinned order on followers' instances via Add/Remove activities
782// (how Mastodon itself federates pins) — pushed to the inbox + processed immediately,
783// unlike the featured COLLECTION which Mastodon caches with sticky StatusPins.
784// Mastodon's Add skips an already-pinned status, so we REMOVE every pin first, wait,
785// then ADD in rank-DESCENDING order (rank 1 added LAST → newest StatusPin → shown first,
786// because Mastodon displays pins newest-first). `alsoRemove` = ids to unpin too.
787export async function resyncFeaturedPins(site, alsoRemove = []) {
788 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
789 if (!base || !site || !site.slug) return;
790 const followers = fStmts().list.all(site.slug);
791 if (!followers.length) return;
792 const inboxes = [...new Set(followers.map((f) => f.shared_inbox || f.inbox).filter(Boolean))];
793 const keys = getOrCreateKeys(site.slug);
794 const me = actorId(base, site.slug);
795 const keyId = `${me}#main-key`;
796 const featured = `${me}/featured`;
797 const AS = 'https://www.w3.org/ns/activitystreams';
798 const note = (id) => noteId(base, id);
799 const pinned = db.prepare(
800 `SELECT id FROM posts WHERE site_id = ? AND status = 'published' AND (fan_only IS NULL OR fan_only = 0)
801 AND pinned IS NOT NULL AND pinned > 0
802 ORDER BY pinned DESC, COALESCE(published_at, created_at) ASC LIMIT 20`
803 ).all(site.id);
804 const removeIds = [...new Set([...pinned.map((p) => p.id), ...alsoRemove])];
805 // 1. Remove every current pin so Mastodon can recreate them in order.
806 for (const id of removeIds) {
807 const rm = { '@context': AS, id: `${me}#rm-${id}-${Date.now()}-${rid()}`, type: 'Remove', actor: me, object: note(id), target: featured, to: [PUBLIC] };
808 for (const inbox of inboxes) deliver(inbox, rm, keyId, keys.private_pem).catch(() => { /* best-effort */ });
809 }
810 if (!pinned.length) { console.log('[AP] unpinned all featured for', site.slug); return; }
811 await new Promise((r) => setTimeout(r, 5000)); // let the Removes land first
812 // 2. Add in rank-DESC order, gaps so each StatusPin gets an increasing created_at.
813 for (const p of pinned) {
814 const add = { '@context': AS, id: `${me}#add-${p.id}-${Date.now()}-${rid()}`, type: 'Add', actor: me, object: note(p.id), target: featured, to: [PUBLIC], cc: [`${me}/followers`] };
815 for (const inbox of inboxes) deliver(inbox, add, keyId, keys.private_pem).catch(() => { /* best-effort */ });
816 await new Promise((r) => setTimeout(r, 2000));
817 }
818 console.log('[AP] resynced', pinned.length, 'featured pins for', site.slug);
819}
820
821// ── outbound replies (Klonkt → fediverse) ─────────────────────────
822const escHtml = (s) => String(s || '').replace(/[<>&]/g, (c) => ({ '<': '&lt;', '>': '&gt;', '&': '&amp;' }[c]));
823const toISO = (v) => { if (!v) return new Date().toISOString(); const s = String(v); const d = new Date(/[TZ]/.test(s) ? s : s.replace(' ', 'T') + 'Z'); return isNaN(d) ? new Date().toISOString() : d.toISOString(); };
824
825// Build one of OUR outbound reply Notes from an ap_outbox row.
826export function buildReplyNote(base, site, row) {
827 const me = actorId(base, site.slug);
828 return {
829 id: noteId(base, row.id),
830 type: 'Note',
831 attributedTo: me,
832 inReplyTo: row.in_reply_to || undefined,
833 content: row.content,
834 url: row.post_slug ? `${base}/${encodeURIComponent(row.post_slug)}` : undefined,
835 published: toISO(row.created_at),
836 to: row.to_actor ? [row.to_actor] : [PUBLIC],
837 cc: [PUBLIC, `${me}/followers`],
838 tag: row.to_actor ? [{ type: 'Mention', href: row.to_actor, name: row.to_handle }] : [],
839 };
840}
841
842// Resolve one of our outbound reply Notes by id (for /ap/notes/:id fallback).
843export function getOutboxNote(base, id) {
844 const row = iStmts().getO.get(id);
845 if (!row) return null;
846 const site = db.prepare('SELECT * FROM sites WHERE slug = ?').get(row.site_slug);
847 if (!site) return null;
848 return buildReplyNote(base, site, row);
849}
850
851// Send a reply FROM this site to a remote actor (in reply to their inbound reply).
852// `parent` = an ap_interactions row (actor_uri, actor_url, actor_handle, object_uri).
853export async function deliverReply(site, { postId, postSlug, parent, text }) {
854 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
855 if (!base || !site || !site.slug || !parent || !String(text || '').trim()) return null;
856 const me = actorId(base, site.slug);
857 const handle = parent.actor_handle || deriveHandle(parent.actor_uri);
858 const body = escHtml(String(text).trim()).replace(/\r?\n/g, '<br>');
859 const mention = parent.actor_uri
860 ? `<a href="${escHtml(parent.actor_url || parent.actor_uri)}" class="u-url mention">${escHtml(handle)}</a> ` : '';
861 const content = `<p>${mention}${body}</p>`;
862 // Dedup: skip if the exact same reply was already sent (double-submit guard).
863 const dup = db.prepare('SELECT 1 FROM ap_outbox WHERE site_slug = ? AND IFNULL(in_reply_to, \'\') = ? AND content = ? LIMIT 1')
864 .get(site.slug, parent.object_uri || '', content);
865 if (dup) { console.log('[AP] outreply skipped (duplicate)'); return { duplicate: true, delivered: 0 }; }
866 const id = crypto.randomUUID();
867 iStmts().insO.run(id, site.slug, postId, postSlug || null, parent.object_uri || null, parent.actor_uri || null, handle, content);
868 const row = iStmts().getO.get(id);
869 const note = buildReplyNote(base, site, row);
870 const create = {
871 '@context': 'https://www.w3.org/ns/activitystreams',
872 id: note.id + '#create', type: 'Create', actor: me,
873 published: note.published, to: note.to, cc: note.cc, object: note,
874 };
875 const keys = getOrCreateKeys(site.slug);
876 const keyId = `${me}#main-key`;
877 const inboxes = new Set();
878 if (parent.actor_uri) {
879 const a = await fetchActor(parent.actor_uri).catch(() => null);
880 if (a) inboxes.add((a.endpoints && a.endpoints.sharedInbox) || a.inbox);
881 }
882 if (parent.threadInbox) inboxes.add(parent.threadInbox); // back-compat (single)
883 (parent.threadInboxes || []).forEach((i) => inboxes.add(i)); // whole ancestor chain
884 for (const f of fStmts().list.all(site.slug)) inboxes.add(f.shared_inbox || f.inbox);
885 inboxes.delete(`${me}/inbox`); // never deliver to ourselves (already in ap_outbox)
886 inboxes.delete(`${base}/ap/inbox`); // (our own shared inbox) → avoids a self-duplicate
887 let delivered = 0;
888 for (const inbox of [...inboxes].filter(Boolean)) {
889 try { const st = await deliver(inbox, create, keyId, keys.private_pem); if (st >= 200 && st < 300) delivered++; } catch { /* best-effort */ }
890 }
891 console.log('[AP] outreply', site.slug, '→', parent.actor_uri, 'delivered', delivered);
892 return { id, content, delivered };
893}
894
895// Resolve a remote post URL (any fediverse/Klonkt post) into a reply target.
896// Returns a parent-shaped object usable by deliverReply(), or null.
897export async function resolveRemoteNote(url) {
898 if (!/^https?:\/\//i.test(String(url || ''))) return null;
899 const note = await fetchActor(url).catch(() => null); // AP GET (content-negotiates)
900 if (!note || !note.id) return null;
901 const att = note.attributedTo;
902 const actorUri = typeof att === 'string' ? att : (att && att.id);
903 if (!actorUri) return null;
904 const actor = await fetchActor(actorUri).catch(() => null);
905 const ai = actorInfo(actor, actorUri);
906 // Is what we're replying to a post (or a comment) on one of OUR posts? If so,
907 // link our reply to that local post so it shows nested in the post thread.
908 const localTgt = findThreadTarget(note.id, (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, ''));
909 // Walk the WHOLE reply chain upward (comment → parent comment → … → root post)
910 // and collect every ancestor author's inbox, so each participant's server —
911 // including the original post's author — receives + threads our reply.
912 const threadInboxes = [];
913 const seenInbox = new Set();
914 let cursor = note.inReplyTo, guard = 0;
915 while (cursor && guard++ < 6) {
916 const url = typeof cursor === 'string' ? cursor : (cursor && cursor.id);
917 if (!url) break;
918 const pn = await fetchActor(url).catch(() => null);
919 if (!pn) break;
920 const pa = typeof pn.attributedTo === 'string' ? pn.attributedTo : (pn.attributedTo && pn.attributedTo.id);
921 if (pa && pa !== actorUri) {
922 const paDoc = await fetchActor(pa).catch(() => null);
923 const inbox = paDoc && ((paDoc.endpoints && paDoc.endpoints.sharedInbox) || paDoc.inbox);
924 if (inbox && !seenInbox.has(inbox)) { seenInbox.add(inbox); threadInboxes.push(inbox); }
925 }
926 cursor = pn.inReplyTo; // climb to the next ancestor
927 }
928 const rawHtml = String(note.content || '').replace(/\[\[(track|album|playlist):[^\]]+\]\]/gi, '');
929 const images = (Array.isArray(note.attachment) ? note.attachment : [])
930 .filter((a) => a && a.url && (!a.mediaType || /^image\//i.test(a.mediaType)))
931 .map((a) => safeUrl(a.url)).filter(Boolean);
932 return {
933 object_uri: safeUrl(note.id) || note.id,
934 actor_uri: actorUri,
935 actor_url: ai.url,
936 actor_handle: ai.handle,
937 actor_name: ai.name,
938 actor_icon: ai.icon,
939 url: note.url || url,
940 content: HtmlSanitizerService.sanitize(rawHtml), // full, sanitized
941 images,
942 threadInboxes, // every ancestor author's inbox
943 localPostId: localTgt ? localTgt.post_id : '', // our post this belongs to (if any)
944 preview: HtmlSanitizerService.toPlainText(note.content || '').slice(0, 240),
945 };
946}
947
948// List a site's own outbound fediverse replies (for the manage/delete view).
949export function listOutbox(siteSlug) {
950 return db.prepare('SELECT id, content, to_handle, in_reply_to, created_at FROM ap_outbox WHERE site_slug = ? ORDER BY created_at DESC').all(siteSlug);
951}
952
953// Delete one of our outbound replies: send Delete(Tombstone) to recipients + remove it.
954export async function deliverOutboxDelete(site, outboxId) {
955 const row = iStmts().getO.get(outboxId);
956 if (!row || row.site_slug !== site.slug) return false;
957 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
958 if (base) {
959 const me = actorId(base, site.slug);
960 const nid = noteId(base, row.id);
961 const del = { '@context': 'https://www.w3.org/ns/activitystreams', id: `${nid}#delete-${Date.now()}-${rid()}`, type: 'Delete', actor: me, to: [PUBLIC], object: { id: nid, type: 'Tombstone' } };
962 const keys = getOrCreateKeys(site.slug);
963 const inboxes = new Set();
964 if (row.to_actor) { const a = await fetchActor(row.to_actor).catch(() => null); if (a) inboxes.add((a.endpoints && a.endpoints.sharedInbox) || a.inbox); }
965 for (const f of fStmts().list.all(site.slug)) inboxes.add(f.shared_inbox || f.inbox);
966 for (const inbox of [...inboxes].filter(Boolean)) { try { await deliver(inbox, del, `${me}#main-key`, keys.private_pem); } catch { /* best-effort */ } }
967 }
968 db.prepare('DELETE FROM ap_outbox WHERE id = ?').run(outboxId);
969 return true;
970}
971
972// ── Fediverse CLIENT: follow accounts + home timeline ─────────────
973// Resolve an @user@domain handle to its actor URL via WebFinger.
974export async function webfingerResolve(handle) {
975 const h = String(handle || '').trim().replace(/^@/, '');
976 const parts = h.split('@');
977 if (parts.length !== 2 || !parts[0] || !parts[1]) return null;
978 const acct = `${parts[0]}@${parts[1]}`;
979 try {
980 const r = await safeFetch(`https://${parts[1]}/.well-known/webfinger?resource=acct:${encodeURIComponent(acct)}`,
981 { headers: { Accept: 'application/jrd+json, application/json' } });
982 if (!r.ok) return null;
983 const jrd = await r.json();
984 const link = (jrd.links || []).find((l) => l.rel === 'self' && /activity\+json|ld\+json/.test(l.type || ''));
985 return safeUrl(link ? link.href : '') || null;
986 } catch { return null; }
987}
988
989let _insFw, _delFw, _listFw, _accFw, _oneFw;
990function fwStmts() {
991 if (!_insFw) {
992 _insFw = db.prepare('INSERT OR REPLACE INTO ap_following (slug, actor_uri, handle, name, icon, url, inbox, follow_id, status, created_at) VALUES (?,?,?,?,?,?,?,?,?,CURRENT_TIMESTAMP)');
993 _delFw = db.prepare('DELETE FROM ap_following WHERE slug = ? AND actor_uri = ?');
994 _listFw = db.prepare('SELECT * FROM ap_following WHERE slug = ? ORDER BY created_at DESC');
995 _accFw = db.prepare("UPDATE ap_following SET status = 'accepted' WHERE follow_id = ?");
996 _oneFw = db.prepare('SELECT * FROM ap_following WHERE slug = ? AND actor_uri = ?');
997 }
998 return { ins: _insFw, del: _delFw, list: _listFw, acc: _accFw, one: _oneFw };
999}
1000export function listFollowing(slug) { return fwStmts().list.all(slug); }
1001
1002let _insTl, _listTl, _delTl;
1003function tlStmts() {
1004 if (!_insTl) {
1005 _insTl = db.prepare('INSERT OR IGNORE INTO ap_timeline (id, slug, author_uri, author_name, author_handle, author_icon, author_url, content, url, published, media_json, created_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,CURRENT_TIMESTAMP)');
1006 _listTl = db.prepare('SELECT * FROM ap_timeline WHERE slug = ? ORDER BY COALESCE(published, created_at) DESC LIMIT ?');
1007 _delTl = db.prepare('DELETE FROM ap_timeline WHERE id = ?');
1008 }
1009 return { ins: _insTl, list: _listTl, del: _delTl };
1010}
1011export function getTimeline(slug, limit) { return tlStmts().list.all(slug, limit || 50); }
1012
1013// Follow a fediverse account by @handle (WebFinger → actor → signed Follow).
1014export async function followActor(site, handle) {
1015 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
1016 if (!base || !site || !site.slug) return { error: 'config' };
1017 const actorUrl = await webfingerResolve(handle);
1018 if (!actorUrl) return { error: 'not_found' };
1019 const actor = await fetchActor(actorUrl).catch(() => null);
1020 if (!actor || !actor.id || !actor.inbox) return { error: 'unreachable' };
1021 const ai = actorInfo(actor, actor.id);
1022 const me = actorId(base, site.slug);
1023 const keys = getOrCreateKeys(site.slug);
1024 const followId = `${me}#follow-${Date.now()}-${rid()}`;
1025 fwStmts().ins.run(site.slug, actor.id, ai.handle, ai.name, ai.icon, ai.url, actor.inbox, followId, 'pending');
1026 const follow = { '@context': 'https://www.w3.org/ns/activitystreams', id: followId, type: 'Follow', actor: me, object: actor.id };
1027 try { await deliver(actor.inbox, follow, `${me}#main-key`, keys.private_pem); }
1028 catch (e) { console.warn('[AP] follow deliver failed:', e.message); }
1029 console.log('[AP] follow', site.slug, '→', actor.id);
1030 return { ok: true, name: ai.name, handle: ai.handle };
1031}
1032
1033export async function unfollowActor(site, actorUri) {
1034 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
1035 const me = actorId(base, site.slug);
1036 const keys = getOrCreateKeys(site.slug);
1037 const row = fwStmts().one.get(site.slug, actorUri);
1038 if (row && row.inbox) {
1039 const undo = { '@context': 'https://www.w3.org/ns/activitystreams', id: `${me}#unfollow-${Date.now()}-${rid()}`, type: 'Undo', actor: me, object: { id: row.follow_id || `${me}#follow`, type: 'Follow', actor: me, object: actorUri } };
1040 try { await deliver(row.inbox, undo, `${me}#main-key`, keys.private_pem); } catch { /* best-effort */ }
1041 }
1042 fwStmts().del.run(site.slug, actorUri);
1043 return { ok: true };
1044}
1045
1046// Send a Like or Announce (boost) on a remote note FROM this site.
1047export async function sendInteraction(site, kind, targetNoteId, authorUri) {
1048 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
1049 if (!base || !site || !site.slug || !targetNoteId) return { error: 'config' };
1050 const type = kind === 'boost' ? 'Announce' : 'Like';
1051 const me = actorId(base, site.slug);
1052 const keys = getOrCreateKeys(site.slug);
1053 const act = {
1054 '@context': 'https://www.w3.org/ns/activitystreams',
1055 id: `${me}#${type.toLowerCase()}-${Date.now()}-${rid()}`,
1056 type, actor: me, object: targetNoteId,
1057 };
1058 if (type === 'Announce') { act.to = [PUBLIC]; act.cc = [`${me}/followers`]; }
1059 const inboxes = new Set();
1060 if (authorUri) { const a = await fetchActor(authorUri).catch(() => null); if (a) inboxes.add((a.endpoints && a.endpoints.sharedInbox) || a.inbox); }
1061 // A boost is public → also deliver to our own followers so it shows for them.
1062 if (type === 'Announce') { for (const f of fStmts().list.all(site.slug)) inboxes.add(f.shared_inbox || f.inbox); }
1063 let delivered = 0;
1064 for (const inbox of [...inboxes].filter(Boolean)) { try { const st = await deliver(inbox, act, `${me}#main-key`, keys.private_pem); if (st >= 200 && st < 300) delivered++; } catch { /* best-effort */ } }
1065 console.log('[AP]', type, site.slug, '→', targetNoteId, 'delivered', delivered);
1066 return { ok: true, delivered };
1067}
1068
1069// Notifications inbox: new followers + replies/likes/boosts on this site's posts.
1070export function getNotifications(slug, limit) {
1071 const out = [];
1072 try {
1073 for (const f of db.prepare('SELECT actor_uri, created_at FROM ap_followers WHERE slug = ? ORDER BY created_at DESC LIMIT 50').all(slug)) {
1074 out.push({ type: 'follow', handle: deriveHandle(f.actor_uri), url: f.actor_uri, created_at: f.created_at });
1075 }
1076 } catch { /* ignore */ }
1077 try {
1078 const rows = db.prepare(`
1079 SELECT i.kind, i.actor_name, i.actor_handle, i.actor_url, i.content, i.created_at,
1080 p.slug AS post_slug, p.title AS post_title
1081 FROM ap_interactions i LEFT JOIN posts p ON p.id = i.post_id
1082 WHERE p.site_id = (SELECT id FROM sites WHERE slug = ?)
1083 ORDER BY i.created_at DESC LIMIT 80
1084 `).all(slug);
1085 for (const r of rows) out.push({
1086 type: r.kind, name: r.actor_name, handle: r.actor_handle, url: r.actor_url,
1087 content: r.content, post_slug: r.post_slug, post_title: r.post_title, created_at: r.created_at,
1088 });
1089 } catch { /* ignore */ }
1090 out.sort((a, b) => new Date(b.created_at) - new Date(a.created_at));
1091 return out.slice(0, limit || 60);
1092}
1093
1094// ── Blocking / defederation ───────────────────────────────────────
1095let _insBl, _delBl, _listBl;
1096function blStmts() {
1097 if (!_insBl) {
1098 _insBl = db.prepare('INSERT OR IGNORE INTO ap_blocks (slug, target, kind, label, created_at) VALUES (?,?,?,?,CURRENT_TIMESTAMP)');
1099 _delBl = db.prepare('DELETE FROM ap_blocks WHERE slug = ? AND target = ?');
1100 _listBl = db.prepare('SELECT * FROM ap_blocks WHERE slug = ? ORDER BY created_at DESC');
1101 }
1102 return { ins: _insBl, del: _delBl, list: _listBl };
1103}
1104export function listBlocks(slug) { return blStmts().list.all(slug); }
1105
1106// True if an actor (or its whole domain) is blocked anywhere on this instance.
1107export function isBlockedAny(actorUri) {
1108 if (!actorUri) return false;
1109 let domain = ''; try { domain = new URL(actorUri).host; } catch { /* ignore */ }
1110 try { return !!db.prepare("SELECT 1 FROM ap_blocks WHERE (kind='actor' AND target=?) OR (kind='domain' AND target=?) LIMIT 1").get(actorUri, domain); }
1111 catch { return false; }
1112}
1113
1114function purgeBlocked(kind, target) {
1115 try {
1116 if (kind === 'domain') {
1117 const like = `%//${target}/%`;
1118 db.prepare('DELETE FROM ap_interactions WHERE actor_uri LIKE ?').run(like);
1119 db.prepare('DELETE FROM ap_timeline WHERE author_uri LIKE ?').run(like);
1120 db.prepare('DELETE FROM ap_followers WHERE actor_uri LIKE ?').run(like);
1121 } else {
1122 db.prepare('DELETE FROM ap_interactions WHERE actor_uri = ?').run(target);
1123 db.prepare('DELETE FROM ap_timeline WHERE author_uri = ?').run(target);
1124 db.prepare('DELETE FROM ap_followers WHERE actor_uri = ?').run(target);
1125 }
1126 } catch { /* best-effort */ }
1127}
1128
1129// Block an actor (@handle or actor URL) or a whole domain; purges their content.
1130export async function blockTarget(site, input) {
1131 const raw = String(input || '').trim();
1132 if (!site || !site.slug || !raw) return { error: 'empty' };
1133 let kind, target, label;
1134 if (/^https?:\/\//i.test(raw)) { kind = 'actor'; target = raw; label = raw; }
1135 else if (raw.includes('@')) {
1136 const actorUrl = await webfingerResolve(raw);
1137 if (!actorUrl) return { error: 'not_found' };
1138 kind = 'actor'; target = actorUrl; label = raw.startsWith('@') ? raw : ('@' + raw);
1139 } else { kind = 'domain'; target = raw.toLowerCase(); label = raw.toLowerCase(); }
1140 blStmts().ins.run(site.slug, target, kind, label);
1141 purgeBlocked(kind, target);
1142 console.log('[AP] block', site.slug, kind, target);
1143 return { ok: true, label };
1144}
1145
1146export function unblock(site, target) { blStmts().del.run(site.slug, target); return { ok: true }; }
1147
1148export default {
1149 getOrCreateKeys, apWants, sendAP, actorId, noteId,
1150 buildActor, buildNote, buildCreate, buildOutbox, buildFollowers, buildFeatured,
1151 followerCount, deliver, fetchActor, verifyRequest, handleInbox, deliverCreate, deliverDelete, deliverUpdate, deliverActorUpdate, resyncFeaturedPins,
1152 getInteractions, getInteractionById, buildReplyNote, getOutboxNote, deliverReply, resolveRemoteNote,
1153 listOutbox, deliverOutboxDelete,
1154 webfingerResolve, followActor, unfollowActor, listFollowing, getTimeline, sendInteraction,
1155 getNotifications, listBlocks, isBlockedAny, blockTarget, unblock,
1156 deliverWithRetry, enqueueDelivery, processDeliveryQueue, startDeliveryWorker,
1157 getReplyUris, markNotificationsSeen, countUnseenNotifications, hasPlayableAudio,
1158};
Note: See TracBrowser for help on using the repository browser.