diff --git a/.env.example b/.env.example index d53d198..7fe7459 100644 --- a/.env.example +++ b/.env.example @@ -67,6 +67,9 @@ DATABASE_TYPE=sqlite # POSTGRES_PASSWORD=your_password_here # POSTGRES_DB=discord_bot +# Redis Configuration (queue + persistent KV store) +# REDIS_URL=redis://localhost:6379 + # PostgreSQL Connection Pool Configuration # POSTGRES_POOL_MIN=2 # POSTGRES_POOL_MAX=10 diff --git a/package.json b/package.json index f6a3c15..9e69d68 100644 --- a/package.json +++ b/package.json @@ -44,6 +44,7 @@ "drizzle-orm": "^0.45.2", "express": "^5.2.1", "helmet": "^8.1.0", + "ioredis": "^5.11.0", "libsodium-wrappers": "^0.8.4", "lucide-react": "^1.16.0", "motion": "^12.40.0", diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index c74d383..5825cda 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -37,7 +37,7 @@ importers: version: 8.20.0 '@vitejs/plugin-react': specifier: ^6.0.2 - version: 6.0.2(vite@8.0.13(@types/node@25.9.0)(esbuild@0.28.0)(jiti@2.7.0)(tsx@4.22.2)) + version: 6.0.2(vite@8.0.13(@types/node@25.9.0)(esbuild@0.28.0)(jiti@2.7.0)(tsx@4.22.2)(yaml@2.9.0)) axios: specifier: ^1.16.1 version: 1.16.1 @@ -62,6 +62,9 @@ importers: helmet: specifier: ^8.1.0 version: 8.1.0 + ioredis: + specifier: ^5.11.0 + version: 5.11.0 libsodium-wrappers: specifier: ^0.8.4 version: 0.8.4 @@ -106,7 +109,7 @@ importers: version: 3.6.0 vite: specifier: ^8.0.13 - version: 8.0.13(@types/node@25.9.0)(esbuild@0.28.0)(jiti@2.7.0)(tsx@4.22.2) + version: 8.0.13(@types/node@25.9.0)(esbuild@0.28.0)(jiti@2.7.0)(tsx@4.22.2)(yaml@2.9.0) winston: specifier: ^3.19.0 version: 3.19.0 @@ -161,7 +164,7 @@ importers: version: 5.9.3 vitest: specifier: latest - version: 4.1.7(@opentelemetry/api@1.9.1)(@types/node@25.9.0)(vite@8.0.13(@types/node@25.9.0)(esbuild@0.28.0)(jiti@2.7.0)(tsx@4.22.2)) + version: 4.1.7(@opentelemetry/api@1.9.1)(@types/node@25.9.0)(vite@8.0.13(@types/node@25.9.0)(esbuild@0.28.0)(jiti@2.7.0)(tsx@4.22.2)(yaml@2.9.0)) vendor/discord-video-stream: dependencies: @@ -1108,6 +1111,9 @@ packages: cpu: [x64] os: [win32] + '@ioredis/commands@1.10.0': + resolution: {integrity: sha512-UmeW7z4LfctwoQ5wkhVzgq8tXkreED2xZGpX+Bg+zA+WJFZCT6c062AfCK/Dfk81xZnnwdhJCUMkitihRaoC2Q==} + '@jest/schemas@29.6.3': resolution: {integrity: sha512-mo5j5X+jIZmJQveBKeS/clAueipV7KgiX1vMgCxam1RNYiqE1w62n0/tJJnHtjW8ZHcQco5gY85jA3mi0L+nSA==} engines: {node: ^14.15.0 || ^16.10.0 || >=18.0.0} @@ -2364,6 +2370,10 @@ packages: resolution: {integrity: sha512-eYm0QWBtUrBWZWG0d386OGAw16Z995PiOVo2B7bjWSbHedGl5e0ZWaq65kOGgUSNesEIDkB9ISbTg/JK9dhCZA==} engines: {node: '>=6'} + cluster-key-slot@1.1.1: + resolution: {integrity: sha512-rwHwUfXL40Chm1r08yrhU3qpUvdVlgkKNeyeGPOxnW8/SyVDvgRaed/Uz54AqWNaTCAThlj6QAs3TZcKI0xDEw==} + engines: {node: '>=0.10.0'} + cmake-ts@1.0.2: resolution: {integrity: sha512-5l++JHE7MxFuyV/OwJf3ek7ZZN1aGPFPM5oUz6AnK5inQAPe4TFXRMz5sA2qg2FRgByPWdqO+gSfIPo8GzoKNQ==} hasBin: true @@ -2538,6 +2548,10 @@ packages: delegates@1.0.0: resolution: {integrity: sha512-bd2L678uiWATM6m5Z1VzNCErI3jiGzt6HGY8OVICs40JQq/HALfbyNJmp0UDakEY4pMMaN0Ly5om/B1VI/+xfQ==} + denque@2.1.0: + resolution: {integrity: sha512-HVQE3AAb/pxF8fQAoiqpvg9i3evqug3hoiwakOyZAwJm+6vZehbkYXZ0l4JxS+I3QxM97v5aaRNhj8v5oBhekw==} + engines: {node: '>=0.10'} + depd@2.0.0: resolution: {integrity: sha512-g7nH6P6dyDioJogAAGprGpCtVImJhpPk/roCzdb3fIh61/s/nPsfR6onyMwkCAR/OlC3yBC0lESvUoQEAssIrw==} engines: {node: '>= 0.8'} @@ -3181,6 +3195,10 @@ packages: int64-buffer@1.1.0: resolution: {integrity: sha512-94smTCQOvigN4d/2R/YDjz8YVG0Sufvv2aAh8P5m42gwhCsDAJqnbNOrxJsrADuAFAA69Q/ptGzxvNcNuIJcvw==} + ioredis@5.11.0: + resolution: {integrity: sha512-EZBErytyVovD8f6pDfG3Kb37N6Y3lmDA9NNj+4+IP13CzzHGeX+OyeRM2Um13khRzoBSzzL+5lVnCX8V2RLeMg==} + engines: {node: '>=12.22.0'} + ip@2.0.1: resolution: {integrity: sha512-lJUL9imLTNi1ZfXT+DU6rBBdbiKGBuay9B6xGSPVjUeQwaH1RIGqef8RZkUtHioLmSNpPR5M4HVKJGm1j8FWVQ==} @@ -4114,6 +4132,14 @@ packages: resolution: {integrity: sha512-6tDA8g98We0zd0GvVeMT9arEOnTw9qM03L9cJXaCjrip1OO764RDBLBfrB4cwzNGDj5OA5ioymC9GkizgWJDUg==} engines: {node: '>=8'} + redis-errors@1.2.0: + resolution: {integrity: sha512-1qny3OExCf0UvUV/5wpYKf2YwPcOqXzkwKKSmKHiE6ZMQs5heeE/c8eXK+PNllPvmjgAbfnsbpkGZWy8cBpn9w==} + engines: {node: '>=4'} + + redis-parser@3.0.0: + resolution: {integrity: sha512-DJnGAeenTdpMEH6uAJRK/uiyEIH9WVsUmoLwzudwGJUwZPp80PDBWPHXSAGNPwNvIXAbe7MSUB1zQFugFml66A==} + engines: {node: '>=4'} + reduce-extract@1.0.0: resolution: {integrity: sha512-QF8vjWx3wnRSL5uFMyCjDeDc5EBMiryoT9tz94VvgjKfzecHAVnqmXAwQDcr7X4JmLc2cjkjFGCVzhMqDjgR9g==} engines: {node: '>=0.10.0'} @@ -4365,6 +4391,9 @@ packages: stackback@0.0.2: resolution: {integrity: sha512-1XMJE5fQo1jGH6Y/7ebnwPOBEkIEnT4QF32d5R1+VXdXveM0IBMJt8zfaxX1P3QhVwrYe+576+jkANtSS2mBbw==} + standard-as-callback@2.1.0: + resolution: {integrity: sha512-qoRRSyROncaz1z0mvYqIE4lCd9p2R90i6GxW3uZv5ucSu8tU7B5HXUP1gG8pVZsYNVaXjk8ClXHPttLyxAL48A==} + statuses@2.0.2: resolution: {integrity: sha512-DvEy55V3DB7uknRo+4iOGT5fP1slR8wQohVdknigZPMpMstaKJQWhwiYBACJE3Ul2pTnATihhBYnRhZQHGBiRw==} engines: {node: '>= 0.8'} @@ -5526,6 +5555,8 @@ snapshots: '@img/sharp-win32-x64@0.34.5': optional: true + '@ioredis/commands@1.10.0': {} + '@jest/schemas@29.6.3': dependencies: '@sinclair/typebox': 0.27.10 @@ -6365,10 +6396,10 @@ snapshots: dependencies: '@types/node': 25.8.0 - '@vitejs/plugin-react@6.0.2(vite@8.0.13(@types/node@25.9.0)(esbuild@0.28.0)(jiti@2.7.0)(tsx@4.22.2))': + '@vitejs/plugin-react@6.0.2(vite@8.0.13(@types/node@25.9.0)(esbuild@0.28.0)(jiti@2.7.0)(tsx@4.22.2)(yaml@2.9.0))': dependencies: '@rolldown/pluginutils': 1.0.1 - vite: 8.0.13(@types/node@25.9.0)(esbuild@0.28.0)(jiti@2.7.0)(tsx@4.22.2) + vite: 8.0.13(@types/node@25.9.0)(esbuild@0.28.0)(jiti@2.7.0)(tsx@4.22.2)(yaml@2.9.0) '@vitest/expect@4.1.7': dependencies: @@ -6379,13 +6410,13 @@ snapshots: chai: 6.2.2 tinyrainbow: 3.1.0 - '@vitest/mocker@4.1.7(vite@8.0.13(@types/node@25.9.0)(esbuild@0.28.0)(jiti@2.7.0)(tsx@4.22.2))': + '@vitest/mocker@4.1.7(vite@8.0.13(@types/node@25.9.0)(esbuild@0.28.0)(jiti@2.7.0)(tsx@4.22.2)(yaml@2.9.0))': dependencies: '@vitest/spy': 4.1.7 estree-walker: 3.0.3 magic-string: 0.30.21 optionalDependencies: - vite: 8.0.13(@types/node@25.9.0)(esbuild@0.28.0)(jiti@2.7.0)(tsx@4.22.2) + vite: 8.0.13(@types/node@25.9.0)(esbuild@0.28.0)(jiti@2.7.0)(tsx@4.22.2)(yaml@2.9.0) '@vitest/pretty-format@4.1.7': dependencies: @@ -6687,6 +6718,8 @@ snapshots: clsx@2.1.1: {} + cluster-key-slot@1.1.1: {} + cmake-ts@1.0.2: {} collect-all@1.0.4: @@ -6846,6 +6879,8 @@ snapshots: delegates@1.0.0: {} + denque@2.1.0: {} + depd@2.0.0: {} deprecation@2.3.1: {} @@ -7550,6 +7585,18 @@ snapshots: int64-buffer@1.1.0: {} + ioredis@5.11.0: + dependencies: + '@ioredis/commands': 1.10.0 + cluster-key-slot: 1.1.1 + debug: 4.4.3 + denque: 2.1.0 + redis-errors: 1.2.0 + redis-parser: 3.0.0 + standard-as-callback: 2.1.0 + transitivePeerDependencies: + - supports-color + ip@2.0.1: {} ipaddr.js@1.9.1: {} @@ -8422,6 +8469,12 @@ snapshots: indent-string: 4.0.0 strip-indent: 3.0.0 + redis-errors@1.2.0: {} + + redis-parser@3.0.0: + dependencies: + redis-errors: 1.2.0 + reduce-extract@1.0.0: dependencies: test-value: 1.1.0 @@ -8711,6 +8764,8 @@ snapshots: stackback@0.0.2: {} + standard-as-callback@2.1.0: {} + statuses@2.0.2: {} std-env@4.1.0: {} @@ -9033,7 +9088,7 @@ snapshots: vary@1.1.2: {} - vite@8.0.13(@types/node@25.9.0)(esbuild@0.28.0)(jiti@2.7.0)(tsx@4.22.2): + vite@8.0.13(@types/node@25.9.0)(esbuild@0.28.0)(jiti@2.7.0)(tsx@4.22.2)(yaml@2.9.0): dependencies: lightningcss: 1.32.0 picomatch: 4.0.4 @@ -9046,11 +9101,12 @@ snapshots: fsevents: 2.3.3 jiti: 2.7.0 tsx: 4.22.2 + yaml: 2.9.0 - vitest@4.1.7(@opentelemetry/api@1.9.1)(@types/node@25.9.0)(vite@8.0.13(@types/node@25.9.0)(esbuild@0.28.0)(jiti@2.7.0)(tsx@4.22.2)): + vitest@4.1.7(@opentelemetry/api@1.9.1)(@types/node@25.9.0)(vite@8.0.13(@types/node@25.9.0)(esbuild@0.28.0)(jiti@2.7.0)(tsx@4.22.2)(yaml@2.9.0)): dependencies: '@vitest/expect': 4.1.7 - '@vitest/mocker': 4.1.7(vite@8.0.13(@types/node@25.9.0)(esbuild@0.28.0)(jiti@2.7.0)(tsx@4.22.2)) + '@vitest/mocker': 4.1.7(vite@8.0.13(@types/node@25.9.0)(esbuild@0.28.0)(jiti@2.7.0)(tsx@4.22.2)(yaml@2.9.0)) '@vitest/pretty-format': 4.1.7 '@vitest/runner': 4.1.7 '@vitest/snapshot': 4.1.7 @@ -9067,7 +9123,7 @@ snapshots: tinyexec: 1.1.2 tinyglobby: 0.2.16 tinyrainbow: 3.1.0 - vite: 8.0.13(@types/node@25.9.0)(esbuild@0.28.0)(jiti@2.7.0)(tsx@4.22.2) + vite: 8.0.13(@types/node@25.9.0)(esbuild@0.28.0)(jiti@2.7.0)(tsx@4.22.2)(yaml@2.9.0) why-is-node-running: 2.3.0 optionalDependencies: '@opentelemetry/api': 1.9.1 diff --git a/src/config.ts b/src/config.ts index aff1002..9ea4970 100644 --- a/src/config.ts +++ b/src/config.ts @@ -165,6 +165,7 @@ const configSchema = z POSTGRES_POOL_MIN: z.coerce.number().int().positive().default(2), POSTGRES_POOL_MAX: z.coerce.number().int().positive().default(10), ADMIN_PASSWORD: z.string().default("admin123"), + REDIS_URL: z.string().min(1).default("redis://localhost:6379"), }) .superRefine((value, ctx) => { if (!value.AI_ANALYSIS_ENABLED) { diff --git a/src/muxer-queue.ts b/src/muxer-queue.ts index 8aa6f43..6eef081 100644 --- a/src/muxer-queue.ts +++ b/src/muxer-queue.ts @@ -1,35 +1,37 @@ -import { and, asc, eq, lt, sql } from "drizzle-orm"; -import { - getDatabase as getDrizzleDatabase, - initializeDatabase, -} from "./database/drizzle.js"; -import { muxerJobsTable, uiStateTable } from "./database/schema.js"; +import Redis from "ioredis"; +import { config } from "./config.js"; import { createChildLogger } from "./logger.js"; 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; +// ── Redis client (lazy singleton) ────────────────────────────────────────── + +let redis: Redis | null = null; + +function getRedis(): Redis { + if (redis !== null) return redis; + + redis = new Redis(config.REDIS_URL, { + maxRetriesPerRequest: 3, + retryStrategy(times) { + if (times > 5) return null; // stop retrying + return Math.min(times * 200, 2000); + }, + lazyConnect: false, + }); + + redis.on("error", (err) => { + logger.error({ err }, "Redis connection error"); + }); + + redis.on("connect", () => { + logger.info({ url: config.REDIS_URL }, "Redis connected"); + }); + + return redis; } -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; -} +// ── Types ────────────────────────────────────────────────────────────────── export interface MuxerJobData { userId: string; @@ -38,56 +40,37 @@ export interface MuxerJobData { outputDir: string; } -interface StoredJobRow { - id: string; - data: string; - status: "pending" | "processing" | "completed" | "failed"; - attempts: number; - maxAttempts: number; - createdAt: number; - updatedAt: number; - error: string | null; +// ── Database backward compatibility ──────────────────────────────────────── + +/** + * @deprecated Use Redis functions directly. Kept for backward compat with old + * tests that import getDatabase from muxer-queue. + */ +export function getDatabase() { + logger.warn( + "getDatabase() is deprecated — queue now uses Redis. Returning a stub.", + ); + return undefined as unknown as never; } -interface JobStatsRow { - status: "pending" | "processing" | "completed" | "failed"; - count: number | string | { count: number | string }; -} +// ── Persistent KV store (replaces SQLite uiState table) ──────────────────── -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(); -} +const KV_PREFIX = "kv:"; 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 { + const r = getRedis(); + const raw = await r.get(`${KV_PREFIX}${key}`); + if (raw === null) return fallback; + return JSON.parse(raw) as T; + } catch (error) { + logger.error( + { key, error: error instanceof Error ? error.message : String(error) }, + "Failed to get persisted value", + ); return fallback; } } @@ -96,45 +79,62 @@ export async function setPersistedValue( key: string, value: unknown, ): Promise { - await initializeDatabase(); - const database = db(); + try { + const r = getRedis(); + await r.set(`${KV_PREFIX}${key}`, JSON.stringify(value)); + } catch (error) { + logger.error( + { key, error: error instanceof Error ? error.message : String(error) }, + "Failed to set persisted value", + ); + throw error; + } +} - 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(), - }, - }); +// ── Job queue ────────────────────────────────────────────────────────────── + +const JOB_PREFIX = "job:"; +const QUEUE_PENDING = "queue:pending"; +const QUEUE_PROCESSING = "queue:processing"; +const QUEUE_COMPLETED = "queue:completed"; +const QUEUE_FAILED = "queue:failed"; + +function queueKey(status: string): string { + switch (status) { + case "pending": + return QUEUE_PENDING; + case "processing": + return QUEUE_PROCESSING; + case "completed": + return QUEUE_COMPLETED; + case "failed": + return QUEUE_FAILED; + default: + return QUEUE_PENDING; + } } export async function enqueueMuxerJob(data: MuxerJobData): Promise { try { - await initializeDatabase(); - const database = db(); - + const r = getRedis(); 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(); + const jobKey = `${JOB_PREFIX}${jobId}`; + + // Use a pipeline for atomicity + const pipeline = r.pipeline(); + pipeline.hset(jobKey, { + data: JSON.stringify(data), + status: "pending", + attempts: "0", + maxAttempts: "3", + createdAt: String(now), + updatedAt: String(now), + error: "", + }); + pipeline.lpush(QUEUE_PENDING, jobId); + await pipeline.exec(); logger.info( { jobId, userId: data.userId, sessionId: data.sessionId }, @@ -154,27 +154,69 @@ export async function enqueueMuxerJob(data: MuxerJobData): Promise { } } -export async function getPendingJobs(): Promise { - await initializeDatabase(); - const database = db(); +export async function getPendingJobs(): Promise< + Array<{ + id: string; + data: string; + status: "pending" | "processing" | "completed" | "failed"; + attempts: number; + maxAttempts: number; + createdAt: number; + updatedAt: number; + error?: string; + }> +> { + try { + const r = getRedis(); - const rows = await database - .select() - .from(muxerJobsTable) - .where(eq(muxerJobsTable.status, "pending")) - .orderBy(asc(muxerJobsTable.createdAt)) - .limit(10); + // Get up to 10 pending job IDs from the left (oldest first) + const jobIds = await r.lrange(QUEUE_PENDING, 0, 9); + if (jobIds.length === 0) return []; - 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, - })); + // Batch fetch all job hashes + const pipeline = r.pipeline(); + for (const id of jobIds) { + pipeline.hgetall(`${JOB_PREFIX}${id}`); + } + const results = await pipeline.exec(); + if (!results) return []; + + const jobs: Array<{ + id: string; + data: string; + status: "pending" | "processing" | "completed" | "failed"; + attempts: number; + maxAttempts: number; + createdAt: number; + updatedAt: number; + error?: string; + }> = []; + + for (let i = 0; i < jobIds.length; i++) { + const [err, fields] = results[i]; + if (err || !fields) continue; + + const raw = fields as Record; + jobs.push({ + id: jobIds[i], + data: raw.data || "", + status: (raw.status as "pending") || "pending", + attempts: Number(raw.attempts) || 0, + maxAttempts: Number(raw.maxAttempts) || 3, + createdAt: Number(raw.createdAt) || 0, + updatedAt: Number(raw.updatedAt) || 0, + error: raw.error || undefined, + }); + } + + return jobs; + } catch (error) { + logger.error( + { error: error instanceof Error ? error.message : String(error) }, + "Failed to get pending jobs", + ); + return []; + } } export async function updateJobStatus( @@ -182,95 +224,135 @@ export async function updateJobStatus( status: "processing" | "completed" | "failed", error?: string, ): Promise { - await initializeDatabase(); - const database = db(); - const now = Date.now(); + try { + const r = getRedis(); + const jobKey = `${JOB_PREFIX}${jobId}`; + 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)); + const exists = await r.exists(jobKey); + if (!exists) { + logger.warn({ jobId }, "Job not found for status update"); + return; + } + + const currentStatus = await r.hget(jobKey, "status"); + + const pipeline = r.pipeline(); + + if (status === "failed") { + pipeline.hincrby(jobKey, "attempts", 1); + pipeline.hset(jobKey, "error", error || ""); + } + + pipeline.hset(jobKey, "status", status); + pipeline.hset(jobKey, "updatedAt", String(now)); + + // Move job ID between queue lists + if (currentStatus) { + pipeline.lrem(queueKey(currentStatus), 0, jobId); + } + pipeline.lpush(queueKey(status), jobId); + + await pipeline.exec(); + + logger.info({ jobId, status, error }, "Job status updated"); + } catch (err) { + logger.error( + { + jobId, + error: err instanceof Error ? err.message : String(err), + }, + "Failed to update job status", + ); + throw err; } - - logger.info({ jobId, status, error }, "Job status updated"); } export async function retryFailedJob(jobId: string): Promise { - await initializeDatabase(); - const database = db(); + try { + const r = getRedis(); + const jobKey = `${JOB_PREFIX}${jobId}`; - const jobs = await database - .select() - .from(muxerJobsTable) - .where(eq(muxerJobsTable.id, jobId)) - .limit(1); + const [attemptsStr, maxAttemptsStr] = await r.hmget( + jobKey, + "attempts", + "maxAttempts", + ); - const job = jobs[0]; + const attempts = Number(attemptsStr) || 0; + const maxAttempts = Number(maxAttemptsStr) || 3; - if (!job) { - logger.warn({ jobId }, "Job not found"); - return false; - } + if (attempts >= maxAttempts) { + logger.warn( + { jobId, attempts, maxAttempts }, + "Max retry attempts reached", + ); + return false; + } - if (job.attempts >= job.maxAttempts) { - logger.warn( - { jobId, attempts: job.attempts, maxAttempts: job.maxAttempts }, - "Max retry attempts reached", + const pipeline = r.pipeline(); + pipeline.hset(jobKey, "status", "pending"); + pipeline.hset(jobKey, "updatedAt", String(Date.now())); + pipeline.lrem(QUEUE_FAILED, 0, jobId); + pipeline.lpush(QUEUE_PENDING, jobId); + await pipeline.exec(); + + logger.info({ jobId, attempt: attempts + 1 }, "Job retried"); + return true; + } catch (err) { + logger.error( + { jobId, error: err instanceof Error ? err.message : String(err) }, + "Failed to retry job", ); 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; + try { + const r = getRedis(); + const cutoffTime = Date.now() - olderThanMs; - const result = await database - .delete(muxerJobsTable) - .where( - and( - eq(muxerJobsTable.status, "completed"), - lt(muxerJobsTable.updatedAt, cutoffTime), - ), + const jobIds = await r.lrange(QUEUE_COMPLETED, 0, -1); + if (jobIds.length === 0) return 0; + + const pipeline = r.pipeline(); + for (const id of jobIds) { + pipeline.hget(`${JOB_PREFIX}${id}`, "updatedAt"); + } + const results = await pipeline.exec(); + if (!results) return 0; + + let deletedCount = 0; + const deletePipeline = r.pipeline(); + + for (let i = 0; i < jobIds.length; i++) { + const [err, updatedAtStr] = results[i]; + if (err) continue; + + const updatedAt = Number(updatedAtStr) || 0; + if (updatedAt < cutoffTime) { + deletePipeline.del(`${JOB_PREFIX}${jobIds[i]}`); + deletePipeline.lrem(QUEUE_COMPLETED, 0, jobIds[i]); + deletedCount++; + } + } + + if (deletedCount > 0) { + await deletePipeline.exec(); + } + + logger.info({ deletedCount }, "Cleaned up completed jobs"); + return deletedCount; + } catch (err) { + logger.error( + { error: err instanceof Error ? err.message : String(err) }, + "Failed to clean up completed jobs", ); - - const deletedCount = - typeof result === "object" && result !== null && "rowsAffected" in result - ? Number(result.rowsAffected) - : 0; - - logger.info({ deletedCount }, "Cleaned up completed jobs"); - - return deletedCount; + return 0; + } } export async function getJobStats(): Promise<{ @@ -279,38 +361,29 @@ export async function getJobStats(): Promise<{ completed: number; failed: number; }> { - await initializeDatabase(); - const database = db(); + try { + const r = getRedis(); + const [pending, processing, completed, failed] = await Promise.all([ + r.llen(QUEUE_PENDING), + r.llen(QUEUE_PROCESSING), + r.llen(QUEUE_COMPLETED), + r.llen(QUEUE_FAILED), + ]); - 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 { pending, processing, completed, failed }; + } catch (err) { + logger.error( + { error: err instanceof Error ? err.message : String(err) }, + "Failed to get job stats", + ); + return { pending: 0, processing: 0, completed: 0, failed: 0 }; } - - return stats; } export async function closeQueue(): Promise { - logger.info("Muxer queue closed"); + if (redis !== null) { + await redis.quit(); + redis = null; + } + logger.info("Muxer queue (Redis) closed"); }