Files
Sravan 3250530596 feat(db): complete async call-site conversion — Phase 3 done, validated on SQLite
The full sync→async conversion is complete and green on the SQLite backend. Every
DB call across the app now awaits the async adapter, so the identical code runs on
Postgres at cutover.

Converted (this commit finishes Phase 3):
- session.js: currentUser/apiKeyFromReq async → 63 route awaits + WS + static.
- routes.js: all ~250 R.* awaited; DTO helpers (namesFor, avatarsFor, buildMsgDTO,
  buildPollDTO, reactionsForMessage, postSystemMessage, pushGroupUpdate,
  issueRefreshToken, provisionFromBizgaze) made async; every `.map(x=>buildDTO(x))`
  restructured to `await Promise.all(...map(async...))` preserving order; `.filter`
  predicates that hit the DB moved to an `asyncFilter` helper; chained
  `R.x.y(...).length/.map/.filter` wrapped as `(await R.x.y(...)).method`; stream
  upload handlers (recording/transcript/attachment) made async.
- calls.js / signaling.js: all call/meeting fns async; leaveMeeting AWAITS
  persistCallHistory + finalizeTranscript BEFORE endCallByRoom (ordering matters —
  fire-and-forget would race the map teardown); WS handle()/cleanup() async with
  .catch guards.
- static.js: authAttachment(Raw) async (the .some carrier check became a loop),
  handleGet async; server.js dispatch catches handler rejections → 500 not a hang.
- media.js backfill, push.js, reminders.js, webhooks.js await their repo calls.

Validation on DB_BACKEND=sqlite: db-smoke 22/22; legacy e2e 80 checks pass with zero
FAILs (throws only at a PRE-EXISTING WS lobby-drift assertion, unrelated). Every
server file `node --check` clean.

Still on the branch — master untouched. Next: Phase 5 (pg backend + ~7 dialect
queries + data migration + Docker Postgres + cutover), then merge.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-24 22:06:27 +05:30

210 lines
9.0 KiB
JavaScript

'use strict';
//
// media.js — web playback renditions for uploaded videos.
//
// WHY THIS EXISTS (measured, not guessed):
// Real uploads on this box were probed at 19.4 Mbps and 19.0 Mbps (1080p screen recordings) and
// 3.6 Mbps (phone portrait). To play a 19 Mbps file the client must SUSTAIN a 19 Mbps download for
// the whole clip; no mobile link does, so the <video> buffer drains every few seconds and you get
// the classic "buffers, plays, buffers, plays". Server disk and CPU were idle throughout — the
// bottleneck is the media, not the delivery.
// Separately, phone MP4s often store `moov` AFTER `mdat` (not faststart), so the player has to
// fetch the tail before it can begin at all.
//
// WHAT WE DO:
// Leave the uploaded bytes untouched — that is what the download button serves, at full quality.
// Alongside it build <id>.web.mp4: longest side capped at 1280, ~2.5 Mbps ceiling, +faststart.
// /stream/<id> prefers that rendition and falls back to the original until it is ready, so a video
// is never unplayable while it transcodes.
//
const fs = require('fs');
const path = require('path');
const { execFile } = require('child_process');
const { UPLOADS_DIR } = require('./config');
const MAX_CONCURRENT = 2; // transcoding is never urgent; leave the cores to the app
const OK_BITRATE = 2500000; // ≤2.5 Mbps streams fine over mobile
const OK_DIMENSION = 1280; // ...provided it isn't oversized as well
const PROBE_TIMEOUT = 20000;
const XCODE_TIMEOUT = 30 * 60 * 1000;
let running = 0;
const queue = [];
const pending = new Set(); // ids queued or transcoding right now
const failed = new Set(); // ffmpeg couldn't handle it — don't retry forever
const webPath = (id) => path.join(UPLOADS_DIR, id + '.web.mp4');
function hasWebRendition(id) {
try { return fs.statSync(webPath(id)).size > 0; } catch (e) { return false; }
}
// Walk the top-level MP4 box headers. If `mdat` comes before `moov` the index lives at the end of
// the file and the player must seek there before it can start — that is what +faststart fixes.
function isFastStart(fp) {
let fd;
try {
fd = fs.openSync(fp, 'r');
const buf = Buffer.alloc(4096);
const read = fs.readSync(fd, buf, 0, 4096, 0);
let off = 0;
while (off + 8 <= read) {
let size = buf.readUInt32BE(off);
const type = buf.toString('latin1', off + 4, off + 8);
if (type === 'moov') return true;
if (type === 'mdat') return false;
if (size === 1) { // 64-bit largesize follows the header
if (off + 16 > read) return true;
size = Number(buf.readBigUInt64BE(off + 8));
} else if (size === 0) return true; // box runs to EOF
if (size < 8) return true; // malformed — leave it alone
off += size;
}
return true; // couldn't tell; assume fine
} catch (e) {
return true;
} finally {
if (fd !== undefined) { try { fs.closeSync(fd); } catch (e) {} }
}
}
function probe(fp, cb) {
execFile('ffprobe', ['-v', 'error', '-select_streams', 'v:0',
'-show_entries', 'stream=width,height,codec_name',
'-show_entries', 'format=bit_rate,duration',
'-of', 'json', fp], { timeout: PROBE_TIMEOUT, maxBuffer: 1 << 20 }, (err, stdout) => {
if (err) return cb(err);
let info;
try { info = JSON.parse(stdout); } catch (e) { return cb(e); }
const s = (info.streams || [])[0], f = info.format || {};
if (!s || !s.width) return cb(new Error('no video stream'));
cb(null, {
width: +s.width, height: +s.height, codec: s.codec_name || '',
bitrate: +f.bit_rate || 0, duration: +f.duration || 0,
});
});
}
// Cap the LONGEST side at 1280 (so 1920x1080 → 1280x720, and portrait 1080x2340 → 591x1280) while
// preserving aspect. force_divisible_by=2 keeps dimensions legal for H.264.
const SCALE = 'scale=w=' + OK_DIMENSION + ':h=' + OK_DIMENSION +
':force_original_aspect_ratio=decrease:force_divisible_by=2';
function buildArgs(src, out, meta) {
const maxDim = Math.max(meta.width, meta.height);
const lightEnough = meta.bitrate > 0 && meta.bitrate <= OK_BITRATE && maxDim <= OK_DIMENSION;
// Already small enough and only the atom order is wrong → remux, no re-encode. Seconds, not minutes.
if (lightEnough && /^(h264|avc1)$/i.test(meta.codec)) {
return ['-y', '-i', src, '-c', 'copy', '-movflags', '+faststart', '-f', 'mp4', out];
}
return ['-y', '-i', src,
'-map', '0:v:0', '-map', '0:a:0?', // audio optional — silent clips exist
'-vf', SCALE,
'-c:v', 'libx264', '-preset', 'veryfast', '-profile:v', 'high', '-level', '4.0',
'-crf', '24', '-maxrate', '2500k', '-bufsize', '5000k',
'-g', '48', '-pix_fmt', 'yuv420p', // 2s keyframes: smooth seeking
'-c:a', 'aac', '-b:a', '128k', '-ac', '2',
'-movflags', '+faststart', '-f', 'mp4', out];
}
function pump() {
while (running < MAX_CONCURRENT && queue.length) {
const id = queue.shift();
running++;
transcode(id, () => { running--; pending.delete(id); pump(); });
}
}
function transcode(id, done) {
const src = path.join(UPLOADS_DIR, id);
const out = webPath(id);
const tmp = out + '.part';
if (!src.startsWith(UPLOADS_DIR)) return done();
fs.stat(src, (e, srcStat) => {
if (e) return done();
probe(src, (perr, meta) => {
if (perr) { failed.add(id); return done(); }
const maxDim = Math.max(meta.width, meta.height);
// Nothing to gain: already light, already sized, already faststart.
if (meta.bitrate > 0 && meta.bitrate <= OK_BITRATE && maxDim <= OK_DIMENSION && isFastStart(src)) {
failed.add(id); // "no rendition needed" — same effect: stop reconsidering it
return done();
}
execFile('ffmpeg', buildArgs(src, tmp, meta), { timeout: XCODE_TIMEOUT, maxBuffer: 1 << 20 }, (xerr) => {
if (xerr) {
try { fs.unlinkSync(tmp); } catch (_) {}
failed.add(id);
return done();
}
let outStat;
try { outStat = fs.statSync(tmp); } catch (_) { failed.add(id); return done(); }
// A rendition bigger than the original helps nobody — throw it away and stream the original.
if (!outStat.size || outStat.size >= srcStat.size) {
try { fs.unlinkSync(tmp); } catch (_) {}
failed.add(id);
return done();
}
try { fs.renameSync(tmp, out); } catch (_) { try { fs.unlinkSync(tmp); } catch (__) {} }
done();
});
});
});
}
// Poster frame, generated to a temp then renamed so a reader never sees a half-written JPEG (the /thumbs
// handler and this can both target the same file). Warming it at upload means the chat bubble shows the
// poster immediately instead of a blank tile while ffmpeg runs on the first view.
function ensureThumb(id) {
const thumb = path.join(UPLOADS_DIR, id + '.thumb.jpg');
const src = path.join(UPLOADS_DIR, id);
if (!src.startsWith(UPLOADS_DIR)) return;
if (fs.existsSync(thumb)) return;
fs.stat(src, (e) => {
if (e) return;
const tmp = thumb + '.part';
execFile('ffmpeg', ['-y', '-ss', '0.5', '-i', src, '-frames:v', '1', '-vf', 'scale=480:-2', '-q:v', '4', tmp],
{ timeout: 15000 }, (err) => {
if (err) { try { fs.unlinkSync(tmp); } catch (_) {} return; }
try { fs.renameSync(tmp, thumb); } catch (_) { try { fs.unlinkSync(tmp); } catch (__) {} }
});
});
}
// Queue a freshly uploaded (or first-played) video. Cheap and idempotent: safe to call on every
// /stream hit, which is also how pre-existing uploads get backfilled.
function ensureWebRendition(id, mime) {
if (!/^video\//.test(mime || '')) return;
ensureThumb(id); // warm the poster so the bubble isn't blank
if (pending.has(id) || failed.has(id) || hasWebRendition(id)) return;
pending.add(id);
queue.push(id);
pump();
}
// Drop derived files when the attachment goes away.
function dropDerived(id) {
for (const p of [webPath(id), webPath(id) + '.part', path.join(UPLOADS_DIR, id + '.thumb.jpg')]) {
try { fs.unlinkSync(p); } catch (e) {}
}
pending.delete(id); failed.delete(id);
}
// One-off catch-up for videos uploaded before this module existed (and after a restore). Renditions
// persist on the data volume, so on a normal restart this finds nothing to do and costs one query.
// Deliberately delayed and rate-limited by the same 2-at-a-time queue — boot must not stall on it.
function backfill() {
setTimeout(async () => {
let rows = [];
try { rows = await require('./repos').attachments.allVideos(); } catch (e) { return; }
let queued = 0;
for (const r of rows) {
if (hasWebRendition(r.id)) continue;
try { if (!fs.statSync(path.join(UPLOADS_DIR, r.id)).size) continue; } catch (e) { continue; }
ensureWebRendition(r.id, r.mime); queued++;
}
if (queued) console.log('[media] backfilling streaming renditions for ' + queued + ' video(s)');
}, 15000).unref();
}
module.exports = { ensureWebRendition, hasWebRendition, webPath, dropDerived, backfill };