source: Klonkt/src/services/ActivityPubService.js@ 356caeb

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

feat(fediverse): auto-boost a followed account ('feature an artist')

Phase 1 of the Cirkels-on-AP rework. A followed account can be marked auto-boost:
their new top-level posts are automatically re-Announced (boosted) to your own
followers. Toggle in the /tijdlijn following list + a checkbox on the follow box.
ap_following.auto_boost column (migrated), setAutoBoost(), hook in handleInbox's
timeline loop, POST /tijdlijn/autoboost. nl/en/de.

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

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