source: Klonkt/src/services/ActivityPubService.js@ 43273c9

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

refactor(fediverse): thread crawl one level per pass, like Mastodon

Match Mastodon's FetchReplies behaviour: crawl ONE hop of replies per pass instead
of a deep BFS. Deeper levels still fill in — a cached reply becomes a seed on the
next crawl, so its replies are pulled on a later view (Mastodon's per-status cascade).
More conservative + polite; same push-first + cache + bounded design.

  • src/services/ActivityPubService.js — THREAD_MAX_DEPTH 3 → 1; comment updated.

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

  • Property mode set to 100644
File size: 124.8 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';
23import AudioEmbedService from './AudioEmbedService.js';
24
25const PUBLIC = 'https://www.w3.org/ns/activitystreams#Public';
26// Full JSON-LD context for every AP object we emit: AS2 core + security (publicKey) + the
27// extension terms we actually use (Mastodon/toot + schema.org), each with a term definition
28// so a strict JSON-LD processor resolves them instead of dropping them → valid AS2/JSON-LD.
29// This is the same context shape Mastodon publishes, so Mastodon sees no change.
30const AP_CONTEXT = [
31 'https://www.w3.org/ns/activitystreams',
32 'https://w3id.org/security/v1',
33 {
34 toot: 'http://joinmastodon.org/ns#',
35 schema: 'http://schema.org#',
36 sensitive: 'as:sensitive',
37 Hashtag: 'as:Hashtag',
38 manuallyApprovesFollowers: 'as:manuallyApprovesFollowers',
39 discoverable: 'toot:discoverable',
40 featured: { '@id': 'toot:featured', '@type': '@id' },
41 PropertyValue: 'schema:PropertyValue',
42 value: 'schema:value',
43 embedUrl: { '@id': 'schema:embedUrl', '@type': '@id' },
44 // Poll (Question) extension: Question/oneOf/anyOf/endTime/closed are AS2 core, but the
45 // per-poll unique-voter count is a Mastodon (toot) term — declare it so the emitted
46 // Question stays valid JSON-LD (a strict processor would otherwise drop votersCount).
47 votersCount: 'toot:votersCount',
48 },
49];
50
51// Short random suffix so two activity ids minted in the same millisecond (e.g.
52// parallel saves) don't collide and get deduped by a receiver.
53const rid = () => crypto.randomBytes(4).toString('hex');
54
55// Keep only http(s) URLs — drops javascript:/data:/etc so a remote actor can't
56// smuggle a dangerous scheme into a stored href/src (rendered in owner-only views).
57const safeUrl = (u) => { const s = String(u == null ? '' : u).trim(); return /^https?:\/\//i.test(s) ? s : ''; };
58
59// ── SSRF guard for outbound fetches ───────────────────────────────
60// Remote URLs (actor/keyId/webfinger/inbox/inReplyTo) are attacker-controlled, so
61// every outbound fetch must refuse hosts that resolve to private/loopback ranges
62// (cloud metadata, internal services) — on the initial host AND each redirect hop.
63function isBlockedIp(ip) {
64 if (!ip) return true;
65 const v = net.isIP(ip);
66 if (v === 4) {
67 const o = ip.split('.').map(Number);
68 return o[0] === 127 || o[0] === 10 || o[0] === 0
69 || (o[0] === 172 && o[1] >= 16 && o[1] <= 31)
70 || (o[0] === 192 && o[1] === 168)
71 || (o[0] === 169 && o[1] === 254)
72 || (o[0] === 100 && o[1] >= 64 && o[1] <= 127); // CGNAT
73 }
74 if (v === 6) {
75 const s = ip.toLowerCase().replace(/^\[|\]$/g, '');
76 return s === '::1' || s === '::' || s.startsWith('fc') || s.startsWith('fd') || s.startsWith('fe80')
77 || s.startsWith('::ffff:127.') || s.startsWith('::ffff:10.') || s.startsWith('::ffff:192.168.')
78 || s.startsWith('::ffff:169.254.') || s.startsWith('::ffff:172.');
79 }
80 return true; // not an IP literal we recognise → refuse
81}
82async function assertPublicHost(hostname) {
83 if (net.isIP(hostname)) { if (isBlockedIp(hostname)) throw new Error('ssrf-blocked-ip'); return; }
84 const addrs = await dns.promises.lookup(hostname, { all: true });
85 if (!addrs.length || addrs.some((a) => isBlockedIp(a.address))) throw new Error('ssrf-blocked-host');
86}
87export async function safeFetch(url, opts = {}, maxRedirects = 3) {
88 let target = url;
89 for (let hop = 0; ; hop++) {
90 const u = new URL(target); // throws on malformed → caller's catch
91 if (u.protocol !== 'https:' && u.protocol !== 'http:') throw new Error('ssrf-bad-scheme');
92 await assertPublicHost(u.hostname);
93 const r = await fetch(target, { ...opts, redirect: 'manual', signal: AbortSignal.timeout(8000) });
94 const loc = (r.status >= 300 && r.status < 400) ? r.headers.get('location') : null;
95 if (loc && hop < maxRedirects) { target = new URL(loc, target).toString(); continue; }
96 return r;
97 }
98}
99const MAX_OUTBOX = 20;
100// Cache-buster for the music listen-link → forces Mastodon to re-crawl a FRESH
101// (square) player card. Bump this whenever the twitter:player card dimensions change.
102const FEDI_CARD_VER = '2';
103
104// ── RSA keys per actor (lazy, cached in DB) ───────────────────────
105// Prepared lazily (NOT at module load) — the ap_keys table is created in
106// initializeDatabase(), which runs after this module is imported.
107let _sel, _ins;
108function keyStmts() {
109 if (!_sel) {
110 _sel = db.prepare('SELECT public_pem, private_pem FROM ap_keys WHERE slug = ?');
111 _ins = db.prepare('INSERT OR IGNORE INTO ap_keys (slug, public_pem, private_pem, created_at) VALUES (?,?,?,CURRENT_TIMESTAMP)');
112 }
113 return { sel: _sel, ins: _ins };
114}
115
116export function getOrCreateKeys(slug) {
117 const { sel, ins } = keyStmts();
118 const row = sel.get(slug);
119 if (row) return row;
120 const { publicKey, privateKey } = crypto.generateKeyPairSync('rsa', {
121 modulusLength: 2048,
122 publicKeyEncoding: { type: 'spki', format: 'pem' },
123 privateKeyEncoding: { type: 'pkcs8', format: 'pem' },
124 });
125 ins.run(slug, publicKey, privateKey);
126 return sel.get(slug) || { public_pem: publicKey, private_pem: privateKey };
127}
128
129// ── content negotiation ───────────────────────────────────────────
130// True when the caller wants ActivityPub JSON rather than the HTML page.
131export function apWants(req) {
132 const a = String(req.headers.accept || '').toLowerCase();
133 return a.includes('application/activity+json') ||
134 (a.includes('application/ld+json') && a.includes('activitystreams'));
135}
136
137const AP_CONTENT_TYPE = 'application/activity+json; charset=utf-8';
138export function sendAP(res, obj) {
139 res.type(AP_CONTENT_TYPE);
140 res.set('Cache-Control', 'public, max-age=120');
141 res.send(JSON.stringify(obj));
142}
143
144// ── document builders ─────────────────────────────────────────────
145export function actorId(base, slug) { return `${base}/ap/users/${encodeURIComponent(slug)}`; }
146export function noteId(base, postId) { return `${base}/ap/notes/${encodeURIComponent(postId)}`; }
147
148export function buildActor(base, site) {
149 const id = actorId(base, site.slug);
150 const keys = getOrCreateKeys(site.slug);
151 const actor = {
152 '@context': AP_CONTEXT,
153 id,
154 type: 'Person',
155 preferredUsername: site.slug,
156 name: site.title || site.slug,
157 summary: site.tagline || site.description || '',
158 url: `${base}/${site.slug === site.primary_slug ? '' : 'user/' + encodeURIComponent(site.slug)}`,
159 manuallyApprovesFollowers: false,
160 discoverable: true,
161 inbox: `${id}/inbox`,
162 outbox: `${id}/outbox`,
163 followers: `${id}/followers`,
164 following: `${id}/following`,
165 featured: `${id}/featured`,
166 endpoints: { sharedInbox: `${base}/ap/inbox` },
167 publicKey: {
168 id: `${id}#main-key`,
169 owner: id,
170 publicKeyPem: keys.public_pem,
171 },
172 };
173 if (site.profile_photo) {
174 const u = /^https?:/.test(site.profile_photo) ? site.profile_photo : `${base}${site.profile_photo.startsWith('/') ? '' : '/'}${site.profile_photo}`;
175 actor.icon = { type: 'Image', url: u };
176 }
177 // Account creation date — shown by Mastodon + read by indexers (additive, standard AS2).
178 if (site.created_at) { try { actor.published = new Date(site.created_at).toISOString(); } catch { /* skip bad date */ } }
179 // Profile links → PropertyValue rows: Mastodon/PeerTube/WordPress-ActivityPub render these as
180 // profile metadata (rel=me enables link-back verification). Additive; ignored by simpler receivers.
181 try {
182 const links = JSON.parse(site.profile_links || '[]');
183 if (Array.isArray(links) && links.length) {
184 const esc = (s) => String(s).replace(/[<>&]/g, (c) => ({ '<': '&lt;', '>': '&gt;', '&': '&amp;' }[c]));
185 const rows = links
186 .filter((l) => l && l.url && /^https?:/i.test(l.url))
187 .map((l) => ({
188 type: 'PropertyValue',
189 name: esc(l.platform || 'Link'),
190 value: `<a href="${esc(l.url).replace(/"/g, '&quot;')}" rel="me nofollow noopener" target="_blank">${esc(String(l.url).replace(/^https?:\/\//, ''))}</a>`,
191 }));
192 if (rows.length) actor.attachment = rows;
193 }
194 } catch { /* skip malformed profile_links */ }
195 return actor;
196}
197
198// Does a post's audio shortcodes reference at least one PLAYABLE (file-backed)
199// track? Link-only tracks (external Spotify/YouTube, media_id NULL) don't count —
200// they have no Klonkt-hosted audio to embed, so no player card / cover-suppression.
201export function hasPlayableAudio(content, siteId) {
202 if (!content || !/\[\[(track|album|playlist):/i.test(content)) return false;
203 try {
204 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; }
205 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; }
206 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; }
207 } catch { /* non-fatal */ }
208 return false;
209}
210
211// A single post as an AS2 Note (the object), and as a Create activity (for outbox/delivery).
212export function buildNote(base, site, post) {
213 const id = noteId(base, post.id);
214 const aId = actorId(base, site.slug);
215 const human = `${base}/${encodeURIComponent(post.slug)}`;
216 // Mastodon ignores a Note's `name`, so put the title INTO the content (bold
217 // first line) — the standard blog→fediverse convention. post.content is
218 // already sanitized HTML; the title is plain text, so escape it.
219 const escTitle = String(post.title || '').replace(/[<>&]/g, (c) => ({ '<': '&lt;', '>': '&gt;', '&': '&amp;' }[c]));
220 const titleHtml = post.title ? `<p><strong>${escTitle}</strong></p>` : '';
221
222 // Images travel as AP `attachment` (Mastodon strips <img> from content). Collect
223 // the cover + any inline <img>, make absolute, then strip <img> from the content
224 // to avoid duplicate rendering on clients that DO keep them.
225 const abs = (u) => !u ? null : (/^https?:/i.test(u) ? u : `${base}${u.startsWith('/') ? '' : '/'}${u}`);
226 const mediaType = (u) => {
227 const e = ((u || '').split('?')[0].match(/\.(\w+)$/) || [])[1];
228 return ({ jpg: 'image/jpeg', jpeg: 'image/jpeg', png: 'image/png', gif: 'image/gif', webp: 'image/webp', avif: 'image/avif', mp4: 'video/mp4', webm: 'video/webm', mov: 'video/quicktime' })[(e || '').toLowerCase()] || 'image/jpeg';
229 };
230 const hadAudio = /\[\[(track|album|playlist):/i.test(post.content || '');
231 const playable = hasPlayableAudio(post.content || '', site && site.id);
232 // A post with an external embed (Spotify/YouTube/SoundCloud/Vimeo/Bandcamp/Apple) should let
233 // Mastodon render the embed's player CARD. Mastodon shows EITHER media attachments OR a link
234 // card, never both — so when the post has an embed link we skip the image attachments so the
235 // card wins. (On Klonkt nothing changes: the cover + the embed player still render.)
236 const hasEmbed = (() => {
237 const c = post.content || '';
238 if (/\[\[embed:/i.test(c)) return true;
239 for (const m of c.matchAll(/https?:\/\/[^\s"'<>]+/gi)) if (AudioEmbedService.detectProvider(m[0])) return true;
240 return false;
241 })();
242 // Link-only tracks (external Spotify/YouTube/SoundCloud, no hosted file): collect their links
243 // so we federate them — Mastodon cards the first (its player), the rest show as clickable links
244 // — instead of a bare "listen on site" link, and we suppress the cover so the card can show.
245 const trackEmbedLinks = (() => {
246 if (playable) return [];
247 const out = [];
248 try {
249 for (const m of (post.content || '').matchAll(/\[\[track:([A-Za-z0-9_-]+)\]\]/g)) {
250 const r = db.prepare('SELECT media_id, link_spotify, link_youtube, link_soundcloud FROM audio_tracks WHERE id = ?').get(m[1]);
251 if (r && !r.media_id) for (const u of [r.link_spotify, r.link_youtube, r.link_soundcloud]) if (u && /^https?:\/\//i.test(u)) out.push(u);
252 }
253 } catch { /* non-fatal */ }
254 return [...new Set(out)].slice(0, 6);
255 })();
256 const noImages = playable || hasEmbed || trackEmbedLinks.length > 0; // suppress images → let the player/embed card show
257 const urls = [];
258 // Posts with PLAYABLE hosted audio suppress image attachments so Mastodon renders
259 // the player CARD (twitter:player) instead of the cover — media attachment and
260 // link/player card are mutually exclusive on Mastodon. Link-only audio (external)
261 // keeps its cover (no player card to show).
262 // An animated cover federates as the muted loop MP4 (→ a Video attachment): animated WebP is
263 // unreliable on Mastodon and its iOS apps; the MP4 plays everywhere. Else the still cover image.
264 if (post.cover_video_url && !noImages) urls.push(abs(post.cover_video_url));
265 else if (post.cover_image_url && !noImages) urls.push(abs(post.cover_image_url));
266 let body = post.content || '';
267 // Only federate inline images we can actually serve: absolute http(s) URLs, or our own
268 // /media/ uploads. A relative path we don't host (e.g. a stale /images/... ref) would 404
269 // and show up as a black tile in Mastodon's attachment grid.
270 if (!noImages) for (const m of body.matchAll(/<img\b[^>]*\bsrc="([^"]+)"[^>]*>/gi)) {
271 const src = m[1];
272 if (/^https?:\/\//i.test(src) || src.startsWith('/media/')) urls.push(abs(src));
273 }
274 body = body.replace(/<img\b[^>]*>/gi, '');
275 // Audio shortcodes: do NOT federate the raw audio file — Klonkt deliberately
276 // gates audio (the /audio/stream URL has friction), and shipping it as an AP
277 // audio attachment would hand Mastodon a plain, downloadable mp3 URL. Instead,
278 // replace the shortcodes with a "🎵 listen on the site" link so the post invites
279 // a click-through to the protected player (discovery without leaking the file).
280 const esc = (s) => String(s == null ? '' : s).replace(/[<>&]/g, (c) => ({ '<': '&lt;', '>': '&gt;', '&': '&amp;' }[c]));
281 const audioLabels = [];
282 try {
283 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); }
284 for (const m of body.matchAll(/\[\[album:([^\]]+)\]\]/g)) audioLabels.push(m[1].trim());
285 } catch { /* non-fatal */ }
286 // fedi_open tracks → real AS2 Audio attachments (the actual file URL, served ungated) so
287 // EVERY client incl. the Mastodon apps plays them inline natively. Gated tracks (default)
288 // stay link/card-only — the file is never exposed for them. Resolve from post.content so a
289 // later body mutation can't affect it.
290 const openAudio = [];
291 if (hadAudio) {
292 const seenA = new Set();
293 const addRow = (r) => {
294 const fn = r.filename || (r.storage_path || '').split('/').pop();
295 if (!fn || seenA.has(fn)) return; seenA.add(fn);
296 openAudio.push({ type: 'Audio', mediaType: r.mime_type || 'audio/mpeg', url: `${base}/audio/stream/${encodeURIComponent(fn)}`, name: r.title || 'Audio' });
297 };
298 const SEL = 'SELECT t.title, m.filename, m.storage_path, m.mime_type FROM audio_tracks t JOIN media m ON m.id = t.media_id WHERE t.fedi_open = 1 AND ';
299 try {
300 for (const mm of (post.content || '').matchAll(/\[\[track:([A-Za-z0-9_-]+)\]\]/g)) { const r = db.prepare(SEL + 't.id = ?').get(mm[1]); if (r) addRow(r); }
301 for (const mm of (post.content || '').matchAll(/\[\[album:([^\]]+)\]\]/g)) for (const r of db.prepare(SEL + 't.site_id = ? AND t.album = ? ORDER BY t.rowid').all(site.id, mm[1].trim())) addRow(r);
302 for (const mm of (post.content || '').matchAll(/\[\[playlist:([A-Za-z0-9_-]+)\]\]/g)) for (const r of db.prepare('SELECT t.title, m.filename, m.storage_path, m.mime_type FROM playlist_tracks pt JOIN audio_tracks t ON t.id = pt.track_id JOIN media m ON m.id = t.media_id WHERE t.fedi_open = 1 AND pt.playlist_id = ? ORDER BY pt.position').all(mm[1])) addRow(r);
303 } catch { /* non-fatal */ }
304 }
305 body = body.replace(/\[\[(track|album|playlist):[^\]]+\]\]/gi, '');
306 // External embeds ([[embed:url]]) → emit the bare URL as a link so Mastodon
307 // renders its OWN preview/player card (YouTube/Spotify/SoundCloud/etc) instead
308 // of federating the raw shortcode text.
309 body = body.replace(/\[\[embed:([^\]]+)\]\]/gi, (mm, raw) => {
310 const u = esc(raw.trim().replace(/&amp;/g, '&'));
311 return `<p><a href="${u}">${u}</a></p>`;
312 });
313 if (hadAudio) {
314 const lbl = audioLabels.length ? esc(audioLabels.slice(0, 4).join(', ')) : '';
315 if (trackEmbedLinks.length) {
316 // Link-only track(s): emit the external link(s). Mastodon cards the first (Spotify → its
317 // player), the rest render as clickable links — the fediverse-native "embed + links".
318 body += `<p>🎵 ${lbl ? `<strong>${lbl}</strong>` : ''}</p>`;
319 for (const u of trackEmbedLinks) { const eu = esc(u); body += `<p><a href="${eu}">${eu}</a></p>`; }
320 } else {
321 // For playable posts, append a version param to the listen-link so Mastodon
322 // sees a NEW card URL and re-crawls it (fresh SQUARE player card) instead of
323 // reusing the cached landscape one. Invisible: the link TEXT stays clean, the
324 // page ignores the param. Bump FEDI_CARD_VER when the card dimensions change.
325 const listenHref = playable ? `${human}?fc=${FEDI_CARD_VER}` : human;
326 body += `<p>🎵 ${lbl ? `<strong>${lbl}</strong> — ` : ''}<a href="${listenHref}">listen on ${esc(site.title || 'the site')}</a></p>`;
327 }
328 }
329 // Klonkt renders post content with white-space:pre-wrap, so raw newlines ARE line
330 // breaks on the site. Mastodon (plain HTML) collapses whitespace and would drop them,
331 // so convert newlines to <br> for the federated copy (content already made with
332 // shift+enter uses <br> and has no \n → this is a no-op there).
333 body = body.replace(/\r?\n/g, '<br>');
334 body = linkHashtags(base, body); // link inline #hashtags in the post body too
335 // Append the tags-field hashtags to the content so Mastodon renders them as clickable
336 // hashtags (a Hashtag that's only in the `tag` array isn't shown inline). CamelCase
337 // multi-word tags; skip any already present inline in the body.
338 {
339 const inlineTags = new Set(hashtagTags(base, body).map((h) => h.name.slice(1).toLowerCase()));
340 const addSeen = new Set();
341 const tagLinks = normalizeTags(post.tags).map(tagParts).filter(Boolean)
342 .filter((p) => !inlineTags.has(p.slug) && !addSeen.has(p.slug) && addSeen.add(p.slug))
343 .map((p) => `<a href="${base}/tag/${encodeURIComponent(p.slug)}" class="mention hashtag" rel="tag">#${p.label}</a>`);
344 if (tagLinks.length) body += `<p>${tagLinks.join(' ')}</p>`;
345 }
346 const seen = new Set();
347 const attachment = urls.filter(Boolean)
348 .filter((u) => { if (seen.has(u)) return false; seen.add(u); return true; })
349 .map((u) => { const mt = mediaType(u); // specific AS2 subtype (Image/Audio/Video) over generic Document
350 const ty = /^image\//i.test(mt) ? 'Image' : /^video\//i.test(mt) ? 'Video' : /^audio\//i.test(mt) ? 'Audio' : 'Document';
351 return { type: ty, mediaType: mt, url: u }; });
352 for (const a of openAudio) attachment.push(a); // fedi_open tracks → native Audio players
353
354 // Inline @user@host mentions: the Mention tag objects + the mentioned actor URIs. Only
355 // present when the content was already mention-linked (deliverCreate/Update resolve them
356 // at send time); a plain buildNote (outbox/notes) yields none.
357 const _mentionTags = mentionTags(body);
358 const _mentionCc = _mentionTags.map((t) => t.href);
359
360 const note = {
361 id,
362 type: 'Note',
363 attributedTo: aId,
364 content: titleHtml + body,
365 url: human,
366 published: new Date(post.published_at || post.created_at || Date.now()).toISOString(),
367 // fan_only = "fans only" → followers-only visibility (delivered to your followers
368 // but not addressed to Public, so Mastodon shows it only to them and can't boost it).
369 to: post.fan_only ? [`${aId}/followers`] : [PUBLIC],
370 // Mentioned actors (from inline @user@host links the caller resolved) are addressed in cc
371 // so Mastodon notifies them; empty unless the content was mention-linked (delivery time).
372 cc: [...new Set([...(post.fan_only ? [] : [`${aId}/followers`]), ..._mentionCc])],
373 tag: [...buildHashtagList(base, post.tags, body), ..._mentionTags],
374 replies: `${id}/replies`,
375 // NSFW → Mastodon-style content warning: sensitive (blurs media) + a summary/spoiler
376 // (hides the whole post behind a "Gevoelige inhoud" button until the reader opens it).
377 sensitive: !!post.nsfw,
378 };
379 if (post.nsfw) note.summary = post.content_warning || 'Gevoelige inhoud';
380 if (attachment.length) note.attachment = attachment;
381 // When the cover attachment is suppressed (hosted audio OR an external embed/link-only track →
382 // so Mastodon shows the player/link card, not media), still expose the cover via AS2 `image` so
383 // card/grid consumers (the Klonkt Cirkel/News feed) can show it. Mastodon ignores a Note's
384 // `image`, so its card is unaffected — but a Klonkt receiver reads it (handleInbox o.image).
385 if (post.cover_image_url && noImages) {
386 const cov = abs(post.cover_image_url);
387 if (cov) note.image = { type: 'Image', mediaType: mediaType(cov), url: cov };
388 }
389 // Experiment (mirrors PeerTube / schema.org `embedUrl`): point at the GATED player page
390 // (/embed) so a client that honours embedUrl can show an inline player WITHOUT ever
391 // getting the audio file — the anti-steal posture is untouched. `embedUrl` is a real
392 // standard field name (not a Klonkt invention); if Mastodon's apps honour it on a Note we
393 // make it JSON-LD-clean with a context term, otherwise it degrades to the player card.
394 if (playable) note.embedUrl = `${base}/embed?post=${encodeURIComponent(post.slug)}`;
395 // A hosted poll → federate as an AS2 Question (options + live tally). Do this last so it
396 // reuses the note's content/addressing/tags, then swaps the type and strips media.
397 const ownPoll = parseOwnPoll(post.poll_json);
398 if (ownPoll) applyPollToNote(note, post.id, ownPoll);
399 return note;
400}
401
402// All reply note URIs on a local post (inbound fediverse replies + our own
403// outbound replies) — backs the Note's `replies` Collection so remote servers
404// can fetch the whole thread.
405export function getReplyUris(base, postId) {
406 const out = [];
407 try {
408 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);
409 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}`);
410 } catch { /* non-fatal */ }
411 return out;
412}
413
414// Notifications "seen" tracking → a real bell badge. Stored per site in app_settings.
415export function markNotificationsSeen(slug) {
416 try {
417 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")
418 .run(`fedi_notif_seen:${slug}`, new Date().toISOString());
419 } catch { /* non-fatal */ }
420}
421export function countUnseenNotifications(slug) {
422 try {
423 const row = db.prepare('SELECT value FROM app_settings WHERE key = ?').get(`fedi_notif_seen:${slug}`);
424 const seen = row ? Date.parse(row.value) : 0;
425 let n = 0;
426 for (const it of getNotifications(slug, 50)) { if (Date.parse(it.created_at) > seen) n++; }
427 return n;
428 } catch { return 0; }
429}
430
431export function buildCreate(base, site, post) {
432 const note = buildNote(base, site, post);
433 return {
434 '@context': AP_CONTEXT,
435 id: note.id + '#create',
436 type: 'Create',
437 actor: actorId(base, site.slug),
438 published: note.published,
439 to: note.to,
440 cc: note.cc,
441 object: note,
442 };
443}
444
445export function buildOutbox(base, site, posts) {
446 const id = `${actorId(base, site.slug)}/outbox`;
447 const items = (posts || []).slice(0, MAX_OUTBOX).map((p) => buildCreate(base, site, p));
448 return {
449 '@context': AP_CONTEXT,
450 id,
451 type: 'OrderedCollection',
452 totalItems: items.length,
453 orderedItems: items,
454 };
455}
456
457export function buildFollowers(base, site, count) {
458 const id = `${actorId(base, site.slug)}/followers`;
459 return {
460 '@context': AP_CONTEXT,
461 id,
462 type: 'OrderedCollection',
463 totalItems: count || 0,
464 orderedItems: [], // hidden for privacy; count only
465 };
466}
467
468// The accounts this site follows — count only, mirroring buildFollowers. The spec lists
469// `following` as a standard actor property; Hubzilla/Friendica + crawlers expect it.
470export function buildFollowing(base, site, count) {
471 const id = `${actorId(base, site.slug)}/following`;
472 return {
473 '@context': AP_CONTEXT,
474 id,
475 type: 'OrderedCollection',
476 totalItems: count || 0,
477 orderedItems: [], // count only
478 };
479}
480
481// Pinned posts → the actor's `featured` collection. Mastodon reads this and shows
482// these as the "Featured" tab (pinned to the profile). Posts come ordered by pin
483// rank; embedded as full Notes so a remote server doesn't need extra fetches.
484export function buildFeatured(base, site, posts) {
485 const id = `${actorId(base, site.slug)}/featured`;
486 const items = (posts || []).map((p) => buildNote(base, site, p));
487 return {
488 '@context': AP_CONTEXT,
489 id,
490 type: 'OrderedCollection',
491 totalItems: items.length,
492 orderedItems: items,
493 };
494}
495
496// ── followers store (lazy stmts) ──────────────────────────────────
497let _insF, _delF, _listF, _cntF;
498function fStmts() {
499 if (!_insF) {
500 _insF = db.prepare('INSERT OR IGNORE INTO ap_followers (slug, actor_uri, inbox, shared_inbox, created_at) VALUES (?,?,?,?,CURRENT_TIMESTAMP)');
501 _delF = db.prepare('DELETE FROM ap_followers WHERE slug = ? AND actor_uri = ?');
502 _listF = db.prepare('SELECT inbox, shared_inbox FROM ap_followers WHERE slug = ?');
503 _cntF = db.prepare('SELECT COUNT(*) n FROM ap_followers WHERE slug = ?');
504 }
505 return { ins: _insF, del: _delF, list: _listF, cnt: _cntF };
506}
507export function followerCount(slug) { return fStmts().cnt.get(slug).n; }
508
509// ── inbound interactions store (replies / likes / boosts) + our outbound replies ──
510let _insI, _delLA, _delReply, _listI, _getI, _insO, _listO, _getO;
511function iStmts() {
512 if (!_insI) {
513 _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)');
514 _delLA = db.prepare('DELETE FROM ap_interactions WHERE kind = ? AND post_id = ? AND actor_uri = ?');
515 _delReply = db.prepare("DELETE FROM ap_interactions WHERE kind = 'reply' AND object_uri = ?");
516 _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');
517 _getI = db.prepare('SELECT * FROM ap_interactions WHERE id = ?');
518 _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)');
519 _listO = db.prepare('SELECT * FROM ap_outbox WHERE post_id = ? ORDER BY created_at ASC');
520 _getO = db.prepare('SELECT * FROM ap_outbox WHERE id = ?');
521 }
522 return { ins: _insI, delLA: _delLA, delReply: _delReply, list: _listI, getI: _getI, insO: _insO, listO: _listO, getO: _getO };
523}
524
525export function getInteractionById(id) { return iStmts().getI.get(id); }
526export function setInteractionBoosted(id, on) {
527 db.prepare('UPDATE ap_interactions SET acted_boost = ? WHERE id = ?').run(on ? 1 : 0, id);
528}
529export function setInteractionLiked(id, on) {
530 db.prepare('UPDATE ap_interactions SET acted_like = ? WHERE id = ?').run(on ? 1 : 0, id);
531}
532// Your like/boost state on a REMOTE post (interact page toggles).
533export function setMyReaction(slug, uri, kind, on) {
534 if (on) db.prepare('INSERT OR IGNORE INTO ap_my_reactions (site_slug, target_uri, kind) VALUES (?,?,?)').run(slug, uri, kind);
535 else db.prepare('DELETE FROM ap_my_reactions WHERE site_slug = ? AND target_uri = ? AND kind = ?').run(slug, uri, kind);
536}
537export function getMyReactions(slug, uri) {
538 const rows = (slug && uri) ? db.prepare('SELECT kind FROM ap_my_reactions WHERE site_slug = ? AND target_uri = ?').all(slug, uri) : [];
539 return { liked: rows.some((r) => r.kind === 'like'), boosted: rows.some((r) => r.kind === 'boost') };
540}
541
542const localPostExists = (id) => { try { return !!db.prepare('SELECT 1 FROM posts WHERE id = ?').get(id); } catch { return false; } };
543// Extract our local post id from a note URL, but only if it's ours (base match).
544function postIdFromNoteUrl(url, base) {
545 const s = String(url || '');
546 if (base && !s.startsWith(base)) return null;
547 const m = s.match(/\/ap\/notes\/([^/?#]+)/);
548 return m ? decodeURIComponent(m[1]) : null;
549}
550function deriveHandle(actorUri) {
551 try { const u = new URL(actorUri); const seg = u.pathname.split('/').filter(Boolean).pop() || ''; return `@${seg}@${u.host}`; } catch { return String(actorUri || ''); }
552}
553function actorInfo(doc, actorUri) {
554 let host = ''; try { host = new URL(actorUri).host; } catch { /* keep empty */ }
555 const handle = doc && doc.preferredUsername ? `@${doc.preferredUsername}@${host}` : deriveHandle(actorUri);
556 const icon = doc && doc.icon ? (doc.icon.url || (Array.isArray(doc.icon) && doc.icon[0] && doc.icon[0].url)) : null;
557 return {
558 name: (doc && (doc.name || doc.preferredUsername)) || handle,
559 handle,
560 url: safeUrl((doc && (doc.url || doc.id)) || actorUri) || null,
561 icon: safeUrl(icon) || null,
562 };
563}
564
565// Given an inReplyTo note URL, find which local post the thread belongs to + the
566// note being replied to (parent), so a reply-to-a-comment can be nested.
567function findThreadTarget(inReplyTo, base) {
568 if (!inReplyTo) return null;
569 const seg = postIdFromNoteUrl(inReplyTo, base); // our /ap/notes/<id> segment (if ours)
570 if (seg && localPostExists(seg)) return { post_id: seg, parent_uri: inReplyTo };
571 if (seg) {
572 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 */ }
573 }
574 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 */ }
575 return null;
576}
577
578// Drop the leading @mention(s) a federated reply carries (the person being replied to),
579// so a comment reads "dope tekening ouwe" instead of "@jason@jasonhacky.nl dope …".
580// Keeps a leading <p> wrapper; handles mention <a> links and plain-text @user@domain.
581export function stripLeadingMentions(html) {
582 if (!html) return html;
583 let s = String(html);
584 s = s.replace(/^(\s*<p[^>]*>)?\s*(?:<a\b[^>]*>\s*@[^<]+<\/a>[  ]*)+/i, (m, p) => p || '');
585 s = s.replace(/^(\s*<p[^>]*>)?\s*(?:@[\w.-]+(?:@[\w.-]+)?[  ]+)+/i, (m, p) => p || '');
586 return s;
587}
588
589// View-ready threaded view of a post's fediverse activity (inbound replies +
590// our outbound replies, nested), plus like/boost counts.
591export function getInteractions(postId, base, site) {
592 const s = iStmts();
593 const rows = s.list.all(postId);
594 const baseClean = (base || process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
595 const postNoteId = baseClean ? `${baseClean}/ap/notes/${postId}` : null;
596 // Our own (outbound) replies show the SITE identity for everyone (not "You").
597 let host = ''; try { host = new URL(baseClean).host; } catch { /* ignore */ }
598 const siteName = (site && (site.title || site.slug)) || '';
599 const siteHandle = (site && site.slug && host) ? `@${site.slug}@${host}` : '';
600 const siteUrl = baseClean ? `${baseClean}/` : '';
601 const siteIcon = (site && site.profile_photo) || null;
602
603 const nodes = [];
604 for (const r of rows) {
605 if (r.kind !== 'reply') continue;
606 nodes.push({
607 noteId: r.object_uri, parent: r.parent_uri || null, mine: false, id: r.id,
608 actor_name: r.actor_name, actor_handle: r.actor_handle, actor_url: r.actor_url,
609 actor_icon: r.actor_icon, content: stripLeadingMentions(r.content), created_at: r.published || r.created_at,
610 acted_boost: !!r.acted_boost, acted_like: !!r.acted_like,
611 children: [],
612 });
613 }
614 for (const o of s.listO.all(postId)) {
615 nodes.push({
616 noteId: baseClean ? `${baseClean}/ap/notes/${o.id}` : o.id, parent: o.in_reply_to || null,
617 mine: true, outboxId: o.id, content: stripLeadingMentions(o.content), created_at: o.created_at,
618 actor_name: siteName, actor_handle: siteHandle, actor_url: siteUrl, actor_icon: siteIcon,
619 children: [],
620 });
621 }
622
623 const byId = new Map(nodes.map((n) => [n.noteId, n]));
624 const isTop = (n) => !n.parent || n.parent === postNoteId || !byId.has(n.parent);
625 const tops = [];
626 for (const n of nodes) {
627 if (isTop(n)) { tops.push(n); continue; }
628 let anc = n, guard = 0;
629 while (!isTop(anc) && guard++ < 12) anc = byId.get(anc.parent);
630 anc.children.push(n);
631 }
632 const byTime = (a, b) => new Date(a.created_at) - new Date(b.created_at);
633 tops.sort(byTime).forEach((t) => t.children.sort(byTime));
634
635 return {
636 thread: tops,
637 likeCount: rows.filter((r) => r.kind === 'like').length,
638 announceCount: rows.filter((r) => r.kind === 'announce').length,
639 total: nodes.length,
640 };
641}
642
643// ── HTTP Signatures + delivery ────────────────────────────────────
644const slugFromActorUrl = (url) => { const m = String(url || '').match(/\/ap\/users\/([^/?#]+)/); return m ? decodeURIComponent(m[1]) : null; };
645
646// Sign + POST an activity to a remote inbox (draft-cavage HTTP Signatures, RSA-SHA256).
647export async function deliver(inboxUrl, bodyObj, keyId, privatePem) {
648 const body = JSON.stringify(bodyObj);
649 const u = new URL(inboxUrl);
650 const date = new Date().toUTCString();
651 const digest = 'SHA-256=' + crypto.createHash('sha256').update(body).digest('base64');
652 const signingString = `(request-target): post ${u.pathname}\nhost: ${u.host}\ndate: ${date}\ndigest: ${digest}`;
653 const signature = crypto.sign('sha256', Buffer.from(signingString), privatePem).toString('base64');
654 const sig = `keyId="${keyId}",algorithm="rsa-sha256",headers="(request-target) host date digest",signature="${signature}"`;
655 const r = await safeFetch(inboxUrl, {
656 method: 'POST',
657 headers: { 'Content-Type': 'application/activity+json', Accept: 'application/activity+json', Date: date, Digest: digest, Signature: sig },
658 body,
659 });
660 return r.status;
661}
662
663export async function fetchActor(url) {
664 try {
665 const r = await safeFetch(url, { headers: { Accept: 'application/activity+json' } });
666 if (!r.ok) return null;
667 const len = Number(r.headers.get('content-length') || 0);
668 if (len > 2_000_000) return null; // refuse oversized actor docs
669 return await r.json();
670 } catch { return null; }
671}
672
673// ── Delivery queue with retries ───────────────────────────────────
674// Outbound deliveries are tried immediately; on failure (down server, timeout,
675// non-2xx) they're queued and retried with backoff so a briefly-offline follower
676// doesn't silently miss the post. The signing key is NOT stored — the worker
677// re-derives it from the actor slug at send time.
678const DELIVERY_MAX_ATTEMPTS = 6;
679const DELIVERY_BACKOFF_MIN = [1, 5, 15, 60, 180, 360];
680let _insDeliv, _dueDeliv, _delDeliv, _bumpDeliv;
681function deliveryStmts() {
682 if (!_insDeliv) {
683 _insDeliv = db.prepare('INSERT INTO ap_delivery (slug, inbox, body, attempts, next_at) VALUES (?,?,?,0,CURRENT_TIMESTAMP)');
684 _dueDeliv = db.prepare("SELECT * FROM ap_delivery WHERE datetime(next_at) <= datetime('now') ORDER BY next_at LIMIT 30");
685 _delDeliv = db.prepare('DELETE FROM ap_delivery WHERE id = ?');
686 _bumpDeliv = db.prepare('UPDATE ap_delivery SET attempts = ?, next_at = ? WHERE id = ?');
687 }
688 return { ins: _insDeliv, due: _dueDeliv, del: _delDeliv, bump: _bumpDeliv };
689}
690export function enqueueDelivery(slug, inbox, activity) {
691 if (!slug || !inbox || !activity) return;
692 try { deliveryStmts().ins.run(slug, inbox, JSON.stringify(activity)); } catch { /* ignore */ }
693}
694// Deliver now; queue for retry if it fails.
695export async function deliverWithRetry(slug, inbox, activity, keyId, privPem) {
696 if (!inbox) return;
697 try { const st = await deliver(inbox, activity, keyId, privPem); if (st >= 200 && st < 300) return; } catch { /* queue below */ }
698 enqueueDelivery(slug, inbox, activity);
699}
700let _processingDeliv = false;
701export async function processDeliveryQueue() {
702 if (_processingDeliv) return; // re-entrancy guard: 30 rows × 8s can exceed the 60s tick → no double-delivery
703 _processingDeliv = true;
704 try {
705 let rows;
706 try { rows = deliveryStmts().due.all(); } catch { return; }
707 if (!rows || !rows.length) return;
708 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
709 for (const row of rows) {
710 let ok = false;
711 try {
712 const keys = getOrCreateKeys(row.slug);
713 const st = await deliver(row.inbox, JSON.parse(row.body), `${actorId(base, row.slug)}#main-key`, keys.private_pem);
714 ok = st >= 200 && st < 300;
715 } catch { ok = false; }
716 if (ok) { deliveryStmts().del.run(row.id); continue; }
717 const attempts = row.attempts + 1;
718 if (attempts >= DELIVERY_MAX_ATTEMPTS) { deliveryStmts().del.run(row.id); console.warn('[AP] delivery gave up after', attempts, 'tries →', row.inbox); continue; }
719 // Index the backoff on the CURRENT attempt count (row.attempts) so the first
720 // retry uses the 1-min tier instead of skipping it.
721 const mins = DELIVERY_BACKOFF_MIN[Math.min(row.attempts, DELIVERY_BACKOFF_MIN.length - 1)];
722 deliveryStmts().bump.run(attempts, new Date(Date.now() + mins * 60000).toISOString(), row.id);
723 }
724 } finally { _processingDeliv = false; }
725}
726let _delivTimer = null;
727export function startDeliveryWorker() {
728 if (_delivTimer) return;
729 _delivTimer = setInterval(() => { processDeliveryQueue().catch(() => {}); }, 60 * 1000);
730 if (_delivTimer.unref) _delivTimer.unref();
731}
732
733// Best-effort verification of an incoming signed request. Returns the sender's
734// actor doc if the signature checks out, else null. (Not gating yet — MVP.)
735// Max clock skew for the signed Date header (replay window). Generous default to tolerate
736// federating servers with drifting clocks; an operator can widen it via env.
737const SIG_MAX_SKEW_MS = (Number(process.env.AP_SIG_MAX_SKEW_MIN) || 60) * 60 * 1000;
738export async function verifyRequest(req) {
739 const sigH = req.headers['signature'];
740 if (!sigH) return null;
741 const p = Object.fromEntries([...sigH.matchAll(/([a-zA-Z]+)="([^"]*)"/g)].map((m) => [m[1], m[2]]));
742 if (!p.keyId || !p.signature) return null;
743 const actor = await fetchActor(p.keyId.split('#')[0]);
744 const pem = actor && actor.publicKey && actor.publicKey.publicKeyPem;
745 if (!pem) return null;
746 const hs = (p.headers || '(request-target) host date').split(/\s+/);
747 // Behind a reverse proxy the raw Host header is the backend bind (e.g. localhost:3000, when
748 // the proxy doesn't preserve it — Apache .htaccess [P] proxying), but the sender signed the
749 // HTTP-Signature over the PUBLIC host. Try each candidate host (the configured PUBLIC_BASE_URL
750 // host, the proxy's X-Forwarded-Host, and the raw Host) and accept if the signature verifies
751 // against any. An attacker can't forge a match (no private key), so this only rescues the
752 // legitimate proxied case. Also normalise a leading double-slash in the request-target.
753 let _pubHost = null;
754 if (process.env.PUBLIC_BASE_URL) { try { _pubHost = new URL(process.env.PUBLIC_BASE_URL).host; } catch { /* ignore */ } }
755 const _hosts = [...new Set([_pubHost, req.headers['x-forwarded-host'], req.headers['host']].filter(Boolean))];
756 const _target = `${req.method.toLowerCase()} ${String(req.originalUrl || '').replace(/^\/{2,}/, '/')}`;
757 const _sig = Buffer.from(p.signature, 'base64');
758 let ok = false;
759 for (const _h of _hosts) {
760 const line = hs.map((x) => x === '(request-target)'
761 ? `(request-target): ${_target}`
762 : x === 'host' ? `host: ${_h}`
763 : `${x}: ${req.headers[x] || ''}`).join('\n');
764 try { if (crypto.verify('sha256', Buffer.from(line), pem, _sig)) { ok = true; break; } } catch { /* try next host */ }
765 }
766 // Replay defence: the Date header must be signed and recent. A captured signed request
767 // replayed later (or with a swapped body) is rejected.
768 if (ok) {
769 if (!hs.includes('date')) ok = false;
770 else {
771 const t = Date.parse(req.headers['date'] || '');
772 if (isNaN(t) || Math.abs(Date.now() - t) > SIG_MAX_SKEW_MS) ok = false;
773 }
774 }
775 // Digest is MANDATORY when the request carries a body: without a signed digest the body
776 // isn't covered by the signature and could be swapped on a replay.
777 if (ok && req.rawBody && req.rawBody.length) {
778 if (!hs.includes('digest')) ok = false;
779 else {
780 const exp = 'SHA-256=' + crypto.createHash('sha256').update(req.rawBody).digest('base64');
781 if (req.headers['digest'] !== exp) ok = false;
782 }
783 }
784 return ok ? actor : null;
785}
786
787// Parse a fediverse poll (an ActivityStreams `Question` — the Mastodon-standard poll form)
788// into our compact shape. `oneOf` = single choice, `anyOf` = multiple; each option is a Note
789// with a `name` and a `replies` collection whose `totalItems` is that option's vote count.
790function parsePoll(o) {
791 if (!o || o.type !== 'Question') return null;
792 const raw = Array.isArray(o.oneOf) ? o.oneOf : (Array.isArray(o.anyOf) ? o.anyOf : null);
793 if (!raw || !raw.length) return null;
794 const options = raw.slice(0, 12).map((opt) => ({
795 name: String((opt && opt.name) || '').slice(0, 300),
796 count: Math.max(0, Number(opt && opt.replies && opt.replies.totalItems) || 0),
797 })).filter((x) => x.name);
798 if (!options.length) return null;
799 const endTime = o.endTime || (typeof o.closed === 'string' ? o.closed : null);
800 const closed = !!o.closed || (endTime ? Date.parse(endTime) <= Date.now() : false);
801 return { multiple: Array.isArray(o.anyOf), options, endTime, closed, voters: Number(o.votersCount) || null, voted: null };
802}
803
804// ── Polls WE host (a local post with a poll) ──────────────────────
805// Parse the poll definition stored on our own post (posts.poll_json). Counts are
806// NOT stored here — they're derived from the poll_votes ballots so a re-render always
807// reflects the authoritative tally.
808export function parseOwnPoll(pollJson) {
809 if (!pollJson) return null;
810 let d; try { d = typeof pollJson === 'string' ? JSON.parse(pollJson) : pollJson; } catch { return null; }
811 if (!d || !Array.isArray(d.options)) return null;
812 const options = d.options.map((o) => ({ name: String((o && o.name != null ? o.name : o) || '').slice(0, 300) })).filter((o) => o.name);
813 if (options.length < 2) return null;
814 const endTime = d.endTime || null;
815 const closed = !!d.closed || (endTime ? Date.parse(endTime) <= Date.now() : false);
816 return { multiple: !!d.multiple, options, endTime, closed };
817}
818
819// Live tally of a hosted poll from its ballots: per-option counts + unique voters.
820export function pollTally(postId) {
821 const counts = {}; let voters = 0;
822 try {
823 for (const r of db.prepare('SELECT choice, COUNT(*) AS n FROM poll_votes WHERE post_id = ? GROUP BY choice').all(postId)) counts[r.choice] = r.n;
824 voters = db.prepare('SELECT COUNT(DISTINCT actor_uri) AS n FROM poll_votes WHERE post_id = ?').get(postId).n || 0;
825 } catch { /* table may not exist yet */ }
826 return { counts, voters };
827}
828
829// Render-ready view of a hosted poll (options with counts + percentages, totals, state).
830// Voting is fediverse-only, so this is display-only on the site.
831export function ownPollView(post) {
832 const poll = parseOwnPoll(post && post.poll_json);
833 if (!poll) return null;
834 const { counts, voters } = pollTally(post.id);
835 const total = Object.values(counts).reduce((a, b) => a + b, 0);
836 const denom = poll.multiple ? voters : total; // multiple-choice %: share of voters (can sum >100%)
837 const options = poll.options.map((o) => {
838 const count = counts[o.name] || 0;
839 return { name: o.name, count, pct: denom ? Math.round((count / denom) * 100) : 0 };
840 });
841 return { multiple: poll.multiple, options, total, voters, endTime: poll.endTime, closed: poll.closed };
842}
843
844// Attach the AS2 Question shape to a note built for a hosted poll. Mastodon renders a
845// status with either media OR a poll (never both), so a poll federates as content +
846// options with no media attachment. oneOf = single choice, anyOf = multiple.
847function applyPollToNote(note, postId, poll) {
848 const { counts, voters } = pollTally(postId);
849 const opts = poll.options.map((o) => ({
850 type: 'Note',
851 name: o.name,
852 replies: { type: 'Collection', totalItems: counts[o.name] || 0 },
853 }));
854 note.type = 'Question';
855 note[poll.multiple ? 'anyOf' : 'oneOf'] = opts;
856 if (poll.endTime) note.endTime = new Date(poll.endTime).toISOString();
857 // Once closed, Mastodon expects a `closed` timestamp (the effective end).
858 if (poll.closed) note.closed = poll.endTime ? new Date(poll.endTime).toISOString() : new Date().toISOString();
859 note.votersCount = voters;
860 delete note.attachment; // media + poll are mutually exclusive on Mastodon
861 delete note.image;
862 return note;
863}
864
865// Record an inbound ballot on one of OUR polls. A vote arrives as a Create(Note) whose
866// `name` is the chosen option and `inReplyTo` is our poll note — the Mastodon-standard
867// vote form. Returns { handled } — handled=true means it was addressed to a poll (so the
868// caller must NOT also store it as a reply), false means "not a poll, fall through".
869function recordPollBallot(postId, actorUri, rawChoice) {
870 const choice = String(rawChoice == null ? '' : rawChoice).slice(0, 300);
871 if (!choice) return { handled: false };
872 let post; try { post = db.prepare('SELECT poll_json FROM posts WHERE id = ?').get(postId); } catch { return { handled: false }; }
873 const poll = post && parseOwnPoll(post.poll_json);
874 if (!poll) return { handled: false }; // not a poll → let the reply logic handle it
875 if (poll.closed) return { handled: true }; // voting closed → drop
876 if (!poll.options.some((o) => o.name === choice)) return { handled: true }; // unknown option → drop
877 try {
878 // Single choice = one ballot per actor: ignore a later/different vote. Multiple choice
879 // allows one ballot per distinct option (the UNIQUE(post,actor,choice) dedupes repeats).
880 if (!poll.multiple && db.prepare('SELECT 1 FROM poll_votes WHERE post_id = ? AND actor_uri = ? LIMIT 1').get(postId, actorUri)) return { handled: true };
881 db.prepare('INSERT OR IGNORE INTO poll_votes (post_id, actor_uri, choice) VALUES (?, ?, ?)').run(postId, actorUri, choice);
882 } catch { return { handled: true }; }
883 schedulePollUpdate(postId);
884 return { handled: true };
885}
886
887// Coalesce a burst of votes into ONE Update(Question) per poll: the first vote schedules a
888// refresh ~15s out; further votes in that window ride the same pending update (which carries
889// the accumulated tally). Non-follower voters re-fetch the Question (live tally) themselves.
890const _pollUpdTimers = new Map();
891function schedulePollUpdate(postId) {
892 if (_pollUpdTimers.has(postId)) return;
893 const t = setTimeout(() => { _pollUpdTimers.delete(postId); deliverPollUpdate(postId).catch(() => { /* best-effort */ }); }, 15000);
894 if (t.unref) t.unref();
895 _pollUpdTimers.set(postId, t);
896}
897
898// Push the fresh poll tally (or closed state) to followers as Update(Question).
899export async function deliverPollUpdate(postId) {
900 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
901 if (!base || !postId) return;
902 let post, site;
903 try {
904 post = db.prepare('SELECT * FROM posts WHERE id = ?').get(postId);
905 if (!post || !post.poll_json) return;
906 site = db.prepare('SELECT * FROM sites WHERE id = ?').get(post.site_id);
907 } catch { return; }
908 if (site) await deliverUpdate(site, post);
909}
910
911// Handle an incoming inbox POST. slugParam = null for the shared /ap/inbox.
912export async function handleInbox(req, slugParam) {
913 const act = req.body || {};
914 const type = act.type;
915 // Real client IP (behind the proxy via `trust proxy`) — logged on dropped/rejected/
916 // ignored inbox hits so an operator can see who is probing their fediverse inbox.
917 const ip = req.ip || (req.connection && req.connection.remoteAddress) || '?';
918 const base = (process.env.PUBLIC_BASE_URL || `${req.protocol}://${req.get('host')}`).replace(/\/+$/, '');
919 const verified = await verifyRequest(req).catch(() => null);
920
921 // ENFORCE HTTP signatures: a data-affecting activity must be signed by the very
922 // actor it claims to be. No valid signature, or signer ≠ actor → reject (no
923 // forged replies/likes/follows/timeline posts). GET/discovery stays open.
924 const claimedActor = typeof act.actor === 'string' ? act.actor : (act.actor && act.actor.id);
925 // Blocked actor/domain → silently drop (202, don't reveal the block).
926 if (claimedActor && isBlockedAny(claimedActor)) { console.log('[AP] inbox dropped (blocked)', claimedActor, 'from', ip); return 202; }
927 const GATED = ['Create', 'Like', 'Announce', 'Follow', 'Delete', 'Undo', 'Accept', 'Reject', 'Add', 'Remove', 'Update'];
928 if (GATED.includes(type)) {
929 if (!verified || !claimedActor || verified.id !== claimedActor) {
930 console.warn('[AP] inbox REJECTED (signature)', type, claimedActor || '?', 'from', ip, verified ? '(signer mismatch)' : '(unsigned/invalid)');
931 return 401;
932 }
933 }
934
935 if (type === 'Follow') {
936 const who = typeof act.actor === 'string' ? act.actor : (act.actor && act.actor.id);
937 const slug = slugParam || slugFromActorUrl(typeof act.object === 'string' ? act.object : (act.object && act.object.id));
938 if (!who || !slug) return 400;
939 const remote = await fetchActor(who);
940 if (!remote || !remote.inbox) return 202; // can't reach them → drop quietly
941 const sharedInbox = (remote.endpoints && remote.endpoints.sharedInbox) || null;
942 fStmts().ins.run(slug, who, remote.inbox, sharedInbox);
943 const me = actorId(base, slug);
944 const keys = getOrCreateKeys(slug);
945 const accept = { '@context': AP_CONTEXT, id: `${me}#accept-${Date.now()}-${rid()}`, type: 'Accept', actor: me, object: act };
946 deliver(remote.inbox, accept, `${me}#main-key`, keys.private_pem).catch((e) => console.warn('[AP] Accept delivery failed:', e.message));
947 // Auto-backfill: send our recent posts as Create so the instance has our history
948 // (Mastodon doesn't fetch history on follow). ONCE PER REMOTE INSTANCE only —
949 // Mastodon dedupes notes per-instance, so re-filling an instance that already has
950 // a follower of ours is wasted work (and won't re-populate the new follower's
951 // timeline anyway). Deliver to the shared inbox (instance-level) when present.
952 // Sync insert+check (no await between) → no interleave race with concurrent Follows.
953 const instanceFilled = sharedInbox &&
954 db.prepare('SELECT 1 FROM ap_followers WHERE slug = ? AND shared_inbox = ? AND actor_uri != ? LIMIT 1')
955 .get(slug, sharedInbox, who);
956 if (!instanceFilled) {
957 backfillNewFollower(base, slug, sharedInbox || remote.inbox).catch(() => { /* best-effort */ });
958 }
959 console.log('[AP] Follow', who, '→', slug, verified ? '(sig ok)' : '(sig unverified)');
960 return 202;
961 }
962 if (type === 'Undo' && act.object) {
963 const who = typeof act.actor === 'string' ? act.actor : (act.actor && act.actor.id);
964 const ot = act.object.type;
965 if (ot === 'Follow') {
966 const obj = act.object.object;
967 const slug = slugParam || slugFromActorUrl(typeof obj === 'string' ? obj : (obj && obj.id));
968 if (who && slug) { fStmts().del.run(slug, who); console.log('[AP] Unfollow', who, '→', slug); }
969 return 202;
970 }
971 if (ot === 'Like' || ot === 'Announce') {
972 const tgt = act.object.object;
973 const pid = postIdFromNoteUrl(typeof tgt === 'string' ? tgt : (tgt && tgt.id), base);
974 if (who && pid) { iStmts().delLA.run(ot.toLowerCase(), pid, who); console.log('[AP] Undo', ot, who, '→', pid); }
975 return 202;
976 }
977 return 202;
978 }
979
980 const actorUri = typeof act.actor === 'string' ? act.actor : (act.actor && act.actor.id);
981 const resolveActor = async (uri) => ((verified && verified.id === uri) ? verified : await fetchActor(uri).catch(() => null));
982 // Activities from our OWN actors are already stored via ap_outbox — don't re-store.
983 const isLocalActor = !!(base && actorUri && actorUri.startsWith(`${base}/ap/users/`));
984
985 // Inbound reply: a Create whose object replies to one of our notes (post OR comment).
986 if (type === 'Create' && act.object && (act.object.type === 'Note' || act.object.type === 'Article' || act.object.type === 'Question')) {
987 const o = act.object;
988 // A poll ballot: a Note carrying a `name` (the chosen option) inReplyTo one of OUR poll
989 // posts. Record it (deduped per actor) BEFORE the reply logic so a vote is never stored
990 // as a comment. recordPollBallot returns handled=false only if the target isn't a poll.
991 if (o.name && o.inReplyTo && actorUri && !isLocalActor) {
992 const seg = postIdFromNoteUrl(o.inReplyTo, base);
993 if (seg && localPostExists(seg)) {
994 const rec = recordPollBallot(seg, actorUri, o.name);
995 if (rec.handled) { console.log('[AP] poll vote', actorUri, '→', seg); return 202; }
996 }
997 }
998 const tgt = findThreadTarget(o.inReplyTo, base);
999 if (tgt && actorUri && !isLocalActor) {
1000 const ai = actorInfo(await resolveActor(actorUri), actorUri);
1001 const html = HtmlSanitizerService.sanitize(o.content || '');
1002 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);
1003 console.log('[AP] reply', actorUri, '→', tgt.post_id);
1004 return 202;
1005 }
1006 // Home timeline (client): a top-level post from an account we follow.
1007 if (actorUri && !isLocalActor && !o.inReplyTo && o.id) {
1008 let subs = []; try { subs = db.prepare('SELECT slug, auto_boost FROM ap_following WHERE actor_uri = ?').all(actorUri); } catch { /* table may not exist yet */ }
1009 if (subs.length) {
1010 const ai = actorInfo(await resolveActor(actorUri), actorUri);
1011 const html = HtmlSanitizerService.sanitize(o.content || '');
1012 const _atts = (Array.isArray(o.attachment) ? o.attachment : []).map((a) => ({ url: safeUrl(a && a.url), type: (a && a.mediaType) || '' })).filter((m) => m.url);
1013 // Fallback cover: a Note's `image` (set when the attachment was suppressed
1014 // for a player-card post, e.g. hosted-audio posts).
1015 if (!_atts.some((m) => !m.type || /image/i.test(m.type)) && o.image) {
1016 const _im = Array.isArray(o.image) ? o.image[0] : o.image;
1017 const _iu = safeUrl(typeof _im === 'string' ? _im : (_im && _im.url));
1018 if (_iu) _atts.push({ url: _iu, type: (_im && _im.mediaType) || 'image/jpeg' });
1019 }
1020 const media = JSON.stringify(_atts);
1021 const poll = parsePoll(o); // a Question (fediverse poll) → cache its options/counts
1022 // "Feature" = show in the Cirkel (local only). We do NOT auto-Announce
1023 // incoming posts to the fediverse — that flooded followers. Boosting to the
1024 // fediverse is only ever a deliberate, manual per-post action (the 🔁 on
1025 // the timeline).
1026 for (const s of subs) {
1027 tlStmts().ins.run(o.id, s.slug, actorUri, ai.name, ai.handle, ai.icon, ai.url, html, o.url || null, o.published || null, media, o.sensitive ? 1 : 0, o.summary || null);
1028 if (poll) { try { db.prepare('UPDATE ap_timeline SET poll_json = ? WHERE id = ? AND slug = ?').run(JSON.stringify(poll), o.id, s.slug); } catch { /* ignore */ } }
1029 }
1030 console.log('[AP] timeline +', actorUri, 'x' + subs.length);
1031 }
1032 }
1033 return 202;
1034 }
1035 // A remote post we cached was edited upstream → refresh our cached copy. This is the
1036 // push-based edit-sync that keeps the Cirkel/timeline fresh without polling (selfHeal
1037 // does it on a version bump; this does it live). Scope to the SIGNING actor so B can't
1038 // edit A's note (the signature gate guarantees claimedActor == the verified signer).
1039 if (type === 'Update' && act.object && (act.object.type === 'Note' || act.object.type === 'Article' || act.object.type === 'Question')) {
1040 const o = act.object;
1041 if (o.id && claimedActor) {
1042 const html = HtmlSanitizerService.sanitize(o.content || '');
1043 const media = mediaFromNote(o);
1044 try {
1045 // Refresh url too (COALESCE keeps the old one if the Update omits it): a remote slug
1046 // rename keeps the same AP id but changes the human url, so without this the cached
1047 // post would keep linking to the old, now-dead URL.
1048 const r = db.prepare('UPDATE ap_timeline SET content = ?, media_json = ?, nsfw = ?, cw = ?, url = COALESCE(?, url) WHERE id = ? AND author_uri = ?')
1049 .run(html, media, o.sensitive ? 1 : 0, o.summary || null, o.url || null, o.id, claimedActor);
1050 if (r.changes) console.log('[AP] timeline update', claimedActor, '→', o.id);
1051 // A poll's Update carries the fresh vote counts / closed state. Refresh per-row so each
1052 // site keeps its own `voted` state while the counts/closed update to the new totals.
1053 const poll = parsePoll(o);
1054 if (poll) {
1055 const rows = db.prepare('SELECT rowid AS rid, poll_json FROM ap_timeline WHERE id = ? AND author_uri = ?').all(o.id, claimedActor);
1056 const upd = db.prepare('UPDATE ap_timeline SET poll_json = ? WHERE rowid = ?');
1057 for (const rw of rows) {
1058 let voted = null; try { voted = rw.poll_json ? (JSON.parse(rw.poll_json).voted || null) : null; } catch { /* ignore */ }
1059 upd.run(JSON.stringify({ ...poll, voted }), rw.rid);
1060 }
1061 }
1062 } catch { /* ignore */ }
1063 // If this note is a cached fediverse reply on one of our posts, refresh its text too.
1064 try { db.prepare('UPDATE ap_interactions SET content = ? WHERE object_uri = ? AND actor_uri = ?').run(html, o.id, claimedActor); } catch { /* ignore */ }
1065 }
1066 return 202;
1067 }
1068 if (type === 'Like' || type === 'Announce') {
1069 const tgt = act.object;
1070 const objUrl = typeof tgt === 'string' ? tgt : (tgt && tgt.id);
1071 const pid = postIdFromNoteUrl(objUrl, base);
1072 if (pid && actorUri && !isLocalActor && localPostExists(pid)) {
1073 const ai = actorInfo(await resolveActor(actorUri), actorUri);
1074 iStmts().ins.run(type.toLowerCase(), pid, '', actorUri, ai.name, ai.handle, ai.url, ai.icon, null, null, null);
1075 console.log('[AP]', type === 'Like' ? 'like' : 'boost', actorUri, '→', pid);
1076 } else if (type === 'Announce' && objUrl && actorUri && !isLocalActor) {
1077 // A boost FROM an account we follow, of a REMOTE post → show it in the News feed.
1078 // We only STORE it for display; we NEVER auto-Announce it onward (anti-feedback-loop:
1079 // re-announcing an incoming Announce would cascade boosts across the network).
1080 let subs = []; try { subs = db.prepare('SELECT slug FROM ap_following WHERE actor_uri = ?').all(actorUri); } catch { /* table may not exist */ }
1081 if (subs.length) {
1082 const bn = await fetchNoteAP(objUrl);
1083 if (bn && bn !== 404 && (bn.type === 'Note' || bn.type === 'Article') && bn.id) {
1084 const origUri = actorUriOf(bn.attributedTo);
1085 // Block completeness: even if you follow the booster, drop a boost whose ORIGINAL
1086 // author is blocked — otherwise a block is bypassed via someone else's boost.
1087 if (origUri && isBlockedAny(origUri)) { console.log('[AP] timeline boost dropped (blocked origin)', origUri, 'via', actorUri); return 202; }
1088 const oai = actorInfo(await resolveActor(origUri), origUri);
1089 const html = HtmlSanitizerService.sanitize(bn.content || '');
1090 const media = mediaFromNote(bn);
1091 const booster = actorInfo(await resolveActor(actorUri), actorUri);
1092 for (const s of subs) {
1093 // published = now → the boost shows as fresh activity at the top (Mastodon shows
1094 // reblogs at reblog-time, not the original's date). INSERT OR IGNORE: if we already
1095 // have the note (e.g. we also follow the author), keep it and DON'T relabel it.
1096 let inserted = false;
1097 try { const r = tlStmts().ins.run(bn.id, s.slug, origUri || '', oai.name, oai.handle, oai.icon, oai.url, html, bn.url || null, new Date().toISOString(), media, bn.sensitive ? 1 : 0, bn.summary || null); inserted = r.changes > 0; } catch { /* ignore */ }
1098 if (inserted) { try { db.prepare('UPDATE ap_timeline SET reblog_name = ?, reblog_handle = ?, reblog_icon = ? WHERE slug = ? AND id = ?').run(booster.name, booster.handle, booster.icon, s.slug, bn.id); } catch { /* ignore */ } }
1099 }
1100 console.log('[AP] timeline boost +', actorUri, 'x' + subs.length);
1101 }
1102 }
1103 }
1104 return 202;
1105 }
1106 if (type === 'Delete') {
1107 // A remote note was deleted upstream → drop it from replies AND the timeline.
1108 // Scope to the SIGNING actor so actor B can't delete actor A's content (the
1109 // signature gate guarantees claimedActor == the verified signer here).
1110 const oid = typeof act.object === 'string' ? act.object : (act.object && act.object.id);
1111 if (oid && claimedActor) {
1112 try { db.prepare('DELETE FROM ap_interactions WHERE object_uri = ? AND actor_uri = ?').run(oid, claimedActor); } catch { /* ignore */ }
1113 try { db.prepare('DELETE FROM ap_timeline WHERE id = ? AND author_uri = ?').run(oid, claimedActor); } catch { /* ignore */ }
1114 // Also clear a boost/like YOU made of this now-deleted remote post (the interact-page
1115 // ap_my_reactions state), so it can't stay stuck as "boosted" on a post that's gone.
1116 // Guard: only when the deleter owns the note's domain (B mustn't clear your reactions
1117 // to A's posts).
1118 try {
1119 let sameHost = false;
1120 try { sameHost = new URL(oid).host === new URL(claimedActor).host; } catch { sameHost = false; }
1121 if (sameHost) db.prepare('DELETE FROM ap_my_reactions WHERE target_uri = ?').run(oid);
1122 } catch { /* ignore */ }
1123 }
1124 return 202;
1125 }
1126 // Accept/Reject of a Follow WE sent (client side).
1127 if (type === 'Accept' && act.object) {
1128 const fid = typeof act.object === 'string' ? act.object : (act.object && act.object.id);
1129 if (fid) { try { fwStmts().acc.run(fid); } catch { /* ignore */ } }
1130 console.log('[AP] follow accepted', actorUri);
1131 return 202;
1132 }
1133 if (type === 'Reject' && act.object) {
1134 const who = actorUri;
1135 if (who && slugParam) { try { fwStmts().del.run(slugParam, who); } catch { /* ignore */ } }
1136 return 202;
1137 }
1138
1139 console.log('[AP] inbox', type || 'unknown', '→', slugParam || 'shared', 'from', ip, '(ignored)');
1140 return 202;
1141}
1142
1143// Deliver a new post as Create(Note) to all followers' inboxes (fire-and-forget).
1144// Needs PUBLIC_BASE_URL (absolute URLs); no-op without followers or base.
1145export async function deliverCreate(site, post) {
1146 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
1147 if (!base || !site || !site.slug) return;
1148 // Resolve inline @user@host mentions → link them in the note + collect their inboxes, so a
1149 // mentioned person is notified even if they don't follow us (Mastodon-standard mention).
1150 const mres = await resolveMentionsInText(base, post.content || '');
1151 const post2 = mres.inboxes.length ? { ...post, content: mres.html } : post;
1152 const followers = fStmts().list.all(site.slug);
1153 const inboxes = [...new Set([...followers.map((f) => f.shared_inbox || f.inbox), ...mres.inboxes].filter(Boolean))];
1154 if (!inboxes.length) return; // no followers and no one mentioned
1155 const keys = getOrCreateKeys(site.slug);
1156 const keyId = `${actorId(base, site.slug)}#main-key`;
1157 const create = buildCreate(base, site, post2);
1158 for (const inbox of inboxes) deliverWithRetry(site.slug, inbox, create, keyId, keys.private_pem);
1159}
1160
1161// On a new Follow, send that follower our most recent posts as Create so their
1162// timeline shows our history (Mastodon does not backfill on follow). Oldest-first
1163// so they sort into the follower's timeline at their original dates.
1164async function backfillNewFollower(base, slug, inbox) {
1165 if (!base || !slug || !inbox) return;
1166 const site = db.prepare('SELECT * FROM sites WHERE slug = ?').get(slug);
1167 if (!site) return;
1168 const recent = db.prepare(
1169 `SELECT id, slug, title, content, cover_image_url, cover_video_url, nsfw, content_warning, published_at, created_at
1170 FROM posts WHERE site_id = ? AND status = 'published' AND (fan_only IS NULL OR fan_only = 0)
1171 ORDER BY COALESCE(published_at, created_at) DESC LIMIT 20`
1172 ).all(site.id).reverse();
1173 if (!recent.length) return;
1174 const keys = getOrCreateKeys(slug);
1175 const keyId = `${actorId(base, slug)}#main-key`;
1176 for (const p of recent) {
1177 try { await deliver(inbox, buildCreate(base, site, p), keyId, keys.private_pem); } catch { /* best-effort */ }
1178 await new Promise((r) => setTimeout(r, 150));
1179 }
1180 console.log('[AP] backfilled', recent.length, 'posts to new follower of', slug);
1181}
1182
1183// Tell followers a post is gone (Delete + Tombstone) so it's removed from their feeds.
1184export async function deliverDelete(site, post) {
1185 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
1186 if (!base || !site || !site.slug || !post || !post.id) return;
1187 const followers = fStmts().list.all(site.slug);
1188 if (!followers.length) return;
1189 const inboxes = [...new Set(followers.map((f) => f.shared_inbox || f.inbox).filter(Boolean))];
1190 const keys = getOrCreateKeys(site.slug);
1191 const me = actorId(base, site.slug);
1192 const nid = noteId(base, post.id);
1193 const del = {
1194 '@context': AP_CONTEXT,
1195 id: `${nid}#delete-${Date.now()}-${rid()}`,
1196 type: 'Delete',
1197 actor: me,
1198 to: [PUBLIC],
1199 object: { id: nid, type: 'Tombstone' },
1200 };
1201 for (const inbox of inboxes) deliverWithRetry(site.slug, inbox, del, `${me}#main-key`, keys.private_pem);
1202}
1203
1204// Tell followers an already-published post changed (Update + edited Note) so
1205// Mastodon refreshes the cached copy (e.g. after fixing content).
1206export async function deliverUpdate(site, post) {
1207 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
1208 if (!base || !site || !site.slug || !post || !post.id) return;
1209 const mres = await resolveMentionsInText(base, post.content || ''); // link mentions + collect inboxes
1210 const post2 = mres.inboxes.length ? { ...post, content: mres.html } : post;
1211 const followers = fStmts().list.all(site.slug);
1212 const inboxes = [...new Set([...followers.map((f) => f.shared_inbox || f.inbox), ...mres.inboxes].filter(Boolean))];
1213 if (!inboxes.length) return;
1214 const keys = getOrCreateKeys(site.slug);
1215 const me = actorId(base, site.slug);
1216 const note = buildNote(base, site, post2);
1217 note.updated = new Date().toISOString();
1218 const update = {
1219 '@context': AP_CONTEXT,
1220 id: `${noteId(base, post.id)}#update-${Date.now()}-${rid()}`,
1221 type: 'Update', actor: me, to: [PUBLIC], cc: note.cc,
1222 object: note,
1223 };
1224 for (const inbox of inboxes) deliverWithRetry(site.slug, inbox, update, `${me}#main-key`, keys.private_pem);
1225}
1226
1227// Tell followers the ACTOR changed (Update + Person) so Mastodon re-processes the
1228// account AND re-fetches the featured (pinned) collection — there is no standard
1229// "featured changed" activity, so this is how a pin/unpin propagates promptly.
1230export async function deliverActorUpdate(site) {
1231 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
1232 if (!base || !site || !site.slug) return;
1233 const followers = fStmts().list.all(site.slug);
1234 if (!followers.length) return;
1235 const inboxes = [...new Set(followers.map((f) => f.shared_inbox || f.inbox).filter(Boolean))];
1236 const keys = getOrCreateKeys(site.slug);
1237 const me = actorId(base, site.slug);
1238 const update = {
1239 '@context': AP_CONTEXT,
1240 id: `${me}#update-${Date.now()}-${rid()}`,
1241 type: 'Update', actor: me, to: [PUBLIC], cc: [`${me}/followers`],
1242 object: buildActor(base, site),
1243 };
1244 for (const inbox of inboxes) deliverWithRetry(site.slug, inbox, update, `${me}#main-key`, keys.private_pem);
1245}
1246
1247// Reliably set the pinned order on followers' instances via Add/Remove activities
1248// (how Mastodon itself federates pins) — pushed to the inbox + processed immediately,
1249// unlike the featured COLLECTION which Mastodon caches with sticky StatusPins.
1250// Mastodon's Add skips an already-pinned status, so we REMOVE every pin first, wait,
1251// then ADD in rank-DESCENDING order (rank 1 added LAST → newest StatusPin → shown first,
1252// because Mastodon displays pins newest-first). `alsoRemove` = ids to unpin too.
1253// Serialize pin-resyncs per site: two concurrent /save calls would otherwise interleave
1254// their Remove -> wait -> Add sequences and scramble the StatusPin order on Mastodon. A
1255// resync already in flight for a site coalesces later requests into ONE rerun after it
1256// finishes (accumulating their extra unpins), so rapid saves don't pile up N full resyncs.
1257const _pinResync = new Map(); // slug -> { promise, pending, pendingRemove:Set, site }
1258export function resyncFeaturedPins(site, alsoRemove = []) {
1259 if (!site || !site.slug) return Promise.resolve();
1260 const slug = site.slug;
1261 const running = _pinResync.get(slug);
1262 if (running) {
1263 running.pending = true;
1264 running.site = site; // use the latest site object on the rerun
1265 for (const id of alsoRemove) running.pendingRemove.add(id);
1266 return running.promise;
1267 }
1268 const state = { promise: null, pending: false, pendingRemove: new Set(), site };
1269 state.promise = (async () => {
1270 let extra = alsoRemove;
1271 for (;;) {
1272 try { await doResyncFeaturedPins(state.site, extra); }
1273 catch (e) { console.warn('[AP] pin resync failed:', e.message); }
1274 if (!state.pending) break;
1275 state.pending = false;
1276 extra = [...state.pendingRemove];
1277 state.pendingRemove = new Set();
1278 }
1279 _pinResync.delete(slug);
1280 })();
1281 _pinResync.set(slug, state);
1282 return state.promise;
1283}
1284
1285// The actual resync work — do NOT call directly; go through resyncFeaturedPins() above so
1286// it stays serialized per site.
1287async function doResyncFeaturedPins(site, alsoRemove = []) {
1288 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
1289 if (!base || !site || !site.slug) return;
1290 const followers = fStmts().list.all(site.slug);
1291 if (!followers.length) return;
1292 const inboxes = [...new Set(followers.map((f) => f.shared_inbox || f.inbox).filter(Boolean))];
1293 const keys = getOrCreateKeys(site.slug);
1294 const me = actorId(base, site.slug);
1295 const keyId = `${me}#main-key`;
1296 const featured = `${me}/featured`;
1297 const note = (id) => noteId(base, id);
1298 const pinned = db.prepare(
1299 `SELECT id FROM posts WHERE site_id = ? AND status = 'published' AND (fan_only IS NULL OR fan_only = 0)
1300 AND pinned IS NOT NULL AND pinned > 0
1301 ORDER BY pinned DESC, COALESCE(published_at, created_at) ASC LIMIT 20`
1302 ).all(site.id);
1303 const removeIds = [...new Set([...pinned.map((p) => p.id), ...alsoRemove])];
1304 // 1. Remove every current pin so Mastodon can recreate them in order.
1305 for (const id of removeIds) {
1306 const rm = { '@context': AP_CONTEXT, id: `${me}#rm-${id}-${Date.now()}-${rid()}`, type: 'Remove', actor: me, object: note(id), target: featured, to: [PUBLIC] };
1307 for (const inbox of inboxes) deliver(inbox, rm, keyId, keys.private_pem).catch(() => { /* best-effort */ });
1308 }
1309 if (!pinned.length) { console.log('[AP] unpinned all featured for', site.slug); return; }
1310 await new Promise((r) => setTimeout(r, 5000)); // let the Removes land first
1311 // 2. Add in rank-DESC order, gaps so each StatusPin gets an increasing created_at.
1312 for (const p of pinned) {
1313 const add = { '@context': AP_CONTEXT, id: `${me}#add-${p.id}-${Date.now()}-${rid()}`, type: 'Add', actor: me, object: note(p.id), target: featured, to: [PUBLIC], cc: [`${me}/followers`] };
1314 for (const inbox of inboxes) deliver(inbox, add, keyId, keys.private_pem).catch(() => { /* best-effort */ });
1315 await new Promise((r) => setTimeout(r, 2000));
1316 }
1317 console.log('[AP] resynced', pinned.length, 'featured pins for', site.slug);
1318}
1319
1320// ── outbound replies (Klonkt → fediverse) ─────────────────────────
1321const escHtml = (s) => String(s || '').replace(/[<>&]/g, (c) => ({ '<': '&lt;', '>': '&gt;', '&': '&amp;' }[c]));
1322const 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(); };
1323
1324// Build one of OUR outbound reply Notes from an ap_outbox row.
1325// Turn #hashtags in reply text into Mastodon-style hashtag links (clickable + federated).
1326function linkHashtags(base, html) {
1327 return String(html || '').replace(/(^|[\s>])#([\p{L}\p{M}\p{N}_]+)/gu, (m, pre, tag) =>
1328 `${pre}<a href="${base}/tag/${encodeURIComponent(tag.toLowerCase())}" class="mention hashtag" rel="tag">#${tag}</a>`);
1329}
1330// Extract the AP Hashtag tag objects from already-linked reply content.
1331function hashtagTags(base, content) {
1332 const tags = [], seen = new Set();
1333 const re = /class="[^"]*\bhashtag\b[^"]*"[^>]*>#([\p{L}\p{M}\p{N}_]+)</giu;
1334 let m;
1335 while ((m = re.exec(content || ''))) {
1336 const k = m[1].toLowerCase();
1337 if (seen.has(k)) continue; seen.add(k);
1338 tags.push({ type: 'Hashtag', href: `${base}/tag/${encodeURIComponent(k)}`, name: '#' + m[1] });
1339 }
1340 return tags;
1341}
1342
1343// Normalise a post's tags field (array, JSON-string, or comma-string) to an array.
1344function normalizeTags(t) {
1345 if (Array.isArray(t)) return t;
1346 if (typeof t === 'string') {
1347 const s = t.trim(); if (!s) return [];
1348 if (s[0] === '[') { try { const a = JSON.parse(s); return Array.isArray(a) ? a : []; } catch { /* fall through */ } }
1349 return s.split(',').map((x) => x.trim()).filter(Boolean);
1350 }
1351 return [];
1352}
1353// A tag → { label, slug }. Multi-word tags become CamelCase (#LiveMusic) for the display
1354// name (Mastodon hashtags can't contain spaces; CamelCase is the accessibility norm); the
1355// slug/href stays lowercase ("livemusic").
1356function tagParts(raw) {
1357 const words = String(raw || '').trim().split(/[\s_]+/).map((w) => w.replace(/[^\p{L}\p{M}\p{N}]/gu, '')).filter(Boolean);
1358 if (!words.length) return null;
1359 const slug = words.join('').toLowerCase();
1360 if (!slug) return null;
1361 const label = words.length > 1 ? words.map((w) => w[0].toUpperCase() + w.slice(1)).join('') : words[0];
1362 return { label, slug };
1363}
1364// Merge a post's tags field + the #hashtags linked inline in its body into one deduped
1365// Hashtag tag list (with hrefs to our /tag page).
1366function buildHashtagList(base, tagsField, content) {
1367 const out = [], seen = new Set();
1368 for (const t of normalizeTags(tagsField)) {
1369 const p = tagParts(t); if (!p || seen.has(p.slug)) continue; seen.add(p.slug);
1370 out.push({ type: 'Hashtag', href: `${base}/tag/${encodeURIComponent(p.slug)}`, name: '#' + p.label });
1371 }
1372 for (const h of hashtagTags(base, content)) {
1373 const k = h.name.slice(1).toLowerCase(); if (seen.has(k)) continue; seen.add(k);
1374 out.push(h);
1375 }
1376 return out;
1377}
1378
1379// Extract Mention tag objects from already-linked content (class="u-url mention").
1380function mentionTags(content) {
1381 const tags = [], seen = new Set();
1382 // The link href is the human profile URL; the actor URI (for the Mention tag) is in data-actor.
1383 const re = /<a href="[^"]*" class="u-url mention" data-actor="([^"]+)">@([^<]+)<\/a>/gi;
1384 let m;
1385 while ((m = re.exec(content || ''))) {
1386 const href = m[1];
1387 if (seen.has(href)) continue; seen.add(href);
1388 tags.push({ type: 'Mention', href, name: '@' + m[2] });
1389 }
1390 return tags;
1391}
1392// Resolve inline @user@domain mentions in reply/post text → link them (href = actor URI)
1393// and collect the mentioned actors' inboxes so they get notified. Best-effort per mention.
1394async function resolveMentionsInText(base, html) {
1395 const inboxes = [];
1396 const handles = new Set();
1397 const re = /(^|[\s>])@([\p{L}\p{M}\p{N}_.-]+@[\p{L}\p{M}\p{N}.-]+)/gu;
1398 let m;
1399 while ((m = re.exec(html || ''))) handles.add(m[2]);
1400 let out = String(html || '');
1401 for (const h of handles) {
1402 let actorUri = null;
1403 try { actorUri = await webfingerResolve('@' + h); } catch { actorUri = null; }
1404 if (!actorUri) continue;
1405 const actor = await fetchActor(actorUri).catch(() => null);
1406 const inbox = actor && ((actor.endpoints && actor.endpoints.sharedInbox) || actor.inbox);
1407 if (inbox) inboxes.push(inbox);
1408 const profileUrl = actorInfo(actor, actorUri).url || actorUri; // human profile page → the link href
1409 const esc = h.replace(/[.*+?^${}()|[\]\\]/g, '\\$&');
1410 out = out.replace(new RegExp('(^|[\\s>])@' + esc + '(?![\\p{L}\\p{M}\\p{N}_.-])', 'gu'),
1411 (full, pre) => `${pre}<a href="${profileUrl}" class="u-url mention" data-actor="${actorUri}">@${h}</a>`);
1412 }
1413 return { html: out, inboxes };
1414}
1415
1416export function buildReplyNote(base, site, row) {
1417 const me = actorId(base, site.slug);
1418 return {
1419 id: noteId(base, row.id),
1420 type: 'Note',
1421 attributedTo: me,
1422 inReplyTo: row.in_reply_to || undefined,
1423 content: row.content,
1424 url: row.post_slug ? `${base}/${encodeURIComponent(row.post_slug)}` : undefined,
1425 published: toISO(row.created_at),
1426 to: row.to_actor ? [row.to_actor] : [PUBLIC],
1427 cc: [PUBLIC, `${me}/followers`],
1428 tag: [
1429 ...mentionTags(row.content),
1430 ...hashtagTags(base, row.content),
1431 ],
1432 };
1433}
1434
1435// Resolve one of our outbound reply Notes by id (for /ap/notes/:id fallback).
1436export function getOutboxNote(base, id) {
1437 const row = iStmts().getO.get(id);
1438 if (!row) return null;
1439 const site = db.prepare('SELECT * FROM sites WHERE slug = ?').get(row.site_slug);
1440 if (!site) return null;
1441 return buildReplyNote(base, site, row);
1442}
1443
1444// Send a reply FROM this site to a remote actor (in reply to their inbound reply).
1445// `parent` = an ap_interactions row (actor_uri, actor_url, actor_handle, object_uri).
1446export async function deliverReply(site, { postId, postSlug, parent, text }) {
1447 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
1448 if (!base || !site || !site.slug || !parent || !String(text || '').trim()) return null;
1449 const me = actorId(base, site.slug);
1450 const handle = parent.actor_handle || deriveHandle(parent.actor_uri);
1451 const dispHandle = handle && handle[0] === '@' ? handle : '@' + (handle || '');
1452 const body = escHtml(String(text).trim()).replace(/\r?\n/g, '<br>');
1453 const mres = await resolveMentionsInText(base, body); // link inline @mentions + collect their inboxes
1454 const mention = parent.actor_uri
1455 ? `<a href="${escHtml(parent.actor_url || parent.actor_uri)}" class="u-url mention" data-actor="${escHtml(parent.actor_uri)}">${escHtml(dispHandle)}</a> ` : '';
1456 const content = `<p>${mention}${linkHashtags(base, mres.html)}</p>`;
1457 // Dedup: skip if the exact same reply was already sent (double-submit guard).
1458 const dup = db.prepare('SELECT 1 FROM ap_outbox WHERE site_slug = ? AND IFNULL(in_reply_to, \'\') = ? AND content = ? LIMIT 1')
1459 .get(site.slug, parent.object_uri || '', content);
1460 if (dup) { console.log('[AP] outreply skipped (duplicate)'); return { duplicate: true, delivered: 0 }; }
1461 const id = crypto.randomUUID();
1462 iStmts().insO.run(id, site.slug, postId, postSlug || null, parent.object_uri || null, parent.actor_uri || null, handle, content);
1463 const row = iStmts().getO.get(id);
1464 const note = buildReplyNote(base, site, row);
1465 const create = {
1466 '@context': AP_CONTEXT,
1467 id: note.id + '#create', type: 'Create', actor: me,
1468 published: note.published, to: note.to, cc: note.cc, object: note,
1469 };
1470 const keys = getOrCreateKeys(site.slug);
1471 const keyId = `${me}#main-key`;
1472 const inboxes = new Set();
1473 if (parent.actor_uri) {
1474 const a = await fetchActor(parent.actor_uri).catch(() => null);
1475 if (a) inboxes.add((a.endpoints && a.endpoints.sharedInbox) || a.inbox);
1476 }
1477 if (parent.threadInbox) inboxes.add(parent.threadInbox); // back-compat (single)
1478 (parent.threadInboxes || []).forEach((i) => inboxes.add(i)); // whole ancestor chain
1479 for (const f of fStmts().list.all(site.slug)) inboxes.add(f.shared_inbox || f.inbox);
1480 mres.inboxes.forEach((i) => inboxes.add(i)); // people @mentioned inline in the reply
1481 inboxes.delete(`${me}/inbox`); // never deliver to ourselves (already in ap_outbox)
1482 inboxes.delete(`${base}/ap/inbox`); // (our own shared inbox) → avoids a self-duplicate
1483 let delivered = 0;
1484 for (const inbox of [...inboxes].filter(Boolean)) {
1485 try { const st = await deliver(inbox, create, keyId, keys.private_pem); if (st >= 200 && st < 300) delivered++; } catch { /* best-effort */ }
1486 }
1487 console.log('[AP] outreply', site.slug, '→', parent.actor_uri, 'delivered', delivered);
1488 return { id, content, delivered };
1489}
1490
1491// attributedTo may be a string, an object {id}, or an ARRAY — e.g. a PeerTube Video is
1492// attributed to [Person (account), Group (channel)]. Pick a usable actor URI (prefer Person).
1493function actorUriOf(att) {
1494 if (!att) return null;
1495 if (typeof att === 'string') return att;
1496 if (Array.isArray(att)) {
1497 const person = att.find((a) => a && typeof a === 'object' && a.type === 'Person' && a.id);
1498 if (person) return person.id;
1499 for (const a of att) { if (typeof a === 'string') return a; if (a && a.id) return a.id; }
1500 return null;
1501 }
1502 return att.id || null;
1503}
1504
1505// Resolve a remote post URL (any fediverse/Klonkt post) into a reply target.
1506// Returns a parent-shaped object usable by deliverReply(), or null.
1507export async function resolveRemoteNote(url) {
1508 if (!/^https?:\/\//i.test(String(url || ''))) return null;
1509 const note = await fetchActor(url).catch(() => null); // AP GET (content-negotiates)
1510 if (!note || !note.id) return null;
1511 const att = note.attributedTo;
1512 const actorUri = actorUriOf(att);
1513 if (!actorUri) return null;
1514 const actor = await fetchActor(actorUri).catch(() => null);
1515 const ai = actorInfo(actor, actorUri);
1516 // Is what we're replying to a post (or a comment) on one of OUR posts? If so,
1517 // link our reply to that local post so it shows nested in the post thread.
1518 const localTgt = findThreadTarget(note.id, (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, ''));
1519 // Walk the WHOLE reply chain upward (comment → parent comment → … → root post)
1520 // and collect every ancestor author's inbox, so each participant's server —
1521 // including the original post's author — receives + threads our reply.
1522 const threadInboxes = [];
1523 const seenInbox = new Set();
1524 let cursor = note.inReplyTo, guard = 0;
1525 while (cursor && guard++ < 6) {
1526 const url = typeof cursor === 'string' ? cursor : (cursor && cursor.id);
1527 if (!url) break;
1528 const pn = await fetchActor(url).catch(() => null);
1529 if (!pn) break;
1530 const pa = actorUriOf(pn.attributedTo);
1531 if (pa && pa !== actorUri) {
1532 const paDoc = await fetchActor(pa).catch(() => null);
1533 const inbox = paDoc && ((paDoc.endpoints && paDoc.endpoints.sharedInbox) || paDoc.inbox);
1534 if (inbox && !seenInbox.has(inbox)) { seenInbox.add(inbox); threadInboxes.push(inbox); }
1535 }
1536 cursor = pn.inReplyTo; // climb to the next ancestor
1537 }
1538 // For non-Note objects (PeerTube Video, Article, …) the meaningful label is `name` (the
1539 // title); prepend it so the reply page shows what you're replying to (sanitize cleans it).
1540 let rawHtml = String(note.content || '').replace(/\[\[(track|album|playlist):[^\]]+\]\]/gi, '');
1541 if (note.name && note.type && note.type !== 'Note') rawHtml = `<p><strong>${note.name}</strong></p>` + rawHtml;
1542 const images = (Array.isArray(note.attachment) ? note.attachment : [])
1543 .filter((a) => a && a.url && (!a.mediaType || /^image\//i.test(a.mediaType)))
1544 .map((a) => safeUrl(a.url)).filter(Boolean);
1545 return {
1546 object_uri: safeUrl(note.id) || note.id,
1547 actor_uri: actorUri,
1548 actor_url: ai.url,
1549 actor_handle: ai.handle,
1550 actor_name: ai.name,
1551 actor_icon: ai.icon,
1552 url: note.url || url,
1553 content: HtmlSanitizerService.sanitize(rawHtml), // full, sanitized
1554 sensitive: !!note.sensitive, // remote CW → blur in the Cirkel
1555 cw: note.summary || '',
1556 images,
1557 threadInboxes, // every ancestor author's inbox
1558 localPostId: localTgt ? localTgt.post_id : '', // our post this belongs to (if any)
1559 poll: parsePoll(note), // a Question → its options/counts (else null)
1560 preview: HtmlSanitizerService.toPlainText(note.content || '').slice(0, 240),
1561 };
1562}
1563
1564// List a site's own outbound fediverse replies (for the manage/delete view).
1565// The plain editable text of a stored reply (unwrap links → their text, <br> → newline)
1566// so the manage view can prefill an edit box; the mention is re-added on save.
1567function outboxEditableText(content) {
1568 return String(content || '')
1569 .replace(/<br\s*\/?>/gi, '\n')
1570 .replace(/<a\b[^>]*>([\s\S]*?)<\/a>/gi, '$1')
1571 .replace(/<[^>]+>/g, '')
1572 .replace(/&lt;/g, '<').replace(/&gt;/g, '>').replace(/&amp;/g, '&')
1573 .trim();
1574}
1575export function listOutbox(siteSlug) {
1576 return db.prepare('SELECT id, content, to_handle, in_reply_to, created_at FROM ap_outbox WHERE site_slug = ? ORDER BY created_at DESC')
1577 .all(siteSlug).map((r) => { const c = stripLeadingMentions(r.content); return { ...r, content: c, editable: outboxEditableText(c) }; });
1578}
1579
1580// Delete one of our outbound replies: send Delete(Tombstone) to recipients + remove it.
1581export async function deliverOutboxDelete(site, outboxId) {
1582 const row = iStmts().getO.get(outboxId);
1583 if (!row || row.site_slug !== site.slug) return false;
1584 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
1585 if (base) {
1586 const me = actorId(base, site.slug);
1587 const nid = noteId(base, row.id);
1588 const del = { '@context': AP_CONTEXT, id: `${nid}#delete-${Date.now()}-${rid()}`, type: 'Delete', actor: me, to: [PUBLIC], object: { id: nid, type: 'Tombstone' } };
1589 const keys = getOrCreateKeys(site.slug);
1590 const inboxes = new Set();
1591 if (row.to_actor) { const a = await fetchActor(row.to_actor).catch(() => null); if (a) inboxes.add((a.endpoints && a.endpoints.sharedInbox) || a.inbox); }
1592 for (const f of fStmts().list.all(site.slug)) inboxes.add(f.shared_inbox || f.inbox);
1593 for (const inbox of [...inboxes].filter(Boolean)) { try { await deliver(inbox, del, `${me}#main-key`, keys.private_pem); } catch { /* best-effort */ } }
1594 }
1595 db.prepare('DELETE FROM ap_outbox WHERE id = ?').run(outboxId);
1596 return true;
1597}
1598
1599// Edit one of our outbound replies: rewrite the stored content (mention re-added + #tags
1600// re-linked) and send an Update(Note) so recipients refresh their cached copy.
1601export async function deliverOutboxUpdate(site, outboxId, newText) {
1602 const row = iStmts().getO.get(outboxId);
1603 if (!row || row.site_slug !== site.slug) return false;
1604 const text = String(newText || '').trim();
1605 if (!text) return false;
1606 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
1607 if (!base) return false;
1608 const me = actorId(base, site.slug);
1609 const toActor = row.to_actor ? await fetchActor(row.to_actor).catch(() => null) : null;
1610 const toProfile = row.to_actor ? (actorInfo(toActor, row.to_actor).url || row.to_actor) : '';
1611 const _h = row.to_handle || deriveHandle(row.to_actor);
1612 const toHandle = _h && _h[0] === '@' ? _h : '@' + (_h || '');
1613 const mention = row.to_actor
1614 ? `<a href="${escHtml(toProfile)}" class="u-url mention" data-actor="${escHtml(row.to_actor)}">${escHtml(toHandle)}</a> ` : '';
1615 const mres = await resolveMentionsInText(base, escHtml(text).replace(/\r?\n/g, '<br>'));
1616 const content = `<p>${mention}${linkHashtags(base, mres.html)}</p>`;
1617 db.prepare('UPDATE ap_outbox SET content = ? WHERE id = ?').run(content, outboxId);
1618 const note = buildReplyNote(base, site, iStmts().getO.get(outboxId));
1619 note.updated = new Date().toISOString();
1620 const update = {
1621 '@context': AP_CONTEXT,
1622 id: `${note.id}#update-${Date.now()}-${rid()}`, type: 'Update', actor: me,
1623 published: note.published, updated: note.updated, to: note.to, cc: note.cc, object: note,
1624 };
1625 const keys = getOrCreateKeys(site.slug);
1626 const inboxes = new Set();
1627 if (toActor) inboxes.add((toActor.endpoints && toActor.endpoints.sharedInbox) || toActor.inbox);
1628 for (const f of fStmts().list.all(site.slug)) inboxes.add(f.shared_inbox || f.inbox);
1629 mres.inboxes.forEach((i) => inboxes.add(i)); // people @mentioned inline in the edit
1630 inboxes.delete(`${me}/inbox`); inboxes.delete(`${base}/ap/inbox`);
1631 let delivered = 0;
1632 for (const inbox of [...inboxes].filter(Boolean)) {
1633 try { const st = await deliver(inbox, update, `${me}#main-key`, keys.private_pem); if (st >= 200 && st < 300) delivered++; } catch { /* best-effort */ }
1634 }
1635 console.log('[AP] outreply edit', site.slug, 'delivered', delivered);
1636 return { ok: true, content, delivered };
1637}
1638
1639// ── Fediverse CLIENT: follow accounts + home timeline ─────────────
1640// Resolve an @user@domain handle to its actor URL via WebFinger.
1641export async function webfingerResolve(handle) {
1642 const h = String(handle || '').trim().replace(/^@/, '');
1643 const parts = h.split('@');
1644 if (parts.length !== 2 || !parts[0] || !parts[1]) return null;
1645 const acct = `${parts[0]}@${parts[1]}`;
1646 try {
1647 const r = await safeFetch(`https://${parts[1]}/.well-known/webfinger?resource=acct:${encodeURIComponent(acct)}`,
1648 { headers: { Accept: 'application/jrd+json, application/json' } });
1649 if (!r.ok) return null;
1650 const jrd = await r.json();
1651 const link = (jrd.links || []).find((l) => l.rel === 'self' && /activity\+json|ld\+json/.test(l.type || ''));
1652 return safeUrl(link ? link.href : '') || null;
1653 } catch { return null; }
1654}
1655
1656let _insFw, _delFw, _listFw, _accFw, _oneFw, _setAB;
1657function fwStmts() {
1658 if (!_insFw) {
1659 _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)');
1660 _delFw = db.prepare('DELETE FROM ap_following WHERE slug = ? AND actor_uri = ?');
1661 _listFw = db.prepare('SELECT * FROM ap_following WHERE slug = ? ORDER BY created_at DESC');
1662 _accFw = db.prepare("UPDATE ap_following SET status = 'accepted' WHERE follow_id = ?");
1663 _oneFw = db.prepare('SELECT * FROM ap_following WHERE slug = ? AND actor_uri = ?');
1664 _setAB = db.prepare('UPDATE ap_following SET auto_boost = ? WHERE slug = ? AND actor_uri = ?');
1665 }
1666 return { ins: _insFw, del: _delFw, list: _listFw, acc: _accFw, one: _oneFw, setAB: _setAB };
1667}
1668export function listFollowing(slug) { return fwStmts().list.all(slug); }
1669
1670// Toggle auto-boost ("feature") on an account we already follow.
1671export function setAutoBoost(slug, actorUri, on) {
1672 try { fwStmts().setAB.run(on ? 1 : 0, slug, actorUri); } catch { /* ignore */ }
1673 // Featuring an account → AP-native catch-up so the Cirkel isn't empty until they next
1674 // post (push doesn't backfill history-before-follow). Fire-and-forget pull, sends nothing.
1675 if (on) backfillFromOutbox(slug, actorUri).catch(() => {});
1676 return { ok: true };
1677}
1678
1679let _insTl, _listTl, _delTl;
1680function tlStmts() {
1681 if (!_insTl) {
1682 _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, nsfw, cw, created_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,CURRENT_TIMESTAMP)');
1683 _listTl = db.prepare('SELECT * FROM ap_timeline WHERE slug = ? ORDER BY COALESCE(published, created_at) DESC LIMIT ?');
1684 _delTl = db.prepare('DELETE FROM ap_timeline WHERE id = ?');
1685 }
1686 return { ins: _insTl, list: _listTl, del: _delTl };
1687}
1688export function getTimeline(slug, limit) { return tlStmts().list.all(slug, limit || 50); }
1689
1690// ── Cirkel = posts from the accounts you auto-boost ("feature an artist") ──
1691let _abCount, _cirkelPosts, _cirkelMembers;
1692export function autoBoostCount(slug) {
1693 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; }
1694}
1695export function getCirkelPosts(slug, limit) {
1696 try {
1697 // Cirkel = posts from featured (auto_boost) accounts + posts you boosted
1698 // (t.boosted), mixed by date. One row per note in ap_timeline → no duplicates.
1699 if (!_cirkelPosts) _cirkelPosts = db.prepare(`
1700 SELECT t.id, t.author_uri, t.author_name, t.author_handle, t.author_icon, t.author_url,
1701 t.content, t.url, t.published, t.media_json, t.boosted, t.nsfw, t.cw
1702 FROM ap_timeline t
1703 LEFT JOIN ap_following f ON f.slug = t.slug AND f.actor_uri = t.author_uri
1704 WHERE t.slug = ? AND (f.auto_boost = 1 OR t.boosted = 1)
1705 ORDER BY COALESCE(t.published, t.created_at) DESC, t.rowid DESC
1706 LIMIT ?`);
1707 return _cirkelPosts.all(slug, limit || 60);
1708 } catch { return []; }
1709}
1710export function getCirkelMembers(slug) {
1711 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 []; }
1712}
1713// Mark a timeline post as boosted so it shows in the Cirkel (mixed by date).
1714let _markBoost, _unmarkBoost, _boostedCount;
1715export function markBoosted(slug, noteId) {
1716 try { if (!_markBoost) _markBoost = db.prepare('UPDATE ap_timeline SET boosted = 1 WHERE slug = ? AND id = ?'); _markBoost.run(slug, noteId); } catch { /* ignore */ }
1717}
1718export function unmarkBoosted(slug, noteId) {
1719 try { if (!_unmarkBoost) _unmarkBoost = db.prepare('UPDATE ap_timeline SET boosted = 0 WHERE slug = ? AND id = ?'); _unmarkBoost.run(slug, noteId); } catch { /* ignore */ }
1720}
1721let _markLike, _unmarkLike;
1722export function markLiked(slug, noteId) {
1723 try { if (!_markLike) _markLike = db.prepare('UPDATE ap_timeline SET liked = 1 WHERE slug = ? AND id = ?'); _markLike.run(slug, noteId); } catch { /* ignore */ }
1724}
1725export function unmarkLiked(slug, noteId) {
1726 try { if (!_unmarkLike) _unmarkLike = db.prepare('UPDATE ap_timeline SET liked = 0 WHERE slug = ? AND id = ?'); _unmarkLike.run(slug, noteId); } catch { /* ignore */ }
1727}
1728export function getTimelineReaction(slug, noteId) {
1729 try { const r = db.prepare('SELECT liked, boosted FROM ap_timeline WHERE slug = ? AND id = ?').get(slug, noteId); return { liked: !!(r && r.liked), boosted: !!(r && r.boosted) }; } catch { return { liked: false, boosted: false }; }
1730}
1731// Boost a REMOTE post that may not be in your timeline (you don't follow the author):
1732// store it in ap_timeline (INSERT OR IGNORE → no dup for followed posts) so it shows in
1733// the Cirkel with a Boost badge, then flag it boosted.
1734export function upsertBoostedNote(slug, note) {
1735 if (!slug || !note || !note.object_uri) return;
1736 const id = note.object_uri;
1737 const media = JSON.stringify((note.images || []).map((u) => ({ url: u, type: 'image/jpeg' })));
1738 try {
1739 tlStmts().ins.run(id, slug, note.actor_uri || '', note.actor_name || '', note.actor_handle || '',
1740 note.actor_icon || '', note.actor_url || '', note.content || '', note.url || null,
1741 new Date().toISOString(), media, note.sensitive ? 1 : 0, note.cw || null);
1742 } catch { /* ignore */ }
1743 markBoosted(slug, id);
1744}
1745export function boostedCount(slug) {
1746 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; }
1747}
1748
1749// Resolve a Klonkt/AP actor URL from a site root: a Klonkt site's root 302s to
1750// /ap/users/<slug> (content negotiation; Location may be relative). Used by
1751// followActor for bare-domain follows.
1752// NB: the old auto-migration of legacy Cirkels (circle_links -> AP follows) was
1753// REMOVED on 2026-06-26 — it auto-sent Follows on boot, which violates "the code
1754// never throws anything into the fediverse automatically" (would surprise-Follow
1755// for some operators at scale). The dead circle_links table stays as harmless dead
1756// data; an operator restores an old cirkel by re-following in /following (their click).
1757async function resolveApActor(siteUrl) {
1758 try {
1759 const r = await fetch(siteUrl, { headers: { Accept: 'application/activity+json' }, redirect: 'manual' });
1760 if (r.status >= 300 && r.status < 400) { const loc = r.headers.get('location'); if (loc) return new URL(loc, siteUrl).href; }
1761 if (r.ok) return siteUrl;
1762 } catch { /* unreachable */ }
1763 return null;
1764}
1765
1766// ── Self-heal: re-sync the fediverse cache (ap_timeline) after a DRASTIC update ──
1767// Runs ONCE per SELFHEAL_VERSION bump — NOT on every boot. Re-fetches each cached
1768// note and refreshes content + media (recovers covers/edits that were delivered
1769// during a flux window, e.g. a fleet-wide update), and drops notes that are gone
1770// (404/410). Bump SELFHEAL_VERSION only on a release that warrants a re-sync.
1771const SELFHEAL_VERSION = 5; // v5: re-fetch so embed/link-only posts pick up the note.image cover in feeds
1772async function fetchNoteAP(url) {
1773 try {
1774 const r = await fetch(url, { headers: { Accept: 'application/activity+json' } });
1775 if (r.status === 404 || r.status === 410) return 404;
1776 if (r.ok) return await r.json();
1777 } catch { /* unreachable */ }
1778 return null;
1779}
1780function mediaFromNote(note) {
1781 const atts = (Array.isArray(note.attachment) ? note.attachment : []).map((a) => ({ url: safeUrl(a && a.url), type: (a && a.mediaType) || '' })).filter((m) => m.url);
1782 if (!atts.some((m) => !m.type || /image/i.test(m.type)) && note.image) {
1783 const im = Array.isArray(note.image) ? note.image[0] : note.image;
1784 const iu = safeUrl(typeof im === 'string' ? im : (im && im.url));
1785 if (iu) atts.push({ url: iu, type: (im && im.mediaType) || 'image/jpeg' });
1786 }
1787 return JSON.stringify(atts);
1788}
1789// A generic SSRF-safe AP GET (collections / pages).
1790async function apGetJson(url) {
1791 try {
1792 const r = await safeFetch(url, { headers: { Accept: 'application/activity+json' } });
1793 if (!r.ok) return null;
1794 const len = Number(r.headers.get('content-length') || 0);
1795 if (len > 3_000_000) return null;
1796 return await r.json();
1797 } catch { return null; }
1798}
1799// AP-native catch-up: pull an actor's standard `outbox` collection and merge their recent
1800// top-level posts into the timeline for `slug`. Push (Create delivery) cannot backfill
1801// history-from-before-you-followed or a delivery that was missed while you were down;
1802// reading the outbox is the spec-conform way to catch up. PULL ONLY — sends nothing.
1803export async function backfillFromOutbox(slug, actorUri, limit = 20) {
1804 try {
1805 if (!slug || !actorUri) return 0;
1806 const actor = await fetchActor(actorUri);
1807 if (!actor || !actor.outbox) return 0;
1808 let page = await apGetJson(typeof actor.outbox === 'string' ? actor.outbox : actor.outbox.id);
1809 let items = (page && (page.orderedItems || page.items)) || [];
1810 if (!items.length && page && page.first) {
1811 page = await apGetJson(typeof page.first === 'string' ? page.first : page.first.id);
1812 items = (page && (page.orderedItems || page.items)) || [];
1813 }
1814 if (!Array.isArray(items) || !items.length) return 0;
1815 const ai = actorInfo(actor, actorUri);
1816 let added = 0;
1817 for (const it of items.slice(0, limit)) {
1818 // Each item is usually a Create wrapping a Note, or sometimes the Note itself.
1819 const o = (it && typeof it.object === 'object' && it.object) ? it.object : it;
1820 if (!o || !o.id) continue;
1821 if (o.type && o.type !== 'Note' && o.type !== 'Article' && o.type !== 'Question') continue; // skip boosts/other
1822 if (o.inReplyTo) continue; // top-level only
1823 const auth = actorUriOf(o.attributedTo);
1824 if (auth && auth !== actorUri) continue; // their OWN posts only
1825 const html = HtmlSanitizerService.sanitize(o.content || '');
1826 const poll = parsePoll(o); // a Question (poll) → carry its options/counts on backfill too
1827 try {
1828 const r = tlStmts().ins.run(o.id, slug, actorUri, ai.name, ai.handle, ai.icon, ai.url, html, o.url || null, o.published || null, mediaFromNote(o), o.sensitive ? 1 : 0, o.summary || null);
1829 if (r && r.changes > 0) added++;
1830 // Set poll_json if this is a poll and we don't already have it (COALESCE preserves a vote).
1831 if (poll) { try { db.prepare('UPDATE ap_timeline SET poll_json = COALESCE(poll_json, ?) WHERE id = ? AND slug = ?').run(JSON.stringify(poll), o.id, slug); } catch { /* ignore */ } }
1832 } catch { /* ignore */ }
1833 }
1834 if (added) console.log('[AP] outbox backfill', actorUri, '→', slug, '+' + added);
1835 return added;
1836 } catch { return 0; }
1837}
1838
1839// ── Remote thread crawl (fill the gaps in a local post's conversation) ────────────
1840// Most replies reach us by delivery, but replies-to-replies that live on other servers and
1841// aren't addressed to us are missed. This pulls the AS2 `replies` collections of the replies
1842// we DO have, caching any newly-found ones in ap_interactions.
1843//
1844// Matches Mastodon's behaviour: ONE level per crawl (like its FetchRepliesService), not a deep
1845// recursive walk. Deeper levels fill in incrementally across crawls — once a fetched reply is
1846// cached it becomes a seed itself, so its own replies are pulled on a later view (Mastodon's
1847// per-status cascade). Bounded + polite (serial), PULL only, and stale-while-revalidate: it
1848// never runs in a page request — the view renders from cache; a stale post kicks off a
1849// background refresh for the NEXT view.
1850const THREAD_TTL_MS = 15 * 60 * 1000; // don't re-crawl a post more than ~4×/hour
1851const THREAD_MAX_DEPTH = 1; // one hop per crawl (like Mastodon); deeper fills in over crawls
1852const THREAD_MAX_FETCHES = 30; // hard cap on remote GETs per crawl (be a good peer)
1853const _crawlingThreads = new Set(); // per-post in-flight lock (no stampede across views)
1854
1855function threadCrawlTs(postId) {
1856 try { const r = db.prepare('SELECT value FROM app_settings WHERE key = ?').get('thread_crawl:' + postId); return r ? (Number(r.value) || 0) : 0; }
1857 catch { return 0; }
1858}
1859function setThreadCrawlTs(postId, ts) {
1860 try { db.prepare('INSERT INTO app_settings (key, value) VALUES (?, ?) ON CONFLICT(key) DO UPDATE SET value = excluded.value').run('thread_crawl:' + postId, String(ts)); }
1861 catch { /* ignore */ }
1862}
1863
1864// Read a note's `replies` (string ref / Collection with `first` / paged CollectionPages) →
1865// child note URIs. Every remote GET goes through `budget` so the whole crawl stays capped.
1866async function collectReplyItems(repliesRef, maxPages, budget) {
1867 const uris = [];
1868 let node = typeof repliesRef === 'string' ? await budget.get(repliesRef) : repliesRef;
1869 if (node && node.first) node = typeof node.first === 'string' ? await budget.get(node.first) : node.first;
1870 let pages = 0;
1871 while (node && pages++ < maxPages) {
1872 for (const it of (node.items || node.orderedItems || [])) {
1873 const u = typeof it === 'string' ? it : (it && it.id);
1874 if (u && /^https?:\/\//i.test(u)) uris.push(u);
1875 }
1876 if (!node.next) break;
1877 node = typeof node.next === 'string' ? await budget.get(node.next) : node.next;
1878 }
1879 return uris;
1880}
1881
1882async function crawlThread(postId) {
1883 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
1884 if (!base) return;
1885 // Seed frontier = the remote reply note URIs we already have; also the dedup set.
1886 let known;
1887 try { known = new Set(db.prepare("SELECT object_uri FROM ap_interactions WHERE post_id = ? AND kind = 'reply' AND object_uri != ''").all(postId).map((r) => r.object_uri)); }
1888 catch { return; }
1889 const seeds = [...known].filter((u) => /^https?:\/\//i.test(u));
1890 if (!seeds.length) return; // nothing remote to expand
1891
1892 let fetches = 0;
1893 const budget = { get: async (u) => { if (fetches >= THREAD_MAX_FETCHES) return null; fetches++; return apGetJson(u); } };
1894 const visited = new Set(); // notes whose replies collection we've already expanded
1895 let frontier = seeds.slice();
1896 let added = 0;
1897
1898 for (let depth = 0; depth < THREAD_MAX_DEPTH && frontier.length && fetches < THREAD_MAX_FETCHES; depth++) {
1899 const nextFrontier = [];
1900 for (const noteUri of frontier) {
1901 if (visited.has(noteUri) || fetches >= THREAD_MAX_FETCHES) continue;
1902 visited.add(noteUri);
1903 const note = await budget.get(noteUri);
1904 if (!note || !note.replies) continue;
1905 const childUris = await collectReplyItems(note.replies, 2, budget);
1906 for (const cu of childUris) {
1907 if (known.has(cu) || fetches >= THREAD_MAX_FETCHES) continue;
1908 known.add(cu);
1909 const child = await budget.get(cu);
1910 if (!child || !child.id || (child.type !== 'Note' && child.type !== 'Article')) continue;
1911 const actorUri = actorUriOf(child.attributedTo);
1912 if (!actorUri || isBlockedAny(actorUri)) continue; // skip blocked authors
1913 const actor = await budget.get(actorUri); // may be null if budget spent → fallback handle
1914 const ai = actorInfo(actor, actorUri);
1915 const html = HtmlSanitizerService.sanitize(child.content || '');
1916 // The child replies to `note` by construction (it's in note's replies collection).
1917 try { iStmts().ins.run('reply', postId, child.id, actorUri, ai.name, ai.handle, ai.url, ai.icon, html, child.published || null, note.id || noteUri); added++; } catch { /* ignore */ }
1918 nextFrontier.push(child.id); // expand this reply's own replies next depth
1919 }
1920 }
1921 frontier = nextFrontier;
1922 }
1923 if (added) console.log('[AP] thread crawl', postId, '+' + added, 'remote replies (' + fetches + ' fetches)');
1924}
1925
1926// Stale-while-revalidate entry point: call from the post view. Renders nothing, blocks nothing —
1927// fires a background crawl only if this post hasn't been crawled within the TTL.
1928export function maybeCrawlThread(postId) {
1929 if (!postId || _crawlingThreads.has(postId)) return;
1930 if (Date.now() - threadCrawlTs(postId) < THREAD_TTL_MS) return;
1931 _crawlingThreads.add(postId);
1932 setThreadCrawlTs(postId, Date.now()); // optimistic mark so concurrent/next views don't re-fire
1933 crawlThread(postId).catch((e) => console.warn('[AP] thread crawl failed:', e && e.message)).finally(() => _crawlingThreads.delete(postId));
1934}
1935
1936let _selfHealing = false;
1937export async function selfHealTimeline() {
1938 if (_selfHealing) return; _selfHealing = true;
1939 try {
1940 let cur = 0;
1941 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; }
1942 if (cur >= SELFHEAL_VERSION) return; // already healed for this version — skip on normal boots
1943 let rows = [];
1944 try { rows = db.prepare('SELECT id, content, media_json, nsfw, cw, url FROM ap_timeline ORDER BY rowid DESC LIMIT 200').all(); } catch { /* no table */ }
1945 let healed = 0;
1946 for (const r of rows) {
1947 try {
1948 const note = await fetchNoteAP(r.id);
1949 if (note === 404) { db.prepare('DELETE FROM ap_timeline WHERE id = ?').run(r.id); healed++; continue; }
1950 if (!note || typeof note !== 'object') continue;
1951 const html = HtmlSanitizerService.sanitize(note.content || '');
1952 const media = mediaFromNote(note);
1953 const nsfw = note.sensitive ? 1 : 0; // re-sync NSFW/sensitive + CW onto already-cached posts
1954 const cw = note.summary || null;
1955 const url = note.url || null; // re-sync the human url (catches a remote slug rename)
1956 if ((html && html !== r.content) || media !== (r.media_json || '[]') || nsfw !== (r.nsfw || 0) || (cw || '') !== (r.cw || '') || (url && url !== r.url)) {
1957 db.prepare('UPDATE ap_timeline SET content = ?, media_json = ?, nsfw = ?, cw = ?, url = COALESCE(?, url) WHERE id = ?').run(html || r.content, media, nsfw, cw, url, r.id);
1958 healed++;
1959 }
1960 } catch { /* per-note best-effort */ }
1961 }
1962 try { db.prepare('INSERT OR REPLACE INTO app_settings (key, value) VALUES (?, ?)').run('selfheal_version', String(SELFHEAL_VERSION)); } catch { /* ignore */ }
1963 if (rows.length) console.log(`[AP] self-heal v${SELFHEAL_VERSION}: ${healed}/${rows.length} timeline notes`);
1964 } catch { /* never block boot */ } finally { _selfHealing = false; }
1965}
1966
1967// Follow a fediverse account by @handle (WebFinger → actor → signed Follow).
1968export async function followActor(site, handle, autoBoost = false) {
1969 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
1970 if (!base || !site || !site.slug) return { error: 'config' };
1971 // Accept any of: a profile/actor URL, an @user@host handle (WebFinger), or a
1972 // bare site domain (site.com) — for a single-actor site (Klonkt etc.) the root
1973 // resolves to its AP actor, so you can follow a site by just its domain.
1974 const s = String(handle || '').trim();
1975 let actorUrl;
1976 if (/^https?:\/\//i.test(s)) actorUrl = safeUrl(s) || null;
1977 else if (s.includes('@')) actorUrl = await webfingerResolve(s);
1978 else if (/^[a-z0-9.-]+\.[a-z]{2,}/i.test(s)) actorUrl = await resolveApActor('https://' + s.replace(/^\/+|\/+$/g, ''));
1979 else actorUrl = null;
1980 if (!actorUrl) return { error: 'not_found' };
1981 const actor = await fetchActor(actorUrl).catch(() => null);
1982 if (!actor || !actor.id || !actor.inbox) return { error: 'unreachable' };
1983 const ai = actorInfo(actor, actor.id);
1984 const me = actorId(base, site.slug);
1985 const keys = getOrCreateKeys(site.slug);
1986 const followId = `${me}#follow-${Date.now()}-${rid()}`;
1987 fwStmts().ins.run(site.slug, actor.id, ai.handle, ai.name, ai.icon, ai.url, actor.inbox, followId, 'pending', autoBoost ? 1 : 0);
1988 const follow = { '@context': AP_CONTEXT, id: followId, type: 'Follow', actor: me, object: actor.id };
1989 // Deliver via the retry queue: a Follow that fails the first attempt (peer down,
1990 // timeout, transient 5xx) is retried with backoff instead of staying stuck on
1991 // 'pending' forever — the Accept can only come back once the Follow lands.
1992 await deliverWithRetry(site.slug, actor.inbox, follow, `${me}#main-key`, keys.private_pem);
1993 console.log('[AP] follow', site.slug, '→', actor.id);
1994 // Follow + feature in one step → backfill their recent posts into the Cirkel right away.
1995 if (autoBoost) backfillFromOutbox(site.slug, actor.id).catch(() => {});
1996 return { ok: true, name: ai.name, handle: ai.handle, actor: actor.id };
1997}
1998
1999// Resolve a profile URL or @handle to a followable remote actor (for the
2000// authorize_interaction "Follow" flow). Returns display fields + inbox, or null
2001// when it isn't a reachable actor (e.g. the input was a post, not a profile).
2002export async function resolveRemoteActor(input) {
2003 const s = String(input || '').trim();
2004 const actorUrl = /^https?:\/\//i.test(s) ? (safeUrl(s) || null) : await webfingerResolve(s);
2005 if (!actorUrl) return null;
2006 const actor = await fetchActor(actorUrl).catch(() => null);
2007 if (!actor || !actor.id || !actor.inbox) return null;
2008 const ai = actorInfo(actor, actor.id);
2009 return { actor_uri: actor.id, actor_name: ai.name, actor_handle: ai.handle, actor_url: ai.url, actor_icon: ai.icon, inbox: actor.inbox };
2010}
2011
2012export async function unfollowActor(site, actorUri) {
2013 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
2014 const me = actorId(base, site.slug);
2015 const keys = getOrCreateKeys(site.slug);
2016 const row = fwStmts().one.get(site.slug, actorUri);
2017 // Undo(Follow) MUST reference the original Follow's real id so the remote can correlate it
2018 // and drop the follow. The old `${me}#follow` fallback never matched anything → the unfollow
2019 // silently failed on the remote. With no stored follow id (legacy row), skip the network Undo
2020 // rather than send an unmatchable one. Deliver durably via the retry queue.
2021 if (row && row.inbox && row.follow_id) {
2022 const undo = { '@context': AP_CONTEXT, id: `${me}/undo/${Date.now()}-${rid()}`, type: 'Undo', actor: me, object: { id: row.follow_id, type: 'Follow', actor: me, object: actorUri } };
2023 deliverWithRetry(site.slug, row.inbox, undo, `${me}#main-key`, keys.private_pem);
2024 } else if (row && row.inbox) {
2025 console.warn('[AP] unfollow', site.slug, '→', actorUri, '— no stored follow id; removed locally only (legacy follow, remote may keep it)');
2026 }
2027 fwStmts().del.run(site.slug, actorUri);
2028 return { ok: true };
2029}
2030
2031// Send a Like or Announce (boost) on a remote note FROM this site.
2032export async function sendInteraction(site, kind, targetNoteId, authorUri) {
2033 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
2034 if (!base || !site || !site.slug || !targetNoteId) return { error: 'config' };
2035 const me = actorId(base, site.slug);
2036 const keys = getOrCreateKeys(site.slug);
2037 // 'unboost' = Undo(Announce): retracts a boost so followers' servers remove the
2038 // reblog (matched on actor+object — no record of the original Announce needed).
2039 const fanout = (kind === 'boost' || kind === 'unboost'); // also goes to our followers
2040 const followersCol = `${me}/followers`;
2041 // Address the original author in cc so their server (Mastodon, WordPress/ActivityPub, …)
2042 // attributes the boost to their post and notifies them — without this, a shared-inbox
2043 // receiver has nothing to route the Announce to. Non-fragment activity ids + a `published`
2044 // stamp keep us aligned with what Mastodon emits.
2045 const audience = authorUri ? [followersCol, authorUri] : [followersCol];
2046 let act;
2047 if (kind === 'unboost' || kind === 'unlike') {
2048 // Undo(Announce) retracts a boost; Undo(Like) un-favourites (matched on actor+object,
2049 // no record of the original activity needed — Mastodon honours both).
2050 const inner = kind === 'unboost' ? 'Announce' : 'Like';
2051 act = {
2052 '@context': AP_CONTEXT,
2053 id: `${me}/undo/${Date.now()}-${rid()}`, type: 'Undo', actor: me,
2054 object: { id: `${me}/${inner.toLowerCase()}/${Date.now()}-${rid()}`, type: inner, actor: me, object: targetNoteId },
2055 };
2056 if (kind === 'unboost') { act.to = [PUBLIC]; act.cc = audience; }
2057 } else {
2058 const type = kind === 'boost' ? 'Announce' : 'Like';
2059 act = {
2060 '@context': AP_CONTEXT,
2061 id: `${me}/${type.toLowerCase()}/${Date.now()}-${rid()}`,
2062 type, actor: me, object: targetNoteId,
2063 };
2064 if (type === 'Announce') { act.published = new Date().toISOString(); act.to = [PUBLIC]; act.cc = audience; }
2065 }
2066 const inboxes = new Set();
2067 // Author first, via their PERSONAL inbox (not the shared one) so a multi-user receiver
2068 // routes the Announce/Like to the right post unambiguously.
2069 if (authorUri) { const a = await fetchActor(authorUri).catch(() => null); if (a) inboxes.add(a.inbox || (a.endpoints && a.endpoints.sharedInbox)); }
2070 if (fanout) { for (const f of fStmts().list.all(site.slug)) inboxes.add(f.shared_inbox || f.inbox); }
2071 // Queue each delivery (immediate attempt + backoff retries on failure via ap_delivery)
2072 // instead of a single fire-and-forget POST, so a transient hiccup at the receiver doesn't
2073 // silently lose the boost — same durability a new post (deliverCreate) already gets.
2074 let queued = 0;
2075 for (const inbox of [...inboxes].filter(Boolean)) { deliverWithRetry(site.slug, inbox, act, `${me}#main-key`, keys.private_pem); queued++; }
2076 console.log('[AP]', kind, site.slug, '→', targetNoteId, 'queued', queued, 'inbox(es)');
2077 return { ok: true, delivered: queued };
2078}
2079
2080// Notifications inbox: new followers + replies/likes/boosts on this site's posts.
2081export function getNotifications(slug, limit) {
2082 const out = [];
2083 try {
2084 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)) {
2085 out.push({ type: 'follow', handle: deriveHandle(f.actor_uri), url: f.actor_uri, created_at: f.created_at });
2086 }
2087 } catch { /* ignore */ }
2088 try {
2089 const rows = db.prepare(`
2090 SELECT i.kind, i.actor_name, i.actor_handle, i.actor_url, i.content, i.created_at,
2091 p.slug AS post_slug, p.title AS post_title
2092 FROM ap_interactions i LEFT JOIN posts p ON p.id = i.post_id
2093 WHERE p.site_id = (SELECT id FROM sites WHERE slug = ?)
2094 ORDER BY i.created_at DESC LIMIT 80
2095 `).all(slug);
2096 for (const r of rows) out.push({
2097 type: r.kind, name: r.actor_name, handle: r.actor_handle, url: r.actor_url,
2098 content: stripLeadingMentions(r.content), post_slug: r.post_slug, post_title: r.post_title, created_at: r.created_at,
2099 });
2100 } catch { /* ignore */ }
2101 out.sort((a, b) => new Date(b.created_at) - new Date(a.created_at));
2102 return out.slice(0, limit || 60);
2103}
2104
2105// ── Blocking / defederation ───────────────────────────────────────
2106let _insBl, _delBl, _listBl;
2107function blStmts() {
2108 if (!_insBl) {
2109 _insBl = db.prepare('INSERT OR IGNORE INTO ap_blocks (slug, target, kind, label, created_at) VALUES (?,?,?,?,CURRENT_TIMESTAMP)');
2110 _delBl = db.prepare('DELETE FROM ap_blocks WHERE slug = ? AND target = ?');
2111 _listBl = db.prepare('SELECT * FROM ap_blocks WHERE slug = ? ORDER BY created_at DESC');
2112 }
2113 return { ins: _insBl, del: _delBl, list: _listBl };
2114}
2115export function listBlocks(slug) { return blStmts().list.all(slug); }
2116
2117// True if an actor (or its whole domain) is blocked anywhere on this instance.
2118// Vote on a remote fediverse poll (a cached Question). A ballot = a Create(Note) carrying only a
2119// `name` (the chosen option) + inReplyTo the Question, addressed to the poll's author — the
2120// Mastodon-standard vote. Records our choice locally + optimistically bumps the counts; the
2121// author's Update(Question) refreshes the authoritative totals when it arrives.
2122export async function voteOnPoll(site, questionId, choices) {
2123 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
2124 if (!base || !site || !site.slug || !questionId) return { error: 'config' };
2125 let row; try { row = db.prepare('SELECT author_uri, poll_json FROM ap_timeline WHERE id = ? AND slug = ? LIMIT 1').get(questionId, site.slug); } catch { /* ignore */ }
2126 if (!row || !row.poll_json) return { error: 'not_found' };
2127 let poll; try { poll = JSON.parse(row.poll_json); } catch { return { error: 'not_found' }; }
2128 if (poll.closed) return { error: 'closed' };
2129 if (poll.voted) return { error: 'already' };
2130 const valid = new Set(poll.options.map((o) => o.name));
2131 const picks = (Array.isArray(choices) ? choices : [choices]).map(String).filter((c) => valid.has(c));
2132 if (!picks.length) return { error: 'invalid' };
2133 const chosen = poll.multiple ? [...new Set(picks)] : [picks[0]];
2134 const me = actorId(base, site.slug);
2135 const keys = getOrCreateKeys(site.slug);
2136 const authorUri = row.author_uri || null;
2137 const author = authorUri ? await fetchActor(authorUri).catch(() => null) : null;
2138 const inbox = author && (author.inbox || (author.endpoints && author.endpoints.sharedInbox));
2139 if (!inbox) return { error: 'unreachable' };
2140 for (const name of chosen) {
2141 const nid = `${me}/votes/${Date.now()}-${rid()}`;
2142 const note = { id: nid, type: 'Note', attributedTo: me, to: authorUri ? [authorUri] : [], name, inReplyTo: questionId, published: new Date().toISOString() };
2143 const create = { '@context': AP_CONTEXT, id: `${nid}/activity`, type: 'Create', actor: me, to: note.to, object: note };
2144 deliverWithRetry(site.slug, inbox, create, `${me}#main-key`, keys.private_pem);
2145 }
2146 // Local optimistic update (authoritative counts arrive via the author's Update(Question)).
2147 poll.voted = poll.multiple ? chosen : chosen[0];
2148 for (const o of poll.options) if (chosen.includes(o.name)) o.count = (o.count || 0) + 1;
2149 if (poll.voters != null) poll.voters += 1;
2150 try { db.prepare('UPDATE ap_timeline SET poll_json = ? WHERE id = ? AND slug = ?').run(JSON.stringify(poll), questionId, site.slug); } catch { /* ignore */ }
2151 return { ok: true };
2152}
2153
2154// Vote on ANY fediverse poll by URL (the interact page) — no timeline cache needed. Fetches
2155// the Question fresh, validates the choice(s), and casts the Mastodon-standard ballot (a
2156// Create(Note) with `name` + inReplyTo) straight to the poll's author. Used for polls you find
2157// by URL, not just ones from accounts you follow (which go through voteOnPoll via /news).
2158export async function voteOnRemotePoll(site, questionUrl, choices) {
2159 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
2160 if (!base || !site || !site.slug || !/^https?:\/\//i.test(String(questionUrl || ''))) return { error: 'config' };
2161 const q = await fetchActor(questionUrl).catch(() => null); // AP GET (SSRF-guarded)
2162 if (!q || q.type !== 'Question' || !q.id) return { error: 'not_found' };
2163 const poll = parsePoll(q);
2164 if (!poll) return { error: 'not_found' };
2165 if (poll.closed) return { error: 'closed' };
2166 const valid = new Set(poll.options.map((o) => o.name));
2167 const picks = (Array.isArray(choices) ? choices : [choices]).map(String).filter((c) => valid.has(c));
2168 if (!picks.length) return { error: 'invalid' };
2169 const chosen = poll.multiple ? [...new Set(picks)] : [picks[0]];
2170 const authorUri = actorUriOf(q.attributedTo);
2171 const author = authorUri ? await fetchActor(authorUri).catch(() => null) : null;
2172 const inbox = author && (author.inbox || (author.endpoints && author.endpoints.sharedInbox));
2173 if (!inbox) return { error: 'unreachable' };
2174 const me = actorId(base, site.slug);
2175 const keys = getOrCreateKeys(site.slug);
2176 for (const name of chosen) {
2177 const nid = `${me}/votes/${Date.now()}-${rid()}`;
2178 const note = { id: nid, type: 'Note', attributedTo: me, to: [authorUri], name, inReplyTo: q.id, published: new Date().toISOString() };
2179 const create = { '@context': AP_CONTEXT, id: `${nid}/activity`, type: 'Create', actor: me, to: note.to, object: note };
2180 deliverWithRetry(site.slug, inbox, create, `${me}#main-key`, keys.private_pem);
2181 }
2182 return { ok: true };
2183}
2184
2185export function isBlockedAny(actorUri) {
2186 if (!actorUri) return false;
2187 let domain = ''; try { domain = new URL(actorUri).host; } catch { /* ignore */ }
2188 try { return !!db.prepare("SELECT 1 FROM ap_blocks WHERE (kind='actor' AND target=?) OR (kind='domain' AND target=?) LIMIT 1").get(actorUri, domain); }
2189 catch { return false; }
2190}
2191
2192function purgeBlocked(kind, target) {
2193 try {
2194 if (kind === 'domain') {
2195 // Exact host match (a URL LIKE over-/under-matches: it misses bare-domain or :port
2196 // actor URIs and can catch look-alikes). Filter by parsed host, same as isBlockedAny.
2197 const purge = (table, col) => {
2198 let rows = [];
2199 try { rows = db.prepare(`SELECT DISTINCT ${col} AS u FROM ${table} WHERE ${col} IS NOT NULL AND ${col} != ''`).all(); } catch { return; }
2200 const del = db.prepare(`DELETE FROM ${table} WHERE ${col} = ?`);
2201 for (const r of rows) { let h = ''; try { h = new URL(r.u).host; } catch { /* skip */ } if (h === target) { try { del.run(r.u); } catch { /* ignore */ } } }
2202 };
2203 purge('ap_interactions', 'actor_uri');
2204 purge('ap_timeline', 'author_uri');
2205 purge('ap_followers', 'actor_uri');
2206 } else {
2207 db.prepare('DELETE FROM ap_interactions WHERE actor_uri = ?').run(target);
2208 db.prepare('DELETE FROM ap_timeline WHERE author_uri = ?').run(target);
2209 db.prepare('DELETE FROM ap_followers WHERE actor_uri = ?').run(target);
2210 }
2211 } catch { /* best-effort */ }
2212}
2213
2214// Block an actor (@handle or actor URL) or a whole domain; purges their content.
2215export async function blockTarget(site, input) {
2216 const raw = String(input || '').trim();
2217 if (!site || !site.slug || !raw) return { error: 'empty' };
2218 let kind, target, label;
2219 if (/^https?:\/\//i.test(raw)) { kind = 'actor'; target = raw; label = raw; }
2220 else if (raw.includes('@')) {
2221 const actorUrl = await webfingerResolve(raw);
2222 if (!actorUrl) return { error: 'not_found' };
2223 kind = 'actor'; target = actorUrl; label = raw.startsWith('@') ? raw : ('@' + raw);
2224 } else { kind = 'domain'; target = raw.toLowerCase(); label = raw.toLowerCase(); }
2225 blStmts().ins.run(site.slug, target, kind, label);
2226 purgeBlocked(kind, target);
2227 console.log('[AP] block', site.slug, kind, target);
2228 return { ok: true, label };
2229}
2230
2231export function unblock(site, target) { blStmts().del.run(site.slug, target); return { ok: true }; }
2232
2233export default {
2234 AP_CONTEXT, getOrCreateKeys, apWants, sendAP, actorId, noteId,
2235 buildActor, buildNote, buildCreate, buildOutbox, buildFollowers, buildFollowing, buildFeatured,
2236 followerCount, deliver, fetchActor, verifyRequest, handleInbox, deliverCreate, deliverDelete, deliverUpdate, deliverActorUpdate, resyncFeaturedPins,
2237 getInteractions, getInteractionById, setInteractionBoosted, setInteractionLiked, setMyReaction, getMyReactions, buildReplyNote, getOutboxNote, deliverReply, resolveRemoteNote,
2238 listOutbox, deliverOutboxDelete, deliverOutboxUpdate,
2239 webfingerResolve, followActor, resolveRemoteActor, unfollowActor, listFollowing, setAutoBoost, backfillFromOutbox, getTimeline, sendInteraction, voteOnPoll, voteOnRemotePoll,
2240 parseOwnPoll, pollTally, ownPollView, deliverPollUpdate, maybeCrawlThread,
2241 autoBoostCount, boostedCount, markBoosted, unmarkBoosted, markLiked, unmarkLiked, getTimelineReaction, upsertBoostedNote, getCirkelPosts, getCirkelMembers, selfHealTimeline,
2242 getNotifications, listBlocks, isBlockedAny, blockTarget, unblock,
2243 deliverWithRetry, enqueueDelivery, processDeliveryQueue, startDeliveryWorker,
2244 getReplyUris, markNotificationsSeen, countUnseenNotifications, hasPlayableAudio,
2245};
Note: See TracBrowser for help on using the repository browser.