Files
youtube-studio-flow/documentation/YouTube Studio Flow/01 - Architecture/05 - Queue System.md
T

6.7 KiB
Raw Blame History

Queue System

BullMQ on Redis. The API process enqueues jobs; the worker process (src/worker.ts) consumes them.

Critical Redis Requirement

Redis must run with --maxmemory-policy noeviction. BullMQ silently loses jobs if Redis uses allkeys-lru. This is pre-configured in infrastructure/docker-compose.yml.

Queue Names

Defined in backend/src/queues/queues.constants.ts:

Constant Queue name
QUEUES.YOUTUBE_SYNC youtube-sync
QUEUES.RENDER render
QUEUES.LINT lint
QUEUES.BULK_METADATA bulk-metadata
QUEUES.IMPORT import
QUEUES.CONFLICT_DETECTION conflict-detection

Processors

youtube-sync.processor.ts

Trigger: POST /videos/:id/sync Job data: { videoId } Behavior:

  1. Calls VideoRenderService to render description and compute hash
  2. Skips push if hash matches lastSyncedHash (no changes)
  3. Calls QuotaService.canSpend(50) — rejects if quota exceeded
  4. Pushes all editable fields to YouTube API (videos.update)
  5. Updates renderedDescription, lastSyncedHash, lastSyncedAt on the video row

render.processor.ts

Trigger: PUT /video-configs/:videoId when the request body includes "autoRender": true Job data: { videoId } Behavior: Renders description without pushing to YouTube. Updates renderedDescription.

The render job is not enqueued automatically on every config save — it is opt-in via the autoRender flag. VideoConfigsService.upsert() checks if (dto.autoRender) and only then calls renderQueue.add(...). The frontend passes autoRender: true when saving from the video editor so the preview updates in the background. Saves that don't need an immediate re-render (e.g. bulk template applies) omit the flag and skip the queue.

lint.processor.ts

Trigger: POST /lint/bulk, POST /lint/channel/:id, POST /lint/team Job data: { videoId } Behavior: Calls LintingService.lintVideo(videoId). Replaces all unresolved LintResult rows and recomputes Video.lintStatus.

Note: POST /lint/videos/:id (single video) is synchronous — it calls LintingService.lintVideo() directly in the HTTP handler and returns results immediately. It does not use this queue. Only bulk operations go through BullMQ.

bulk-metadata.processor.ts

Trigger: Bulk job confirmed by user Job data: { bulkJobId } Behavior: Processes each BulkJobItem in sequence. Updates video fields, saves rollback snapshots, updates job status counters.

import.processor.ts

Trigger: CSV import committed (POST /imports/csv/commit) Job data: { importJobId } Behavior: Calls ImportsService.executeCommit(), which marks the ImportJob as committed and writes an audit log entry. Does not create or update any Video rows. The preview step does not persist the validated rows, so the processor has no data to act on. CSV import is structurally incomplete — see the backlog.

conflict-detection.processor.ts

Trigger: BullMQ repeatable job registered on worker startup by ConflictDetectionScheduler. Cron pattern comes from CONFLICT_DETECTION_CRON (default 0 3 * * *). Registration is gated by the global CONFLICT_DETECTION_ENABLED env var. Job data: {} Behavior:

  1. Loads all teams with conflictDetectionEnabled: true
  2. For each team, selects videos where remoteConflict = false, youtubeDeletedAt = null, and lastSyncedAt < now - minAgeDays, stalest-first, capped at batchSize
  3. Hands the selected video IDs to YouTubeSyncService.detectConflictsForVideos(), which groups by channelId (needed for OAuth) and issues one videos.list call per batch of up to 50 IDs. Cost is 1 quota unit per batch, not per video (YouTube quota is per method-call, invariant to the number of parts or IDs). QuotaService.canSpend(1) is checked before each batch.
  4. Per video within a batch: on hash mismatch, writes pendingRemoteSnapshot + pendingRemoteDescription and sets remoteConflict = true. On hash match with a stale flag, clears the flag and pending fields — self-heals when the creator reverts an out-of-band edit.
  5. Stops the entire sweep as soon as quota is exhausted; resumes on the next cron tick.

Team.conflictDetectionBatchSize (1500) and conflictDetectionMinAgeDays (≥0) are per-team knobs. Because the API cost is per batch (up to 50 videos each), a batchSize of 500 consumes ~10 units per team per run — not 500. Users resolve detected conflicts via POST /videos/:id/accept-remote (adopt remote, zero extra quota) or POST /videos/:id/sync (push local, overwrite remote).

Frontend Job Completion Detection

There is no push mechanism (no SSE, no WebSocket, no EventEmitter). The frontend detects job completion entirely through polling via TanStack Query's refetchInterval.

Sync queue status — dynamic polling

The <Header> component (components/shared/Header.tsx) polls GET /youtube-sync/queue-status continuously with an interval that adapts to queue state:

Queue state Poll interval
Active jobs running 3 s
Jobs waiting 8 s
Idle 30 s

This drives the sync-in-progress indicator in the header. When the queue drains, the indicator clears within one polling cycle.

Video editor — scheduled invalidations

After POST /videos/:id/sync is enqueued, the video editor schedules two forced refetches of ['video', id] — at 3 s and 8 s — to pick up the updated hasPendingChanges once the processor finishes. If the job takes longer than 8 s, the UI shows stale "push pending" status until the next sync-status poll triggers a broader refresh.

Bulk jobs page — fixed polling

/bulk-jobs polls GET /bulk-jobs every 5 s via refetchInterval: 5000. Job status transitions (pending → processing → completed) are reflected within one polling cycle.

Linting page — fixed polling + immediate invalidation

/linting polls GET /lint/results every 30 s. After enqueuing a bulk lint operation, the mutation's onSuccess immediately invalidates ['lintResults'] and ['videos'] for a faster first update.

Implication

Because completion detection is polling-based, the UI does not reflect job results in real time — there is always a latency of up to one polling interval. For sync jobs that take longer than 8 s, the video editor in particular may lag. If tighter feedback is needed in future, SSE on the sync queue endpoint would be the natural addition.


Job ID Deduplication

  • Lint-once: jobId: lint-{videoId} — BullMQ ignores duplicate job IDs, so submitting the same video twice before the job runs only creates one job
  • Forced rerun: jobId: lint-{videoId}-{timestamp} — always creates a new job
  • 02 - Backend
  • 05 - Local Setup (Redis setup)