source: Klonkt/src/services/ActivityPubService.js@ 4c12783

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

perf(ap): backfill new followers once per remote instance, not per follower

Backfill fired 20 Create deliveries on every Follow. Mastodon dedupes notes
per-instance, so the 2nd+ follower from the same instance got the same 20 posts
re-sent for nothing (and old posts don't re-enter a new follower's timeline
anyway). Now we skip backfill when another follower already represents that
instance (same shared_inbox), and deliver to the shared inbox when present.

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

  • Property mode set to 100644
File size: 60.9 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 const sharedInbox = (remote.endpoints && remote.endpoints.sharedInbox) || null;
588 fStmts().ins.run(slug, who, remote.inbox, sharedInbox);
589 const me = actorId(base, slug);
590 const keys = getOrCreateKeys(slug);
591 const accept = { '@context': 'https://www.w3.org/ns/activitystreams', id: `${me}#accept-${Date.now()}-${rid()}`, type: 'Accept', actor: me, object: act };
592 deliver(remote.inbox, accept, `${me}#main-key`, keys.private_pem).catch((e) => console.warn('[AP] Accept delivery failed:', e.message));
593 // Auto-backfill: send our recent posts as Create so the instance has our history
594 // (Mastodon doesn't fetch history on follow). ONCE PER REMOTE INSTANCE only —
595 // Mastodon dedupes notes per-instance, so re-filling an instance that already has
596 // a follower of ours is wasted work (and won't re-populate the new follower's
597 // timeline anyway). Deliver to the shared inbox (instance-level) when present.
598 // Sync insert+check (no await between) → no interleave race with concurrent Follows.
599 const instanceFilled = sharedInbox &&
600 db.prepare('SELECT 1 FROM ap_followers WHERE slug = ? AND shared_inbox = ? AND actor_uri != ? LIMIT 1')
601 .get(slug, sharedInbox, who);
602 if (!instanceFilled) {
603 backfillNewFollower(base, slug, sharedInbox || remote.inbox).catch(() => { /* best-effort */ });
604 }
605 console.log('[AP] Follow', who, '→', slug, verified ? '(sig ok)' : '(sig unverified)');
606 return 202;
607 }
608 if (type === 'Undo' && act.object) {
609 const who = typeof act.actor === 'string' ? act.actor : (act.actor && act.actor.id);
610 const ot = act.object.type;
611 if (ot === 'Follow') {
612 const obj = act.object.object;
613 const slug = slugParam || slugFromActorUrl(typeof obj === 'string' ? obj : (obj && obj.id));
614 if (who && slug) { fStmts().del.run(slug, who); console.log('[AP] Unfollow', who, '→', slug); }
615 return 202;
616 }
617 if (ot === 'Like' || ot === 'Announce') {
618 const tgt = act.object.object;
619 const pid = postIdFromNoteUrl(typeof tgt === 'string' ? tgt : (tgt && tgt.id), base);
620 if (who && pid) { iStmts().delLA.run(ot.toLowerCase(), pid, who); console.log('[AP] Undo', ot, who, '→', pid); }
621 return 202;
622 }
623 return 202;
624 }
625
626 const actorUri = typeof act.actor === 'string' ? act.actor : (act.actor && act.actor.id);
627 const resolveActor = async (uri) => ((verified && verified.id === uri) ? verified : await fetchActor(uri).catch(() => null));
628 // Activities from our OWN actors are already stored via ap_outbox — don't re-store.
629 const isLocalActor = !!(base && actorUri && actorUri.startsWith(`${base}/ap/users/`));
630
631 // Inbound reply: a Create whose object replies to one of our notes (post OR comment).
632 if (type === 'Create' && act.object && (act.object.type === 'Note' || act.object.type === 'Article')) {
633 const o = act.object;
634 const tgt = findThreadTarget(o.inReplyTo, base);
635 if (tgt && actorUri && !isLocalActor) {
636 const ai = actorInfo(await resolveActor(actorUri), actorUri);
637 const html = HtmlSanitizerService.sanitize(o.content || '');
638 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);
639 console.log('[AP] reply', actorUri, '→', tgt.post_id);
640 return 202;
641 }
642 // Home timeline (client): a top-level post from an account we follow.
643 if (actorUri && !isLocalActor && !o.inReplyTo && o.id) {
644 let subs = []; try { subs = db.prepare('SELECT slug FROM ap_following WHERE actor_uri = ?').all(actorUri); } catch { /* table may not exist yet */ }
645 if (subs.length) {
646 const ai = actorInfo(await resolveActor(actorUri), actorUri);
647 const html = HtmlSanitizerService.sanitize(o.content || '');
648 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));
649 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);
650 console.log('[AP] timeline +', actorUri, 'x' + subs.length);
651 }
652 }
653 return 202;
654 }
655 if (type === 'Like' || type === 'Announce') {
656 const tgt = act.object;
657 const pid = postIdFromNoteUrl(typeof tgt === 'string' ? tgt : (tgt && tgt.id), base);
658 if (pid && actorUri && !isLocalActor && localPostExists(pid)) {
659 const ai = actorInfo(await resolveActor(actorUri), actorUri);
660 iStmts().ins.run(type.toLowerCase(), pid, '', actorUri, ai.name, ai.handle, ai.url, ai.icon, null, null, null);
661 console.log('[AP]', type === 'Like' ? 'like' : 'boost', actorUri, '→', pid);
662 }
663 return 202;
664 }
665 if (type === 'Delete') {
666 // A remote note was deleted upstream → drop it from replies AND the timeline.
667 // Scope to the SIGNING actor so actor B can't delete actor A's content (the
668 // signature gate guarantees claimedActor == the verified signer here).
669 const oid = typeof act.object === 'string' ? act.object : (act.object && act.object.id);
670 if (oid && claimedActor) {
671 try { db.prepare('DELETE FROM ap_interactions WHERE object_uri = ? AND actor_uri = ?').run(oid, claimedActor); } catch { /* ignore */ }
672 try { db.prepare('DELETE FROM ap_timeline WHERE id = ? AND author_uri = ?').run(oid, claimedActor); } catch { /* ignore */ }
673 }
674 return 202;
675 }
676 // Accept/Reject of a Follow WE sent (client side).
677 if (type === 'Accept' && act.object) {
678 const fid = typeof act.object === 'string' ? act.object : (act.object && act.object.id);
679 if (fid) { try { fwStmts().acc.run(fid); } catch { /* ignore */ } }
680 console.log('[AP] follow accepted', actorUri);
681 return 202;
682 }
683 if (type === 'Reject' && act.object) {
684 const who = actorUri;
685 if (who && slugParam) { try { fwStmts().del.run(slugParam, who); } catch { /* ignore */ } }
686 return 202;
687 }
688
689 console.log('[AP] inbox', type || 'unknown', '→', slugParam || 'shared', '(ignored)');
690 return 202;
691}
692
693// Deliver a new post as Create(Note) to all followers' inboxes (fire-and-forget).
694// Needs PUBLIC_BASE_URL (absolute URLs); no-op without followers or base.
695export async function deliverCreate(site, post) {
696 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
697 if (!base || !site || !site.slug) return;
698 const followers = fStmts().list.all(site.slug);
699 if (!followers.length) return;
700 const inboxes = [...new Set(followers.map((f) => f.shared_inbox || f.inbox).filter(Boolean))];
701 const keys = getOrCreateKeys(site.slug);
702 const keyId = `${actorId(base, site.slug)}#main-key`;
703 const create = buildCreate(base, site, post);
704 for (const inbox of inboxes) deliverWithRetry(site.slug, inbox, create, keyId, keys.private_pem);
705}
706
707// On a new Follow, send that follower our most recent posts as Create so their
708// timeline shows our history (Mastodon does not backfill on follow). Oldest-first
709// so they sort into the follower's timeline at their original dates.
710async function backfillNewFollower(base, slug, inbox) {
711 if (!base || !slug || !inbox) return;
712 const site = db.prepare('SELECT * FROM sites WHERE slug = ?').get(slug);
713 if (!site) return;
714 const recent = db.prepare(
715 `SELECT id, slug, title, content, cover_image_url, published_at, created_at
716 FROM posts WHERE site_id = ? AND status = 'published' AND (fan_only IS NULL OR fan_only = 0)
717 ORDER BY COALESCE(published_at, created_at) DESC LIMIT 20`
718 ).all(site.id).reverse();
719 if (!recent.length) return;
720 const keys = getOrCreateKeys(slug);
721 const keyId = `${actorId(base, slug)}#main-key`;
722 for (const p of recent) {
723 try { await deliver(inbox, buildCreate(base, site, p), keyId, keys.private_pem); } catch { /* best-effort */ }
724 await new Promise((r) => setTimeout(r, 150));
725 }
726 console.log('[AP] backfilled', recent.length, 'posts to new follower of', slug);
727}
728
729// Tell followers a post is gone (Delete + Tombstone) so it's removed from their feeds.
730export async function deliverDelete(site, post) {
731 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
732 if (!base || !site || !site.slug || !post || !post.id) return;
733 const followers = fStmts().list.all(site.slug);
734 if (!followers.length) return;
735 const inboxes = [...new Set(followers.map((f) => f.shared_inbox || f.inbox).filter(Boolean))];
736 const keys = getOrCreateKeys(site.slug);
737 const me = actorId(base, site.slug);
738 const nid = noteId(base, post.id);
739 const del = {
740 '@context': 'https://www.w3.org/ns/activitystreams',
741 id: `${nid}#delete-${Date.now()}-${rid()}`,
742 type: 'Delete',
743 actor: me,
744 to: [PUBLIC],
745 object: { id: nid, type: 'Tombstone' },
746 };
747 for (const inbox of inboxes) deliverWithRetry(site.slug, inbox, del, `${me}#main-key`, keys.private_pem);
748}
749
750// Tell followers an already-published post changed (Update + edited Note) so
751// Mastodon refreshes the cached copy (e.g. after fixing content).
752export async function deliverUpdate(site, post) {
753 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
754 if (!base || !site || !site.slug || !post || !post.id) return;
755 const followers = fStmts().list.all(site.slug);
756 if (!followers.length) return;
757 const inboxes = [...new Set(followers.map((f) => f.shared_inbox || f.inbox).filter(Boolean))];
758 const keys = getOrCreateKeys(site.slug);
759 const me = actorId(base, site.slug);
760 const note = buildNote(base, site, post);
761 note.updated = new Date().toISOString();
762 const update = {
763 '@context': 'https://www.w3.org/ns/activitystreams',
764 id: `${noteId(base, post.id)}#update-${Date.now()}-${rid()}`,
765 type: 'Update', actor: me, to: [PUBLIC], cc: [`${me}/followers`],
766 object: note,
767 };
768 for (const inbox of inboxes) deliverWithRetry(site.slug, inbox, update, `${me}#main-key`, keys.private_pem);
769}
770
771// Tell followers the ACTOR changed (Update + Person) so Mastodon re-processes the
772// account AND re-fetches the featured (pinned) collection — there is no standard
773// "featured changed" activity, so this is how a pin/unpin propagates promptly.
774export async function deliverActorUpdate(site) {
775 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
776 if (!base || !site || !site.slug) return;
777 const followers = fStmts().list.all(site.slug);
778 if (!followers.length) return;
779 const inboxes = [...new Set(followers.map((f) => f.shared_inbox || f.inbox).filter(Boolean))];
780 const keys = getOrCreateKeys(site.slug);
781 const me = actorId(base, site.slug);
782 const update = {
783 '@context': ['https://www.w3.org/ns/activitystreams', 'https://w3id.org/security/v1'],
784 id: `${me}#update-${Date.now()}-${rid()}`,
785 type: 'Update', actor: me, to: [PUBLIC], cc: [`${me}/followers`],
786 object: buildActor(base, site),
787 };
788 for (const inbox of inboxes) deliverWithRetry(site.slug, inbox, update, `${me}#main-key`, keys.private_pem);
789}
790
791// Reliably set the pinned order on followers' instances via Add/Remove activities
792// (how Mastodon itself federates pins) — pushed to the inbox + processed immediately,
793// unlike the featured COLLECTION which Mastodon caches with sticky StatusPins.
794// Mastodon's Add skips an already-pinned status, so we REMOVE every pin first, wait,
795// then ADD in rank-DESCENDING order (rank 1 added LAST → newest StatusPin → shown first,
796// because Mastodon displays pins newest-first). `alsoRemove` = ids to unpin too.
797export async function resyncFeaturedPins(site, alsoRemove = []) {
798 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
799 if (!base || !site || !site.slug) return;
800 const followers = fStmts().list.all(site.slug);
801 if (!followers.length) return;
802 const inboxes = [...new Set(followers.map((f) => f.shared_inbox || f.inbox).filter(Boolean))];
803 const keys = getOrCreateKeys(site.slug);
804 const me = actorId(base, site.slug);
805 const keyId = `${me}#main-key`;
806 const featured = `${me}/featured`;
807 const AS = 'https://www.w3.org/ns/activitystreams';
808 const note = (id) => noteId(base, id);
809 const pinned = db.prepare(
810 `SELECT id FROM posts WHERE site_id = ? AND status = 'published' AND (fan_only IS NULL OR fan_only = 0)
811 AND pinned IS NOT NULL AND pinned > 0
812 ORDER BY pinned DESC, COALESCE(published_at, created_at) ASC LIMIT 20`
813 ).all(site.id);
814 const removeIds = [...new Set([...pinned.map((p) => p.id), ...alsoRemove])];
815 // 1. Remove every current pin so Mastodon can recreate them in order.
816 for (const id of removeIds) {
817 const rm = { '@context': AS, id: `${me}#rm-${id}-${Date.now()}-${rid()}`, type: 'Remove', actor: me, object: note(id), target: featured, to: [PUBLIC] };
818 for (const inbox of inboxes) deliver(inbox, rm, keyId, keys.private_pem).catch(() => { /* best-effort */ });
819 }
820 if (!pinned.length) { console.log('[AP] unpinned all featured for', site.slug); return; }
821 await new Promise((r) => setTimeout(r, 5000)); // let the Removes land first
822 // 2. Add in rank-DESC order, gaps so each StatusPin gets an increasing created_at.
823 for (const p of pinned) {
824 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`] };
825 for (const inbox of inboxes) deliver(inbox, add, keyId, keys.private_pem).catch(() => { /* best-effort */ });
826 await new Promise((r) => setTimeout(r, 2000));
827 }
828 console.log('[AP] resynced', pinned.length, 'featured pins for', site.slug);
829}
830
831// ── outbound replies (Klonkt → fediverse) ─────────────────────────
832const escHtml = (s) => String(s || '').replace(/[<>&]/g, (c) => ({ '<': '&lt;', '>': '&gt;', '&': '&amp;' }[c]));
833const 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(); };
834
835// Build one of OUR outbound reply Notes from an ap_outbox row.
836export function buildReplyNote(base, site, row) {
837 const me = actorId(base, site.slug);
838 return {
839 id: noteId(base, row.id),
840 type: 'Note',
841 attributedTo: me,
842 inReplyTo: row.in_reply_to || undefined,
843 content: row.content,
844 url: row.post_slug ? `${base}/${encodeURIComponent(row.post_slug)}` : undefined,
845 published: toISO(row.created_at),
846 to: row.to_actor ? [row.to_actor] : [PUBLIC],
847 cc: [PUBLIC, `${me}/followers`],
848 tag: row.to_actor ? [{ type: 'Mention', href: row.to_actor, name: row.to_handle }] : [],
849 };
850}
851
852// Resolve one of our outbound reply Notes by id (for /ap/notes/:id fallback).
853export function getOutboxNote(base, id) {
854 const row = iStmts().getO.get(id);
855 if (!row) return null;
856 const site = db.prepare('SELECT * FROM sites WHERE slug = ?').get(row.site_slug);
857 if (!site) return null;
858 return buildReplyNote(base, site, row);
859}
860
861// Send a reply FROM this site to a remote actor (in reply to their inbound reply).
862// `parent` = an ap_interactions row (actor_uri, actor_url, actor_handle, object_uri).
863export async function deliverReply(site, { postId, postSlug, parent, text }) {
864 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
865 if (!base || !site || !site.slug || !parent || !String(text || '').trim()) return null;
866 const me = actorId(base, site.slug);
867 const handle = parent.actor_handle || deriveHandle(parent.actor_uri);
868 const body = escHtml(String(text).trim()).replace(/\r?\n/g, '<br>');
869 const mention = parent.actor_uri
870 ? `<a href="${escHtml(parent.actor_url || parent.actor_uri)}" class="u-url mention">${escHtml(handle)}</a> ` : '';
871 const content = `<p>${mention}${body}</p>`;
872 // Dedup: skip if the exact same reply was already sent (double-submit guard).
873 const dup = db.prepare('SELECT 1 FROM ap_outbox WHERE site_slug = ? AND IFNULL(in_reply_to, \'\') = ? AND content = ? LIMIT 1')
874 .get(site.slug, parent.object_uri || '', content);
875 if (dup) { console.log('[AP] outreply skipped (duplicate)'); return { duplicate: true, delivered: 0 }; }
876 const id = crypto.randomUUID();
877 iStmts().insO.run(id, site.slug, postId, postSlug || null, parent.object_uri || null, parent.actor_uri || null, handle, content);
878 const row = iStmts().getO.get(id);
879 const note = buildReplyNote(base, site, row);
880 const create = {
881 '@context': 'https://www.w3.org/ns/activitystreams',
882 id: note.id + '#create', type: 'Create', actor: me,
883 published: note.published, to: note.to, cc: note.cc, object: note,
884 };
885 const keys = getOrCreateKeys(site.slug);
886 const keyId = `${me}#main-key`;
887 const inboxes = new Set();
888 if (parent.actor_uri) {
889 const a = await fetchActor(parent.actor_uri).catch(() => null);
890 if (a) inboxes.add((a.endpoints && a.endpoints.sharedInbox) || a.inbox);
891 }
892 if (parent.threadInbox) inboxes.add(parent.threadInbox); // back-compat (single)
893 (parent.threadInboxes || []).forEach((i) => inboxes.add(i)); // whole ancestor chain
894 for (const f of fStmts().list.all(site.slug)) inboxes.add(f.shared_inbox || f.inbox);
895 inboxes.delete(`${me}/inbox`); // never deliver to ourselves (already in ap_outbox)
896 inboxes.delete(`${base}/ap/inbox`); // (our own shared inbox) → avoids a self-duplicate
897 let delivered = 0;
898 for (const inbox of [...inboxes].filter(Boolean)) {
899 try { const st = await deliver(inbox, create, keyId, keys.private_pem); if (st >= 200 && st < 300) delivered++; } catch { /* best-effort */ }
900 }
901 console.log('[AP] outreply', site.slug, '→', parent.actor_uri, 'delivered', delivered);
902 return { id, content, delivered };
903}
904
905// Resolve a remote post URL (any fediverse/Klonkt post) into a reply target.
906// Returns a parent-shaped object usable by deliverReply(), or null.
907export async function resolveRemoteNote(url) {
908 if (!/^https?:\/\//i.test(String(url || ''))) return null;
909 const note = await fetchActor(url).catch(() => null); // AP GET (content-negotiates)
910 if (!note || !note.id) return null;
911 const att = note.attributedTo;
912 const actorUri = typeof att === 'string' ? att : (att && att.id);
913 if (!actorUri) return null;
914 const actor = await fetchActor(actorUri).catch(() => null);
915 const ai = actorInfo(actor, actorUri);
916 // Is what we're replying to a post (or a comment) on one of OUR posts? If so,
917 // link our reply to that local post so it shows nested in the post thread.
918 const localTgt = findThreadTarget(note.id, (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, ''));
919 // Walk the WHOLE reply chain upward (comment → parent comment → … → root post)
920 // and collect every ancestor author's inbox, so each participant's server —
921 // including the original post's author — receives + threads our reply.
922 const threadInboxes = [];
923 const seenInbox = new Set();
924 let cursor = note.inReplyTo, guard = 0;
925 while (cursor && guard++ < 6) {
926 const url = typeof cursor === 'string' ? cursor : (cursor && cursor.id);
927 if (!url) break;
928 const pn = await fetchActor(url).catch(() => null);
929 if (!pn) break;
930 const pa = typeof pn.attributedTo === 'string' ? pn.attributedTo : (pn.attributedTo && pn.attributedTo.id);
931 if (pa && pa !== actorUri) {
932 const paDoc = await fetchActor(pa).catch(() => null);
933 const inbox = paDoc && ((paDoc.endpoints && paDoc.endpoints.sharedInbox) || paDoc.inbox);
934 if (inbox && !seenInbox.has(inbox)) { seenInbox.add(inbox); threadInboxes.push(inbox); }
935 }
936 cursor = pn.inReplyTo; // climb to the next ancestor
937 }
938 const rawHtml = String(note.content || '').replace(/\[\[(track|album|playlist):[^\]]+\]\]/gi, '');
939 const images = (Array.isArray(note.attachment) ? note.attachment : [])
940 .filter((a) => a && a.url && (!a.mediaType || /^image\//i.test(a.mediaType)))
941 .map((a) => safeUrl(a.url)).filter(Boolean);
942 return {
943 object_uri: safeUrl(note.id) || note.id,
944 actor_uri: actorUri,
945 actor_url: ai.url,
946 actor_handle: ai.handle,
947 actor_name: ai.name,
948 actor_icon: ai.icon,
949 url: note.url || url,
950 content: HtmlSanitizerService.sanitize(rawHtml), // full, sanitized
951 images,
952 threadInboxes, // every ancestor author's inbox
953 localPostId: localTgt ? localTgt.post_id : '', // our post this belongs to (if any)
954 preview: HtmlSanitizerService.toPlainText(note.content || '').slice(0, 240),
955 };
956}
957
958// List a site's own outbound fediverse replies (for the manage/delete view).
959export function listOutbox(siteSlug) {
960 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);
961}
962
963// Delete one of our outbound replies: send Delete(Tombstone) to recipients + remove it.
964export async function deliverOutboxDelete(site, outboxId) {
965 const row = iStmts().getO.get(outboxId);
966 if (!row || row.site_slug !== site.slug) return false;
967 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
968 if (base) {
969 const me = actorId(base, site.slug);
970 const nid = noteId(base, row.id);
971 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' } };
972 const keys = getOrCreateKeys(site.slug);
973 const inboxes = new Set();
974 if (row.to_actor) { const a = await fetchActor(row.to_actor).catch(() => null); if (a) inboxes.add((a.endpoints && a.endpoints.sharedInbox) || a.inbox); }
975 for (const f of fStmts().list.all(site.slug)) inboxes.add(f.shared_inbox || f.inbox);
976 for (const inbox of [...inboxes].filter(Boolean)) { try { await deliver(inbox, del, `${me}#main-key`, keys.private_pem); } catch { /* best-effort */ } }
977 }
978 db.prepare('DELETE FROM ap_outbox WHERE id = ?').run(outboxId);
979 return true;
980}
981
982// ── Fediverse CLIENT: follow accounts + home timeline ─────────────
983// Resolve an @user@domain handle to its actor URL via WebFinger.
984export async function webfingerResolve(handle) {
985 const h = String(handle || '').trim().replace(/^@/, '');
986 const parts = h.split('@');
987 if (parts.length !== 2 || !parts[0] || !parts[1]) return null;
988 const acct = `${parts[0]}@${parts[1]}`;
989 try {
990 const r = await safeFetch(`https://${parts[1]}/.well-known/webfinger?resource=acct:${encodeURIComponent(acct)}`,
991 { headers: { Accept: 'application/jrd+json, application/json' } });
992 if (!r.ok) return null;
993 const jrd = await r.json();
994 const link = (jrd.links || []).find((l) => l.rel === 'self' && /activity\+json|ld\+json/.test(l.type || ''));
995 return safeUrl(link ? link.href : '') || null;
996 } catch { return null; }
997}
998
999let _insFw, _delFw, _listFw, _accFw, _oneFw;
1000function fwStmts() {
1001 if (!_insFw) {
1002 _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)');
1003 _delFw = db.prepare('DELETE FROM ap_following WHERE slug = ? AND actor_uri = ?');
1004 _listFw = db.prepare('SELECT * FROM ap_following WHERE slug = ? ORDER BY created_at DESC');
1005 _accFw = db.prepare("UPDATE ap_following SET status = 'accepted' WHERE follow_id = ?");
1006 _oneFw = db.prepare('SELECT * FROM ap_following WHERE slug = ? AND actor_uri = ?');
1007 }
1008 return { ins: _insFw, del: _delFw, list: _listFw, acc: _accFw, one: _oneFw };
1009}
1010export function listFollowing(slug) { return fwStmts().list.all(slug); }
1011
1012let _insTl, _listTl, _delTl;
1013function tlStmts() {
1014 if (!_insTl) {
1015 _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)');
1016 _listTl = db.prepare('SELECT * FROM ap_timeline WHERE slug = ? ORDER BY COALESCE(published, created_at) DESC LIMIT ?');
1017 _delTl = db.prepare('DELETE FROM ap_timeline WHERE id = ?');
1018 }
1019 return { ins: _insTl, list: _listTl, del: _delTl };
1020}
1021export function getTimeline(slug, limit) { return tlStmts().list.all(slug, limit || 50); }
1022
1023// Follow a fediverse account by @handle (WebFinger → actor → signed Follow).
1024export async function followActor(site, handle) {
1025 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
1026 if (!base || !site || !site.slug) return { error: 'config' };
1027 const actorUrl = await webfingerResolve(handle);
1028 if (!actorUrl) return { error: 'not_found' };
1029 const actor = await fetchActor(actorUrl).catch(() => null);
1030 if (!actor || !actor.id || !actor.inbox) return { error: 'unreachable' };
1031 const ai = actorInfo(actor, actor.id);
1032 const me = actorId(base, site.slug);
1033 const keys = getOrCreateKeys(site.slug);
1034 const followId = `${me}#follow-${Date.now()}-${rid()}`;
1035 fwStmts().ins.run(site.slug, actor.id, ai.handle, ai.name, ai.icon, ai.url, actor.inbox, followId, 'pending');
1036 const follow = { '@context': 'https://www.w3.org/ns/activitystreams', id: followId, type: 'Follow', actor: me, object: actor.id };
1037 try { await deliver(actor.inbox, follow, `${me}#main-key`, keys.private_pem); }
1038 catch (e) { console.warn('[AP] follow deliver failed:', e.message); }
1039 console.log('[AP] follow', site.slug, '→', actor.id);
1040 return { ok: true, name: ai.name, handle: ai.handle };
1041}
1042
1043export async function unfollowActor(site, actorUri) {
1044 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
1045 const me = actorId(base, site.slug);
1046 const keys = getOrCreateKeys(site.slug);
1047 const row = fwStmts().one.get(site.slug, actorUri);
1048 if (row && row.inbox) {
1049 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 } };
1050 try { await deliver(row.inbox, undo, `${me}#main-key`, keys.private_pem); } catch { /* best-effort */ }
1051 }
1052 fwStmts().del.run(site.slug, actorUri);
1053 return { ok: true };
1054}
1055
1056// Send a Like or Announce (boost) on a remote note FROM this site.
1057export async function sendInteraction(site, kind, targetNoteId, authorUri) {
1058 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
1059 if (!base || !site || !site.slug || !targetNoteId) return { error: 'config' };
1060 const type = kind === 'boost' ? 'Announce' : 'Like';
1061 const me = actorId(base, site.slug);
1062 const keys = getOrCreateKeys(site.slug);
1063 const act = {
1064 '@context': 'https://www.w3.org/ns/activitystreams',
1065 id: `${me}#${type.toLowerCase()}-${Date.now()}-${rid()}`,
1066 type, actor: me, object: targetNoteId,
1067 };
1068 if (type === 'Announce') { act.to = [PUBLIC]; act.cc = [`${me}/followers`]; }
1069 const inboxes = new Set();
1070 if (authorUri) { const a = await fetchActor(authorUri).catch(() => null); if (a) inboxes.add((a.endpoints && a.endpoints.sharedInbox) || a.inbox); }
1071 // A boost is public → also deliver to our own followers so it shows for them.
1072 if (type === 'Announce') { for (const f of fStmts().list.all(site.slug)) inboxes.add(f.shared_inbox || f.inbox); }
1073 let delivered = 0;
1074 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 */ } }
1075 console.log('[AP]', type, site.slug, '→', targetNoteId, 'delivered', delivered);
1076 return { ok: true, delivered };
1077}
1078
1079// Notifications inbox: new followers + replies/likes/boosts on this site's posts.
1080export function getNotifications(slug, limit) {
1081 const out = [];
1082 try {
1083 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)) {
1084 out.push({ type: 'follow', handle: deriveHandle(f.actor_uri), url: f.actor_uri, created_at: f.created_at });
1085 }
1086 } catch { /* ignore */ }
1087 try {
1088 const rows = db.prepare(`
1089 SELECT i.kind, i.actor_name, i.actor_handle, i.actor_url, i.content, i.created_at,
1090 p.slug AS post_slug, p.title AS post_title
1091 FROM ap_interactions i LEFT JOIN posts p ON p.id = i.post_id
1092 WHERE p.site_id = (SELECT id FROM sites WHERE slug = ?)
1093 ORDER BY i.created_at DESC LIMIT 80
1094 `).all(slug);
1095 for (const r of rows) out.push({
1096 type: r.kind, name: r.actor_name, handle: r.actor_handle, url: r.actor_url,
1097 content: r.content, post_slug: r.post_slug, post_title: r.post_title, created_at: r.created_at,
1098 });
1099 } catch { /* ignore */ }
1100 out.sort((a, b) => new Date(b.created_at) - new Date(a.created_at));
1101 return out.slice(0, limit || 60);
1102}
1103
1104// ── Blocking / defederation ───────────────────────────────────────
1105let _insBl, _delBl, _listBl;
1106function blStmts() {
1107 if (!_insBl) {
1108 _insBl = db.prepare('INSERT OR IGNORE INTO ap_blocks (slug, target, kind, label, created_at) VALUES (?,?,?,?,CURRENT_TIMESTAMP)');
1109 _delBl = db.prepare('DELETE FROM ap_blocks WHERE slug = ? AND target = ?');
1110 _listBl = db.prepare('SELECT * FROM ap_blocks WHERE slug = ? ORDER BY created_at DESC');
1111 }
1112 return { ins: _insBl, del: _delBl, list: _listBl };
1113}
1114export function listBlocks(slug) { return blStmts().list.all(slug); }
1115
1116// True if an actor (or its whole domain) is blocked anywhere on this instance.
1117export function isBlockedAny(actorUri) {
1118 if (!actorUri) return false;
1119 let domain = ''; try { domain = new URL(actorUri).host; } catch { /* ignore */ }
1120 try { return !!db.prepare("SELECT 1 FROM ap_blocks WHERE (kind='actor' AND target=?) OR (kind='domain' AND target=?) LIMIT 1").get(actorUri, domain); }
1121 catch { return false; }
1122}
1123
1124function purgeBlocked(kind, target) {
1125 try {
1126 if (kind === 'domain') {
1127 const like = `%//${target}/%`;
1128 db.prepare('DELETE FROM ap_interactions WHERE actor_uri LIKE ?').run(like);
1129 db.prepare('DELETE FROM ap_timeline WHERE author_uri LIKE ?').run(like);
1130 db.prepare('DELETE FROM ap_followers WHERE actor_uri LIKE ?').run(like);
1131 } else {
1132 db.prepare('DELETE FROM ap_interactions WHERE actor_uri = ?').run(target);
1133 db.prepare('DELETE FROM ap_timeline WHERE author_uri = ?').run(target);
1134 db.prepare('DELETE FROM ap_followers WHERE actor_uri = ?').run(target);
1135 }
1136 } catch { /* best-effort */ }
1137}
1138
1139// Block an actor (@handle or actor URL) or a whole domain; purges their content.
1140export async function blockTarget(site, input) {
1141 const raw = String(input || '').trim();
1142 if (!site || !site.slug || !raw) return { error: 'empty' };
1143 let kind, target, label;
1144 if (/^https?:\/\//i.test(raw)) { kind = 'actor'; target = raw; label = raw; }
1145 else if (raw.includes('@')) {
1146 const actorUrl = await webfingerResolve(raw);
1147 if (!actorUrl) return { error: 'not_found' };
1148 kind = 'actor'; target = actorUrl; label = raw.startsWith('@') ? raw : ('@' + raw);
1149 } else { kind = 'domain'; target = raw.toLowerCase(); label = raw.toLowerCase(); }
1150 blStmts().ins.run(site.slug, target, kind, label);
1151 purgeBlocked(kind, target);
1152 console.log('[AP] block', site.slug, kind, target);
1153 return { ok: true, label };
1154}
1155
1156export function unblock(site, target) { blStmts().del.run(site.slug, target); return { ok: true }; }
1157
1158export default {
1159 getOrCreateKeys, apWants, sendAP, actorId, noteId,
1160 buildActor, buildNote, buildCreate, buildOutbox, buildFollowers, buildFeatured,
1161 followerCount, deliver, fetchActor, verifyRequest, handleInbox, deliverCreate, deliverDelete, deliverUpdate, deliverActorUpdate, resyncFeaturedPins,
1162 getInteractions, getInteractionById, buildReplyNote, getOutboxNote, deliverReply, resolveRemoteNote,
1163 listOutbox, deliverOutboxDelete,
1164 webfingerResolve, followActor, unfollowActor, listFollowing, getTimeline, sendInteraction,
1165 getNotifications, listBlocks, isBlockedAny, blockTarget, unblock,
1166 deliverWithRetry, enqueueDelivery, processDeliveryQueue, startDeliveryWorker,
1167 getReplyUris, markNotificationsSeen, countUnseenNotifications, hasPlayableAudio,
1168};
Note: See TracBrowser for help on using the repository browser.