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

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

feat(fediverse): NodeInfo 2.1 + per-note replies collection

NodeInfo (/.well-known/nodeinfo + /nodeinfo/2.1) so fediverse tools recognise the
instance (software 'klonkt', user/post counts). Notes now carry a 'replies'
OrderedCollection (/ap/notes/:id/replies) listing inbound + our own reply note
URIs, so remote servers can fetch the whole thread.

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

  • Property mode set to 100644
File size: 46.9 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 db from '../config/database.js';
20import HtmlSanitizerService from './HtmlSanitizerService.js';
21
22const PUBLIC = 'https://www.w3.org/ns/activitystreams#Public';
23const MAX_OUTBOX = 20;
24
25// ── RSA keys per actor (lazy, cached in DB) ───────────────────────
26// Prepared lazily (NOT at module load) — the ap_keys table is created in
27// initializeDatabase(), which runs after this module is imported.
28let _sel, _ins;
29function keyStmts() {
30 if (!_sel) {
31 _sel = db.prepare('SELECT public_pem, private_pem FROM ap_keys WHERE slug = ?');
32 _ins = db.prepare('INSERT OR IGNORE INTO ap_keys (slug, public_pem, private_pem, created_at) VALUES (?,?,?,CURRENT_TIMESTAMP)');
33 }
34 return { sel: _sel, ins: _ins };
35}
36
37export function getOrCreateKeys(slug) {
38 const { sel, ins } = keyStmts();
39 const row = sel.get(slug);
40 if (row) return row;
41 const { publicKey, privateKey } = crypto.generateKeyPairSync('rsa', {
42 modulusLength: 2048,
43 publicKeyEncoding: { type: 'spki', format: 'pem' },
44 privateKeyEncoding: { type: 'pkcs8', format: 'pem' },
45 });
46 ins.run(slug, publicKey, privateKey);
47 return sel.get(slug) || { public_pem: publicKey, private_pem: privateKey };
48}
49
50// ── content negotiation ───────────────────────────────────────────
51// True when the caller wants ActivityPub JSON rather than the HTML page.
52export function apWants(req) {
53 const a = String(req.headers.accept || '').toLowerCase();
54 return a.includes('application/activity+json') ||
55 (a.includes('application/ld+json') && a.includes('activitystreams'));
56}
57
58const AP_CONTENT_TYPE = 'application/activity+json; charset=utf-8';
59export function sendAP(res, obj) {
60 res.type(AP_CONTENT_TYPE);
61 res.set('Cache-Control', 'public, max-age=120');
62 res.send(JSON.stringify(obj));
63}
64
65// ── document builders ─────────────────────────────────────────────
66export function actorId(base, slug) { return `${base}/ap/users/${encodeURIComponent(slug)}`; }
67export function noteId(base, postId) { return `${base}/ap/notes/${encodeURIComponent(postId)}`; }
68
69export function buildActor(base, site) {
70 const id = actorId(base, site.slug);
71 const keys = getOrCreateKeys(site.slug);
72 const actor = {
73 '@context': ['https://www.w3.org/ns/activitystreams', 'https://w3id.org/security/v1'],
74 id,
75 type: 'Person',
76 preferredUsername: site.slug,
77 name: site.title || site.slug,
78 summary: site.tagline || site.description || '',
79 url: `${base}/${site.slug === site.primary_slug ? '' : 'user/' + encodeURIComponent(site.slug)}`,
80 manuallyApprovesFollowers: false,
81 discoverable: true,
82 inbox: `${id}/inbox`,
83 outbox: `${id}/outbox`,
84 followers: `${id}/followers`,
85 endpoints: { sharedInbox: `${base}/ap/inbox` },
86 publicKey: {
87 id: `${id}#main-key`,
88 owner: id,
89 publicKeyPem: keys.public_pem,
90 },
91 };
92 if (site.profile_photo) {
93 const u = /^https?:/.test(site.profile_photo) ? site.profile_photo : `${base}${site.profile_photo.startsWith('/') ? '' : '/'}${site.profile_photo}`;
94 actor.icon = { type: 'Image', url: u };
95 }
96 return actor;
97}
98
99// A single post as an AS2 Note (the object), and as a Create activity (for outbox/delivery).
100export function buildNote(base, site, post) {
101 const id = noteId(base, post.id);
102 const aId = actorId(base, site.slug);
103 const human = `${base}/${encodeURIComponent(post.slug)}`;
104 // Mastodon ignores a Note's `name`, so put the title INTO the content (bold
105 // first line) — the standard blog→fediverse convention. post.content is
106 // already sanitized HTML; the title is plain text, so escape it.
107 const escTitle = String(post.title || '').replace(/[<>&]/g, (c) => ({ '<': '&lt;', '>': '&gt;', '&': '&amp;' }[c]));
108 const titleHtml = post.title ? `<p><strong>${escTitle}</strong></p>` : '';
109
110 // Images travel as AP `attachment` (Mastodon strips <img> from content). Collect
111 // the cover + any inline <img>, make absolute, then strip <img> from the content
112 // to avoid duplicate rendering on clients that DO keep them.
113 const abs = (u) => !u ? null : (/^https?:/i.test(u) ? u : `${base}${u.startsWith('/') ? '' : '/'}${u}`);
114 const mediaType = (u) => {
115 const e = ((u || '').split('?')[0].match(/\.(\w+)$/) || [])[1];
116 return ({ jpg: 'image/jpeg', jpeg: 'image/jpeg', png: 'image/png', gif: 'image/gif', webp: 'image/webp', avif: 'image/avif' })[(e || '').toLowerCase()] || 'image/jpeg';
117 };
118 const urls = [];
119 if (post.cover_image_url) urls.push(abs(post.cover_image_url));
120 let body = post.content || '';
121 for (const m of body.matchAll(/<img\b[^>]*\bsrc="([^"]+)"[^>]*>/gi)) urls.push(abs(m[1]));
122 body = body.replace(/<img\b[^>]*>/gi, '');
123 // Audio shortcodes: do NOT federate the raw audio file — Klonkt deliberately
124 // gates audio (the /audio/stream URL has friction), and shipping it as an AP
125 // audio attachment would hand Mastodon a plain, downloadable mp3 URL. Instead,
126 // replace the shortcodes with a "🎵 listen on the site" link so the post invites
127 // a click-through to the protected player (discovery without leaking the file).
128 const esc = (s) => String(s == null ? '' : s).replace(/[<>&]/g, (c) => ({ '<': '&lt;', '>': '&gt;', '&': '&amp;' }[c]));
129 const audioLabels = [];
130 try {
131 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); }
132 for (const m of body.matchAll(/\[\[album:([^\]]+)\]\]/g)) audioLabels.push(m[1].trim());
133 } catch { /* non-fatal */ }
134 const hadAudio = /\[\[(track|album|playlist):/i.test(body);
135 body = body.replace(/\[\[(track|album|playlist):[^\]]+\]\]/gi, '');
136 if (hadAudio) {
137 const lbl = audioLabels.length ? esc(audioLabels.slice(0, 4).join(', ')) : '';
138 body += `<p>🎵 ${lbl ? `<strong>${lbl}</strong> — ` : ''}<a href="${human}">listen on ${esc(site.title || 'the site')}</a></p>`;
139 }
140 const seen = new Set();
141 const attachment = urls.filter(Boolean)
142 .filter((u) => { if (seen.has(u)) return false; seen.add(u); return true; })
143 .map((u) => ({ type: 'Document', mediaType: mediaType(u), url: u }));
144
145 const note = {
146 id,
147 type: 'Note',
148 attributedTo: aId,
149 content: titleHtml + body,
150 url: human,
151 published: new Date(post.published_at || post.created_at || Date.now()).toISOString(),
152 to: [PUBLIC],
153 cc: [`${aId}/followers`],
154 tag: Array.isArray(post.tags) ? post.tags.map((t) => ({ type: 'Hashtag', name: '#' + String(t).replace(/\s+/g, '') })) : [],
155 replies: `${id}/replies`,
156 };
157 if (attachment.length) note.attachment = attachment;
158 return note;
159}
160
161// All reply note URIs on a local post (inbound fediverse replies + our own
162// outbound replies) — backs the Note's `replies` Collection so remote servers
163// can fetch the whole thread.
164export function getReplyUris(base, postId) {
165 const out = [];
166 try {
167 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);
168 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}`);
169 } catch { /* non-fatal */ }
170 return out;
171}
172
173export function buildCreate(base, site, post) {
174 const note = buildNote(base, site, post);
175 return {
176 '@context': 'https://www.w3.org/ns/activitystreams',
177 id: note.id + '#create',
178 type: 'Create',
179 actor: actorId(base, site.slug),
180 published: note.published,
181 to: note.to,
182 cc: note.cc,
183 object: note,
184 };
185}
186
187export function buildOutbox(base, site, posts) {
188 const id = `${actorId(base, site.slug)}/outbox`;
189 const items = (posts || []).slice(0, MAX_OUTBOX).map((p) => buildCreate(base, site, p));
190 return {
191 '@context': 'https://www.w3.org/ns/activitystreams',
192 id,
193 type: 'OrderedCollection',
194 totalItems: items.length,
195 orderedItems: items,
196 };
197}
198
199export function buildFollowers(base, site, count) {
200 const id = `${actorId(base, site.slug)}/followers`;
201 return {
202 '@context': 'https://www.w3.org/ns/activitystreams',
203 id,
204 type: 'OrderedCollection',
205 totalItems: count || 0,
206 orderedItems: [], // hidden for privacy; count only
207 };
208}
209
210// ── followers store (lazy stmts) ──────────────────────────────────
211let _insF, _delF, _listF, _cntF;
212function fStmts() {
213 if (!_insF) {
214 _insF = db.prepare('INSERT OR IGNORE INTO ap_followers (slug, actor_uri, inbox, shared_inbox, created_at) VALUES (?,?,?,?,CURRENT_TIMESTAMP)');
215 _delF = db.prepare('DELETE FROM ap_followers WHERE slug = ? AND actor_uri = ?');
216 _listF = db.prepare('SELECT inbox, shared_inbox FROM ap_followers WHERE slug = ?');
217 _cntF = db.prepare('SELECT COUNT(*) n FROM ap_followers WHERE slug = ?');
218 }
219 return { ins: _insF, del: _delF, list: _listF, cnt: _cntF };
220}
221export function followerCount(slug) { return fStmts().cnt.get(slug).n; }
222
223// ── inbound interactions store (replies / likes / boosts) + our outbound replies ──
224let _insI, _delLA, _delReply, _listI, _getI, _insO, _listO, _getO;
225function iStmts() {
226 if (!_insI) {
227 _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)');
228 _delLA = db.prepare('DELETE FROM ap_interactions WHERE kind = ? AND post_id = ? AND actor_uri = ?');
229 _delReply = db.prepare("DELETE FROM ap_interactions WHERE kind = 'reply' AND object_uri = ?");
230 _listI = db.prepare('SELECT id, kind, object_uri, parent_uri, actor_uri, actor_name, actor_handle, actor_url, actor_icon, content, published, created_at FROM ap_interactions WHERE post_id = ? ORDER BY created_at ASC');
231 _getI = db.prepare('SELECT * FROM ap_interactions WHERE id = ?');
232 _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)');
233 _listO = db.prepare('SELECT * FROM ap_outbox WHERE post_id = ? ORDER BY created_at ASC');
234 _getO = db.prepare('SELECT * FROM ap_outbox WHERE id = ?');
235 }
236 return { ins: _insI, delLA: _delLA, delReply: _delReply, list: _listI, getI: _getI, insO: _insO, listO: _listO, getO: _getO };
237}
238
239export function getInteractionById(id) { return iStmts().getI.get(id); }
240
241const localPostExists = (id) => { try { return !!db.prepare('SELECT 1 FROM posts WHERE id = ?').get(id); } catch { return false; } };
242// Extract our local post id from a note URL, but only if it's ours (base match).
243function postIdFromNoteUrl(url, base) {
244 const s = String(url || '');
245 if (base && !s.startsWith(base)) return null;
246 const m = s.match(/\/ap\/notes\/([^/?#]+)/);
247 return m ? decodeURIComponent(m[1]) : null;
248}
249function deriveHandle(actorUri) {
250 try { const u = new URL(actorUri); const seg = u.pathname.split('/').filter(Boolean).pop() || ''; return `@${seg}@${u.host}`; } catch { return String(actorUri || ''); }
251}
252function actorInfo(doc, actorUri) {
253 let host = ''; try { host = new URL(actorUri).host; } catch { /* keep empty */ }
254 const handle = doc && doc.preferredUsername ? `@${doc.preferredUsername}@${host}` : deriveHandle(actorUri);
255 const icon = doc && doc.icon ? (doc.icon.url || (Array.isArray(doc.icon) && doc.icon[0] && doc.icon[0].url)) : null;
256 return {
257 name: (doc && (doc.name || doc.preferredUsername)) || handle,
258 handle,
259 url: (doc && (doc.url || doc.id)) || actorUri,
260 icon: icon || null,
261 };
262}
263
264// Given an inReplyTo note URL, find which local post the thread belongs to + the
265// note being replied to (parent), so a reply-to-a-comment can be nested.
266function findThreadTarget(inReplyTo, base) {
267 if (!inReplyTo) return null;
268 const seg = postIdFromNoteUrl(inReplyTo, base); // our /ap/notes/<id> segment (if ours)
269 if (seg && localPostExists(seg)) return { post_id: seg, parent_uri: inReplyTo };
270 if (seg) {
271 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 */ }
272 }
273 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 */ }
274 return null;
275}
276
277// View-ready threaded view of a post's fediverse activity (inbound replies +
278// our outbound replies, nested), plus like/boost counts.
279export function getInteractions(postId, base, site) {
280 const s = iStmts();
281 const rows = s.list.all(postId);
282 const baseClean = (base || process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
283 const postNoteId = baseClean ? `${baseClean}/ap/notes/${postId}` : null;
284 // Our own (outbound) replies show the SITE identity for everyone (not "You").
285 let host = ''; try { host = new URL(baseClean).host; } catch { /* ignore */ }
286 const siteName = (site && (site.title || site.slug)) || '';
287 const siteHandle = (site && site.slug && host) ? `@${site.slug}@${host}` : '';
288 const siteUrl = baseClean ? `${baseClean}/` : '';
289 const siteIcon = (site && site.profile_photo) || null;
290
291 const nodes = [];
292 for (const r of rows) {
293 if (r.kind !== 'reply') continue;
294 nodes.push({
295 noteId: r.object_uri, parent: r.parent_uri || null, mine: false,
296 actor_name: r.actor_name, actor_handle: r.actor_handle, actor_url: r.actor_url,
297 actor_icon: r.actor_icon, content: r.content, created_at: r.published || r.created_at,
298 children: [],
299 });
300 }
301 for (const o of s.listO.all(postId)) {
302 nodes.push({
303 noteId: baseClean ? `${baseClean}/ap/notes/${o.id}` : o.id, parent: o.in_reply_to || null,
304 mine: true, outboxId: o.id, content: o.content, created_at: o.created_at,
305 actor_name: siteName, actor_handle: siteHandle, actor_url: siteUrl, actor_icon: siteIcon,
306 children: [],
307 });
308 }
309
310 const byId = new Map(nodes.map((n) => [n.noteId, n]));
311 const isTop = (n) => !n.parent || n.parent === postNoteId || !byId.has(n.parent);
312 const tops = [];
313 for (const n of nodes) {
314 if (isTop(n)) { tops.push(n); continue; }
315 let anc = n, guard = 0;
316 while (!isTop(anc) && guard++ < 12) anc = byId.get(anc.parent);
317 anc.children.push(n);
318 }
319 const byTime = (a, b) => new Date(a.created_at) - new Date(b.created_at);
320 tops.sort(byTime).forEach((t) => t.children.sort(byTime));
321
322 return {
323 thread: tops,
324 likeCount: rows.filter((r) => r.kind === 'like').length,
325 announceCount: rows.filter((r) => r.kind === 'announce').length,
326 total: nodes.length,
327 };
328}
329
330// ── HTTP Signatures + delivery ────────────────────────────────────
331const slugFromActorUrl = (url) => { const m = String(url || '').match(/\/ap\/users\/([^/?#]+)/); return m ? decodeURIComponent(m[1]) : null; };
332
333// Sign + POST an activity to a remote inbox (draft-cavage HTTP Signatures, RSA-SHA256).
334export async function deliver(inboxUrl, bodyObj, keyId, privatePem) {
335 const body = JSON.stringify(bodyObj);
336 const u = new URL(inboxUrl);
337 const date = new Date().toUTCString();
338 const digest = 'SHA-256=' + crypto.createHash('sha256').update(body).digest('base64');
339 const signingString = `(request-target): post ${u.pathname}\nhost: ${u.host}\ndate: ${date}\ndigest: ${digest}`;
340 const signature = crypto.sign('sha256', Buffer.from(signingString), privatePem).toString('base64');
341 const sig = `keyId="${keyId}",algorithm="rsa-sha256",headers="(request-target) host date digest",signature="${signature}"`;
342 const r = await fetch(inboxUrl, {
343 method: 'POST',
344 headers: { 'Content-Type': 'application/activity+json', Accept: 'application/activity+json', Date: date, Digest: digest, Signature: sig },
345 body,
346 signal: AbortSignal.timeout(8000),
347 });
348 return r.status;
349}
350
351export async function fetchActor(url) {
352 try {
353 const r = await fetch(url, { headers: { Accept: 'application/activity+json' }, redirect: 'follow', signal: AbortSignal.timeout(8000) });
354 if (!r.ok) return null;
355 return await r.json();
356 } catch { return null; }
357}
358
359// ── Delivery queue with retries ───────────────────────────────────
360// Outbound deliveries are tried immediately; on failure (down server, timeout,
361// non-2xx) they're queued and retried with backoff so a briefly-offline follower
362// doesn't silently miss the post. The signing key is NOT stored — the worker
363// re-derives it from the actor slug at send time.
364const DELIVERY_MAX_ATTEMPTS = 6;
365const DELIVERY_BACKOFF_MIN = [1, 5, 15, 60, 180, 360];
366let _insDeliv, _dueDeliv, _delDeliv, _bumpDeliv;
367function deliveryStmts() {
368 if (!_insDeliv) {
369 _insDeliv = db.prepare('INSERT INTO ap_delivery (slug, inbox, body, attempts, next_at) VALUES (?,?,?,0,CURRENT_TIMESTAMP)');
370 _dueDeliv = db.prepare("SELECT * FROM ap_delivery WHERE datetime(next_at) <= datetime('now') ORDER BY next_at LIMIT 30");
371 _delDeliv = db.prepare('DELETE FROM ap_delivery WHERE id = ?');
372 _bumpDeliv = db.prepare('UPDATE ap_delivery SET attempts = ?, next_at = ? WHERE id = ?');
373 }
374 return { ins: _insDeliv, due: _dueDeliv, del: _delDeliv, bump: _bumpDeliv };
375}
376export function enqueueDelivery(slug, inbox, activity) {
377 if (!slug || !inbox || !activity) return;
378 try { deliveryStmts().ins.run(slug, inbox, JSON.stringify(activity)); } catch { /* ignore */ }
379}
380// Deliver now; queue for retry if it fails.
381export async function deliverWithRetry(slug, inbox, activity, keyId, privPem) {
382 if (!inbox) return;
383 try { const st = await deliver(inbox, activity, keyId, privPem); if (st >= 200 && st < 300) return; } catch { /* queue below */ }
384 enqueueDelivery(slug, inbox, activity);
385}
386export async function processDeliveryQueue() {
387 let rows;
388 try { rows = deliveryStmts().due.all(); } catch { return; }
389 if (!rows || !rows.length) return;
390 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
391 for (const row of rows) {
392 let ok = false;
393 try {
394 const keys = getOrCreateKeys(row.slug);
395 const st = await deliver(row.inbox, JSON.parse(row.body), `${actorId(base, row.slug)}#main-key`, keys.private_pem);
396 ok = st >= 200 && st < 300;
397 } catch { ok = false; }
398 if (ok) { deliveryStmts().del.run(row.id); continue; }
399 const attempts = row.attempts + 1;
400 if (attempts >= DELIVERY_MAX_ATTEMPTS) { deliveryStmts().del.run(row.id); console.warn('[AP] delivery gave up after', attempts, 'tries →', row.inbox); continue; }
401 const mins = DELIVERY_BACKOFF_MIN[Math.min(attempts, DELIVERY_BACKOFF_MIN.length - 1)];
402 deliveryStmts().bump.run(attempts, new Date(Date.now() + mins * 60000).toISOString(), row.id);
403 }
404}
405let _delivTimer = null;
406export function startDeliveryWorker() {
407 if (_delivTimer) return;
408 _delivTimer = setInterval(() => { processDeliveryQueue().catch(() => {}); }, 60 * 1000);
409 if (_delivTimer.unref) _delivTimer.unref();
410}
411
412// Best-effort verification of an incoming signed request. Returns the sender's
413// actor doc if the signature checks out, else null. (Not gating yet — MVP.)
414export async function verifyRequest(req) {
415 const sigH = req.headers['signature'];
416 if (!sigH) return null;
417 const p = Object.fromEntries([...sigH.matchAll(/([a-zA-Z]+)="([^"]*)"/g)].map((m) => [m[1], m[2]]));
418 if (!p.keyId || !p.signature) return null;
419 const actor = await fetchActor(p.keyId.split('#')[0]);
420 const pem = actor && actor.publicKey && actor.publicKey.publicKeyPem;
421 if (!pem) return null;
422 const hs = (p.headers || '(request-target) host date').split(/\s+/);
423 const line = hs.map((h) => h === '(request-target)'
424 ? `(request-target): ${req.method.toLowerCase()} ${req.originalUrl}`
425 : `${h}: ${req.headers[h] || ''}`).join('\n');
426 let ok = false;
427 try { ok = crypto.verify('sha256', Buffer.from(line), pem, Buffer.from(p.signature, 'base64')); } catch { ok = false; }
428 if (ok && hs.includes('digest') && req.rawBody) {
429 const exp = 'SHA-256=' + crypto.createHash('sha256').update(req.rawBody).digest('base64');
430 if (req.headers['digest'] !== exp) ok = false;
431 }
432 return ok ? actor : null;
433}
434
435// Handle an incoming inbox POST. slugParam = null for the shared /ap/inbox.
436export async function handleInbox(req, slugParam) {
437 const act = req.body || {};
438 const type = act.type;
439 const base = (process.env.PUBLIC_BASE_URL || `${req.protocol}://${req.get('host')}`).replace(/\/+$/, '');
440 const verified = await verifyRequest(req).catch(() => null);
441
442 // ENFORCE HTTP signatures: a data-affecting activity must be signed by the very
443 // actor it claims to be. No valid signature, or signer ≠ actor → reject (no
444 // forged replies/likes/follows/timeline posts). GET/discovery stays open.
445 const claimedActor = typeof act.actor === 'string' ? act.actor : (act.actor && act.actor.id);
446 // Blocked actor/domain → silently drop (202, don't reveal the block).
447 if (claimedActor && isBlockedAny(claimedActor)) { console.log('[AP] inbox dropped (blocked)', claimedActor); return 202; }
448 const GATED = ['Create', 'Like', 'Announce', 'Follow', 'Delete', 'Undo', 'Accept', 'Reject'];
449 if (GATED.includes(type)) {
450 if (!verified || !claimedActor || verified.id !== claimedActor) {
451 console.warn('[AP] inbox REJECTED (signature)', type, claimedActor || '?', verified ? '(signer mismatch)' : '(unsigned/invalid)');
452 return 401;
453 }
454 }
455
456 if (type === 'Follow') {
457 const who = typeof act.actor === 'string' ? act.actor : (act.actor && act.actor.id);
458 const slug = slugParam || slugFromActorUrl(typeof act.object === 'string' ? act.object : (act.object && act.object.id));
459 if (!who || !slug) return 400;
460 const remote = await fetchActor(who);
461 if (!remote || !remote.inbox) return 202; // can't reach them → drop quietly
462 fStmts().ins.run(slug, who, remote.inbox, (remote.endpoints && remote.endpoints.sharedInbox) || null);
463 const me = actorId(base, slug);
464 const keys = getOrCreateKeys(slug);
465 const accept = { '@context': 'https://www.w3.org/ns/activitystreams', id: `${me}#accept-${Date.now()}`, type: 'Accept', actor: me, object: act };
466 deliver(remote.inbox, accept, `${me}#main-key`, keys.private_pem).catch((e) => console.warn('[AP] Accept delivery failed:', e.message));
467 console.log('[AP] Follow', who, '→', slug, verified ? '(sig ok)' : '(sig unverified)');
468 return 202;
469 }
470 if (type === 'Undo' && act.object) {
471 const who = typeof act.actor === 'string' ? act.actor : (act.actor && act.actor.id);
472 const ot = act.object.type;
473 if (ot === 'Follow') {
474 const obj = act.object.object;
475 const slug = slugParam || slugFromActorUrl(typeof obj === 'string' ? obj : (obj && obj.id));
476 if (who && slug) { fStmts().del.run(slug, who); console.log('[AP] Unfollow', who, '→', slug); }
477 return 202;
478 }
479 if (ot === 'Like' || ot === 'Announce') {
480 const tgt = act.object.object;
481 const pid = postIdFromNoteUrl(typeof tgt === 'string' ? tgt : (tgt && tgt.id), base);
482 if (who && pid) { iStmts().delLA.run(ot.toLowerCase(), pid, who); console.log('[AP] Undo', ot, who, '→', pid); }
483 return 202;
484 }
485 return 202;
486 }
487
488 const actorUri = typeof act.actor === 'string' ? act.actor : (act.actor && act.actor.id);
489 const resolveActor = async (uri) => ((verified && verified.id === uri) ? verified : await fetchActor(uri).catch(() => null));
490 // Activities from our OWN actors are already stored via ap_outbox — don't re-store.
491 const isLocalActor = !!(base && actorUri && actorUri.startsWith(`${base}/ap/users/`));
492
493 // Inbound reply: a Create whose object replies to one of our notes (post OR comment).
494 if (type === 'Create' && act.object && (act.object.type === 'Note' || act.object.type === 'Article')) {
495 const o = act.object;
496 const tgt = findThreadTarget(o.inReplyTo, base);
497 if (tgt && actorUri && !isLocalActor) {
498 const ai = actorInfo(await resolveActor(actorUri), actorUri);
499 const html = HtmlSanitizerService.sanitize(o.content || '');
500 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);
501 console.log('[AP] reply', actorUri, '→', tgt.post_id);
502 return 202;
503 }
504 // Home timeline (client): a top-level post from an account we follow.
505 if (actorUri && !isLocalActor && !o.inReplyTo && o.id) {
506 let subs = []; try { subs = db.prepare('SELECT slug FROM ap_following WHERE actor_uri = ?').all(actorUri); } catch { /* table may not exist yet */ }
507 if (subs.length) {
508 const ai = actorInfo(await resolveActor(actorUri), actorUri);
509 const html = HtmlSanitizerService.sanitize(o.content || '');
510 const media = JSON.stringify((Array.isArray(o.attachment) ? o.attachment : []).filter((a) => a && a.url).map((a) => ({ url: a.url, type: a.mediaType || '' })));
511 for (const s of subs) tlStmts().ins.run(o.id, s.slug, actorUri, ai.name, ai.handle, ai.icon, ai.url, html, o.url || null, o.published || null, media);
512 console.log('[AP] timeline +', actorUri, 'x' + subs.length);
513 }
514 }
515 return 202;
516 }
517 if (type === 'Like' || type === 'Announce') {
518 const tgt = act.object;
519 const pid = postIdFromNoteUrl(typeof tgt === 'string' ? tgt : (tgt && tgt.id), base);
520 if (pid && actorUri && !isLocalActor && localPostExists(pid)) {
521 const ai = actorInfo(await resolveActor(actorUri), actorUri);
522 iStmts().ins.run(type.toLowerCase(), pid, '', actorUri, ai.name, ai.handle, ai.url, ai.icon, null, null, null);
523 console.log('[AP]', type === 'Like' ? 'like' : 'boost', actorUri, '→', pid);
524 }
525 return 202;
526 }
527 if (type === 'Delete') {
528 // A remote note was deleted upstream → drop it from replies AND the timeline.
529 const oid = typeof act.object === 'string' ? act.object : (act.object && act.object.id);
530 if (oid) { iStmts().delReply.run(oid); try { tlStmts().del.run(oid); } catch { /* ignore */ } }
531 return 202;
532 }
533 // Accept/Reject of a Follow WE sent (client side).
534 if (type === 'Accept' && act.object) {
535 const fid = typeof act.object === 'string' ? act.object : (act.object && act.object.id);
536 if (fid) { try { fwStmts().acc.run(fid); } catch { /* ignore */ } }
537 console.log('[AP] follow accepted', actorUri);
538 return 202;
539 }
540 if (type === 'Reject' && act.object) {
541 const who = actorUri;
542 if (who && slugParam) { try { fwStmts().del.run(slugParam, who); } catch { /* ignore */ } }
543 return 202;
544 }
545
546 console.log('[AP] inbox', type || 'unknown', '→', slugParam || 'shared', '(ignored)');
547 return 202;
548}
549
550// Deliver a new post as Create(Note) to all followers' inboxes (fire-and-forget).
551// Needs PUBLIC_BASE_URL (absolute URLs); no-op without followers or base.
552export async function deliverCreate(site, post) {
553 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
554 if (!base || !site || !site.slug) return;
555 const followers = fStmts().list.all(site.slug);
556 if (!followers.length) return;
557 const inboxes = [...new Set(followers.map((f) => f.shared_inbox || f.inbox).filter(Boolean))];
558 const keys = getOrCreateKeys(site.slug);
559 const keyId = `${actorId(base, site.slug)}#main-key`;
560 const create = buildCreate(base, site, post);
561 for (const inbox of inboxes) deliverWithRetry(site.slug, inbox, create, keyId, keys.private_pem);
562}
563
564// Tell followers a post is gone (Delete + Tombstone) so it's removed from their feeds.
565export async function deliverDelete(site, post) {
566 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
567 if (!base || !site || !site.slug || !post || !post.id) return;
568 const followers = fStmts().list.all(site.slug);
569 if (!followers.length) return;
570 const inboxes = [...new Set(followers.map((f) => f.shared_inbox || f.inbox).filter(Boolean))];
571 const keys = getOrCreateKeys(site.slug);
572 const me = actorId(base, site.slug);
573 const nid = noteId(base, post.id);
574 const del = {
575 '@context': 'https://www.w3.org/ns/activitystreams',
576 id: `${nid}#delete-${Date.now()}`,
577 type: 'Delete',
578 actor: me,
579 to: [PUBLIC],
580 object: { id: nid, type: 'Tombstone' },
581 };
582 for (const inbox of inboxes) deliverWithRetry(site.slug, inbox, del, `${me}#main-key`, keys.private_pem);
583}
584
585// ── outbound replies (Klonkt → fediverse) ─────────────────────────
586const escHtml = (s) => String(s || '').replace(/[<>&]/g, (c) => ({ '<': '&lt;', '>': '&gt;', '&': '&amp;' }[c]));
587const 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(); };
588
589// Build one of OUR outbound reply Notes from an ap_outbox row.
590export function buildReplyNote(base, site, row) {
591 const me = actorId(base, site.slug);
592 return {
593 id: noteId(base, row.id),
594 type: 'Note',
595 attributedTo: me,
596 inReplyTo: row.in_reply_to || undefined,
597 content: row.content,
598 url: row.post_slug ? `${base}/${encodeURIComponent(row.post_slug)}` : undefined,
599 published: toISO(row.created_at),
600 to: row.to_actor ? [row.to_actor] : [PUBLIC],
601 cc: [PUBLIC, `${me}/followers`],
602 tag: row.to_actor ? [{ type: 'Mention', href: row.to_actor, name: row.to_handle }] : [],
603 };
604}
605
606// Resolve one of our outbound reply Notes by id (for /ap/notes/:id fallback).
607export function getOutboxNote(base, id) {
608 const row = iStmts().getO.get(id);
609 if (!row) return null;
610 const site = db.prepare('SELECT * FROM sites WHERE slug = ?').get(row.site_slug);
611 if (!site) return null;
612 return buildReplyNote(base, site, row);
613}
614
615// Send a reply FROM this site to a remote actor (in reply to their inbound reply).
616// `parent` = an ap_interactions row (actor_uri, actor_url, actor_handle, object_uri).
617export async function deliverReply(site, { postId, postSlug, parent, text }) {
618 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
619 if (!base || !site || !site.slug || !parent || !String(text || '').trim()) return null;
620 const me = actorId(base, site.slug);
621 const handle = parent.actor_handle || deriveHandle(parent.actor_uri);
622 const body = escHtml(String(text).trim()).replace(/\r?\n/g, '<br>');
623 const mention = parent.actor_uri
624 ? `<a href="${escHtml(parent.actor_url || parent.actor_uri)}" class="u-url mention">${escHtml(handle)}</a> ` : '';
625 const content = `<p>${mention}${body}</p>`;
626 // Dedup: skip if the exact same reply was already sent (double-submit guard).
627 const dup = db.prepare('SELECT 1 FROM ap_outbox WHERE site_slug = ? AND IFNULL(in_reply_to, \'\') = ? AND content = ? LIMIT 1')
628 .get(site.slug, parent.object_uri || '', content);
629 if (dup) { console.log('[AP] outreply skipped (duplicate)'); return { duplicate: true, delivered: 0 }; }
630 const id = crypto.randomUUID();
631 iStmts().insO.run(id, site.slug, postId, postSlug || null, parent.object_uri || null, parent.actor_uri || null, handle, content);
632 const row = iStmts().getO.get(id);
633 const note = buildReplyNote(base, site, row);
634 const create = {
635 '@context': 'https://www.w3.org/ns/activitystreams',
636 id: note.id + '#create', type: 'Create', actor: me,
637 published: note.published, to: note.to, cc: note.cc, object: note,
638 };
639 const keys = getOrCreateKeys(site.slug);
640 const keyId = `${me}#main-key`;
641 const inboxes = new Set();
642 if (parent.actor_uri) {
643 const a = await fetchActor(parent.actor_uri).catch(() => null);
644 if (a) inboxes.add((a.endpoints && a.endpoints.sharedInbox) || a.inbox);
645 }
646 if (parent.threadInbox) inboxes.add(parent.threadInbox); // back-compat (single)
647 (parent.threadInboxes || []).forEach((i) => inboxes.add(i)); // whole ancestor chain
648 for (const f of fStmts().list.all(site.slug)) inboxes.add(f.shared_inbox || f.inbox);
649 inboxes.delete(`${me}/inbox`); // never deliver to ourselves (already in ap_outbox)
650 inboxes.delete(`${base}/ap/inbox`); // (our own shared inbox) → avoids a self-duplicate
651 let delivered = 0;
652 for (const inbox of [...inboxes].filter(Boolean)) {
653 try { const st = await deliver(inbox, create, keyId, keys.private_pem); if (st >= 200 && st < 300) delivered++; } catch { /* best-effort */ }
654 }
655 console.log('[AP] outreply', site.slug, '→', parent.actor_uri, 'delivered', delivered);
656 return { id, content, delivered };
657}
658
659// Resolve a remote post URL (any fediverse/Klonkt post) into a reply target.
660// Returns a parent-shaped object usable by deliverReply(), or null.
661export async function resolveRemoteNote(url) {
662 if (!/^https?:\/\//i.test(String(url || ''))) return null;
663 const note = await fetchActor(url).catch(() => null); // AP GET (content-negotiates)
664 if (!note || !note.id) return null;
665 const att = note.attributedTo;
666 const actorUri = typeof att === 'string' ? att : (att && att.id);
667 if (!actorUri) return null;
668 const actor = await fetchActor(actorUri).catch(() => null);
669 const ai = actorInfo(actor, actorUri);
670 // Is what we're replying to a post (or a comment) on one of OUR posts? If so,
671 // link our reply to that local post so it shows nested in the post thread.
672 const localTgt = findThreadTarget(note.id, (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, ''));
673 // Walk the WHOLE reply chain upward (comment → parent comment → … → root post)
674 // and collect every ancestor author's inbox, so each participant's server —
675 // including the original post's author — receives + threads our reply.
676 const threadInboxes = [];
677 const seenInbox = new Set();
678 let cursor = note.inReplyTo, guard = 0;
679 while (cursor && guard++ < 6) {
680 const url = typeof cursor === 'string' ? cursor : (cursor && cursor.id);
681 if (!url) break;
682 const pn = await fetchActor(url).catch(() => null);
683 if (!pn) break;
684 const pa = typeof pn.attributedTo === 'string' ? pn.attributedTo : (pn.attributedTo && pn.attributedTo.id);
685 if (pa && pa !== actorUri) {
686 const paDoc = await fetchActor(pa).catch(() => null);
687 const inbox = paDoc && ((paDoc.endpoints && paDoc.endpoints.sharedInbox) || paDoc.inbox);
688 if (inbox && !seenInbox.has(inbox)) { seenInbox.add(inbox); threadInboxes.push(inbox); }
689 }
690 cursor = pn.inReplyTo; // climb to the next ancestor
691 }
692 const rawHtml = String(note.content || '').replace(/\[\[(track|album|playlist):[^\]]+\]\]/gi, '');
693 const images = (Array.isArray(note.attachment) ? note.attachment : [])
694 .filter((a) => a && a.url && (!a.mediaType || /^image\//i.test(a.mediaType)))
695 .map((a) => a.url);
696 return {
697 object_uri: note.id,
698 actor_uri: actorUri,
699 actor_url: ai.url,
700 actor_handle: ai.handle,
701 actor_name: ai.name,
702 actor_icon: ai.icon,
703 url: note.url || url,
704 content: HtmlSanitizerService.sanitize(rawHtml), // full, sanitized
705 images,
706 threadInboxes, // every ancestor author's inbox
707 localPostId: localTgt ? localTgt.post_id : '', // our post this belongs to (if any)
708 preview: HtmlSanitizerService.toPlainText(note.content || '').slice(0, 240),
709 };
710}
711
712// List a site's own outbound fediverse replies (for the manage/delete view).
713export function listOutbox(siteSlug) {
714 return db.prepare('SELECT id, content, to_handle, in_reply_to, created_at FROM ap_outbox WHERE site_slug = ? ORDER BY created_at DESC').all(siteSlug);
715}
716
717// Delete one of our outbound replies: send Delete(Tombstone) to recipients + remove it.
718export async function deliverOutboxDelete(site, outboxId) {
719 const row = iStmts().getO.get(outboxId);
720 if (!row || row.site_slug !== site.slug) return false;
721 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
722 if (base) {
723 const me = actorId(base, site.slug);
724 const nid = noteId(base, row.id);
725 const del = { '@context': 'https://www.w3.org/ns/activitystreams', id: `${nid}#delete-${Date.now()}`, type: 'Delete', actor: me, to: [PUBLIC], object: { id: nid, type: 'Tombstone' } };
726 const keys = getOrCreateKeys(site.slug);
727 const inboxes = new Set();
728 if (row.to_actor) { const a = await fetchActor(row.to_actor).catch(() => null); if (a) inboxes.add((a.endpoints && a.endpoints.sharedInbox) || a.inbox); }
729 for (const f of fStmts().list.all(site.slug)) inboxes.add(f.shared_inbox || f.inbox);
730 for (const inbox of [...inboxes].filter(Boolean)) { try { await deliver(inbox, del, `${me}#main-key`, keys.private_pem); } catch { /* best-effort */ } }
731 }
732 db.prepare('DELETE FROM ap_outbox WHERE id = ?').run(outboxId);
733 return true;
734}
735
736// ── Fediverse CLIENT: follow accounts + home timeline ─────────────
737// Resolve an @user@domain handle to its actor URL via WebFinger.
738export async function webfingerResolve(handle) {
739 const h = String(handle || '').trim().replace(/^@/, '');
740 const parts = h.split('@');
741 if (parts.length !== 2 || !parts[0] || !parts[1]) return null;
742 const acct = `${parts[0]}@${parts[1]}`;
743 try {
744 const r = await fetch(`https://${parts[1]}/.well-known/webfinger?resource=acct:${encodeURIComponent(acct)}`,
745 { headers: { Accept: 'application/jrd+json, application/json' }, redirect: 'follow', signal: AbortSignal.timeout(8000) });
746 if (!r.ok) return null;
747 const jrd = await r.json();
748 const link = (jrd.links || []).find((l) => l.rel === 'self' && /activity\+json|ld\+json/.test(l.type || ''));
749 return link ? link.href : null;
750 } catch { return null; }
751}
752
753let _insFw, _delFw, _listFw, _accFw, _oneFw;
754function fwStmts() {
755 if (!_insFw) {
756 _insFw = db.prepare('INSERT OR REPLACE INTO ap_following (slug, actor_uri, handle, name, icon, url, inbox, follow_id, status, created_at) VALUES (?,?,?,?,?,?,?,?,?,CURRENT_TIMESTAMP)');
757 _delFw = db.prepare('DELETE FROM ap_following WHERE slug = ? AND actor_uri = ?');
758 _listFw = db.prepare('SELECT * FROM ap_following WHERE slug = ? ORDER BY created_at DESC');
759 _accFw = db.prepare("UPDATE ap_following SET status = 'accepted' WHERE follow_id = ?");
760 _oneFw = db.prepare('SELECT * FROM ap_following WHERE slug = ? AND actor_uri = ?');
761 }
762 return { ins: _insFw, del: _delFw, list: _listFw, acc: _accFw, one: _oneFw };
763}
764export function listFollowing(slug) { return fwStmts().list.all(slug); }
765
766let _insTl, _listTl, _delTl;
767function tlStmts() {
768 if (!_insTl) {
769 _insTl = db.prepare('INSERT OR IGNORE INTO ap_timeline (id, slug, author_uri, author_name, author_handle, author_icon, author_url, content, url, published, media_json, created_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,CURRENT_TIMESTAMP)');
770 _listTl = db.prepare('SELECT * FROM ap_timeline WHERE slug = ? ORDER BY COALESCE(published, created_at) DESC LIMIT ?');
771 _delTl = db.prepare('DELETE FROM ap_timeline WHERE id = ?');
772 }
773 return { ins: _insTl, list: _listTl, del: _delTl };
774}
775export function getTimeline(slug, limit) { return tlStmts().list.all(slug, limit || 50); }
776
777// Follow a fediverse account by @handle (WebFinger → actor → signed Follow).
778export async function followActor(site, handle) {
779 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
780 if (!base || !site || !site.slug) return { error: 'config' };
781 const actorUrl = await webfingerResolve(handle);
782 if (!actorUrl) return { error: 'not_found' };
783 const actor = await fetchActor(actorUrl).catch(() => null);
784 if (!actor || !actor.id || !actor.inbox) return { error: 'unreachable' };
785 const ai = actorInfo(actor, actor.id);
786 const me = actorId(base, site.slug);
787 const keys = getOrCreateKeys(site.slug);
788 const followId = `${me}#follow-${Date.now()}`;
789 fwStmts().ins.run(site.slug, actor.id, ai.handle, ai.name, ai.icon, ai.url, actor.inbox, followId, 'pending');
790 const follow = { '@context': 'https://www.w3.org/ns/activitystreams', id: followId, type: 'Follow', actor: me, object: actor.id };
791 try { await deliver(actor.inbox, follow, `${me}#main-key`, keys.private_pem); }
792 catch (e) { console.warn('[AP] follow deliver failed:', e.message); }
793 console.log('[AP] follow', site.slug, '→', actor.id);
794 return { ok: true, name: ai.name, handle: ai.handle };
795}
796
797export async function unfollowActor(site, actorUri) {
798 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
799 const me = actorId(base, site.slug);
800 const keys = getOrCreateKeys(site.slug);
801 const row = fwStmts().one.get(site.slug, actorUri);
802 if (row && row.inbox) {
803 const undo = { '@context': 'https://www.w3.org/ns/activitystreams', id: `${me}#unfollow-${Date.now()}`, type: 'Undo', actor: me, object: { id: row.follow_id || `${me}#follow`, type: 'Follow', actor: me, object: actorUri } };
804 try { await deliver(row.inbox, undo, `${me}#main-key`, keys.private_pem); } catch { /* best-effort */ }
805 }
806 fwStmts().del.run(site.slug, actorUri);
807 return { ok: true };
808}
809
810// Send a Like or Announce (boost) on a remote note FROM this site.
811export async function sendInteraction(site, kind, targetNoteId, authorUri) {
812 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
813 if (!base || !site || !site.slug || !targetNoteId) return { error: 'config' };
814 const type = kind === 'boost' ? 'Announce' : 'Like';
815 const me = actorId(base, site.slug);
816 const keys = getOrCreateKeys(site.slug);
817 const act = {
818 '@context': 'https://www.w3.org/ns/activitystreams',
819 id: `${me}#${type.toLowerCase()}-${Date.now()}`,
820 type, actor: me, object: targetNoteId,
821 };
822 if (type === 'Announce') { act.to = [PUBLIC]; act.cc = [`${me}/followers`]; }
823 const inboxes = new Set();
824 if (authorUri) { const a = await fetchActor(authorUri).catch(() => null); if (a) inboxes.add((a.endpoints && a.endpoints.sharedInbox) || a.inbox); }
825 // A boost is public → also deliver to our own followers so it shows for them.
826 if (type === 'Announce') { for (const f of fStmts().list.all(site.slug)) inboxes.add(f.shared_inbox || f.inbox); }
827 let delivered = 0;
828 for (const inbox of [...inboxes].filter(Boolean)) { try { const st = await deliver(inbox, act, `${me}#main-key`, keys.private_pem); if (st >= 200 && st < 300) delivered++; } catch { /* best-effort */ } }
829 console.log('[AP]', type, site.slug, '→', targetNoteId, 'delivered', delivered);
830 return { ok: true, delivered };
831}
832
833// Notifications inbox: new followers + replies/likes/boosts on this site's posts.
834export function getNotifications(slug, limit) {
835 const out = [];
836 try {
837 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)) {
838 out.push({ type: 'follow', handle: deriveHandle(f.actor_uri), url: f.actor_uri, created_at: f.created_at });
839 }
840 } catch { /* ignore */ }
841 try {
842 const rows = db.prepare(`
843 SELECT i.kind, i.actor_name, i.actor_handle, i.actor_url, i.content, i.created_at,
844 p.slug AS post_slug, p.title AS post_title
845 FROM ap_interactions i LEFT JOIN posts p ON p.id = i.post_id
846 WHERE p.site_id = (SELECT id FROM sites WHERE slug = ?)
847 ORDER BY i.created_at DESC LIMIT 80
848 `).all(slug);
849 for (const r of rows) out.push({
850 type: r.kind, name: r.actor_name, handle: r.actor_handle, url: r.actor_url,
851 content: r.content, post_slug: r.post_slug, post_title: r.post_title, created_at: r.created_at,
852 });
853 } catch { /* ignore */ }
854 out.sort((a, b) => new Date(b.created_at) - new Date(a.created_at));
855 return out.slice(0, limit || 60);
856}
857
858// ── Blocking / defederation ───────────────────────────────────────
859let _insBl, _delBl, _listBl;
860function blStmts() {
861 if (!_insBl) {
862 _insBl = db.prepare('INSERT OR IGNORE INTO ap_blocks (slug, target, kind, label, created_at) VALUES (?,?,?,?,CURRENT_TIMESTAMP)');
863 _delBl = db.prepare('DELETE FROM ap_blocks WHERE slug = ? AND target = ?');
864 _listBl = db.prepare('SELECT * FROM ap_blocks WHERE slug = ? ORDER BY created_at DESC');
865 }
866 return { ins: _insBl, del: _delBl, list: _listBl };
867}
868export function listBlocks(slug) { return blStmts().list.all(slug); }
869
870// True if an actor (or its whole domain) is blocked anywhere on this instance.
871export function isBlockedAny(actorUri) {
872 if (!actorUri) return false;
873 let domain = ''; try { domain = new URL(actorUri).host; } catch { /* ignore */ }
874 try { return !!db.prepare("SELECT 1 FROM ap_blocks WHERE (kind='actor' AND target=?) OR (kind='domain' AND target=?) LIMIT 1").get(actorUri, domain); }
875 catch { return false; }
876}
877
878function purgeBlocked(kind, target) {
879 try {
880 if (kind === 'domain') {
881 const like = `%//${target}/%`;
882 db.prepare('DELETE FROM ap_interactions WHERE actor_uri LIKE ?').run(like);
883 db.prepare('DELETE FROM ap_timeline WHERE author_uri LIKE ?').run(like);
884 db.prepare('DELETE FROM ap_followers WHERE actor_uri LIKE ?').run(like);
885 } else {
886 db.prepare('DELETE FROM ap_interactions WHERE actor_uri = ?').run(target);
887 db.prepare('DELETE FROM ap_timeline WHERE author_uri = ?').run(target);
888 db.prepare('DELETE FROM ap_followers WHERE actor_uri = ?').run(target);
889 }
890 } catch { /* best-effort */ }
891}
892
893// Block an actor (@handle or actor URL) or a whole domain; purges their content.
894export async function blockTarget(site, input) {
895 const raw = String(input || '').trim();
896 if (!site || !site.slug || !raw) return { error: 'empty' };
897 let kind, target, label;
898 if (/^https?:\/\//i.test(raw)) { kind = 'actor'; target = raw; label = raw; }
899 else if (raw.includes('@')) {
900 const actorUrl = await webfingerResolve(raw);
901 if (!actorUrl) return { error: 'not_found' };
902 kind = 'actor'; target = actorUrl; label = raw.startsWith('@') ? raw : ('@' + raw);
903 } else { kind = 'domain'; target = raw.toLowerCase(); label = raw.toLowerCase(); }
904 blStmts().ins.run(site.slug, target, kind, label);
905 purgeBlocked(kind, target);
906 console.log('[AP] block', site.slug, kind, target);
907 return { ok: true, label };
908}
909
910export function unblock(site, target) { blStmts().del.run(site.slug, target); return { ok: true }; }
911
912export default {
913 getOrCreateKeys, apWants, sendAP, actorId, noteId,
914 buildActor, buildNote, buildCreate, buildOutbox, buildFollowers,
915 followerCount, deliver, fetchActor, verifyRequest, handleInbox, deliverCreate, deliverDelete,
916 getInteractions, getInteractionById, buildReplyNote, getOutboxNote, deliverReply, resolveRemoteNote,
917 listOutbox, deliverOutboxDelete,
918 webfingerResolve, followActor, unfollowActor, listFollowing, getTimeline, sendInteraction,
919 getNotifications, listBlocks, isBlockedAny, blockTarget, unblock,
920 deliverWithRetry, enqueueDelivery, processDeliveryQueue, startDeliveryWorker,
921 getReplyUris,
922};
Note: See TracBrowser for help on using the repository browser.