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

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

feat(fedi): hide the leading @mention in displayed comments

A federated reply carries a leading @user@domain mention (the person replied to).
Strip it for display (Mastodon does the same) so a comment reads 'dope tekening ouwe'
not '@jason@… dope …'. Applied to inbound + outbound reply nodes, the
My-replies list, and the notifications preview. stripLeadingMentions keeps the <p>
wrapper and handles both <a>mention</a> links and plain-text mentions.

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