didi-lot2-backend/backend/services/data-layer/didiQueue/INDEX.md
2026-07-10 03:39:53 -07:00

13 KiB

didiQueue - Index

Coada de mesaje RabbitMQ pentru procesarea asincrona a analizelor. Primeste task-uri de analiza de la agent-v3, le distribuie la workeri pe componente, si colecteaza rezultatele intr-un agregator de verdict. Nu contine cod custom -- doar container RabbitMQ cu script de initializare.

Productia ruleaza pe RabbitMQ LOCAL (container staging-dataLayer-rabbitmq). Decizie: stabilitate + zero dependinte externe. Clusterul RabbitMQ RAG (managed extern) ramane configurat ca fallback de urgenta pentru HA, activabil cu redis-switch.sh cluster --rabbit, nu este folosit operational acum.

Productie (LOCAL — activ):

  • Imagine: rabbitmq:3.12-management-alpine
  • Container: staging-dataLayer-rabbitmq
  • Port AMQP: 5672 (expus pe host: 0.0.0.0:5672->5672)
  • Port Management UI: 15672 (expus pe host: 0.0.0.0:15672->15672)
  • Credentiale: admin / rabbitmq123 (din .env)
  • Vhost: /
  • Retea: didi-network

Fallback HA (cluster RAG — disponibil dar inactiv):

  • Endpoint: 10.11.50.100:16672 (HAProxy VIP)
  • Vhost: /didi
  • User: didi
  • Parola: din .cluster-credentials.env (gitignored)

Ce face

  1. Primeste task-uri de analiza -- publicate de agent-v3 dispatcher
  2. Distribuie catre workeri -- fiecare componenta (techniques, ai_tampered, claims, domain) are cozi separate, plus media-preprocess pentru audio/video
  3. Prioritizeaza dupa plan -- plan 1 (free) = prioritate 1, plan 6 (enterprise) = prioritate 10
  4. Colecteaza rezultate -- workerii publica in coada de rezultate, agregatorul face fan-in
  5. DLQ -- mesajele care esueaza dupa 3 incercari merg in dead-letter queue
  6. TTL -- mesajele expira dupa 24 ore

Nota HIL Moderation: tabela moderation_queue (schema bos_analysis) este o tabela PostgreSQL pentru starea de review uman, NU o coada AMQP. RabbitMQ ramane folosit doar pentru dispatch-ul analizei asincrone (5 familii de cozi componente x 6 plan tiers = 30 cozi: media_preprocess, techniques, ai_tampered, claims, domain — plus coada analysis.results + DLQ). HIL nu introduce cozi noi.


Topologie cozi

Exchange

Nume Tip Durabil Scop
analysis topic da Ruteaza task-uri si rezultate

Cozi componente (30 total = 5 cozi x 6 plan types)

analysis.media_preprocess.1  analysis.media_preprocess.2  ... analysis.media_preprocess.6
analysis.techniques.1        analysis.techniques.2        ... analysis.techniques.6
analysis.ai_tampered.1       analysis.ai_tampered.2       ... analysis.ai_tampered.6
analysis.claims.1            analysis.claims.2            ... analysis.claims.6
analysis.domain.1            analysis.domain.2            ... analysis.domain.6

Configurare per coada:

  • Durabil: da
  • Max prioritate: 10
  • Dead-letter exchange: '' (default)
  • Dead-letter routing key: analysis_dlq
  • Message TTL: 86,400,000 ms (24 ore)

Coada rezultate (fan-in)

Coada Bindings Scop
analysis.results analysis.results.techniques, analysis.results.ai_tampered, analysis.results.claims, analysis.results.domain Colecteaza rezultate de la toti workerii

Dead-letter queue

Coada Scop
analysis_dlq Mesaje care au esuat dupa 3 retry-uri

Prioritati per plan

Plan Type Prioritate Tip utilizator
1 1 freemium
2 2 starter
3 4 basic
4 6 pro
5 8 business
6 10 enterprise

Mesajele cu prioritate mai mare sunt procesate primele din coada.


Format mesaje

Task message (Dispatcher -> Worker)

Publicat de agent-v3 dispatcher in cozile de componente.

{
  sessionId: "uuid",
  component: "techniques" | "ai_tampered" | "claims" | "domain",
  planType: 1-6,
  priority: 1-10,
  input: {
    content: "text de analizat",
    url: "URL optional (video/domain)",
    mediaPath: "cale MinIO optional (audio/video)",
    inputType: "text" | "url" | "image" | "audio" | "video"
  },
  userId: "string optional",
  userEmail: "string optional",
  timestamp: 1695312000000,
  retryCount: 0
}

AMQP properties: persistent=true, contentType=application/json, headers={sessionId, component, planType}

Result message (Worker -> Aggregator)

Publicat de worker in coada analysis.results.

{
  sessionId: "uuid",
  component: "techniques" | "ai_tampered" | "claims" | "domain",
  success: true | false,
  score: 0-100,
  data: { ... rezultat flat componenta ... },
  error: "mesaj eroare daca success=false",
  processingTime: 3500,
  timestamp: 1695312003500
}

Cine publica mesaje

Cine Ce publica In ce coada Logica in fisier
agent-v3 dispatcher Task-uri de analiza analysis.{component}.{planType} agent-v3/src/queue/dispatcher.ts
Workeri componente Rezultate analiza analysis.results.{component} agent-v3/src/queue/workers/component-worker.ts

Cine consuma mesaje

Cine Din ce coada Ce face Replici Docker
worker-media-preprocess analysis.media_preprocess.1-6 Download yt-dlp + ffmpeg cadre + Whisper + Vision OCR, cache in Redis, dispatch task-uri componente 2
worker-techniques analysis.techniques.1-6 Ruleaza TechniquesV3Executor 2 (prefetch 5)
worker-ai-tampered analysis.ai_tampered.1-6 Ruleaza AITamperedExecutor 2 (prefetch 5)
worker-claims analysis.claims.1-6 Ruleaza ClaimsExecutor 3 (prefetch 3)
worker-domain analysis.domain.1-6 Ruleaza analyzeDomain() 2 (prefetch 10)
verdict-aggregator analysis.results Fan-in + VerdictCalculator 2 (prefetch 10)

Claims are 3 replici (nu 2) pentru ca e cel mai lent (cautare web per claim). Domain are prefetch 10 pentru ca e cel mai rapid (analiza locala, fara LLM). Media-preprocess este nou (din 2026-03): centralizeaza download/transcribe/vision pentru audio+video, inlocuind logica per-worker. Workerii componente citesc media procesata din Redis (TTL 1h, chei agent:media:{sessionId}:transcript, agent:media:{sessionId}:vision:misinformation, etc.).


Fluxul complet async

Client POST /api/v3/pipeline/analyze-async { text, plan_type: 4 }
  |
  v
agent-v3 Dispatcher
  |-- Salveaza SessionState in Redis (status: processing)
  |-- Salveaza sesiune initiala in Redis + PostgreSQL
  |-- Publica 4 task-uri in RabbitMQ:
  |     analysis.techniques.4 (prioritate 6)
  |     analysis.ai_tampered.4 (prioritate 6)
  |     analysis.claims.4 (prioritate 6)
  |     analysis.domain.4 (prioritate 6)
  |
  v
Response 202: { session_id, poll_url, result_url }

--- In paralel, 4 workeri proceseaza ---

Worker Techniques (consuma din analysis.techniques.4)
  |-- Achizitioneaza lock Redis (300s TTL)
  |-- Ruleaza TechniquesV3Executor (screening -> deep analysis)
  |-- Publica rezultat in analysis.results.techniques
  |-- ACK mesaj

Worker AI-Tampered (consuma din analysis.ai_tampered.4)
  |-- Ruleaza AITamperedExecutor
  |-- Publica rezultat in analysis.results.ai_tampered

Worker Claims (consuma din analysis.claims.4)
  |-- Ruleaza ClaimsExecutor (extrage + verifica prin web)
  |-- Publica rezultat in analysis.results.claims

Worker Domain (consuma din analysis.domain.4)
  |-- Ruleaza analyzeDomain()
  |-- Publica rezultat in analysis.results.domain

--- Agregatorul colecteaza ---

Verdict Aggregator (consuma din analysis.results)
  |-- Primeste rezultat componenta
  |-- Achizitioneaza lock Redis (30s TTL)
  |-- Actualizeaza SessionState in Redis (completedComponents++)
  |-- Daca toate 4 componente gata:
  |     |-- VerdictCalculator.calculate() (functie pura)
  |     |-- VerdictExplanation.generate() (LLM, RO+EN)
  |     |-- PersistService.persist() (Redis + PostgreSQL)
  |-- ACK mesaj

--- Clientul polleaza ---

GET /api/v3/pipeline/{sessionId}/queue-status
  -> { progress: 75%, completed_components: ["techniques", "ai_tampered", "domain"] }

GET /api/v3/pipeline/{sessionId}/result
  -> AnalysisSession completa (cand status=completed)

Retry si error handling

Situatie Actiune Rezultat
Procesare reusita channel.ack(msg) Mesaj sters din coada
Eroare retryable + retryCount < 3 channel.nack(msg, false, true) Mesaj pus inapoi in coada
Eroare retryable + retryCount >= 3 channel.nack(msg, false, false) Mesaj trimis in analysis_dlq
Eroare non-retryable channel.nack(msg, false, false) Mesaj trimis in analysis_dlq
RabbitMQ indisponibil Fallback sync agent-v3 ruleaza analiza sincrona

Lock-uri Redis previn procesarea dubla:

  • Lock componenta: didi:queue:lock:{sessionId}:{component} (TTL 300s)
  • Lock agregator: didi:queue:lock:aggregator:{sessionId} (TTL 30s)

Procesare media in workeri

Workerii proceseaza media inainte de analiza text:

Input type Ce face workerul Logica in
text Nimic, trimite direct la executor component-worker.ts
audio Transcriere via M17/Groq/OpenAI shared/media/transcription.ts
video Download yt-dlp + ffmpeg cadre + transcriere shared/media/video-processor.ts
image Extragere text via vision cascade (pentru techniques/claims) shared/media/vision.ts

Timeout-uri worker:

  • Video: 600,000 ms (10 minute)
  • Default: 120,000 ms (2 minute)

Conexiune RabbitMQ (din agent-v3)

Fisier: agent-v3/src/queue/connection.ts + agent-v3/src/shared/queue/constants.ts

  • Lazy initialization (conectare la prima utilizare)
  • Doua canale: regular (consume) + confirm (publish cu confirmare)
  • Auto-recovery la deconectare (CONNECTION_RETRY_DELAY = 2s pentru failover HA pe cluster)
  • Graceful shutdown pe SIGTERM/SIGINT (stop consume, close channels, close connection)
  • URL building foloseste encodeURIComponent() pentru vhost (/didi -> %2Fdidi in AMQP URI)
# Productie (LOCAL — activ)
URL: amqp://admin:rabbitmq123@staging-dataLayer-rabbitmq:5672/   (vhost `/`)

# Fallback HA (cluster RAG, HAProxy VIP — disponibil dar inactiv)
URL: amqp://didi:<password>@10.11.50.100:16672/%2Fdidi

Switch intre cluster si local

Script: backend/services/orchestration-layer/scripts/redis-switch.sh (denumirea istorica este redis-switch, dar suporta si RabbitMQ via flag --rabbit).

redis-switch.sh {cluster|local|status} [redis|rabbit|both]

Verificare topologie: agent-v3/scripts/verify-rabbitmq-cluster.ts (verifica privilegii vhost + topologia celor 30 cozi + exchange + DLQ).


Fisiere in directorul didiQueue

init-queues.sh          -- Creeaza exchange + coada legacy singulara + DLQ + binding + policy (idempotent)
.env                    -- Credentiale + porturi
.env.example            -- Template
README.md               -- Documentatie

Zero cod custom. Doar container RabbitMQ standard cu management plugin.

Nota: init-queues.sh creeaza topologia legacy cu o singura coada. Topologia actuala cu 30 cozi (5 componente x 6 planuri: media_preprocess, techniques, ai_tampered, claims, domain) + coada results + DLQ este creata dinamic de workerii agent-v3 la startup (vezi agent-v3/src/queue/connection.ts si constants.ts).


Fisiere cod integrare (in agent-v3)

Fisier Rol
agent-v3/src/queue/connection.ts Manager conexiune RabbitMQ (lazy, auto-recovery)
agent-v3/src/shared/queue/constants.ts Nume exchange/cozi, prioritati, config workeri
agent-v3/src/queue/dispatcher.ts Publica task-uri in cozi componente
agent-v3/src/queue/aggregator.ts Consuma rezultate, calculeaza verdict, persista
agent-v3/src/queue/workers/component-worker.ts Worker generic (lock, procesare, publish result, ack/nack)
agent-v3/src/worker-entrypoints/techniques.ts Entry point Docker worker techniques
agent-v3/src/worker-entrypoints/ai-tampered.ts Entry point Docker worker ai-tampered
agent-v3/src/worker-entrypoints/claims.ts Entry point Docker worker claims
agent-v3/src/worker-entrypoints/domain.ts Entry point Docker worker domain
agent-v3/src/worker-entrypoints/aggregator.ts Entry point Docker verdict aggregator

Configurare Docker

# din data-layer/docker-compose.yml
staging-dataLayer-rabbitmq:
  image: rabbitmq:3.12-management-alpine
  container_name: staging-dataLayer-rabbitmq
  restart: unless-stopped
  environment:
    RABBITMQ_DEFAULT_USER: admin
    RABBITMQ_DEFAULT_PASS: rabbitmq123
    RABBITMQ_DEFAULT_VHOST: /
  ports:
    - "5672:5672"       # AMQP (expus pe host)
    - "15672:15672"     # Management UI (expus pe host)
  volumes:
    - didi-staging-rabbitmq-data:/var/lib/rabbitmq
  networks:
    - didi-network
  healthcheck:
    test: rabbitmq-diagnostics -q ping
    interval: 30s
    timeout: 10s
    retries: 3
    start_period: 60s

Workerii sunt definiti in agent-v3/docker-compose.yml (vezi agent-v3/INDEX.md pentru detalii replici).


Ce NU face

  • Nu are cod custom (container RabbitMQ standard)
  • Nu are clustering (instanta singulara)
  • Nu are mirroring/quorum queues (nu e HA)
  • Nu are SSL/TLS (AMQP plain text intern)
  • Nu are ACL per serviciu (toti folosesc userul admin)
  • Nu are delayed message plugin (retry prin requeue nativ)
  • Nu are shovel/federation (nu transfera mesaje intre brokeri)