source: Klonkt/src/services/ActivityPubService.js@ c745659

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

feat(fedi): comment boost is a toggle with state + coloured action icons

Inbound fediverse comments: the owner like/boost/reply controls are now coloured SVG
icons (gold star / green repeat / accent arrow), matching the News feed. Boost is a
real toggle — ap_interactions.acted_boost remembers what you boosted, the button shows
an 'on' state, and clicking again retracts it (Undo Announce). Like stays fire-and-forget
(no unlike in AS).

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