From 5751ef63d53bcb0d671be59e0a944cae60f002fb Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Gerg=C5=91=20M=C3=B3ricz?= Date: Wed, 19 Nov 2025 15:15:37 +0100 Subject: [PATCH] fix(api): concurrency endpoints/reporting --- apps/api/src/controllers/v0/admin/cclog.ts | 6 +++--- apps/api/src/controllers/v0/admin/metrics.ts | 8 ++++---- apps/api/src/controllers/v1/concurrency-check.ts | 4 ++-- apps/api/src/controllers/v1/queue-status.ts | 4 ++-- apps/api/src/controllers/v2/concurrency-check.ts | 4 ++-- apps/api/src/controllers/v2/queue-status.ts | 4 ++-- 6 files changed, 15 insertions(+), 15 deletions(-) diff --git a/apps/api/src/controllers/v0/admin/cclog.ts b/apps/api/src/controllers/v0/admin/cclog.ts index 76d2c415d..06ed24a6a 100644 --- a/apps/api/src/controllers/v0/admin/cclog.ts +++ b/apps/api/src/controllers/v0/admin/cclog.ts @@ -1,4 +1,4 @@ -import { redisEvictConnection } from "../../../services/redis"; +import { getRedisConnection } from "../../../services/queue-service"; import { supabase_service } from "../../../services/supabase"; import { logger as _logger } from "../../../lib/logger"; import { Request, Response } from "express"; @@ -10,7 +10,7 @@ async function cclog() { let cursor = 0; do { - const result = await redisEvictConnection.scan( + const result = await getRedisConnection().scan( cursor, "MATCH", "concurrency-limiter:*", @@ -31,7 +31,7 @@ async function cclog() { for (const x of usable) { const at = new Date(); - const concurrency = await redisEvictConnection.zrangebyscore( + const concurrency = await getRedisConnection().zrangebyscore( x, Date.now(), Infinity, diff --git a/apps/api/src/controllers/v0/admin/metrics.ts b/apps/api/src/controllers/v0/admin/metrics.ts index 47abb5bcc..4174d2372 100644 --- a/apps/api/src/controllers/v0/admin/metrics.ts +++ b/apps/api/src/controllers/v0/admin/metrics.ts @@ -1,12 +1,12 @@ import type { Request, Response } from "express"; -import { redisEvictConnection } from "../../../services/redis"; +import { getRedisConnection } from "../../../services/queue-service"; import { nuqGetLocalMetrics, scrapeQueue } from "../../../services/worker/nuq"; export async function metricsController(_: Request, res: Response) { let cursor: string = "0"; const metrics: Record = {}; do { - const res = await redisEvictConnection.sscan( + const res = await getRedisConnection().sscan( "concurrency-limit-queues", cursor, ); @@ -15,10 +15,10 @@ export async function metricsController(_: Request, res: Response) { const keys = res[1]; for (const key of keys) { - const jobCount = await redisEvictConnection.zcard(key); + const jobCount = await getRedisConnection().zcard(key); if (jobCount === 0) { - await redisEvictConnection.srem("concurrency-limit-queues", key); + await getRedisConnection().srem("concurrency-limit-queues", key); } else { const teamId = key.split(":")[1]; metrics[teamId] = jobCount; diff --git a/apps/api/src/controllers/v1/concurrency-check.ts b/apps/api/src/controllers/v1/concurrency-check.ts index 517aed515..b25fedbc9 100644 --- a/apps/api/src/controllers/v1/concurrency-check.ts +++ b/apps/api/src/controllers/v1/concurrency-check.ts @@ -4,7 +4,7 @@ import { RequestWithAuth, } from "./types"; import { Response } from "express"; -import { redisEvictConnection } from "../../../src/services/redis"; +import { getRedisConnection } from "../../../src/services/queue-service"; // Basically just middleware and error wrapping export async function concurrencyCheckController( @@ -13,7 +13,7 @@ export async function concurrencyCheckController( ) { const concurrencyLimiterKey = "concurrency-limiter:" + req.auth.team_id; const now = Date.now(); - const activeJobsOfTeam = await redisEvictConnection.zrangebyscore( + const activeJobsOfTeam = await getRedisConnection().zrangebyscore( concurrencyLimiterKey, now, Infinity, diff --git a/apps/api/src/controllers/v1/queue-status.ts b/apps/api/src/controllers/v1/queue-status.ts index 63e0157fc..dcec60b99 100644 --- a/apps/api/src/controllers/v1/queue-status.ts +++ b/apps/api/src/controllers/v1/queue-status.ts @@ -2,7 +2,7 @@ import { RateLimiterMode } from "../../types"; import { getACUCTeam } from "../auth"; import { AuthCreditUsageChunkFromTeam, RequestWithAuth } from "./types"; import { Response } from "express"; -import { redisEvictConnection } from "../../services/redis"; +import { getRedisConnection } from "../../services/queue-service"; import { cleanOldConcurrencyLimitedJobs, cleanOldConcurrencyLimitEntries, @@ -47,7 +47,7 @@ export async function queueStatusController( await cleanOldConcurrencyLimitedJobs(req.auth.team_id); const queuedJobsOfTeam = await getConcurrencyQueueJobsCount(req.auth.team_id); - const mostRecentSuccess = await redisEvictConnection.get( + const mostRecentSuccess = await getRedisConnection().get( "most-recent-success:" + req.auth.team_id, ); diff --git a/apps/api/src/controllers/v2/concurrency-check.ts b/apps/api/src/controllers/v2/concurrency-check.ts index 17edca1c4..4431ac41f 100644 --- a/apps/api/src/controllers/v2/concurrency-check.ts +++ b/apps/api/src/controllers/v2/concurrency-check.ts @@ -5,7 +5,7 @@ import { } from "./types"; import { AuthCreditUsageChunkFromTeam } from "../v1/types"; import { Response } from "express"; -import { redisEvictConnection } from "../../../src/services/redis"; +import { getRedisConnection } from "../../../src/services/queue-service"; import { getACUCTeam } from "../auth"; import { RateLimiterMode } from "../../types"; @@ -40,7 +40,7 @@ export async function concurrencyCheckController( const concurrencyLimiterKey = "concurrency-limiter:" + req.auth.team_id; const now = Date.now(); - const activeJobsOfTeam = await redisEvictConnection.zrangebyscore( + const activeJobsOfTeam = await getRedisConnection().zrangebyscore( concurrencyLimiterKey, now, Infinity, diff --git a/apps/api/src/controllers/v2/queue-status.ts b/apps/api/src/controllers/v2/queue-status.ts index e52542dce..1322e544d 100644 --- a/apps/api/src/controllers/v2/queue-status.ts +++ b/apps/api/src/controllers/v2/queue-status.ts @@ -3,7 +3,7 @@ import { getACUCTeam } from "../auth"; import { RequestWithAuth } from "./types"; import { AuthCreditUsageChunkFromTeam } from "../v1/types"; import { Response } from "express"; -import { redisEvictConnection } from "../../services/redis"; +import { getRedisConnection } from "../../services/queue-service"; import { cleanOldConcurrencyLimitedJobs, cleanOldConcurrencyLimitEntries, @@ -48,7 +48,7 @@ export async function queueStatusController( await cleanOldConcurrencyLimitedJobs(req.auth.team_id); const queuedJobsOfTeam = await getConcurrencyQueueJobsCount(req.auth.team_id); - const mostRecentSuccess = await redisEvictConnection.get( + const mostRecentSuccess = await getRedisConnection().get( "most-recent-success:" + req.auth.team_id, );