source: Klonkt/src/services/ActivityPubService.js@ 4c47eff

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

refactor(cirkel): remove the legacy circle auto-migration (no auto-fediverse)

autoMigrateCircles auto-sent Follows on boot to re-establish legacy circle_links
as AP follows. That violates the rule 'the code never throws anything into the
fediverse automatically' — at scale it would surprise-Follow on behalf of operators
who never asked. Removed the boot call, the function, its export, and the manual
scripts/migrate-circles.mjs (bulk Follows). resolveApActor stays (used by bare-domain
follows). Dead circle_links table is left as harmless dead data; an operator restores
an old cirkel by re-following in /volgend (their own click). We can do this clean
removal now precisely because we're still small-scale.

Co-Authored-By: Claude <noreply@…>

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