Changeset 4f322bc in Klonkt
- Timestamp:
- 08/07/2026 08:49:58 AM (5 weeks ago)
- Branches:
- main
- Children:
- a6d703a
- Parents:
- 5861849
- git-author:
- Robin <roboburr@…> (08/07/2026 08:49:56 AM)
- git-committer:
- roboburr <roboburr@…> (08/07/2026 08:49:58 AM)
- Files:
-
- 1 added
- 2 edited
-
src/routes/activitypub.js (modified) (2 diffs)
-
src/services/ActivityPubService.js (modified) (2 diffs)
-
test/feed-wait.test.js (added)
Legend:
- Unmodified
- Added
- Removed
-
src/routes/activitypub.js
r5861849 r4f322bc 268 268 // they follow) as Create(Note) items, so an app (Shaer) can build a unified 269 269 // feed. Anyone else gets 403; the inbox stays write-only for the public. 270 router.get('/ap/users/:slug/inbox', (req, res) => {270 router.get('/ap/users/:slug/inbox', async (req, res) => { 271 271 const auth = OAuth.verifyBearer(req.headers.authorization); 272 272 if (!auth || auth.site.slug !== req.params.slug) return res.status(403).end(); 273 273 const base = baseUrl(req); 274 // Wachten is een UITBREIDING van deze lezing, geen tweede endpoint (shaer-n05). 275 // Geef `since` (de shaer:cursor van je vorige antwoord) en `wait` mee, en het 276 // antwoord blijft hangen tot er iets is of de tijd om is. Zonder die twee 277 // gedraagt de route zich exact zoals altijd. 278 // 279 // Bewust hetzelfde antwoord in plaats van een "er is nieuws"-seintje: dan 280 // hoeft er niets nieuws geparsed te worden, is er geen tweede beschrijving van 281 // de kaartvorm die uit de pas kan lopen, en scheelt het de client een tweede 282 // ronde. 283 const wachtS = Math.min(Math.max(parseInt(req.query.wait, 10) || 0, 0), 50); 284 if (req.query.since && wachtS > 0) { 285 const afbreken = new AbortController(); 286 res.on('close', () => afbreken.abort()); // client hing op: niet doorgaan met wachten 287 await AP.waitForFeedChange(auth.site.slug, { 288 since: String(req.query.since), waitMs: wachtS * 1000, signal: afbreken.signal, 289 }); 290 if (res.writableEnded || afbreken.signal.aborted) return undefined; 291 } 274 292 // Gated feature (FEP-633c): may this account see EXTERNAL embeds? A ward's 275 293 // world outside the fediverse is the guardians' call. The gate is applied … … 444 462 'shaer:externalLinks': playbackAllowed, 445 463 }, 464 // Het merk van wat hierin zit. Geef hem terug als `since` om op het 465 // volgende te wachten. NA het samenstellen bepaald, zodat hij precies dekt 466 // wat je in handen hebt en niet iets dat er ondertussen bij kwam. 467 'shaer:cursor': AP.feedCursor(auth.site.slug), 446 468 totalItems: items.length, 447 469 orderedItems: items, 448 470 }); 471 return undefined; 449 472 }); 450 473 -
src/services/ActivityPubService.js
r5861849 r4f322bc 3485 3485 } 3486 3486 3487 /** 3488 * Een merk voor "is er iets veranderd aan wat de inbox-lezing zou opleveren?" 3489 * (shaer-n05). 3490 * 3491 * Alle VIER de poten die de inbox samenvoegt tellen mee -- tijdlijn, berichten, 3492 * antwoorden op je eigen posts, en wat je zelf verstuurde. Zou er een ontbreken, 3493 * dan blijft een wachtende client slapen terwijl er wel degelijk iets is 3494 * bijgekomen, en dat is erger dan niet wachten: het lijkt te werken. 3495 * 3496 * rowid en niet een tijdstempel: rowid loopt strikt op per invoeging, terwijl 3497 * twee dingen in dezelfde seconde kunnen aankomen en een `published` van een 3498 * andere server niet te vertrouwen is. 3499 * 3500 * Ondoorzichtig voor de client. Hij krijgt hem terug en geeft hem ongewijzigd 3501 * mee; de vorm mag veranderen zonder dat dat iets breekt. 3502 */ 3503 export function feedCursor(slug) { 3504 const max = (sql, ...args) => { try { const r = db.prepare(sql).get(...args); return (r && r.n) || 0; } catch { return 0; } }; 3505 const t = max('SELECT MAX(rowid) AS n FROM ap_timeline WHERE slug = ?', slug); 3506 const m = max('SELECT MAX(rowid) AS n FROM ap_mentions WHERE slug = ?', slug); 3507 const o = max('SELECT MAX(rowid) AS n FROM ap_outbox WHERE site_slug = ?', slug); 3508 const r = max(`SELECT MAX(i.rowid) AS n FROM ap_interactions i 3509 JOIN posts p ON p.id = i.post_id 3510 JOIN sites s ON s.id = p.site_id 3511 WHERE s.slug = ? AND i.kind = 'reply'`, slug); 3512 return `${t}.${m}.${o}.${r}`; 3513 } 3514 3515 // Zoveel clients mogen er tegelijk op EEN account staan wachten. Een client met 3516 // een kapotte herverbind-lus mag de instance niet vastzetten; de overtolligen 3517 // krijgen gewoon meteen antwoord in plaats van een fout. 3518 const FEED_WAIT_MAX = 4; 3519 const _wachters = new Map(); 3520 3521 /** 3522 * Wacht tot de inbox-lezing iets anders zou opleveren dan bij `since`. 3523 * 3524 * Bewust met een interne tik en niet met een gebeurtenis-emitter. Een emitter 3525 * moet op ELKE plek worden aangeroepen waar er iets bijkomt, en de plek die je 3526 * vergeet is precies de melding die nooit aankomt. Twee tot vier MAX(rowid)- 3527 * queries per seconde is niets, en dit kan niets missen. Prijs: hooguit een tik 3528 * vertraging. 3529 */ 3530 export async function waitForFeedChange(slug, opts = {}) { 3531 const tickMs = Math.max(50, opts.tickMs || 1000); 3532 const waitMs = Math.max(0, opts.waitMs || 0); 3533 const since = String(opts.since || ''); 3534 let cursor = feedCursor(slug); 3535 // Geen sinds, al iets veranderd, of niet willen wachten: meteen antwoorden. 3536 if (!since || since !== cursor || !waitMs) return { cursor, changed: !!since && since !== cursor, waited: false }; 3537 3538 const bezet = _wachters.get(slug) || 0; 3539 if (bezet >= FEED_WAIT_MAX) return { cursor, changed: false, waited: false, busy: true }; 3540 _wachters.set(slug, bezet + 1); 3541 try { 3542 const einde = Date.now() + waitMs; 3543 while (Date.now() < einde) { 3544 if (opts.signal && opts.signal.aborted) break; // client hing op 3545 const rest = Math.min(tickMs, einde - Date.now()); 3546 await new Promise((r) => setTimeout(r, rest)); 3547 cursor = feedCursor(slug); 3548 if (cursor !== since) return { cursor, changed: true, waited: true }; 3549 } 3550 return { cursor, changed: false, waited: true }; 3551 } finally { 3552 const n = (_wachters.get(slug) || 1) - 1; 3553 if (n > 0) _wachters.set(slug, n); else _wachters.delete(slug); 3554 } 3555 } 3556 3487 3557 export function getDirectMessages(slug, limit) { 3488 3558 try { … … 5254 5324 buildActor, buildNote, buildCreate, buildOutbox, buildFollowers, buildFollowing, buildFeatured, 5255 5325 followerCount, deliver, fetchActor, verifyRequest, handleInbox, deliverCreate, deliverDelete, deliverUpdate, deliverActorUpdate, resyncFeaturedPins, 5326 feedCursor, waitForFeedChange, 5256 5327 getInteractions, getInteractionById, setInteractionBoosted, setInteractionLiked, buildReplyNote, getOutboxNote, getSentNotes, deliverReply, resolveRemoteNote, noteAudience, mayReadNote, 5257 5328 listOutbox, deliverOutboxDelete, deliverOutboxUpdate, deliverDirectNote,
Note:
See TracChangeset
for help on using the changeset viewer.
![(please configure the [header_logo] section in trac.ini)](/chrome/site/your_project_logo.png)