source: Klonkt/src/services/ActivityPubService.js@ 00a2bcd

main
Last change on this file since 00a2bcd was cd6a008, checked in by Robin Genis <roboburr@…>, 2 months ago

fix(fedi): preserve newlines when federating to Mastodon

Klonkt renders post content with white-space:pre-wrap, so raw \n are line breaks on the
site, but Mastodon collapses whitespace and dropped them (a poem arrived as one block).
buildNote now converts newlines to <br> for the federated copy. Content authored with
shift+enter already uses <br> (no \n) → no-op there.

  • Property mode set to 100644
File size: 73.7 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 // Klonkt renders post content with white-space:pre-wrap, so raw newlines ARE line
226 // breaks on the site. Mastodon (plain HTML) collapses whitespace and would drop them,
227 // so convert newlines to <br> for the federated copy (content already made with
228 // shift+enter uses <br> and has no \n → this is a no-op there).
229 body = body.replace(/\r?\n/g, '<br>');
230 const seen = new Set();
231 const attachment = urls.filter(Boolean)
232 .filter((u) => { if (seen.has(u)) return false; seen.add(u); return true; })
233 .map((u) => ({ type: 'Document', mediaType: mediaType(u), url: u }));
234
235 const note = {
236 id,
237 type: 'Note',
238 attributedTo: aId,
239 content: titleHtml + body,
240 url: human,
241 published: new Date(post.published_at || post.created_at || Date.now()).toISOString(),
242 to: [PUBLIC],
243 cc: [`${aId}/followers`],
244 tag: Array.isArray(post.tags) ? post.tags.map((t) => ({ type: 'Hashtag', name: '#' + String(t).replace(/\s+/g, '') })) : [],
245 replies: `${id}/replies`,
246 };
247 if (attachment.length) note.attachment = attachment;
248 // Playable-audio posts suppress the cover attachment (player card). Still expose
249 // the cover via AS2 `image` so card/grid consumers (the Klonkt Cirkel) can show
250 // it — Mastodon ignores a Note's `image`, so the player card is unaffected.
251 if (post.cover_image_url && playable) {
252 const cov = abs(post.cover_image_url);
253 if (cov) note.image = { type: 'Image', mediaType: mediaType(cov), url: cov };
254 }
255 return note;
256}
257
258// All reply note URIs on a local post (inbound fediverse replies + our own
259// outbound replies) — backs the Note's `replies` Collection so remote servers
260// can fetch the whole thread.
261export function getReplyUris(base, postId) {
262 const out = [];
263 try {
264 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);
265 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}`);
266 } catch { /* non-fatal */ }
267 return out;
268}
269
270// Notifications "seen" tracking → a real bell badge. Stored per site in app_settings.
271export function markNotificationsSeen(slug) {
272 try {
273 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")
274 .run(`fedi_notif_seen:${slug}`, new Date().toISOString());
275 } catch { /* non-fatal */ }
276}
277export function countUnseenNotifications(slug) {
278 try {
279 const row = db.prepare('SELECT value FROM app_settings WHERE key = ?').get(`fedi_notif_seen:${slug}`);
280 const seen = row ? Date.parse(row.value) : 0;
281 let n = 0;
282 for (const it of getNotifications(slug, 50)) { if (Date.parse(it.created_at) > seen) n++; }
283 return n;
284 } catch { return 0; }
285}
286
287export function buildCreate(base, site, post) {
288 const note = buildNote(base, site, post);
289 return {
290 '@context': 'https://www.w3.org/ns/activitystreams',
291 id: note.id + '#create',
292 type: 'Create',
293 actor: actorId(base, site.slug),
294 published: note.published,
295 to: note.to,
296 cc: note.cc,
297 object: note,
298 };
299}
300
301export function buildOutbox(base, site, posts) {
302 const id = `${actorId(base, site.slug)}/outbox`;
303 const items = (posts || []).slice(0, MAX_OUTBOX).map((p) => buildCreate(base, site, p));
304 return {
305 '@context': 'https://www.w3.org/ns/activitystreams',
306 id,
307 type: 'OrderedCollection',
308 totalItems: items.length,
309 orderedItems: items,
310 };
311}
312
313export function buildFollowers(base, site, count) {
314 const id = `${actorId(base, site.slug)}/followers`;
315 return {
316 '@context': 'https://www.w3.org/ns/activitystreams',
317 id,
318 type: 'OrderedCollection',
319 totalItems: count || 0,
320 orderedItems: [], // hidden for privacy; count only
321 };
322}
323
324// Pinned posts → the actor's `featured` collection. Mastodon reads this and shows
325// these as the "Featured" tab (pinned to the profile). Posts come ordered by pin
326// rank; embedded as full Notes so a remote server doesn't need extra fetches.
327export function buildFeatured(base, site, posts) {
328 const id = `${actorId(base, site.slug)}/featured`;
329 const items = (posts || []).map((p) => buildNote(base, site, p));
330 return {
331 '@context': 'https://www.w3.org/ns/activitystreams',
332 id,
333 type: 'OrderedCollection',
334 totalItems: items.length,
335 orderedItems: items,
336 };
337}
338
339// ── followers store (lazy stmts) ──────────────────────────────────
340let _insF, _delF, _listF, _cntF;
341function fStmts() {
342 if (!_insF) {
343 _insF = db.prepare('INSERT OR IGNORE INTO ap_followers (slug, actor_uri, inbox, shared_inbox, created_at) VALUES (?,?,?,?,CURRENT_TIMESTAMP)');
344 _delF = db.prepare('DELETE FROM ap_followers WHERE slug = ? AND actor_uri = ?');
345 _listF = db.prepare('SELECT inbox, shared_inbox FROM ap_followers WHERE slug = ?');
346 _cntF = db.prepare('SELECT COUNT(*) n FROM ap_followers WHERE slug = ?');
347 }
348 return { ins: _insF, del: _delF, list: _listF, cnt: _cntF };
349}
350export function followerCount(slug) { return fStmts().cnt.get(slug).n; }
351
352// ── inbound interactions store (replies / likes / boosts) + our outbound replies ──
353let _insI, _delLA, _delReply, _listI, _getI, _insO, _listO, _getO;
354function iStmts() {
355 if (!_insI) {
356 _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)');
357 _delLA = db.prepare('DELETE FROM ap_interactions WHERE kind = ? AND post_id = ? AND actor_uri = ?');
358 _delReply = db.prepare("DELETE FROM ap_interactions WHERE kind = 'reply' AND object_uri = ?");
359 _listI = db.prepare('SELECT id, kind, object_uri, parent_uri, actor_uri, actor_name, actor_handle, actor_url, actor_icon, content, published, created_at, acted_boost, acted_like FROM ap_interactions WHERE post_id = ? ORDER BY created_at ASC');
360 _getI = db.prepare('SELECT * FROM ap_interactions WHERE id = ?');
361 _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)');
362 _listO = db.prepare('SELECT * FROM ap_outbox WHERE post_id = ? ORDER BY created_at ASC');
363 _getO = db.prepare('SELECT * FROM ap_outbox WHERE id = ?');
364 }
365 return { ins: _insI, delLA: _delLA, delReply: _delReply, list: _listI, getI: _getI, insO: _insO, listO: _listO, getO: _getO };
366}
367
368export function getInteractionById(id) { return iStmts().getI.get(id); }
369export function setInteractionBoosted(id, on) {
370 db.prepare('UPDATE ap_interactions SET acted_boost = ? WHERE id = ?').run(on ? 1 : 0, id);
371}
372export function setInteractionLiked(id, on) {
373 db.prepare('UPDATE ap_interactions SET acted_like = ? WHERE id = ?').run(on ? 1 : 0, id);
374}
375// Your like/boost state on a REMOTE post (interact page toggles).
376export function setMyReaction(slug, uri, kind, on) {
377 if (on) db.prepare('INSERT OR IGNORE INTO ap_my_reactions (site_slug, target_uri, kind) VALUES (?,?,?)').run(slug, uri, kind);
378 else db.prepare('DELETE FROM ap_my_reactions WHERE site_slug = ? AND target_uri = ? AND kind = ?').run(slug, uri, kind);
379}
380export function getMyReactions(slug, uri) {
381 const rows = (slug && uri) ? db.prepare('SELECT kind FROM ap_my_reactions WHERE site_slug = ? AND target_uri = ?').all(slug, uri) : [];
382 return { liked: rows.some((r) => r.kind === 'like'), boosted: rows.some((r) => r.kind === 'boost') };
383}
384
385const localPostExists = (id) => { try { return !!db.prepare('SELECT 1 FROM posts WHERE id = ?').get(id); } catch { return false; } };
386// Extract our local post id from a note URL, but only if it's ours (base match).
387function postIdFromNoteUrl(url, base) {
388 const s = String(url || '');
389 if (base && !s.startsWith(base)) return null;
390 const m = s.match(/\/ap\/notes\/([^/?#]+)/);
391 return m ? decodeURIComponent(m[1]) : null;
392}
393function deriveHandle(actorUri) {
394 try { const u = new URL(actorUri); const seg = u.pathname.split('/').filter(Boolean).pop() || ''; return `@${seg}@${u.host}`; } catch { return String(actorUri || ''); }
395}
396function actorInfo(doc, actorUri) {
397 let host = ''; try { host = new URL(actorUri).host; } catch { /* keep empty */ }
398 const handle = doc && doc.preferredUsername ? `@${doc.preferredUsername}@${host}` : deriveHandle(actorUri);
399 const icon = doc && doc.icon ? (doc.icon.url || (Array.isArray(doc.icon) && doc.icon[0] && doc.icon[0].url)) : null;
400 return {
401 name: (doc && (doc.name || doc.preferredUsername)) || handle,
402 handle,
403 url: safeUrl((doc && (doc.url || doc.id)) || actorUri) || null,
404 icon: safeUrl(icon) || null,
405 };
406}
407
408// Given an inReplyTo note URL, find which local post the thread belongs to + the
409// note being replied to (parent), so a reply-to-a-comment can be nested.
410function findThreadTarget(inReplyTo, base) {
411 if (!inReplyTo) return null;
412 const seg = postIdFromNoteUrl(inReplyTo, base); // our /ap/notes/<id> segment (if ours)
413 if (seg && localPostExists(seg)) return { post_id: seg, parent_uri: inReplyTo };
414 if (seg) {
415 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 */ }
416 }
417 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 */ }
418 return null;
419}
420
421// Drop the leading @mention(s) a federated reply carries (the person being replied to),
422// so a comment reads "dope tekening ouwe" instead of "@jason@jasonhacky.nl dope …".
423// Keeps a leading <p> wrapper; handles mention <a> links and plain-text @user@domain.
424export function stripLeadingMentions(html) {
425 if (!html) return html;
426 let s = String(html);
427 s = s.replace(/^(\s*<p[^>]*>)?\s*(?:<a\b[^>]*>\s*@[^<]+<\/a>[  ]*)+/i, (m, p) => p || '');
428 s = s.replace(/^(\s*<p[^>]*>)?\s*(?:@[\w.-]+(?:@[\w.-]+)?[  ]+)+/i, (m, p) => p || '');
429 return s;
430}
431
432// View-ready threaded view of a post's fediverse activity (inbound replies +
433// our outbound replies, nested), plus like/boost counts.
434export function getInteractions(postId, base, site) {
435 const s = iStmts();
436 const rows = s.list.all(postId);
437 const baseClean = (base || process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
438 const postNoteId = baseClean ? `${baseClean}/ap/notes/${postId}` : null;
439 // Our own (outbound) replies show the SITE identity for everyone (not "You").
440 let host = ''; try { host = new URL(baseClean).host; } catch { /* ignore */ }
441 const siteName = (site && (site.title || site.slug)) || '';
442 const siteHandle = (site && site.slug && host) ? `@${site.slug}@${host}` : '';
443 const siteUrl = baseClean ? `${baseClean}/` : '';
444 const siteIcon = (site && site.profile_photo) || null;
445
446 const nodes = [];
447 for (const r of rows) {
448 if (r.kind !== 'reply') continue;
449 nodes.push({
450 noteId: r.object_uri, parent: r.parent_uri || null, mine: false, id: r.id,
451 actor_name: r.actor_name, actor_handle: r.actor_handle, actor_url: r.actor_url,
452 actor_icon: r.actor_icon, content: stripLeadingMentions(r.content), created_at: r.published || r.created_at,
453 acted_boost: !!r.acted_boost, acted_like: !!r.acted_like,
454 children: [],
455 });
456 }
457 for (const o of s.listO.all(postId)) {
458 nodes.push({
459 noteId: baseClean ? `${baseClean}/ap/notes/${o.id}` : o.id, parent: o.in_reply_to || null,
460 mine: true, outboxId: o.id, content: stripLeadingMentions(o.content), created_at: o.created_at,
461 actor_name: siteName, actor_handle: siteHandle, actor_url: siteUrl, actor_icon: siteIcon,
462 children: [],
463 });
464 }
465
466 const byId = new Map(nodes.map((n) => [n.noteId, n]));
467 const isTop = (n) => !n.parent || n.parent === postNoteId || !byId.has(n.parent);
468 const tops = [];
469 for (const n of nodes) {
470 if (isTop(n)) { tops.push(n); continue; }
471 let anc = n, guard = 0;
472 while (!isTop(anc) && guard++ < 12) anc = byId.get(anc.parent);
473 anc.children.push(n);
474 }
475 const byTime = (a, b) => new Date(a.created_at) - new Date(b.created_at);
476 tops.sort(byTime).forEach((t) => t.children.sort(byTime));
477
478 return {
479 thread: tops,
480 likeCount: rows.filter((r) => r.kind === 'like').length,
481 announceCount: rows.filter((r) => r.kind === 'announce').length,
482 total: nodes.length,
483 };
484}
485
486// ── HTTP Signatures + delivery ────────────────────────────────────
487const slugFromActorUrl = (url) => { const m = String(url || '').match(/\/ap\/users\/([^/?#]+)/); return m ? decodeURIComponent(m[1]) : null; };
488
489// Sign + POST an activity to a remote inbox (draft-cavage HTTP Signatures, RSA-SHA256).
490export async function deliver(inboxUrl, bodyObj, keyId, privatePem) {
491 const body = JSON.stringify(bodyObj);
492 const u = new URL(inboxUrl);
493 const date = new Date().toUTCString();
494 const digest = 'SHA-256=' + crypto.createHash('sha256').update(body).digest('base64');
495 const signingString = `(request-target): post ${u.pathname}\nhost: ${u.host}\ndate: ${date}\ndigest: ${digest}`;
496 const signature = crypto.sign('sha256', Buffer.from(signingString), privatePem).toString('base64');
497 const sig = `keyId="${keyId}",algorithm="rsa-sha256",headers="(request-target) host date digest",signature="${signature}"`;
498 const r = await safeFetch(inboxUrl, {
499 method: 'POST',
500 headers: { 'Content-Type': 'application/activity+json', Accept: 'application/activity+json', Date: date, Digest: digest, Signature: sig },
501 body,
502 });
503 return r.status;
504}
505
506export async function fetchActor(url) {
507 try {
508 const r = await safeFetch(url, { headers: { Accept: 'application/activity+json' } });
509 if (!r.ok) return null;
510 const len = Number(r.headers.get('content-length') || 0);
511 if (len > 2_000_000) return null; // refuse oversized actor docs
512 return await r.json();
513 } catch { return null; }
514}
515
516// ── Delivery queue with retries ───────────────────────────────────
517// Outbound deliveries are tried immediately; on failure (down server, timeout,
518// non-2xx) they're queued and retried with backoff so a briefly-offline follower
519// doesn't silently miss the post. The signing key is NOT stored — the worker
520// re-derives it from the actor slug at send time.
521const DELIVERY_MAX_ATTEMPTS = 6;
522const DELIVERY_BACKOFF_MIN = [1, 5, 15, 60, 180, 360];
523let _insDeliv, _dueDeliv, _delDeliv, _bumpDeliv;
524function deliveryStmts() {
525 if (!_insDeliv) {
526 _insDeliv = db.prepare('INSERT INTO ap_delivery (slug, inbox, body, attempts, next_at) VALUES (?,?,?,0,CURRENT_TIMESTAMP)');
527 _dueDeliv = db.prepare("SELECT * FROM ap_delivery WHERE datetime(next_at) <= datetime('now') ORDER BY next_at LIMIT 30");
528 _delDeliv = db.prepare('DELETE FROM ap_delivery WHERE id = ?');
529 _bumpDeliv = db.prepare('UPDATE ap_delivery SET attempts = ?, next_at = ? WHERE id = ?');
530 }
531 return { ins: _insDeliv, due: _dueDeliv, del: _delDeliv, bump: _bumpDeliv };
532}
533export function enqueueDelivery(slug, inbox, activity) {
534 if (!slug || !inbox || !activity) return;
535 try { deliveryStmts().ins.run(slug, inbox, JSON.stringify(activity)); } catch { /* ignore */ }
536}
537// Deliver now; queue for retry if it fails.
538export async function deliverWithRetry(slug, inbox, activity, keyId, privPem) {
539 if (!inbox) return;
540 try { const st = await deliver(inbox, activity, keyId, privPem); if (st >= 200 && st < 300) return; } catch { /* queue below */ }
541 enqueueDelivery(slug, inbox, activity);
542}
543let _processingDeliv = false;
544export async function processDeliveryQueue() {
545 if (_processingDeliv) return; // re-entrancy guard: 30 rows × 8s can exceed the 60s tick → no double-delivery
546 _processingDeliv = true;
547 try {
548 let rows;
549 try { rows = deliveryStmts().due.all(); } catch { return; }
550 if (!rows || !rows.length) return;
551 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
552 for (const row of rows) {
553 let ok = false;
554 try {
555 const keys = getOrCreateKeys(row.slug);
556 const st = await deliver(row.inbox, JSON.parse(row.body), `${actorId(base, row.slug)}#main-key`, keys.private_pem);
557 ok = st >= 200 && st < 300;
558 } catch { ok = false; }
559 if (ok) { deliveryStmts().del.run(row.id); continue; }
560 const attempts = row.attempts + 1;
561 if (attempts >= DELIVERY_MAX_ATTEMPTS) { deliveryStmts().del.run(row.id); console.warn('[AP] delivery gave up after', attempts, 'tries →', row.inbox); continue; }
562 // Index the backoff on the CURRENT attempt count (row.attempts) so the first
563 // retry uses the 1-min tier instead of skipping it.
564 const mins = DELIVERY_BACKOFF_MIN[Math.min(row.attempts, DELIVERY_BACKOFF_MIN.length - 1)];
565 deliveryStmts().bump.run(attempts, new Date(Date.now() + mins * 60000).toISOString(), row.id);
566 }
567 } finally { _processingDeliv = false; }
568}
569let _delivTimer = null;
570export function startDeliveryWorker() {
571 if (_delivTimer) return;
572 _delivTimer = setInterval(() => { processDeliveryQueue().catch(() => {}); }, 60 * 1000);
573 if (_delivTimer.unref) _delivTimer.unref();
574}
575
576// Best-effort verification of an incoming signed request. Returns the sender's
577// actor doc if the signature checks out, else null. (Not gating yet — MVP.)
578export async function verifyRequest(req) {
579 const sigH = req.headers['signature'];
580 if (!sigH) return null;
581 const p = Object.fromEntries([...sigH.matchAll(/([a-zA-Z]+)="([^"]*)"/g)].map((m) => [m[1], m[2]]));
582 if (!p.keyId || !p.signature) return null;
583 const actor = await fetchActor(p.keyId.split('#')[0]);
584 const pem = actor && actor.publicKey && actor.publicKey.publicKeyPem;
585 if (!pem) return null;
586 const hs = (p.headers || '(request-target) host date').split(/\s+/);
587 const line = hs.map((h) => h === '(request-target)'
588 ? `(request-target): ${req.method.toLowerCase()} ${req.originalUrl}`
589 : `${h}: ${req.headers[h] || ''}`).join('\n');
590 let ok = false;
591 try { ok = crypto.verify('sha256', Buffer.from(line), pem, Buffer.from(p.signature, 'base64')); } catch { ok = false; }
592 if (ok && hs.includes('digest') && req.rawBody) {
593 const exp = 'SHA-256=' + crypto.createHash('sha256').update(req.rawBody).digest('base64');
594 if (req.headers['digest'] !== exp) ok = false;
595 }
596 return ok ? actor : null;
597}
598
599// Handle an incoming inbox POST. slugParam = null for the shared /ap/inbox.
600export async function handleInbox(req, slugParam) {
601 const act = req.body || {};
602 const type = act.type;
603 // Real client IP (behind the proxy via `trust proxy`) — logged on dropped/rejected/
604 // ignored inbox hits so an operator can see who is probing their fediverse inbox.
605 const ip = req.ip || (req.connection && req.connection.remoteAddress) || '?';
606 const base = (process.env.PUBLIC_BASE_URL || `${req.protocol}://${req.get('host')}`).replace(/\/+$/, '');
607 const verified = await verifyRequest(req).catch(() => null);
608
609 // ENFORCE HTTP signatures: a data-affecting activity must be signed by the very
610 // actor it claims to be. No valid signature, or signer ≠ actor → reject (no
611 // forged replies/likes/follows/timeline posts). GET/discovery stays open.
612 const claimedActor = typeof act.actor === 'string' ? act.actor : (act.actor && act.actor.id);
613 // Blocked actor/domain → silently drop (202, don't reveal the block).
614 if (claimedActor && isBlockedAny(claimedActor)) { console.log('[AP] inbox dropped (blocked)', claimedActor, 'from', ip); return 202; }
615 const GATED = ['Create', 'Like', 'Announce', 'Follow', 'Delete', 'Undo', 'Accept', 'Reject', 'Add', 'Remove', 'Update'];
616 if (GATED.includes(type)) {
617 if (!verified || !claimedActor || verified.id !== claimedActor) {
618 console.warn('[AP] inbox REJECTED (signature)', type, claimedActor || '?', 'from', ip, verified ? '(signer mismatch)' : '(unsigned/invalid)');
619 return 401;
620 }
621 }
622
623 if (type === 'Follow') {
624 const who = typeof act.actor === 'string' ? act.actor : (act.actor && act.actor.id);
625 const slug = slugParam || slugFromActorUrl(typeof act.object === 'string' ? act.object : (act.object && act.object.id));
626 if (!who || !slug) return 400;
627 const remote = await fetchActor(who);
628 if (!remote || !remote.inbox) return 202; // can't reach them → drop quietly
629 const sharedInbox = (remote.endpoints && remote.endpoints.sharedInbox) || null;
630 fStmts().ins.run(slug, who, remote.inbox, sharedInbox);
631 const me = actorId(base, slug);
632 const keys = getOrCreateKeys(slug);
633 const accept = { '@context': 'https://www.w3.org/ns/activitystreams', id: `${me}#accept-${Date.now()}-${rid()}`, type: 'Accept', actor: me, object: act };
634 deliver(remote.inbox, accept, `${me}#main-key`, keys.private_pem).catch((e) => console.warn('[AP] Accept delivery failed:', e.message));
635 // Auto-backfill: send our recent posts as Create so the instance has our history
636 // (Mastodon doesn't fetch history on follow). ONCE PER REMOTE INSTANCE only —
637 // Mastodon dedupes notes per-instance, so re-filling an instance that already has
638 // a follower of ours is wasted work (and won't re-populate the new follower's
639 // timeline anyway). Deliver to the shared inbox (instance-level) when present.
640 // Sync insert+check (no await between) → no interleave race with concurrent Follows.
641 const instanceFilled = sharedInbox &&
642 db.prepare('SELECT 1 FROM ap_followers WHERE slug = ? AND shared_inbox = ? AND actor_uri != ? LIMIT 1')
643 .get(slug, sharedInbox, who);
644 if (!instanceFilled) {
645 backfillNewFollower(base, slug, sharedInbox || remote.inbox).catch(() => { /* best-effort */ });
646 }
647 console.log('[AP] Follow', who, '→', slug, verified ? '(sig ok)' : '(sig unverified)');
648 return 202;
649 }
650 if (type === 'Undo' && act.object) {
651 const who = typeof act.actor === 'string' ? act.actor : (act.actor && act.actor.id);
652 const ot = act.object.type;
653 if (ot === 'Follow') {
654 const obj = act.object.object;
655 const slug = slugParam || slugFromActorUrl(typeof obj === 'string' ? obj : (obj && obj.id));
656 if (who && slug) { fStmts().del.run(slug, who); console.log('[AP] Unfollow', who, '→', slug); }
657 return 202;
658 }
659 if (ot === 'Like' || ot === 'Announce') {
660 const tgt = act.object.object;
661 const pid = postIdFromNoteUrl(typeof tgt === 'string' ? tgt : (tgt && tgt.id), base);
662 if (who && pid) { iStmts().delLA.run(ot.toLowerCase(), pid, who); console.log('[AP] Undo', ot, who, '→', pid); }
663 return 202;
664 }
665 return 202;
666 }
667
668 const actorUri = typeof act.actor === 'string' ? act.actor : (act.actor && act.actor.id);
669 const resolveActor = async (uri) => ((verified && verified.id === uri) ? verified : await fetchActor(uri).catch(() => null));
670 // Activities from our OWN actors are already stored via ap_outbox — don't re-store.
671 const isLocalActor = !!(base && actorUri && actorUri.startsWith(`${base}/ap/users/`));
672
673 // Inbound reply: a Create whose object replies to one of our notes (post OR comment).
674 if (type === 'Create' && act.object && (act.object.type === 'Note' || act.object.type === 'Article')) {
675 const o = act.object;
676 const tgt = findThreadTarget(o.inReplyTo, base);
677 if (tgt && actorUri && !isLocalActor) {
678 const ai = actorInfo(await resolveActor(actorUri), actorUri);
679 const html = HtmlSanitizerService.sanitize(o.content || '');
680 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);
681 console.log('[AP] reply', actorUri, '→', tgt.post_id);
682 return 202;
683 }
684 // Home timeline (client): a top-level post from an account we follow.
685 if (actorUri && !isLocalActor && !o.inReplyTo && o.id) {
686 let subs = []; try { subs = db.prepare('SELECT slug, auto_boost FROM ap_following WHERE actor_uri = ?').all(actorUri); } catch { /* table may not exist yet */ }
687 if (subs.length) {
688 const ai = actorInfo(await resolveActor(actorUri), actorUri);
689 const html = HtmlSanitizerService.sanitize(o.content || '');
690 const _atts = (Array.isArray(o.attachment) ? o.attachment : []).map((a) => ({ url: safeUrl(a && a.url), type: (a && a.mediaType) || '' })).filter((m) => m.url);
691 // Fallback cover: a Note's `image` (set when the attachment was suppressed
692 // for a player-card post, e.g. hosted-audio posts).
693 if (!_atts.some((m) => !m.type || /image/i.test(m.type)) && o.image) {
694 const _im = Array.isArray(o.image) ? o.image[0] : o.image;
695 const _iu = safeUrl(typeof _im === 'string' ? _im : (_im && _im.url));
696 if (_iu) _atts.push({ url: _iu, type: (_im && _im.mediaType) || 'image/jpeg' });
697 }
698 const media = JSON.stringify(_atts);
699 // "Feature" = show in the Cirkel (local only). We do NOT auto-Announce
700 // incoming posts to the fediverse — that flooded followers. Boosting to the
701 // fediverse is only ever a deliberate, manual per-post action (the 🔁 on
702 // the timeline).
703 for (const s of subs) {
704 tlStmts().ins.run(o.id, s.slug, actorUri, ai.name, ai.handle, ai.icon, ai.url, html, o.url || null, o.published || null, media);
705 }
706 console.log('[AP] timeline +', actorUri, 'x' + subs.length);
707 }
708 }
709 return 202;
710 }
711 if (type === 'Like' || type === 'Announce') {
712 const tgt = act.object;
713 const pid = postIdFromNoteUrl(typeof tgt === 'string' ? tgt : (tgt && tgt.id), base);
714 if (pid && actorUri && !isLocalActor && localPostExists(pid)) {
715 const ai = actorInfo(await resolveActor(actorUri), actorUri);
716 iStmts().ins.run(type.toLowerCase(), pid, '', actorUri, ai.name, ai.handle, ai.url, ai.icon, null, null, null);
717 console.log('[AP]', type === 'Like' ? 'like' : 'boost', actorUri, '→', pid);
718 }
719 return 202;
720 }
721 if (type === 'Delete') {
722 // A remote note was deleted upstream → drop it from replies AND the timeline.
723 // Scope to the SIGNING actor so actor B can't delete actor A's content (the
724 // signature gate guarantees claimedActor == the verified signer here).
725 const oid = typeof act.object === 'string' ? act.object : (act.object && act.object.id);
726 if (oid && claimedActor) {
727 try { db.prepare('DELETE FROM ap_interactions WHERE object_uri = ? AND actor_uri = ?').run(oid, claimedActor); } catch { /* ignore */ }
728 try { db.prepare('DELETE FROM ap_timeline WHERE id = ? AND author_uri = ?').run(oid, claimedActor); } catch { /* ignore */ }
729 }
730 return 202;
731 }
732 // Accept/Reject of a Follow WE sent (client side).
733 if (type === 'Accept' && act.object) {
734 const fid = typeof act.object === 'string' ? act.object : (act.object && act.object.id);
735 if (fid) { try { fwStmts().acc.run(fid); } catch { /* ignore */ } }
736 console.log('[AP] follow accepted', actorUri);
737 return 202;
738 }
739 if (type === 'Reject' && act.object) {
740 const who = actorUri;
741 if (who && slugParam) { try { fwStmts().del.run(slugParam, who); } catch { /* ignore */ } }
742 return 202;
743 }
744
745 console.log('[AP] inbox', type || 'unknown', '→', slugParam || 'shared', 'from', ip, '(ignored)');
746 return 202;
747}
748
749// Deliver a new post as Create(Note) to all followers' inboxes (fire-and-forget).
750// Needs PUBLIC_BASE_URL (absolute URLs); no-op without followers or base.
751export async function deliverCreate(site, post) {
752 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
753 if (!base || !site || !site.slug) return;
754 const followers = fStmts().list.all(site.slug);
755 if (!followers.length) return;
756 const inboxes = [...new Set(followers.map((f) => f.shared_inbox || f.inbox).filter(Boolean))];
757 const keys = getOrCreateKeys(site.slug);
758 const keyId = `${actorId(base, site.slug)}#main-key`;
759 const create = buildCreate(base, site, post);
760 for (const inbox of inboxes) deliverWithRetry(site.slug, inbox, create, keyId, keys.private_pem);
761}
762
763// On a new Follow, send that follower our most recent posts as Create so their
764// timeline shows our history (Mastodon does not backfill on follow). Oldest-first
765// so they sort into the follower's timeline at their original dates.
766async function backfillNewFollower(base, slug, inbox) {
767 if (!base || !slug || !inbox) return;
768 const site = db.prepare('SELECT * FROM sites WHERE slug = ?').get(slug);
769 if (!site) return;
770 const recent = db.prepare(
771 `SELECT id, slug, title, content, cover_image_url, published_at, created_at
772 FROM posts WHERE site_id = ? AND status = 'published' AND (fan_only IS NULL OR fan_only = 0)
773 ORDER BY COALESCE(published_at, created_at) DESC LIMIT 20`
774 ).all(site.id).reverse();
775 if (!recent.length) return;
776 const keys = getOrCreateKeys(slug);
777 const keyId = `${actorId(base, slug)}#main-key`;
778 for (const p of recent) {
779 try { await deliver(inbox, buildCreate(base, site, p), keyId, keys.private_pem); } catch { /* best-effort */ }
780 await new Promise((r) => setTimeout(r, 150));
781 }
782 console.log('[AP] backfilled', recent.length, 'posts to new follower of', slug);
783}
784
785// Tell followers a post is gone (Delete + Tombstone) so it's removed from their feeds.
786export async function deliverDelete(site, post) {
787 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
788 if (!base || !site || !site.slug || !post || !post.id) return;
789 const followers = fStmts().list.all(site.slug);
790 if (!followers.length) return;
791 const inboxes = [...new Set(followers.map((f) => f.shared_inbox || f.inbox).filter(Boolean))];
792 const keys = getOrCreateKeys(site.slug);
793 const me = actorId(base, site.slug);
794 const nid = noteId(base, post.id);
795 const del = {
796 '@context': 'https://www.w3.org/ns/activitystreams',
797 id: `${nid}#delete-${Date.now()}-${rid()}`,
798 type: 'Delete',
799 actor: me,
800 to: [PUBLIC],
801 object: { id: nid, type: 'Tombstone' },
802 };
803 for (const inbox of inboxes) deliverWithRetry(site.slug, inbox, del, `${me}#main-key`, keys.private_pem);
804}
805
806// Tell followers an already-published post changed (Update + edited Note) so
807// Mastodon refreshes the cached copy (e.g. after fixing content).
808export async function deliverUpdate(site, post) {
809 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
810 if (!base || !site || !site.slug || !post || !post.id) return;
811 const followers = fStmts().list.all(site.slug);
812 if (!followers.length) return;
813 const inboxes = [...new Set(followers.map((f) => f.shared_inbox || f.inbox).filter(Boolean))];
814 const keys = getOrCreateKeys(site.slug);
815 const me = actorId(base, site.slug);
816 const note = buildNote(base, site, post);
817 note.updated = new Date().toISOString();
818 const update = {
819 '@context': 'https://www.w3.org/ns/activitystreams',
820 id: `${noteId(base, post.id)}#update-${Date.now()}-${rid()}`,
821 type: 'Update', actor: me, to: [PUBLIC], cc: [`${me}/followers`],
822 object: note,
823 };
824 for (const inbox of inboxes) deliverWithRetry(site.slug, inbox, update, `${me}#main-key`, keys.private_pem);
825}
826
827// Tell followers the ACTOR changed (Update + Person) so Mastodon re-processes the
828// account AND re-fetches the featured (pinned) collection — there is no standard
829// "featured changed" activity, so this is how a pin/unpin propagates promptly.
830export async function deliverActorUpdate(site) {
831 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
832 if (!base || !site || !site.slug) return;
833 const followers = fStmts().list.all(site.slug);
834 if (!followers.length) return;
835 const inboxes = [...new Set(followers.map((f) => f.shared_inbox || f.inbox).filter(Boolean))];
836 const keys = getOrCreateKeys(site.slug);
837 const me = actorId(base, site.slug);
838 const update = {
839 '@context': ['https://www.w3.org/ns/activitystreams', 'https://w3id.org/security/v1'],
840 id: `${me}#update-${Date.now()}-${rid()}`,
841 type: 'Update', actor: me, to: [PUBLIC], cc: [`${me}/followers`],
842 object: buildActor(base, site),
843 };
844 for (const inbox of inboxes) deliverWithRetry(site.slug, inbox, update, `${me}#main-key`, keys.private_pem);
845}
846
847// Reliably set the pinned order on followers' instances via Add/Remove activities
848// (how Mastodon itself federates pins) — pushed to the inbox + processed immediately,
849// unlike the featured COLLECTION which Mastodon caches with sticky StatusPins.
850// Mastodon's Add skips an already-pinned status, so we REMOVE every pin first, wait,
851// then ADD in rank-DESCENDING order (rank 1 added LAST → newest StatusPin → shown first,
852// because Mastodon displays pins newest-first). `alsoRemove` = ids to unpin too.
853export async function resyncFeaturedPins(site, alsoRemove = []) {
854 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
855 if (!base || !site || !site.slug) return;
856 const followers = fStmts().list.all(site.slug);
857 if (!followers.length) return;
858 const inboxes = [...new Set(followers.map((f) => f.shared_inbox || f.inbox).filter(Boolean))];
859 const keys = getOrCreateKeys(site.slug);
860 const me = actorId(base, site.slug);
861 const keyId = `${me}#main-key`;
862 const featured = `${me}/featured`;
863 const AS = 'https://www.w3.org/ns/activitystreams';
864 const note = (id) => noteId(base, id);
865 const pinned = db.prepare(
866 `SELECT id FROM posts WHERE site_id = ? AND status = 'published' AND (fan_only IS NULL OR fan_only = 0)
867 AND pinned IS NOT NULL AND pinned > 0
868 ORDER BY pinned DESC, COALESCE(published_at, created_at) ASC LIMIT 20`
869 ).all(site.id);
870 const removeIds = [...new Set([...pinned.map((p) => p.id), ...alsoRemove])];
871 // 1. Remove every current pin so Mastodon can recreate them in order.
872 for (const id of removeIds) {
873 const rm = { '@context': AS, id: `${me}#rm-${id}-${Date.now()}-${rid()}`, type: 'Remove', actor: me, object: note(id), target: featured, to: [PUBLIC] };
874 for (const inbox of inboxes) deliver(inbox, rm, keyId, keys.private_pem).catch(() => { /* best-effort */ });
875 }
876 if (!pinned.length) { console.log('[AP] unpinned all featured for', site.slug); return; }
877 await new Promise((r) => setTimeout(r, 5000)); // let the Removes land first
878 // 2. Add in rank-DESC order, gaps so each StatusPin gets an increasing created_at.
879 for (const p of pinned) {
880 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`] };
881 for (const inbox of inboxes) deliver(inbox, add, keyId, keys.private_pem).catch(() => { /* best-effort */ });
882 await new Promise((r) => setTimeout(r, 2000));
883 }
884 console.log('[AP] resynced', pinned.length, 'featured pins for', site.slug);
885}
886
887// ── outbound replies (Klonkt → fediverse) ─────────────────────────
888const escHtml = (s) => String(s || '').replace(/[<>&]/g, (c) => ({ '<': '&lt;', '>': '&gt;', '&': '&amp;' }[c]));
889const 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(); };
890
891// Build one of OUR outbound reply Notes from an ap_outbox row.
892export function buildReplyNote(base, site, row) {
893 const me = actorId(base, site.slug);
894 return {
895 id: noteId(base, row.id),
896 type: 'Note',
897 attributedTo: me,
898 inReplyTo: row.in_reply_to || undefined,
899 content: row.content,
900 url: row.post_slug ? `${base}/${encodeURIComponent(row.post_slug)}` : undefined,
901 published: toISO(row.created_at),
902 to: row.to_actor ? [row.to_actor] : [PUBLIC],
903 cc: [PUBLIC, `${me}/followers`],
904 tag: row.to_actor ? [{ type: 'Mention', href: row.to_actor, name: row.to_handle }] : [],
905 };
906}
907
908// Resolve one of our outbound reply Notes by id (for /ap/notes/:id fallback).
909export function getOutboxNote(base, id) {
910 const row = iStmts().getO.get(id);
911 if (!row) return null;
912 const site = db.prepare('SELECT * FROM sites WHERE slug = ?').get(row.site_slug);
913 if (!site) return null;
914 return buildReplyNote(base, site, row);
915}
916
917// Send a reply FROM this site to a remote actor (in reply to their inbound reply).
918// `parent` = an ap_interactions row (actor_uri, actor_url, actor_handle, object_uri).
919export async function deliverReply(site, { postId, postSlug, parent, text }) {
920 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
921 if (!base || !site || !site.slug || !parent || !String(text || '').trim()) return null;
922 const me = actorId(base, site.slug);
923 const handle = parent.actor_handle || deriveHandle(parent.actor_uri);
924 const body = escHtml(String(text).trim()).replace(/\r?\n/g, '<br>');
925 const mention = parent.actor_uri
926 ? `<a href="${escHtml(parent.actor_url || parent.actor_uri)}" class="u-url mention">${escHtml(handle)}</a> ` : '';
927 const content = `<p>${mention}${body}</p>`;
928 // Dedup: skip if the exact same reply was already sent (double-submit guard).
929 const dup = db.prepare('SELECT 1 FROM ap_outbox WHERE site_slug = ? AND IFNULL(in_reply_to, \'\') = ? AND content = ? LIMIT 1')
930 .get(site.slug, parent.object_uri || '', content);
931 if (dup) { console.log('[AP] outreply skipped (duplicate)'); return { duplicate: true, delivered: 0 }; }
932 const id = crypto.randomUUID();
933 iStmts().insO.run(id, site.slug, postId, postSlug || null, parent.object_uri || null, parent.actor_uri || null, handle, content);
934 const row = iStmts().getO.get(id);
935 const note = buildReplyNote(base, site, row);
936 const create = {
937 '@context': 'https://www.w3.org/ns/activitystreams',
938 id: note.id + '#create', type: 'Create', actor: me,
939 published: note.published, to: note.to, cc: note.cc, object: note,
940 };
941 const keys = getOrCreateKeys(site.slug);
942 const keyId = `${me}#main-key`;
943 const inboxes = new Set();
944 if (parent.actor_uri) {
945 const a = await fetchActor(parent.actor_uri).catch(() => null);
946 if (a) inboxes.add((a.endpoints && a.endpoints.sharedInbox) || a.inbox);
947 }
948 if (parent.threadInbox) inboxes.add(parent.threadInbox); // back-compat (single)
949 (parent.threadInboxes || []).forEach((i) => inboxes.add(i)); // whole ancestor chain
950 for (const f of fStmts().list.all(site.slug)) inboxes.add(f.shared_inbox || f.inbox);
951 inboxes.delete(`${me}/inbox`); // never deliver to ourselves (already in ap_outbox)
952 inboxes.delete(`${base}/ap/inbox`); // (our own shared inbox) → avoids a self-duplicate
953 let delivered = 0;
954 for (const inbox of [...inboxes].filter(Boolean)) {
955 try { const st = await deliver(inbox, create, keyId, keys.private_pem); if (st >= 200 && st < 300) delivered++; } catch { /* best-effort */ }
956 }
957 console.log('[AP] outreply', site.slug, '→', parent.actor_uri, 'delivered', delivered);
958 return { id, content, delivered };
959}
960
961// Resolve a remote post URL (any fediverse/Klonkt post) into a reply target.
962// Returns a parent-shaped object usable by deliverReply(), or null.
963export async function resolveRemoteNote(url) {
964 if (!/^https?:\/\//i.test(String(url || ''))) return null;
965 const note = await fetchActor(url).catch(() => null); // AP GET (content-negotiates)
966 if (!note || !note.id) return null;
967 const att = note.attributedTo;
968 const actorUri = typeof att === 'string' ? att : (att && att.id);
969 if (!actorUri) return null;
970 const actor = await fetchActor(actorUri).catch(() => null);
971 const ai = actorInfo(actor, actorUri);
972 // Is what we're replying to a post (or a comment) on one of OUR posts? If so,
973 // link our reply to that local post so it shows nested in the post thread.
974 const localTgt = findThreadTarget(note.id, (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, ''));
975 // Walk the WHOLE reply chain upward (comment → parent comment → … → root post)
976 // and collect every ancestor author's inbox, so each participant's server —
977 // including the original post's author — receives + threads our reply.
978 const threadInboxes = [];
979 const seenInbox = new Set();
980 let cursor = note.inReplyTo, guard = 0;
981 while (cursor && guard++ < 6) {
982 const url = typeof cursor === 'string' ? cursor : (cursor && cursor.id);
983 if (!url) break;
984 const pn = await fetchActor(url).catch(() => null);
985 if (!pn) break;
986 const pa = typeof pn.attributedTo === 'string' ? pn.attributedTo : (pn.attributedTo && pn.attributedTo.id);
987 if (pa && pa !== actorUri) {
988 const paDoc = await fetchActor(pa).catch(() => null);
989 const inbox = paDoc && ((paDoc.endpoints && paDoc.endpoints.sharedInbox) || paDoc.inbox);
990 if (inbox && !seenInbox.has(inbox)) { seenInbox.add(inbox); threadInboxes.push(inbox); }
991 }
992 cursor = pn.inReplyTo; // climb to the next ancestor
993 }
994 const rawHtml = String(note.content || '').replace(/\[\[(track|album|playlist):[^\]]+\]\]/gi, '');
995 const images = (Array.isArray(note.attachment) ? note.attachment : [])
996 .filter((a) => a && a.url && (!a.mediaType || /^image\//i.test(a.mediaType)))
997 .map((a) => safeUrl(a.url)).filter(Boolean);
998 return {
999 object_uri: safeUrl(note.id) || note.id,
1000 actor_uri: actorUri,
1001 actor_url: ai.url,
1002 actor_handle: ai.handle,
1003 actor_name: ai.name,
1004 actor_icon: ai.icon,
1005 url: note.url || url,
1006 content: HtmlSanitizerService.sanitize(rawHtml), // full, sanitized
1007 images,
1008 threadInboxes, // every ancestor author's inbox
1009 localPostId: localTgt ? localTgt.post_id : '', // our post this belongs to (if any)
1010 preview: HtmlSanitizerService.toPlainText(note.content || '').slice(0, 240),
1011 };
1012}
1013
1014// List a site's own outbound fediverse replies (for the manage/delete view).
1015export function listOutbox(siteSlug) {
1016 return db.prepare('SELECT id, content, to_handle, in_reply_to, created_at FROM ap_outbox WHERE site_slug = ? ORDER BY created_at DESC')
1017 .all(siteSlug).map((r) => ({ ...r, content: stripLeadingMentions(r.content) }));
1018}
1019
1020// Delete one of our outbound replies: send Delete(Tombstone) to recipients + remove it.
1021export async function deliverOutboxDelete(site, outboxId) {
1022 const row = iStmts().getO.get(outboxId);
1023 if (!row || row.site_slug !== site.slug) return false;
1024 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
1025 if (base) {
1026 const me = actorId(base, site.slug);
1027 const nid = noteId(base, row.id);
1028 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' } };
1029 const keys = getOrCreateKeys(site.slug);
1030 const inboxes = new Set();
1031 if (row.to_actor) { const a = await fetchActor(row.to_actor).catch(() => null); if (a) inboxes.add((a.endpoints && a.endpoints.sharedInbox) || a.inbox); }
1032 for (const f of fStmts().list.all(site.slug)) inboxes.add(f.shared_inbox || f.inbox);
1033 for (const inbox of [...inboxes].filter(Boolean)) { try { await deliver(inbox, del, `${me}#main-key`, keys.private_pem); } catch { /* best-effort */ } }
1034 }
1035 db.prepare('DELETE FROM ap_outbox WHERE id = ?').run(outboxId);
1036 return true;
1037}
1038
1039// ── Fediverse CLIENT: follow accounts + home timeline ─────────────
1040// Resolve an @user@domain handle to its actor URL via WebFinger.
1041export async function webfingerResolve(handle) {
1042 const h = String(handle || '').trim().replace(/^@/, '');
1043 const parts = h.split('@');
1044 if (parts.length !== 2 || !parts[0] || !parts[1]) return null;
1045 const acct = `${parts[0]}@${parts[1]}`;
1046 try {
1047 const r = await safeFetch(`https://${parts[1]}/.well-known/webfinger?resource=acct:${encodeURIComponent(acct)}`,
1048 { headers: { Accept: 'application/jrd+json, application/json' } });
1049 if (!r.ok) return null;
1050 const jrd = await r.json();
1051 const link = (jrd.links || []).find((l) => l.rel === 'self' && /activity\+json|ld\+json/.test(l.type || ''));
1052 return safeUrl(link ? link.href : '') || null;
1053 } catch { return null; }
1054}
1055
1056let _insFw, _delFw, _listFw, _accFw, _oneFw, _setAB;
1057function fwStmts() {
1058 if (!_insFw) {
1059 _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)');
1060 _delFw = db.prepare('DELETE FROM ap_following WHERE slug = ? AND actor_uri = ?');
1061 _listFw = db.prepare('SELECT * FROM ap_following WHERE slug = ? ORDER BY created_at DESC');
1062 _accFw = db.prepare("UPDATE ap_following SET status = 'accepted' WHERE follow_id = ?");
1063 _oneFw = db.prepare('SELECT * FROM ap_following WHERE slug = ? AND actor_uri = ?');
1064 _setAB = db.prepare('UPDATE ap_following SET auto_boost = ? WHERE slug = ? AND actor_uri = ?');
1065 }
1066 return { ins: _insFw, del: _delFw, list: _listFw, acc: _accFw, one: _oneFw, setAB: _setAB };
1067}
1068export function listFollowing(slug) { return fwStmts().list.all(slug); }
1069
1070// Toggle auto-boost ("feature") on an account we already follow.
1071export function setAutoBoost(slug, actorUri, on) {
1072 try { fwStmts().setAB.run(on ? 1 : 0, slug, actorUri); } catch { /* ignore */ }
1073 return { ok: true };
1074}
1075
1076let _insTl, _listTl, _delTl;
1077function tlStmts() {
1078 if (!_insTl) {
1079 _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)');
1080 _listTl = db.prepare('SELECT * FROM ap_timeline WHERE slug = ? ORDER BY COALESCE(published, created_at) DESC LIMIT ?');
1081 _delTl = db.prepare('DELETE FROM ap_timeline WHERE id = ?');
1082 }
1083 return { ins: _insTl, list: _listTl, del: _delTl };
1084}
1085export function getTimeline(slug, limit) { return tlStmts().list.all(slug, limit || 50); }
1086
1087// ── Cirkel = posts from the accounts you auto-boost ("feature an artist") ──
1088let _abCount, _cirkelPosts, _cirkelMembers;
1089export function autoBoostCount(slug) {
1090 try { if (!_abCount) _abCount = db.prepare('SELECT COUNT(*) AS n FROM ap_following WHERE slug = ? AND auto_boost = 1'); return _abCount.get(slug).n; } catch { return 0; }
1091}
1092export function getCirkelPosts(slug, limit) {
1093 try {
1094 // Cirkel = posts from featured (auto_boost) accounts + posts you boosted
1095 // (t.boosted), mixed by date. One row per note in ap_timeline → no duplicates.
1096 if (!_cirkelPosts) _cirkelPosts = db.prepare(`
1097 SELECT t.id, t.author_uri, t.author_name, t.author_handle, t.author_icon, t.author_url,
1098 t.content, t.url, t.published, t.media_json, t.boosted
1099 FROM ap_timeline t
1100 LEFT JOIN ap_following f ON f.slug = t.slug AND f.actor_uri = t.author_uri
1101 WHERE t.slug = ? AND (f.auto_boost = 1 OR t.boosted = 1)
1102 ORDER BY COALESCE(t.published, t.created_at) DESC, t.rowid DESC
1103 LIMIT ?`);
1104 return _cirkelPosts.all(slug, limit || 60);
1105 } catch { return []; }
1106}
1107export function getCirkelMembers(slug) {
1108 try { if (!_cirkelMembers) _cirkelMembers = db.prepare('SELECT name, url, icon FROM ap_following WHERE slug = ? AND auto_boost = 1 ORDER BY name'); return _cirkelMembers.all(slug); } catch { return []; }
1109}
1110// Mark a timeline post as boosted so it shows in the Cirkel (mixed by date).
1111let _markBoost, _unmarkBoost, _boostedCount;
1112export function markBoosted(slug, noteId) {
1113 try { if (!_markBoost) _markBoost = db.prepare('UPDATE ap_timeline SET boosted = 1 WHERE slug = ? AND id = ?'); _markBoost.run(slug, noteId); } catch { /* ignore */ }
1114}
1115export function unmarkBoosted(slug, noteId) {
1116 try { if (!_unmarkBoost) _unmarkBoost = db.prepare('UPDATE ap_timeline SET boosted = 0 WHERE slug = ? AND id = ?'); _unmarkBoost.run(slug, noteId); } catch { /* ignore */ }
1117}
1118let _markLike, _unmarkLike;
1119export function markLiked(slug, noteId) {
1120 try { if (!_markLike) _markLike = db.prepare('UPDATE ap_timeline SET liked = 1 WHERE slug = ? AND id = ?'); _markLike.run(slug, noteId); } catch { /* ignore */ }
1121}
1122export function unmarkLiked(slug, noteId) {
1123 try { if (!_unmarkLike) _unmarkLike = db.prepare('UPDATE ap_timeline SET liked = 0 WHERE slug = ? AND id = ?'); _unmarkLike.run(slug, noteId); } catch { /* ignore */ }
1124}
1125export function getTimelineReaction(slug, noteId) {
1126 try { const r = db.prepare('SELECT liked, boosted FROM ap_timeline WHERE slug = ? AND id = ?').get(slug, noteId); return { liked: !!(r && r.liked), boosted: !!(r && r.boosted) }; } catch { return { liked: false, boosted: false }; }
1127}
1128export function boostedCount(slug) {
1129 try { if (!_boostedCount) _boostedCount = db.prepare('SELECT COUNT(*) AS n FROM ap_timeline WHERE slug = ? AND boosted = 1'); return _boostedCount.get(slug).n; } catch { return 0; }
1130}
1131
1132// Resolve a Klonkt/AP actor URL from a site root: a Klonkt site's root 302s to
1133// /ap/users/<slug> (content negotiation; Location may be relative). Used by
1134// followActor for bare-domain follows.
1135// NB: the old auto-migration of legacy Cirkels (circle_links -> AP follows) was
1136// REMOVED on 2026-06-26 — it auto-sent Follows on boot, which violates "the code
1137// never throws anything into the fediverse automatically" (would surprise-Follow
1138// for some operators at scale). The dead circle_links table stays as harmless dead
1139// data; an operator restores an old cirkel by re-following in /following (their click).
1140async function resolveApActor(siteUrl) {
1141 try {
1142 const r = await fetch(siteUrl, { headers: { Accept: 'application/activity+json' }, redirect: 'manual' });
1143 if (r.status >= 300 && r.status < 400) { const loc = r.headers.get('location'); if (loc) return new URL(loc, siteUrl).href; }
1144 if (r.ok) return siteUrl;
1145 } catch { /* unreachable */ }
1146 return null;
1147}
1148
1149// ── Self-heal: re-sync the fediverse cache (ap_timeline) after a DRASTIC update ──
1150// Runs ONCE per SELFHEAL_VERSION bump — NOT on every boot. Re-fetches each cached
1151// note and refreshes content + media (recovers covers/edits that were delivered
1152// during a flux window, e.g. a fleet-wide update), and drops notes that are gone
1153// (404/410). Bump SELFHEAL_VERSION only on a release that warrants a re-sync.
1154const SELFHEAL_VERSION = 1;
1155async function fetchNoteAP(url) {
1156 try {
1157 const r = await fetch(url, { headers: { Accept: 'application/activity+json' } });
1158 if (r.status === 404 || r.status === 410) return 404;
1159 if (r.ok) return await r.json();
1160 } catch { /* unreachable */ }
1161 return null;
1162}
1163function mediaFromNote(note) {
1164 const atts = (Array.isArray(note.attachment) ? note.attachment : []).map((a) => ({ url: safeUrl(a && a.url), type: (a && a.mediaType) || '' })).filter((m) => m.url);
1165 if (!atts.some((m) => !m.type || /image/i.test(m.type)) && note.image) {
1166 const im = Array.isArray(note.image) ? note.image[0] : note.image;
1167 const iu = safeUrl(typeof im === 'string' ? im : (im && im.url));
1168 if (iu) atts.push({ url: iu, type: (im && im.mediaType) || 'image/jpeg' });
1169 }
1170 return JSON.stringify(atts);
1171}
1172let _selfHealing = false;
1173export async function selfHealTimeline() {
1174 if (_selfHealing) return; _selfHealing = true;
1175 try {
1176 let cur = 0;
1177 try { const r = db.prepare('SELECT value FROM app_settings WHERE key = ?').get('selfheal_version'); cur = r ? (parseInt(r.value, 10) || 0) : 0; } catch { return; }
1178 if (cur >= SELFHEAL_VERSION) return; // already healed for this version — skip on normal boots
1179 let rows = [];
1180 try { rows = db.prepare('SELECT id, content, media_json FROM ap_timeline ORDER BY rowid DESC LIMIT 200').all(); } catch { /* no table */ }
1181 let healed = 0;
1182 for (const r of rows) {
1183 try {
1184 const note = await fetchNoteAP(r.id);
1185 if (note === 404) { db.prepare('DELETE FROM ap_timeline WHERE id = ?').run(r.id); healed++; continue; }
1186 if (!note || typeof note !== 'object') continue;
1187 const html = HtmlSanitizerService.sanitize(note.content || '');
1188 const media = mediaFromNote(note);
1189 if ((html && html !== r.content) || media !== (r.media_json || '[]')) {
1190 db.prepare('UPDATE ap_timeline SET content = ?, media_json = ? WHERE id = ?').run(html || r.content, media, r.id);
1191 healed++;
1192 }
1193 } catch { /* per-note best-effort */ }
1194 }
1195 try { db.prepare('INSERT OR REPLACE INTO app_settings (key, value) VALUES (?, ?)').run('selfheal_version', String(SELFHEAL_VERSION)); } catch { /* ignore */ }
1196 if (rows.length) console.log(`[AP] self-heal v${SELFHEAL_VERSION}: ${healed}/${rows.length} timeline notes`);
1197 } catch { /* never block boot */ } finally { _selfHealing = false; }
1198}
1199
1200// Follow a fediverse account by @handle (WebFinger → actor → signed Follow).
1201export async function followActor(site, handle, autoBoost = false) {
1202 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
1203 if (!base || !site || !site.slug) return { error: 'config' };
1204 // Accept any of: a profile/actor URL, an @user@host handle (WebFinger), or a
1205 // bare site domain (site.com) — for a single-actor site (Klonkt etc.) the root
1206 // resolves to its AP actor, so you can follow a site by just its domain.
1207 const s = String(handle || '').trim();
1208 let actorUrl;
1209 if (/^https?:\/\//i.test(s)) actorUrl = safeUrl(s) || null;
1210 else if (s.includes('@')) actorUrl = await webfingerResolve(s);
1211 else if (/^[a-z0-9.-]+\.[a-z]{2,}/i.test(s)) actorUrl = await resolveApActor('https://' + s.replace(/^\/+|\/+$/g, ''));
1212 else actorUrl = null;
1213 if (!actorUrl) return { error: 'not_found' };
1214 const actor = await fetchActor(actorUrl).catch(() => null);
1215 if (!actor || !actor.id || !actor.inbox) return { error: 'unreachable' };
1216 const ai = actorInfo(actor, actor.id);
1217 const me = actorId(base, site.slug);
1218 const keys = getOrCreateKeys(site.slug);
1219 const followId = `${me}#follow-${Date.now()}-${rid()}`;
1220 fwStmts().ins.run(site.slug, actor.id, ai.handle, ai.name, ai.icon, ai.url, actor.inbox, followId, 'pending', autoBoost ? 1 : 0);
1221 const follow = { '@context': 'https://www.w3.org/ns/activitystreams', id: followId, type: 'Follow', actor: me, object: actor.id };
1222 try { await deliver(actor.inbox, follow, `${me}#main-key`, keys.private_pem); }
1223 catch (e) { console.warn('[AP] follow deliver failed:', e.message); }
1224 console.log('[AP] follow', site.slug, '→', actor.id);
1225 return { ok: true, name: ai.name, handle: ai.handle, actor: actor.id };
1226}
1227
1228// Resolve a profile URL or @handle to a followable remote actor (for the
1229// authorize_interaction "Follow" flow). Returns display fields + inbox, or null
1230// when it isn't a reachable actor (e.g. the input was a post, not a profile).
1231export async function resolveRemoteActor(input) {
1232 const s = String(input || '').trim();
1233 const actorUrl = /^https?:\/\//i.test(s) ? (safeUrl(s) || null) : await webfingerResolve(s);
1234 if (!actorUrl) return null;
1235 const actor = await fetchActor(actorUrl).catch(() => null);
1236 if (!actor || !actor.id || !actor.inbox) return null;
1237 const ai = actorInfo(actor, actor.id);
1238 return { actor_uri: actor.id, actor_name: ai.name, actor_handle: ai.handle, actor_url: ai.url, actor_icon: ai.icon, inbox: actor.inbox };
1239}
1240
1241export async function unfollowActor(site, actorUri) {
1242 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
1243 const me = actorId(base, site.slug);
1244 const keys = getOrCreateKeys(site.slug);
1245 const row = fwStmts().one.get(site.slug, actorUri);
1246 if (row && row.inbox) {
1247 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 } };
1248 try { await deliver(row.inbox, undo, `${me}#main-key`, keys.private_pem); } catch { /* best-effort */ }
1249 }
1250 fwStmts().del.run(site.slug, actorUri);
1251 return { ok: true };
1252}
1253
1254// Send a Like or Announce (boost) on a remote note FROM this site.
1255export async function sendInteraction(site, kind, targetNoteId, authorUri) {
1256 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
1257 if (!base || !site || !site.slug || !targetNoteId) return { error: 'config' };
1258 const me = actorId(base, site.slug);
1259 const keys = getOrCreateKeys(site.slug);
1260 // 'unboost' = Undo(Announce): retracts a boost so followers' servers remove the
1261 // reblog (matched on actor+object — no record of the original Announce needed).
1262 const fanout = (kind === 'boost' || kind === 'unboost'); // also goes to our followers
1263 let act;
1264 if (kind === 'unboost' || kind === 'unlike') {
1265 // Undo(Announce) retracts a boost; Undo(Like) un-favourites (matched on actor+object,
1266 // no record of the original activity needed — Mastodon honours both).
1267 const inner = kind === 'unboost' ? 'Announce' : 'Like';
1268 act = {
1269 '@context': 'https://www.w3.org/ns/activitystreams',
1270 id: `${me}#undo-${Date.now()}-${rid()}`, type: 'Undo', actor: me,
1271 object: { id: `${me}#${inner.toLowerCase()}-${Date.now()}-${rid()}`, type: inner, actor: me, object: targetNoteId },
1272 };
1273 if (kind === 'unboost') { act.to = [PUBLIC]; act.cc = [`${me}/followers`]; }
1274 } else {
1275 const type = kind === 'boost' ? 'Announce' : 'Like';
1276 act = {
1277 '@context': 'https://www.w3.org/ns/activitystreams',
1278 id: `${me}#${type.toLowerCase()}-${Date.now()}-${rid()}`,
1279 type, actor: me, object: targetNoteId,
1280 };
1281 if (type === 'Announce') { act.to = [PUBLIC]; act.cc = [`${me}/followers`]; }
1282 }
1283 const inboxes = new Set();
1284 if (authorUri) { const a = await fetchActor(authorUri).catch(() => null); if (a) inboxes.add((a.endpoints && a.endpoints.sharedInbox) || a.inbox); }
1285 if (fanout) { for (const f of fStmts().list.all(site.slug)) inboxes.add(f.shared_inbox || f.inbox); }
1286 let delivered = 0;
1287 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 */ } }
1288 console.log('[AP]', kind, site.slug, '→', targetNoteId, 'delivered', delivered);
1289 return { ok: true, delivered };
1290}
1291
1292// Notifications inbox: new followers + replies/likes/boosts on this site's posts.
1293export function getNotifications(slug, limit) {
1294 const out = [];
1295 try {
1296 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)) {
1297 out.push({ type: 'follow', handle: deriveHandle(f.actor_uri), url: f.actor_uri, created_at: f.created_at });
1298 }
1299 } catch { /* ignore */ }
1300 try {
1301 const rows = db.prepare(`
1302 SELECT i.kind, i.actor_name, i.actor_handle, i.actor_url, i.content, i.created_at,
1303 p.slug AS post_slug, p.title AS post_title
1304 FROM ap_interactions i LEFT JOIN posts p ON p.id = i.post_id
1305 WHERE p.site_id = (SELECT id FROM sites WHERE slug = ?)
1306 ORDER BY i.created_at DESC LIMIT 80
1307 `).all(slug);
1308 for (const r of rows) out.push({
1309 type: r.kind, name: r.actor_name, handle: r.actor_handle, url: r.actor_url,
1310 content: stripLeadingMentions(r.content), post_slug: r.post_slug, post_title: r.post_title, created_at: r.created_at,
1311 });
1312 } catch { /* ignore */ }
1313 out.sort((a, b) => new Date(b.created_at) - new Date(a.created_at));
1314 return out.slice(0, limit || 60);
1315}
1316
1317// ── Blocking / defederation ───────────────────────────────────────
1318let _insBl, _delBl, _listBl;
1319function blStmts() {
1320 if (!_insBl) {
1321 _insBl = db.prepare('INSERT OR IGNORE INTO ap_blocks (slug, target, kind, label, created_at) VALUES (?,?,?,?,CURRENT_TIMESTAMP)');
1322 _delBl = db.prepare('DELETE FROM ap_blocks WHERE slug = ? AND target = ?');
1323 _listBl = db.prepare('SELECT * FROM ap_blocks WHERE slug = ? ORDER BY created_at DESC');
1324 }
1325 return { ins: _insBl, del: _delBl, list: _listBl };
1326}
1327export function listBlocks(slug) { return blStmts().list.all(slug); }
1328
1329// True if an actor (or its whole domain) is blocked anywhere on this instance.
1330export function isBlockedAny(actorUri) {
1331 if (!actorUri) return false;
1332 let domain = ''; try { domain = new URL(actorUri).host; } catch { /* ignore */ }
1333 try { return !!db.prepare("SELECT 1 FROM ap_blocks WHERE (kind='actor' AND target=?) OR (kind='domain' AND target=?) LIMIT 1").get(actorUri, domain); }
1334 catch { return false; }
1335}
1336
1337function purgeBlocked(kind, target) {
1338 try {
1339 if (kind === 'domain') {
1340 const like = `%//${target}/%`;
1341 db.prepare('DELETE FROM ap_interactions WHERE actor_uri LIKE ?').run(like);
1342 db.prepare('DELETE FROM ap_timeline WHERE author_uri LIKE ?').run(like);
1343 db.prepare('DELETE FROM ap_followers WHERE actor_uri LIKE ?').run(like);
1344 } else {
1345 db.prepare('DELETE FROM ap_interactions WHERE actor_uri = ?').run(target);
1346 db.prepare('DELETE FROM ap_timeline WHERE author_uri = ?').run(target);
1347 db.prepare('DELETE FROM ap_followers WHERE actor_uri = ?').run(target);
1348 }
1349 } catch { /* best-effort */ }
1350}
1351
1352// Block an actor (@handle or actor URL) or a whole domain; purges their content.
1353export async function blockTarget(site, input) {
1354 const raw = String(input || '').trim();
1355 if (!site || !site.slug || !raw) return { error: 'empty' };
1356 let kind, target, label;
1357 if (/^https?:\/\//i.test(raw)) { kind = 'actor'; target = raw; label = raw; }
1358 else if (raw.includes('@')) {
1359 const actorUrl = await webfingerResolve(raw);
1360 if (!actorUrl) return { error: 'not_found' };
1361 kind = 'actor'; target = actorUrl; label = raw.startsWith('@') ? raw : ('@' + raw);
1362 } else { kind = 'domain'; target = raw.toLowerCase(); label = raw.toLowerCase(); }
1363 blStmts().ins.run(site.slug, target, kind, label);
1364 purgeBlocked(kind, target);
1365 console.log('[AP] block', site.slug, kind, target);
1366 return { ok: true, label };
1367}
1368
1369export function unblock(site, target) { blStmts().del.run(site.slug, target); return { ok: true }; }
1370
1371export default {
1372 getOrCreateKeys, apWants, sendAP, actorId, noteId,
1373 buildActor, buildNote, buildCreate, buildOutbox, buildFollowers, buildFeatured,
1374 followerCount, deliver, fetchActor, verifyRequest, handleInbox, deliverCreate, deliverDelete, deliverUpdate, deliverActorUpdate, resyncFeaturedPins,
1375 getInteractions, getInteractionById, setInteractionBoosted, setInteractionLiked, setMyReaction, getMyReactions, buildReplyNote, getOutboxNote, deliverReply, resolveRemoteNote,
1376 listOutbox, deliverOutboxDelete,
1377 webfingerResolve, followActor, resolveRemoteActor, unfollowActor, listFollowing, setAutoBoost, getTimeline, sendInteraction,
1378 autoBoostCount, boostedCount, markBoosted, unmarkBoosted, markLiked, unmarkLiked, getTimelineReaction, getCirkelPosts, getCirkelMembers, selfHealTimeline,
1379 getNotifications, listBlocks, isBlockedAny, blockTarget, unblock,
1380 deliverWithRetry, enqueueDelivery, processDeliveryQueue, startDeliveryWorker,
1381 getReplyUris, markNotificationsSeen, countUnseenNotifications, hasPlayableAudio,
1382};
Note: See TracBrowser for help on using the repository browser.