source: Klonkt/src/services/ActivityPubService.js@ 74d61e6

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

feat(fedi): boosting a non-followed post shows it in the Cirkel

The interact-page boost now stores the remote post in ap_timeline (INSERT OR IGNORE, no
dup for followed posts) before flagging it boosted, so a boost of someone you don't
follow still surfaces in the Cirkel with a Boost badge. upsertBoostedNote(slug, note).

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