356 lines
13 KiB
Markdown
356 lines
13 KiB
Markdown
# 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
|
|
|
|
```yaml
|
|
# 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)
|