source: Klonkt/src/services/ActivityPubService.js@ 9d34855

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

fix(news): feed like is a toggle too (like/unlike), like boost

The feed star was fire-and-forget with no state. Now ap_timeline.liked remembers it,
the star shows an 'on' state, and a second click retracts via /news/unlike (Undo Like)
— mirroring the existing boost/unboost toggle. markLiked/unmarkLiked added.

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