source: Klonkt/src/services/ActivityPubService.js@ 3289a64

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

feat(fedi): comment like is a toggle too (Undo Like / un-favourite)

Symmetric to boost: ap_interactions.acted_like remembers a liked comment, the star
shows an 'on' state, and clicking again retracts it via Undo(Like) (Mastodon honours it).
sendInteraction gains an 'unlike' branch (author-only, no fanout).

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