186 lines
7.2 KiB
JavaScript
186 lines
7.2 KiB
JavaScript
// 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<? AND attempts>=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<?)) AND attempts<3 ORDER BY created_at LIMIT 1) RETURNING *`,
|
|
args: [randomUUID(), time + LEASE, time, time],
|
|
})
|
|
).rows[0];
|
|
return c.json({
|
|
ok: true,
|
|
job: row
|
|
? {
|
|
...publicJob(row),
|
|
lease: row.lease_token,
|
|
audioPath: row.video_id.startsWith('upl_')
|
|
? `/api/uploads/${row.video_id}`
|
|
: `/api/media/${row.video_id}?a=1`,
|
|
}
|
|
: null,
|
|
});
|
|
});
|
|
app.post('/api/piano-worker/:id', workerAuth, async (c) => {
|
|
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 };
|
|
}
|