diff --git a/blob-wanter.js b/blob-wanter.js new file mode 100755 index 0000000..d62c784 --- /dev/null +++ b/blob-wanter.js @@ -0,0 +1,214 @@ +#!/usr/bin/env node +// +// blob-wanter.js — issue standing blobs.want() calls for blobs referenced by +// the feeds this node replicates. +// +// WHY THIS EXISTS +// --------------- +// Nothing in ssb-server scans replicated messages and fetches the blobs they +// mention. `blobs.want` is called from exactly one place in the whole install +// (ssb-ws/blobs.js), and only in response to an inbound HTTP request. +// +// So how did friends' images ever arrive? Via ssb-blobs `sympathy` (default 3): +// when a peer publishes a blob it calls blobs.push(), broadcasting a pretend +// "want" at hop -1; we adopt that want out of sympathy and fetch. That is the +// SAME code path that let 83 strangers fill this disk with 31k blobs we never +// asked for. Turning sympathy off without a replacement would also stop our own +// friends' images arriving. +// +// This is the replacement, and it is strictly better for our situation: +// +// * want() has no timeout and no expiry (ssb-blobs inject.js:312), and +// createWantStream ships the whole want map to every peer on connect. +// So an explicit want is a STANDING ORDER: the moment a feed that has been +// dark for months reappears, our want goes out and the blob transfers. +// Sympathy only fires if that peer happens to re-announce. +// +// * It wants ONLY blobs our own feeds reference. Strangers get nothing. +// +// IMPORTANT: sbot's want map is in-memory (`var want = {}`), so it is lost on +// every sbot restart -- which on this box means every OOM kill. This process +// exits when its sbot connection drops and is restarted by its run loop, which +// re-runs the backfill and re-arms every want. The two lifetimes are deliberately +// tied together. +// +// Usage: +// ./blob-wanter.js # backfill, then follow the log live +// ./blob-wanter.js --dry-run # report what it would want, change nothing +// ./blob-wanter.js --once # backfill only, then exit + +var fs = require('fs') +var path = require('path') +var pull = require('pull-stream') +var ssbKeys = require('ssb-keys') + +var DRY_RUN = process.argv.indexOf('--dry-run') !== -1 +var ONCE = process.argv.indexOf('--once') !== -1 + +// sha256 of the empty string. ssb-blobs special-cases this and never stores it, +// so asking for it is pointless noise. +var EMPTY_HASH = '&47DEQpj8HBSa+/TImW+5JCeuQeRkm5NMpJWZG3hSuFU=.sha256' +var BLOB_RE = /&[A-Za-z0-9+/]{43}=\.sha256/g +// If we ever want more than this, something is wrong -- log loudly rather than +// silently opening tens of thousands of muxrpc requests on a 469MB box. +var SANITY_CAP = 2000 + +function log () { + var msg = '[blob-wanter] ' + new Date().toISOString() + ' ' + + Array.prototype.slice.call(arguments).join(' ') + console.log(msg) +} + +function bail (why) { + console.error('[blob-wanter] ' + new Date().toISOString() + ' ' + why + + ' — exiting so the run loop restarts us') + process.exit(1) +} + +// Pull every blob ref out of a message, including inside private messages we +// can unbox. Regex over the serialised content catches mentions, markdown +// image links, and ad-hoc fields alike. +function blobRefsIn (msg, keys) { + var content = msg && msg.value && msg.value.content + if (!content) return [] + + if (typeof content === 'string') { + if (!/\.box\d*$/.test(content)) return [] // encrypted but not for us + try { + var unboxed = ssbKeys.unbox(content, keys) + if (!unboxed) return [] + content = unboxed + } catch (e) { + return [] + } + } + + var found + try { + found = JSON.stringify(content).match(BLOB_RE) + } catch (e) { + return [] + } + return found || [] +} + +require('ssb-client')(function (err, sbot, config) { + if (err) return bail(String(err.message || err)) + + sbot.on('closed', function () { bail('sbot connection closed') }) + + var keys = ssbKeys.loadOrCreateSync(path.join(config.path, 'secret')) + + // `seen` must be marked SYNCHRONOUSLY on entry to consider(). Deduping inside + // the async has() callback would let the same ref spawn many in-flight has() + // calls, and would push duplicate callbacks into ssb-blobs' waiting[id] array + // (a real leak, since those callbacks are never freed until the blob arrives). + var seen = Object.create(null) + var wanted = Object.create(null) + var stats = { scanned: 0, refs: 0, already: 0, wanted: 0, arrived: 0 } + + // has() is async, so the backfill stream ends long before the lookups finish. + // Track them so --once/--dry-run can report real numbers instead of zeros. + var pending = 0 + var onDrained = null + function settle () { + if (pending === 0 && onDrained) { var f = onDrained; onDrained = null; f() } + } + + function consider (ref, why) { + if (ref === EMPTY_HASH) return + if (seen[ref]) return + seen[ref] = true + stats.refs++ + + pending++ + sbot.blobs.has(ref, function (err, has) { + pending-- + if (err) { log('has() failed for', ref, '-', err.message || err); return settle() } + if (has) { stats.already++; return settle() } + + if (Object.keys(wanted).length >= SANITY_CAP) { + log('WARN hit sanity cap of', SANITY_CAP, 'outstanding wants — not adding', ref) + return settle() + } + + wanted[ref] = true + stats.wanted++ + + if (DRY_RUN) { + log('DRY-RUN would want', ref, '(' + why + ')') + return settle() + } + + log('want', ref, '(' + why + ')') + // The callback fires only if/when the blob actually arrives. If the peer + // holding it never appears it simply never fires -- that is the standing + // order working as intended, not a leak. + sbot.blobs.want(ref, function (err) { + if (err) return log('want failed for', ref, '-', err.message || err) + stats.arrived++ + log('ARRIVED', ref) + }) + settle() + }) + } + + function handle (msg, why) { + if (!msg || !msg.value) return + stats.scanned++ + blobRefsIn(msg, keys).forEach(function (ref) { consider(ref, why) }) + } + + log('connected to sbot' + (DRY_RUN ? ' (DRY RUN — nothing will be wanted)' : '')) + log('backfilling from the local log…') + + pull( + sbot.createLogStream({ keys: true, values: true }), + pull.drain( + function (msg) { handle(msg, 'backfill') }, + function (err) { + if (err) return bail('backfill stream failed: ' + (err.message || err)) + + // Wait for the in-flight has() lookups before reporting or exiting. + onDrained = function () { + log('backfill done —', + 'messages=' + stats.scanned, + 'unique_blobrefs=' + stats.refs, + 'already_have=' + stats.already, + 'newly_wanted=' + stats.wanted) + + if (ONCE || DRY_RUN) { + log('exiting (' + (DRY_RUN ? '--dry-run' : '--once') + ')') + return process.exit(0) + } + startLive() + } + settle() + } + ) + ) + + function startLive () { + log('following the log live; standing wants re-arm on every restart') + pull( + sbot.createLogStream({ keys: true, values: true, live: true, old: false }), + pull.drain( + function (msg) { handle(msg, 'live') }, + function (err) { + bail('live stream ended' + (err ? ': ' + (err.message || err) : '')) + } + ) + ) + + // Periodic heartbeat so the log shows the process is alive and what it is + // holding, without needing to attach to the screen. + setInterval(function () { + log('stats —', + 'msgs=' + stats.scanned, + 'unique_refs=' + stats.refs, + 'have=' + stats.already, + 'standing_wants=' + stats.wanted, + 'arrived=' + stats.arrived) + }, 3600000).unref() + } +}) diff --git a/run-blob-wanter.sh b/run-blob-wanter.sh new file mode 100755 index 0000000..397bd66 --- /dev/null +++ b/run-blob-wanter.sh @@ -0,0 +1,13 @@ +#!/bin/bash +# run-blob-wanter.sh +# Same shape as run-server.sh / run-ssb-viewer.sh: keep it alive, and let the +# process exit cleanly when its sbot connection dies. +# +# Restarting on sbot loss is not just tolerated here, it is REQUIRED: sbot's +# want map is in-memory and is wiped on every sbot restart. Coming back up +# re-runs the backfill and re-arms every standing want. +while true; do + NODE_OPTIONS="--dns-result-order=ipv4first --max-old-space-size=128" ./blob-wanter.js + echo "Restarting blob-wanter at $(date)" + sleep 5 +done