// Optional leased piano jobs; the heavy model always lives in a separate worker. import { randomUUID, timingSafeEqual } from 'node:crypto'; import { createRequire } from 'node:module'; import { existsSync } from 'node:fs'; import { db as sql } from './db.js'; const { normalizeNotes } = createRequire(import.meta.url)( existsSync(new URL('../frontend/piano-core.js', import.meta.url)) ? '../frontend/piano-core.js' : './public/piano-core.js', ); const ID = /^(?:[\w-]{11}|upl_[a-f0-9]{12})$/; const LEASE = 120000; const sourceKey = (row) => String(row?.sha256 || row?.gen || row?.path || row?.created_at || 'original'); export function registerPianoRoutes( app, { db, adminAuth, enabled = false, workerToken = '', now = Date.now }, ) { enabled = enabled && workerToken.length >= 24; const ready = sql.execute(`CREATE TABLE IF NOT EXISTS piano_jobs ( video_id TEXT PRIMARY KEY,status TEXT NOT NULL,source_key TEXT NOT NULL,notes TEXT,error TEXT, attempts INTEGER NOT NULL DEFAULT 0,lease_token TEXT,lease_until INTEGER NOT NULL DEFAULT 0, created_at INTEGER NOT NULL,updated_at INTEGER NOT NULL)`); const source = async (id) => id.startsWith('upl_') ? db.getUpload(id) : db.getMedia(id).then((row) => (row?.status === 'ready' ? row : null)); const get = async (id) => (await sql.execute({ sql: 'SELECT * FROM piano_jobs WHERE video_id=?', args: [id] })).rows[0]; const publicJob = (row) => row ? { videoId: row.video_id, status: row.status, error: row.error, notes: row.notes ? JSON.parse(row.notes) : null, updatedAt: Number(row.updated_at), } : null; async function workerAuth(c, next) { const token = (c.req.header('authorization') || '').replace(/^Bearer\s+/i, ''); const a = Buffer.from(token), b = Buffer.from(workerToken); if (!enabled || a.length !== b.length || !timingSafeEqual(a, b)) return c.json({ ok: false, error: 'Unauthorized worker' }, 401); await ready; return next(); } app.get('/api/media/:id/piano', async (c) => { await ready; const id = c.req.param('id'); if (!ID.test(id)) return c.json({ ok: false, error: 'Invalid video id' }, 400); const row = await get(id), media = await source(id); const stale = row && media && sourceKey(media) !== row.source_key; return c.json( { ok: true, enabled, job: stale ? { videoId: id, status: 'stale', notes: null, error: 'The source changed. Transcribe it again.', } : publicJob(row), }, 200, { 'Cache-Control': 'no-store' }, ); }); app.post('/api/media/:id/piano', adminAuth, async (c) => { await ready; if (!enabled) return c.json({ ok: false, error: 'The optional piano worker is off.' }, 503); const id = c.req.param('id'); if (!ID.test(id)) return c.json({ ok: false, error: 'Invalid video id' }, 400); const media = await source(id); if (!media) return c.json({ ok: false, error: 'Save this song on the server first.' }, 409); if (Number(media.duration) > 900) return c.json({ ok: false, error: 'Piano jobs are limited to 15 minutes.' }, 422); const old = await get(id); if ( old && ['queued', 'running', 'ready'].includes(old.status) && old.source_key === sourceKey(media) ) return c.json({ ok: true, job: publicJob(old) }); const active = Number( (await sql.execute("SELECT COUNT(*) n FROM piano_jobs WHERE status IN ('queued','running')")) .rows[0].n, ); if (active >= 10) return c.json({ ok: false, error: 'Piano queue is full.' }, 429); await sql.execute({ sql: `INSERT INTO piano_jobs(video_id,status,source_key,created_at,updated_at) SELECT ?,'queued',?,?,? WHERE (SELECT COUNT(*) FROM piano_jobs WHERE status IN ('queued','running'))<10 ON CONFLICT(video_id) DO UPDATE SET status='queued',source_key=excluded.source_key,notes=NULL,error=NULL,attempts=0,lease_token=NULL,lease_until=0,updated_at=excluded.updated_at`, args: [id, sourceKey(media), now(), now()], }); const queued = await get(id); if ( !queued || !['queued', 'running'].includes(queued.status) || queued.source_key !== sourceKey(media) ) return c.json({ ok: false, error: 'Piano queue is full.' }, 429); return c.json({ ok: true, job: publicJob(queued) }, 202); }); app.post('/api/piano-worker/claim', workerAuth, async (c) => { const time = now(); await sql.execute({ sql: "UPDATE piano_jobs SET status='failed',error='Worker stopped repeatedly',updated_at=? WHERE status='running' AND lease_until=3", args: [time, time], }); const row = ( await sql.execute({ sql: `UPDATE piano_jobs SET status='running',attempts=attempts+1,lease_token=?,lease_until=?,updated_at=? WHERE video_id=(SELECT video_id FROM piano_jobs WHERE (status='queued' OR (status='running' AND lease_until { const id = c.req.param('id'); if (!ID.test(id)) return c.json({ ok: false, error: 'Invalid video id' }, 400); const text = await c.req.text(); if (text.length > 4 * 1048576) return c.json({ ok: false, error: 'Result too large' }, 413); let body; try { body = JSON.parse(text); } catch { return c.json({ ok: false, error: 'Invalid JSON' }, 400); } const row = await get(id); if ( !row || row.status !== 'running' || row.lease_token !== body.lease || Number(row.lease_until) < now() ) return c.json({ ok: false, error: 'Job lease expired' }, 409); if (body.action === 'heartbeat') { await sql.execute({ sql: 'UPDATE piano_jobs SET lease_until=?,updated_at=? WHERE video_id=? AND lease_token=?', args: [now() + LEASE, now(), id, body.lease], }); return c.json({ ok: true }); } let notes = null, status = 'failed'; if (body.action === 'complete') { try { notes = normalizeNotes(body.notes); } catch (error) { return c.json({ ok: false, error: error.message }, 422); } status = 'ready'; } else if (body.action !== 'fail') return c.json({ ok: false, error: 'Unknown worker action' }, 400); await sql.execute({ sql: 'UPDATE piano_jobs SET status=?,notes=?,error=?,updated_at=?,lease_token=NULL,lease_until=0 WHERE video_id=? AND lease_token=?', args: [ status, notes ? JSON.stringify(notes) : null, status === 'failed' ? String(body.error || 'Transcription failed').slice(0, 500) : null, now(), id, body.lease, ], }); return c.json({ ok: true }); }); return { ready }; }