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

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

feat(security): enforce HTTP signatures on the inbox

Data-affecting activities (Create/Like/Announce/Follow/Delete/Undo/Accept/Reject)
must now carry a valid HTTP signature whose signer matches the claimed actor —
else 401. Stops forged replies/likes/follows/timeline posts. Discovery (GET) open.

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

  • Property mode set to 100644
File size: 39.5 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 // Strip Klonkt audio shortcodes ([[track:…]] etc.) — they'd federate raw as
124 // ugly text (audio federation itself is a later phase).
125 body = body.replace(/\[\[(track|album|playlist):[^\]]+\]\]/gi, '');
126 const seen = new Set();
127 const attachment = urls.filter(Boolean)
128 .filter((u) => { if (seen.has(u)) return false; seen.add(u); return true; })
129 .map((u) => ({ type: 'Document', mediaType: mediaType(u), url: u }));
130
131 const note = {
132 id,
133 type: 'Note',
134 attributedTo: aId,
135 content: titleHtml + body,
136 url: human,
137 published: new Date(post.published_at || post.created_at || Date.now()).toISOString(),
138 to: [PUBLIC],
139 cc: [`${aId}/followers`],
140 tag: Array.isArray(post.tags) ? post.tags.map((t) => ({ type: 'Hashtag', name: '#' + String(t).replace(/\s+/g, '') })) : [],
141 };
142 if (attachment.length) note.attachment = attachment;
143 return note;
144}
145
146export function buildCreate(base, site, post) {
147 const note = buildNote(base, site, post);
148 return {
149 '@context': 'https://www.w3.org/ns/activitystreams',
150 id: note.id + '#create',
151 type: 'Create',
152 actor: actorId(base, site.slug),
153 published: note.published,
154 to: note.to,
155 cc: note.cc,
156 object: note,
157 };
158}
159
160export function buildOutbox(base, site, posts) {
161 const id = `${actorId(base, site.slug)}/outbox`;
162 const items = (posts || []).slice(0, MAX_OUTBOX).map((p) => buildCreate(base, site, p));
163 return {
164 '@context': 'https://www.w3.org/ns/activitystreams',
165 id,
166 type: 'OrderedCollection',
167 totalItems: items.length,
168 orderedItems: items,
169 };
170}
171
172export function buildFollowers(base, site, count) {
173 const id = `${actorId(base, site.slug)}/followers`;
174 return {
175 '@context': 'https://www.w3.org/ns/activitystreams',
176 id,
177 type: 'OrderedCollection',
178 totalItems: count || 0,
179 orderedItems: [], // hidden for privacy; count only
180 };
181}
182
183// ── followers store (lazy stmts) ──────────────────────────────────
184let _insF, _delF, _listF, _cntF;
185function fStmts() {
186 if (!_insF) {
187 _insF = db.prepare('INSERT OR IGNORE INTO ap_followers (slug, actor_uri, inbox, shared_inbox, created_at) VALUES (?,?,?,?,CURRENT_TIMESTAMP)');
188 _delF = db.prepare('DELETE FROM ap_followers WHERE slug = ? AND actor_uri = ?');
189 _listF = db.prepare('SELECT inbox, shared_inbox FROM ap_followers WHERE slug = ?');
190 _cntF = db.prepare('SELECT COUNT(*) n FROM ap_followers WHERE slug = ?');
191 }
192 return { ins: _insF, del: _delF, list: _listF, cnt: _cntF };
193}
194export function followerCount(slug) { return fStmts().cnt.get(slug).n; }
195
196// ── inbound interactions store (replies / likes / boosts) + our outbound replies ──
197let _insI, _delLA, _delReply, _listI, _getI, _insO, _listO, _getO;
198function iStmts() {
199 if (!_insI) {
200 _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)');
201 _delLA = db.prepare('DELETE FROM ap_interactions WHERE kind = ? AND post_id = ? AND actor_uri = ?');
202 _delReply = db.prepare("DELETE FROM ap_interactions WHERE kind = 'reply' AND object_uri = ?");
203 _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');
204 _getI = db.prepare('SELECT * FROM ap_interactions WHERE id = ?');
205 _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)');
206 _listO = db.prepare('SELECT * FROM ap_outbox WHERE post_id = ? ORDER BY created_at ASC');
207 _getO = db.prepare('SELECT * FROM ap_outbox WHERE id = ?');
208 }
209 return { ins: _insI, delLA: _delLA, delReply: _delReply, list: _listI, getI: _getI, insO: _insO, listO: _listO, getO: _getO };
210}
211
212export function getInteractionById(id) { return iStmts().getI.get(id); }
213
214const localPostExists = (id) => { try { return !!db.prepare('SELECT 1 FROM posts WHERE id = ?').get(id); } catch { return false; } };
215// Extract our local post id from a note URL, but only if it's ours (base match).
216function postIdFromNoteUrl(url, base) {
217 const s = String(url || '');
218 if (base && !s.startsWith(base)) return null;
219 const m = s.match(/\/ap\/notes\/([^/?#]+)/);
220 return m ? decodeURIComponent(m[1]) : null;
221}
222function deriveHandle(actorUri) {
223 try { const u = new URL(actorUri); const seg = u.pathname.split('/').filter(Boolean).pop() || ''; return `@${seg}@${u.host}`; } catch { return String(actorUri || ''); }
224}
225function actorInfo(doc, actorUri) {
226 let host = ''; try { host = new URL(actorUri).host; } catch { /* keep empty */ }
227 const handle = doc && doc.preferredUsername ? `@${doc.preferredUsername}@${host}` : deriveHandle(actorUri);
228 const icon = doc && doc.icon ? (doc.icon.url || (Array.isArray(doc.icon) && doc.icon[0] && doc.icon[0].url)) : null;
229 return {
230 name: (doc && (doc.name || doc.preferredUsername)) || handle,
231 handle,
232 url: (doc && (doc.url || doc.id)) || actorUri,
233 icon: icon || null,
234 };
235}
236
237// Given an inReplyTo note URL, find which local post the thread belongs to + the
238// note being replied to (parent), so a reply-to-a-comment can be nested.
239function findThreadTarget(inReplyTo, base) {
240 if (!inReplyTo) return null;
241 const seg = postIdFromNoteUrl(inReplyTo, base); // our /ap/notes/<id> segment (if ours)
242 if (seg && localPostExists(seg)) return { post_id: seg, parent_uri: inReplyTo };
243 if (seg) {
244 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 */ }
245 }
246 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 */ }
247 return null;
248}
249
250// View-ready threaded view of a post's fediverse activity (inbound replies +
251// our outbound replies, nested), plus like/boost counts.
252export function getInteractions(postId, base, site) {
253 const s = iStmts();
254 const rows = s.list.all(postId);
255 const baseClean = (base || process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
256 const postNoteId = baseClean ? `${baseClean}/ap/notes/${postId}` : null;
257 // Our own (outbound) replies show the SITE identity for everyone (not "You").
258 let host = ''; try { host = new URL(baseClean).host; } catch { /* ignore */ }
259 const siteName = (site && (site.title || site.slug)) || '';
260 const siteHandle = (site && site.slug && host) ? `@${site.slug}@${host}` : '';
261 const siteUrl = baseClean ? `${baseClean}/` : '';
262 const siteIcon = (site && site.profile_photo) || null;
263
264 const nodes = [];
265 for (const r of rows) {
266 if (r.kind !== 'reply') continue;
267 nodes.push({
268 noteId: r.object_uri, parent: r.parent_uri || null, mine: false,
269 actor_name: r.actor_name, actor_handle: r.actor_handle, actor_url: r.actor_url,
270 actor_icon: r.actor_icon, content: r.content, created_at: r.published || r.created_at,
271 children: [],
272 });
273 }
274 for (const o of s.listO.all(postId)) {
275 nodes.push({
276 noteId: baseClean ? `${baseClean}/ap/notes/${o.id}` : o.id, parent: o.in_reply_to || null,
277 mine: true, outboxId: o.id, content: o.content, created_at: o.created_at,
278 actor_name: siteName, actor_handle: siteHandle, actor_url: siteUrl, actor_icon: siteIcon,
279 children: [],
280 });
281 }
282
283 const byId = new Map(nodes.map((n) => [n.noteId, n]));
284 const isTop = (n) => !n.parent || n.parent === postNoteId || !byId.has(n.parent);
285 const tops = [];
286 for (const n of nodes) {
287 if (isTop(n)) { tops.push(n); continue; }
288 let anc = n, guard = 0;
289 while (!isTop(anc) && guard++ < 12) anc = byId.get(anc.parent);
290 anc.children.push(n);
291 }
292 const byTime = (a, b) => new Date(a.created_at) - new Date(b.created_at);
293 tops.sort(byTime).forEach((t) => t.children.sort(byTime));
294
295 return {
296 thread: tops,
297 likeCount: rows.filter((r) => r.kind === 'like').length,
298 announceCount: rows.filter((r) => r.kind === 'announce').length,
299 total: nodes.length,
300 };
301}
302
303// ── HTTP Signatures + delivery ────────────────────────────────────
304const slugFromActorUrl = (url) => { const m = String(url || '').match(/\/ap\/users\/([^/?#]+)/); return m ? decodeURIComponent(m[1]) : null; };
305
306// Sign + POST an activity to a remote inbox (draft-cavage HTTP Signatures, RSA-SHA256).
307export async function deliver(inboxUrl, bodyObj, keyId, privatePem) {
308 const body = JSON.stringify(bodyObj);
309 const u = new URL(inboxUrl);
310 const date = new Date().toUTCString();
311 const digest = 'SHA-256=' + crypto.createHash('sha256').update(body).digest('base64');
312 const signingString = `(request-target): post ${u.pathname}\nhost: ${u.host}\ndate: ${date}\ndigest: ${digest}`;
313 const signature = crypto.sign('sha256', Buffer.from(signingString), privatePem).toString('base64');
314 const sig = `keyId="${keyId}",algorithm="rsa-sha256",headers="(request-target) host date digest",signature="${signature}"`;
315 const r = await fetch(inboxUrl, {
316 method: 'POST',
317 headers: { 'Content-Type': 'application/activity+json', Accept: 'application/activity+json', Date: date, Digest: digest, Signature: sig },
318 body,
319 signal: AbortSignal.timeout(8000),
320 });
321 return r.status;
322}
323
324export async function fetchActor(url) {
325 try {
326 const r = await fetch(url, { headers: { Accept: 'application/activity+json' }, redirect: 'follow', signal: AbortSignal.timeout(8000) });
327 if (!r.ok) return null;
328 return await r.json();
329 } catch { return null; }
330}
331
332// Best-effort verification of an incoming signed request. Returns the sender's
333// actor doc if the signature checks out, else null. (Not gating yet — MVP.)
334export async function verifyRequest(req) {
335 const sigH = req.headers['signature'];
336 if (!sigH) return null;
337 const p = Object.fromEntries([...sigH.matchAll(/([a-zA-Z]+)="([^"]*)"/g)].map((m) => [m[1], m[2]]));
338 if (!p.keyId || !p.signature) return null;
339 const actor = await fetchActor(p.keyId.split('#')[0]);
340 const pem = actor && actor.publicKey && actor.publicKey.publicKeyPem;
341 if (!pem) return null;
342 const hs = (p.headers || '(request-target) host date').split(/\s+/);
343 const line = hs.map((h) => h === '(request-target)'
344 ? `(request-target): ${req.method.toLowerCase()} ${req.originalUrl}`
345 : `${h}: ${req.headers[h] || ''}`).join('\n');
346 let ok = false;
347 try { ok = crypto.verify('sha256', Buffer.from(line), pem, Buffer.from(p.signature, 'base64')); } catch { ok = false; }
348 if (ok && hs.includes('digest') && req.rawBody) {
349 const exp = 'SHA-256=' + crypto.createHash('sha256').update(req.rawBody).digest('base64');
350 if (req.headers['digest'] !== exp) ok = false;
351 }
352 return ok ? actor : null;
353}
354
355// Handle an incoming inbox POST. slugParam = null for the shared /ap/inbox.
356export async function handleInbox(req, slugParam) {
357 const act = req.body || {};
358 const type = act.type;
359 const base = (process.env.PUBLIC_BASE_URL || `${req.protocol}://${req.get('host')}`).replace(/\/+$/, '');
360 const verified = await verifyRequest(req).catch(() => null);
361
362 // ENFORCE HTTP signatures: a data-affecting activity must be signed by the very
363 // actor it claims to be. No valid signature, or signer ≠ actor → reject (no
364 // forged replies/likes/follows/timeline posts). GET/discovery stays open.
365 const claimedActor = typeof act.actor === 'string' ? act.actor : (act.actor && act.actor.id);
366 const GATED = ['Create', 'Like', 'Announce', 'Follow', 'Delete', 'Undo', 'Accept', 'Reject'];
367 if (GATED.includes(type)) {
368 if (!verified || !claimedActor || verified.id !== claimedActor) {
369 console.warn('[AP] inbox REJECTED (signature)', type, claimedActor || '?', verified ? '(signer mismatch)' : '(unsigned/invalid)');
370 return 401;
371 }
372 }
373
374 if (type === 'Follow') {
375 const who = typeof act.actor === 'string' ? act.actor : (act.actor && act.actor.id);
376 const slug = slugParam || slugFromActorUrl(typeof act.object === 'string' ? act.object : (act.object && act.object.id));
377 if (!who || !slug) return 400;
378 const remote = await fetchActor(who);
379 if (!remote || !remote.inbox) return 202; // can't reach them → drop quietly
380 fStmts().ins.run(slug, who, remote.inbox, (remote.endpoints && remote.endpoints.sharedInbox) || null);
381 const me = actorId(base, slug);
382 const keys = getOrCreateKeys(slug);
383 const accept = { '@context': 'https://www.w3.org/ns/activitystreams', id: `${me}#accept-${Date.now()}`, type: 'Accept', actor: me, object: act };
384 deliver(remote.inbox, accept, `${me}#main-key`, keys.private_pem).catch((e) => console.warn('[AP] Accept delivery failed:', e.message));
385 console.log('[AP] Follow', who, '→', slug, verified ? '(sig ok)' : '(sig unverified)');
386 return 202;
387 }
388 if (type === 'Undo' && act.object) {
389 const who = typeof act.actor === 'string' ? act.actor : (act.actor && act.actor.id);
390 const ot = act.object.type;
391 if (ot === 'Follow') {
392 const obj = act.object.object;
393 const slug = slugParam || slugFromActorUrl(typeof obj === 'string' ? obj : (obj && obj.id));
394 if (who && slug) { fStmts().del.run(slug, who); console.log('[AP] Unfollow', who, '→', slug); }
395 return 202;
396 }
397 if (ot === 'Like' || ot === 'Announce') {
398 const tgt = act.object.object;
399 const pid = postIdFromNoteUrl(typeof tgt === 'string' ? tgt : (tgt && tgt.id), base);
400 if (who && pid) { iStmts().delLA.run(ot.toLowerCase(), pid, who); console.log('[AP] Undo', ot, who, '→', pid); }
401 return 202;
402 }
403 return 202;
404 }
405
406 const actorUri = typeof act.actor === 'string' ? act.actor : (act.actor && act.actor.id);
407 const resolveActor = async (uri) => ((verified && verified.id === uri) ? verified : await fetchActor(uri).catch(() => null));
408 // Activities from our OWN actors are already stored via ap_outbox — don't re-store.
409 const isLocalActor = !!(base && actorUri && actorUri.startsWith(`${base}/ap/users/`));
410
411 // Inbound reply: a Create whose object replies to one of our notes (post OR comment).
412 if (type === 'Create' && act.object && (act.object.type === 'Note' || act.object.type === 'Article')) {
413 const o = act.object;
414 const tgt = findThreadTarget(o.inReplyTo, base);
415 if (tgt && actorUri && !isLocalActor) {
416 const ai = actorInfo(await resolveActor(actorUri), actorUri);
417 const html = HtmlSanitizerService.sanitize(o.content || '');
418 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);
419 console.log('[AP] reply', actorUri, '→', tgt.post_id);
420 return 202;
421 }
422 // Home timeline (client): a top-level post from an account we follow.
423 if (actorUri && !isLocalActor && !o.inReplyTo && o.id) {
424 let subs = []; try { subs = db.prepare('SELECT slug FROM ap_following WHERE actor_uri = ?').all(actorUri); } catch { /* table may not exist yet */ }
425 if (subs.length) {
426 const ai = actorInfo(await resolveActor(actorUri), actorUri);
427 const html = HtmlSanitizerService.sanitize(o.content || '');
428 const media = JSON.stringify((Array.isArray(o.attachment) ? o.attachment : []).filter((a) => a && a.url).map((a) => ({ url: a.url, type: a.mediaType || '' })));
429 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);
430 console.log('[AP] timeline +', actorUri, 'x' + subs.length);
431 }
432 }
433 return 202;
434 }
435 if (type === 'Like' || type === 'Announce') {
436 const tgt = act.object;
437 const pid = postIdFromNoteUrl(typeof tgt === 'string' ? tgt : (tgt && tgt.id), base);
438 if (pid && actorUri && !isLocalActor && localPostExists(pid)) {
439 const ai = actorInfo(await resolveActor(actorUri), actorUri);
440 iStmts().ins.run(type.toLowerCase(), pid, '', actorUri, ai.name, ai.handle, ai.url, ai.icon, null, null, null);
441 console.log('[AP]', type === 'Like' ? 'like' : 'boost', actorUri, '→', pid);
442 }
443 return 202;
444 }
445 if (type === 'Delete') {
446 // A remote note was deleted upstream → drop it from replies AND the timeline.
447 const oid = typeof act.object === 'string' ? act.object : (act.object && act.object.id);
448 if (oid) { iStmts().delReply.run(oid); try { tlStmts().del.run(oid); } catch { /* ignore */ } }
449 return 202;
450 }
451 // Accept/Reject of a Follow WE sent (client side).
452 if (type === 'Accept' && act.object) {
453 const fid = typeof act.object === 'string' ? act.object : (act.object && act.object.id);
454 if (fid) { try { fwStmts().acc.run(fid); } catch { /* ignore */ } }
455 console.log('[AP] follow accepted', actorUri);
456 return 202;
457 }
458 if (type === 'Reject' && act.object) {
459 const who = actorUri;
460 if (who && slugParam) { try { fwStmts().del.run(slugParam, who); } catch { /* ignore */ } }
461 return 202;
462 }
463
464 console.log('[AP] inbox', type || 'unknown', '→', slugParam || 'shared', '(ignored)');
465 return 202;
466}
467
468// Deliver a new post as Create(Note) to all followers' inboxes (fire-and-forget).
469// Needs PUBLIC_BASE_URL (absolute URLs); no-op without followers or base.
470export async function deliverCreate(site, post) {
471 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
472 if (!base || !site || !site.slug) return;
473 const followers = fStmts().list.all(site.slug);
474 if (!followers.length) return;
475 const inboxes = [...new Set(followers.map((f) => f.shared_inbox || f.inbox).filter(Boolean))];
476 const keys = getOrCreateKeys(site.slug);
477 const keyId = `${actorId(base, site.slug)}#main-key`;
478 const create = buildCreate(base, site, post);
479 for (const inbox of inboxes) deliver(inbox, create, keyId, keys.private_pem).catch(() => { /* best-effort */ });
480}
481
482// Tell followers a post is gone (Delete + Tombstone) so it's removed from their feeds.
483export async function deliverDelete(site, post) {
484 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
485 if (!base || !site || !site.slug || !post || !post.id) return;
486 const followers = fStmts().list.all(site.slug);
487 if (!followers.length) return;
488 const inboxes = [...new Set(followers.map((f) => f.shared_inbox || f.inbox).filter(Boolean))];
489 const keys = getOrCreateKeys(site.slug);
490 const me = actorId(base, site.slug);
491 const nid = noteId(base, post.id);
492 const del = {
493 '@context': 'https://www.w3.org/ns/activitystreams',
494 id: `${nid}#delete-${Date.now()}`,
495 type: 'Delete',
496 actor: me,
497 to: [PUBLIC],
498 object: { id: nid, type: 'Tombstone' },
499 };
500 for (const inbox of inboxes) deliver(inbox, del, `${me}#main-key`, keys.private_pem).catch(() => { /* best-effort */ });
501}
502
503// ── outbound replies (Klonkt → fediverse) ─────────────────────────
504const escHtml = (s) => String(s || '').replace(/[<>&]/g, (c) => ({ '<': '&lt;', '>': '&gt;', '&': '&amp;' }[c]));
505const 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(); };
506
507// Build one of OUR outbound reply Notes from an ap_outbox row.
508export function buildReplyNote(base, site, row) {
509 const me = actorId(base, site.slug);
510 return {
511 id: noteId(base, row.id),
512 type: 'Note',
513 attributedTo: me,
514 inReplyTo: row.in_reply_to || undefined,
515 content: row.content,
516 url: row.post_slug ? `${base}/${encodeURIComponent(row.post_slug)}` : undefined,
517 published: toISO(row.created_at),
518 to: row.to_actor ? [row.to_actor] : [PUBLIC],
519 cc: [PUBLIC, `${me}/followers`],
520 tag: row.to_actor ? [{ type: 'Mention', href: row.to_actor, name: row.to_handle }] : [],
521 };
522}
523
524// Resolve one of our outbound reply Notes by id (for /ap/notes/:id fallback).
525export function getOutboxNote(base, id) {
526 const row = iStmts().getO.get(id);
527 if (!row) return null;
528 const site = db.prepare('SELECT * FROM sites WHERE slug = ?').get(row.site_slug);
529 if (!site) return null;
530 return buildReplyNote(base, site, row);
531}
532
533// Send a reply FROM this site to a remote actor (in reply to their inbound reply).
534// `parent` = an ap_interactions row (actor_uri, actor_url, actor_handle, object_uri).
535export async function deliverReply(site, { postId, postSlug, parent, text }) {
536 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
537 if (!base || !site || !site.slug || !parent || !String(text || '').trim()) return null;
538 const me = actorId(base, site.slug);
539 const handle = parent.actor_handle || deriveHandle(parent.actor_uri);
540 const body = escHtml(String(text).trim()).replace(/\r?\n/g, '<br>');
541 const mention = parent.actor_uri
542 ? `<a href="${escHtml(parent.actor_url || parent.actor_uri)}" class="u-url mention">${escHtml(handle)}</a> ` : '';
543 const content = `<p>${mention}${body}</p>`;
544 // Dedup: skip if the exact same reply was already sent (double-submit guard).
545 const dup = db.prepare('SELECT 1 FROM ap_outbox WHERE site_slug = ? AND IFNULL(in_reply_to, \'\') = ? AND content = ? LIMIT 1')
546 .get(site.slug, parent.object_uri || '', content);
547 if (dup) { console.log('[AP] outreply skipped (duplicate)'); return { duplicate: true, delivered: 0 }; }
548 const id = crypto.randomUUID();
549 iStmts().insO.run(id, site.slug, postId, postSlug || null, parent.object_uri || null, parent.actor_uri || null, handle, content);
550 const row = iStmts().getO.get(id);
551 const note = buildReplyNote(base, site, row);
552 const create = {
553 '@context': 'https://www.w3.org/ns/activitystreams',
554 id: note.id + '#create', type: 'Create', actor: me,
555 published: note.published, to: note.to, cc: note.cc, object: note,
556 };
557 const keys = getOrCreateKeys(site.slug);
558 const keyId = `${me}#main-key`;
559 const inboxes = new Set();
560 if (parent.actor_uri) {
561 const a = await fetchActor(parent.actor_uri).catch(() => null);
562 if (a) inboxes.add((a.endpoints && a.endpoints.sharedInbox) || a.inbox);
563 }
564 if (parent.threadInbox) inboxes.add(parent.threadInbox); // back-compat (single)
565 (parent.threadInboxes || []).forEach((i) => inboxes.add(i)); // whole ancestor chain
566 for (const f of fStmts().list.all(site.slug)) inboxes.add(f.shared_inbox || f.inbox);
567 inboxes.delete(`${me}/inbox`); // never deliver to ourselves (already in ap_outbox)
568 inboxes.delete(`${base}/ap/inbox`); // (our own shared inbox) → avoids a self-duplicate
569 let delivered = 0;
570 for (const inbox of [...inboxes].filter(Boolean)) {
571 try { const st = await deliver(inbox, create, keyId, keys.private_pem); if (st >= 200 && st < 300) delivered++; } catch { /* best-effort */ }
572 }
573 console.log('[AP] outreply', site.slug, '→', parent.actor_uri, 'delivered', delivered);
574 return { id, content, delivered };
575}
576
577// Resolve a remote post URL (any fediverse/Klonkt post) into a reply target.
578// Returns a parent-shaped object usable by deliverReply(), or null.
579export async function resolveRemoteNote(url) {
580 if (!/^https?:\/\//i.test(String(url || ''))) return null;
581 const note = await fetchActor(url).catch(() => null); // AP GET (content-negotiates)
582 if (!note || !note.id) return null;
583 const att = note.attributedTo;
584 const actorUri = typeof att === 'string' ? att : (att && att.id);
585 if (!actorUri) return null;
586 const actor = await fetchActor(actorUri).catch(() => null);
587 const ai = actorInfo(actor, actorUri);
588 // Is what we're replying to a post (or a comment) on one of OUR posts? If so,
589 // link our reply to that local post so it shows nested in the post thread.
590 const localTgt = findThreadTarget(note.id, (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, ''));
591 // Walk the WHOLE reply chain upward (comment → parent comment → … → root post)
592 // and collect every ancestor author's inbox, so each participant's server —
593 // including the original post's author — receives + threads our reply.
594 const threadInboxes = [];
595 const seenInbox = new Set();
596 let cursor = note.inReplyTo, guard = 0;
597 while (cursor && guard++ < 6) {
598 const url = typeof cursor === 'string' ? cursor : (cursor && cursor.id);
599 if (!url) break;
600 const pn = await fetchActor(url).catch(() => null);
601 if (!pn) break;
602 const pa = typeof pn.attributedTo === 'string' ? pn.attributedTo : (pn.attributedTo && pn.attributedTo.id);
603 if (pa && pa !== actorUri) {
604 const paDoc = await fetchActor(pa).catch(() => null);
605 const inbox = paDoc && ((paDoc.endpoints && paDoc.endpoints.sharedInbox) || paDoc.inbox);
606 if (inbox && !seenInbox.has(inbox)) { seenInbox.add(inbox); threadInboxes.push(inbox); }
607 }
608 cursor = pn.inReplyTo; // climb to the next ancestor
609 }
610 const rawHtml = String(note.content || '').replace(/\[\[(track|album|playlist):[^\]]+\]\]/gi, '');
611 const images = (Array.isArray(note.attachment) ? note.attachment : [])
612 .filter((a) => a && a.url && (!a.mediaType || /^image\//i.test(a.mediaType)))
613 .map((a) => a.url);
614 return {
615 object_uri: note.id,
616 actor_uri: actorUri,
617 actor_url: ai.url,
618 actor_handle: ai.handle,
619 actor_name: ai.name,
620 actor_icon: ai.icon,
621 url: note.url || url,
622 content: HtmlSanitizerService.sanitize(rawHtml), // full, sanitized
623 images,
624 threadInboxes, // every ancestor author's inbox
625 localPostId: localTgt ? localTgt.post_id : '', // our post this belongs to (if any)
626 preview: HtmlSanitizerService.toPlainText(note.content || '').slice(0, 240),
627 };
628}
629
630// List a site's own outbound fediverse replies (for the manage/delete view).
631export function listOutbox(siteSlug) {
632 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);
633}
634
635// Delete one of our outbound replies: send Delete(Tombstone) to recipients + remove it.
636export async function deliverOutboxDelete(site, outboxId) {
637 const row = iStmts().getO.get(outboxId);
638 if (!row || row.site_slug !== site.slug) return false;
639 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
640 if (base) {
641 const me = actorId(base, site.slug);
642 const nid = noteId(base, row.id);
643 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' } };
644 const keys = getOrCreateKeys(site.slug);
645 const inboxes = new Set();
646 if (row.to_actor) { const a = await fetchActor(row.to_actor).catch(() => null); if (a) inboxes.add((a.endpoints && a.endpoints.sharedInbox) || a.inbox); }
647 for (const f of fStmts().list.all(site.slug)) inboxes.add(f.shared_inbox || f.inbox);
648 for (const inbox of [...inboxes].filter(Boolean)) { try { await deliver(inbox, del, `${me}#main-key`, keys.private_pem); } catch { /* best-effort */ } }
649 }
650 db.prepare('DELETE FROM ap_outbox WHERE id = ?').run(outboxId);
651 return true;
652}
653
654// ── Fediverse CLIENT: follow accounts + home timeline ─────────────
655// Resolve an @user@domain handle to its actor URL via WebFinger.
656export async function webfingerResolve(handle) {
657 const h = String(handle || '').trim().replace(/^@/, '');
658 const parts = h.split('@');
659 if (parts.length !== 2 || !parts[0] || !parts[1]) return null;
660 const acct = `${parts[0]}@${parts[1]}`;
661 try {
662 const r = await fetch(`https://${parts[1]}/.well-known/webfinger?resource=acct:${encodeURIComponent(acct)}`,
663 { headers: { Accept: 'application/jrd+json, application/json' }, redirect: 'follow', signal: AbortSignal.timeout(8000) });
664 if (!r.ok) return null;
665 const jrd = await r.json();
666 const link = (jrd.links || []).find((l) => l.rel === 'self' && /activity\+json|ld\+json/.test(l.type || ''));
667 return link ? link.href : null;
668 } catch { return null; }
669}
670
671let _insFw, _delFw, _listFw, _accFw, _oneFw;
672function fwStmts() {
673 if (!_insFw) {
674 _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)');
675 _delFw = db.prepare('DELETE FROM ap_following WHERE slug = ? AND actor_uri = ?');
676 _listFw = db.prepare('SELECT * FROM ap_following WHERE slug = ? ORDER BY created_at DESC');
677 _accFw = db.prepare("UPDATE ap_following SET status = 'accepted' WHERE follow_id = ?");
678 _oneFw = db.prepare('SELECT * FROM ap_following WHERE slug = ? AND actor_uri = ?');
679 }
680 return { ins: _insFw, del: _delFw, list: _listFw, acc: _accFw, one: _oneFw };
681}
682export function listFollowing(slug) { return fwStmts().list.all(slug); }
683
684let _insTl, _listTl, _delTl;
685function tlStmts() {
686 if (!_insTl) {
687 _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)');
688 _listTl = db.prepare('SELECT * FROM ap_timeline WHERE slug = ? ORDER BY COALESCE(published, created_at) DESC LIMIT ?');
689 _delTl = db.prepare('DELETE FROM ap_timeline WHERE id = ?');
690 }
691 return { ins: _insTl, list: _listTl, del: _delTl };
692}
693export function getTimeline(slug, limit) { return tlStmts().list.all(slug, limit || 50); }
694
695// Follow a fediverse account by @handle (WebFinger → actor → signed Follow).
696export async function followActor(site, handle) {
697 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
698 if (!base || !site || !site.slug) return { error: 'config' };
699 const actorUrl = await webfingerResolve(handle);
700 if (!actorUrl) return { error: 'not_found' };
701 const actor = await fetchActor(actorUrl).catch(() => null);
702 if (!actor || !actor.id || !actor.inbox) return { error: 'unreachable' };
703 const ai = actorInfo(actor, actor.id);
704 const me = actorId(base, site.slug);
705 const keys = getOrCreateKeys(site.slug);
706 const followId = `${me}#follow-${Date.now()}`;
707 fwStmts().ins.run(site.slug, actor.id, ai.handle, ai.name, ai.icon, ai.url, actor.inbox, followId, 'pending');
708 const follow = { '@context': 'https://www.w3.org/ns/activitystreams', id: followId, type: 'Follow', actor: me, object: actor.id };
709 try { await deliver(actor.inbox, follow, `${me}#main-key`, keys.private_pem); }
710 catch (e) { console.warn('[AP] follow deliver failed:', e.message); }
711 console.log('[AP] follow', site.slug, '→', actor.id);
712 return { ok: true, name: ai.name, handle: ai.handle };
713}
714
715export async function unfollowActor(site, actorUri) {
716 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
717 const me = actorId(base, site.slug);
718 const keys = getOrCreateKeys(site.slug);
719 const row = fwStmts().one.get(site.slug, actorUri);
720 if (row && row.inbox) {
721 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 } };
722 try { await deliver(row.inbox, undo, `${me}#main-key`, keys.private_pem); } catch { /* best-effort */ }
723 }
724 fwStmts().del.run(site.slug, actorUri);
725 return { ok: true };
726}
727
728// Send a Like or Announce (boost) on a remote note FROM this site.
729export async function sendInteraction(site, kind, targetNoteId, authorUri) {
730 const base = (process.env.PUBLIC_BASE_URL || '').replace(/\/+$/, '');
731 if (!base || !site || !site.slug || !targetNoteId) return { error: 'config' };
732 const type = kind === 'boost' ? 'Announce' : 'Like';
733 const me = actorId(base, site.slug);
734 const keys = getOrCreateKeys(site.slug);
735 const act = {
736 '@context': 'https://www.w3.org/ns/activitystreams',
737 id: `${me}#${type.toLowerCase()}-${Date.now()}`,
738 type, actor: me, object: targetNoteId,
739 };
740 if (type === 'Announce') { act.to = [PUBLIC]; act.cc = [`${me}/followers`]; }
741 const inboxes = new Set();
742 if (authorUri) { const a = await fetchActor(authorUri).catch(() => null); if (a) inboxes.add((a.endpoints && a.endpoints.sharedInbox) || a.inbox); }
743 // A boost is public → also deliver to our own followers so it shows for them.
744 if (type === 'Announce') { for (const f of fStmts().list.all(site.slug)) inboxes.add(f.shared_inbox || f.inbox); }
745 let delivered = 0;
746 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 */ } }
747 console.log('[AP]', type, site.slug, '→', targetNoteId, 'delivered', delivered);
748 return { ok: true, delivered };
749}
750
751// Notifications inbox: new followers + replies/likes/boosts on this site's posts.
752export function getNotifications(slug, limit) {
753 const out = [];
754 try {
755 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)) {
756 out.push({ type: 'follow', handle: deriveHandle(f.actor_uri), url: f.actor_uri, created_at: f.created_at });
757 }
758 } catch { /* ignore */ }
759 try {
760 const rows = db.prepare(`
761 SELECT i.kind, i.actor_name, i.actor_handle, i.actor_url, i.content, i.created_at,
762 p.slug AS post_slug, p.title AS post_title
763 FROM ap_interactions i LEFT JOIN posts p ON p.id = i.post_id
764 WHERE p.site_id = (SELECT id FROM sites WHERE slug = ?)
765 ORDER BY i.created_at DESC LIMIT 80
766 `).all(slug);
767 for (const r of rows) out.push({
768 type: r.kind, name: r.actor_name, handle: r.actor_handle, url: r.actor_url,
769 content: r.content, post_slug: r.post_slug, post_title: r.post_title, created_at: r.created_at,
770 });
771 } catch { /* ignore */ }
772 out.sort((a, b) => new Date(b.created_at) - new Date(a.created_at));
773 return out.slice(0, limit || 60);
774}
775
776export default {
777 getOrCreateKeys, apWants, sendAP, actorId, noteId,
778 buildActor, buildNote, buildCreate, buildOutbox, buildFollowers,
779 followerCount, deliver, fetchActor, verifyRequest, handleInbox, deliverCreate, deliverDelete,
780 getInteractions, getInteractionById, buildReplyNote, getOutboxNote, deliverReply, resolveRemoteNote,
781 listOutbox, deliverOutboxDelete,
782 webfingerResolve, followActor, unfollowActor, listFollowing, getTimeline, sendInteraction,
783 getNotifications,
784};
Note: See TracBrowser for help on using the repository browser.