import { and, asc, eq, lt, sql } from "drizzle-orm"; import { getDatabase as getDrizzleDatabase, initializeDatabase, } from "./database/drizzle"; import { muxerJobsTable, uiStateTable } from "./database/schema"; import { createChildLogger } from "./logger"; const logger = createChildLogger("muxer-queue"); interface QueryBuilder extends PromiseLike { from(...args: unknown[]): QueryBuilder; where(...args: unknown[]): QueryBuilder; orderBy(...args: unknown[]): QueryBuilder; limit(...args: unknown[]): QueryBuilder; values(...args: unknown[]): QueryBuilder; onConflictDoNothing(...args: unknown[]): QueryBuilder; onConflictDoUpdate(...args: unknown[]): QueryBuilder; set(...args: unknown[]): QueryBuilder; groupBy(...args: unknown[]): QueryBuilder; } export interface SqliteDatabase { select(...args: unknown[]): QueryBuilder; insert(...args: unknown[]): QueryBuilder; update(...args: unknown[]): QueryBuilder; delete(...args: unknown[]): QueryBuilder; } function db(): SqliteDatabase { return getDrizzleDatabase() as unknown as SqliteDatabase; } export interface MuxerJobData { userId: string; sessionId: string; recordingsDir: string; outputDir: string; } interface StoredJobRow { id: string; data: string; status: "pending" | "processing" | "completed" | "failed"; attempts: number; maxAttempts: number; createdAt: number; updatedAt: number; error: string | null; } interface JobStatsRow { status: "pending" | "processing" | "completed" | "failed"; count: number | string | { count: number | string }; } interface StoredJob { id: string; data: string; status: "pending" | "processing" | "completed" | "failed"; attempts: number; maxAttempts: number; createdAt: number; updatedAt: number; error?: string; } // Export getDatabase for backward compatibility with webserver.ts export function getDatabase(): SqliteDatabase { return db(); } export async function getPersistedValue( key: string, fallback: T, ): Promise { await initializeDatabase(); const database = db(); const row = await database .select>() .from(uiStateTable) .where(eq(uiStateTable.key, key)) .limit(1); if (!row || row.length === 0) return fallback; try { return JSON.parse(row[0].value) as T; } catch { return fallback; } } export async function setPersistedValue( key: string, value: unknown, ): Promise { await initializeDatabase(); const database = db(); await database .insert(uiStateTable) .values({ key, value: JSON.stringify(value), updated_at: Date.now(), }) .onConflictDoUpdate({ target: uiStateTable.key, set: { value: JSON.stringify(value), updated_at: Date.now(), }, }); } export async function enqueueMuxerJob(data: MuxerJobData): Promise { try { await initializeDatabase(); const database = db(); const jobId = `${data.userId}-${data.sessionId}`; const now = Date.now(); await database .insert(muxerJobsTable) .values({ id: jobId, data: JSON.stringify(data), status: "pending", attempts: 0, maxAttempts: 3, createdAt: now, updatedAt: now, }) .onConflictDoNothing(); logger.info( { jobId, userId: data.userId, sessionId: data.sessionId }, "Muxer job enqueued", ); return jobId; } catch (error) { logger.error( { userId: data.userId, error: error instanceof Error ? error.message : String(error), }, "Failed to enqueue muxer job", ); throw error; } } export async function getPendingJobs(): Promise { await initializeDatabase(); const database = db(); const rows = await database .select() .from(muxerJobsTable) .where(eq(muxerJobsTable.status, "pending")) .orderBy(asc(muxerJobsTable.createdAt)) .limit(10); return rows.map((row) => ({ id: row.id, data: row.data, status: row.status as "pending" | "processing" | "completed" | "failed", attempts: row.attempts, maxAttempts: row.maxAttempts, createdAt: row.createdAt, updatedAt: row.updatedAt, error: row.error || undefined, })); } export async function updateJobStatus( jobId: string, status: "processing" | "completed" | "failed", error?: string, ): Promise { await initializeDatabase(); const database = db(); const now = Date.now(); if (status === "failed") { await database .update(muxerJobsTable) .set({ status, attempts: sql`${muxerJobsTable.attempts} + 1`, updatedAt: now, error: error || null, }) .where(eq(muxerJobsTable.id, jobId)); } else { await database .update(muxerJobsTable) .set({ status, updatedAt: now, }) .where(eq(muxerJobsTable.id, jobId)); } logger.info({ jobId, status, error }, "Job status updated"); } export async function retryFailedJob(jobId: string): Promise { await initializeDatabase(); const database = db(); const jobs = await database .select() .from(muxerJobsTable) .where(eq(muxerJobsTable.id, jobId)) .limit(1); const job = jobs[0]; if (!job) { logger.warn({ jobId }, "Job not found"); return false; } if (job.attempts >= job.maxAttempts) { logger.warn( { jobId, attempts: job.attempts, maxAttempts: job.maxAttempts }, "Max retry attempts reached", ); return false; } await database .update(muxerJobsTable) .set({ status: "pending", updatedAt: Date.now(), }) .where(eq(muxerJobsTable.id, jobId)); logger.info({ jobId, attempt: job.attempts + 1 }, "Job retried"); return true; } export async function cleanupCompletedJobs( olderThanMs: number = 24 * 60 * 60 * 1000, ): Promise { await initializeDatabase(); const database = db(); const cutoffTime = Date.now() - olderThanMs; const result = await database .delete(muxerJobsTable) .where( and( eq(muxerJobsTable.status, "completed"), lt(muxerJobsTable.updatedAt, cutoffTime), ), ); const deletedCount = typeof result === "object" && result !== null && "rowsAffected" in result ? Number(result.rowsAffected) : 0; logger.info({ deletedCount }, "Cleaned up completed jobs"); return deletedCount; } export async function getJobStats(): Promise<{ pending: number; processing: number; completed: number; failed: number; }> { await initializeDatabase(); const database = db(); const rows = await database .select({ status: muxerJobsTable.status, count: sql`COUNT(*)`, }) .from(muxerJobsTable) .groupBy(muxerJobsTable.status); const stats = { pending: 0, processing: 0, completed: 0, failed: 0, }; for (const row of rows) { const count = typeof row.count === "object" && "count" in row.count ? Number((row.count as { count: number | string }).count) : Number(row.count); if (row.status === "pending") stats.pending = count; else if (row.status === "processing") stats.processing = count; else if (row.status === "completed") stats.completed = count; else if (row.status === "failed") stats.failed = count; } return stats; } export async function closeQueue(): Promise { logger.info("Muxer queue closed"); }