33800 Docs

← Retour

Audit Passe 2 — ai-orchestrator

Data flow, gestion d'etat, propagation d'erreurs, coherence inter-fonctions

Date : 15-02-2026 Status : IMPLEMENTEE Auditeur : Claude Opus 4.6 Fichiers audites :

Methode : Lecture integrale de chaque fichier, ligne par ligne. Tracage des data flows entree-traitement-sortie, des chemins d'erreur, et de la coherence d'etat Redis/PostgreSQL/memoire.


TABLE DES MATIERES

  1. FINDINGS CRITIQUES (BUG CERTAIN)
  2. FINDINGS IMPORTANTS (BUG PROBABLE)
  3. FINDINGS MODEREES (RISQUE THEORIQUE)
  4. ANALYSE DES DATA FLOWS
  5. ANALYSE DE LA GESTION D'ETAT
  6. ANALYSE DE LA PROPAGATION D'ERREURS
  7. SECURITE
  8. SYNTHESE ET RECOMMANDATIONS

1. FINDINGS CRITIQUES (BUG CERTAIN)

F01 — Cloud job tasks sont fire-and-forget, jamais trackees dans gpu_slots

Classification : BUG CERTAIN Fichier : main.py, lignes 4267-4271 Impact : Fuite de taches asyncio, aucun suivi, aucune limite de concurrence

# Ligne 4268-4271
if tool_id in CLOUD_TOOLS:
    task = asyncio.create_task(process_single_job(job_id, tool_id))
    logger.info(f"worker_dispatch_cloud job={job_id[:8]} tool={tool_id}")
    continue

La tache task est creee mais jamais stockee dans aucune structure (ni gpu_slots, ni une liste dediee). Consequences :

  1. Aucune limite de concurrence : si 100 jobs Anthropic sont en queue, 100 requetes HTTP partent en parallele instantanement.
  2. Au shutdown (lignes 4318-4327), les taches cloud ne sont PAS annulees (seules celles dans gpu_slots le sont).
  3. Pas de watchdog : un job cloud bloque (timeout Anthropic de 120s) n'est jamais detecte par le watchdog idle (lignes 4193-4222, qui ne regarde que gpu_slots).
  4. L'objet task est garbage-collecte sans attente, ce qui peut provoquer un warning asyncio "Task was destroyed but it is pending".

Scenario : Un burst de 50 jobs Anthropic arrive. 50 requetes HTTP partent simultanement. Le budget check (ligne 2900) est fait au debut de chaque job, mais le cout n'est pas enregistre en DB avant la fin du job. Donc tous les 50 passent le budget check et sont executes, meme si les 10 premiers suffisent a epuiser le budget.


F02 — Race condition sur le budget Anthropic (TOCTOU)

Classification : BUG CERTAIN Fichier : main.py, lignes 2895-2905 (check_anthropic_budget) et 4044-4071 (cost update) Impact : Depassement de budget garanti sous charge concurrente

Flow :

  1. check_anthropic_budget() lit SUM(cost_usd) en DB pour les jobs completed du mois (ligne 2887-2891)
  2. Le job s'execute (peut prendre 1-120 secondes)
  3. Le cout est ecrit en DB seulement a la fin, dans db_update_job_status("completed", cost_usd=...) (ligne 4056-4071)

Pendant l'etape 2, d'autres jobs Anthropic font le meme check et voient le meme budget restant. Combine avec F01 (pas de limite de concurrence cloud), un burst de N jobs simultanees peut depasser le budget de N fois le cout unitaire.

Exemple concret : Budget = 120$, deja depense = 119$. 10 jobs arrivent en meme temps. Chacun voit 119$ < 120$, passe le check, depense ~1$ chacun. Total final = 129$ au lieu de 120$.


F03 — Variable is_cloud redefinie dans le scope, masque la variable du try externe

Classification : BUG CERTAIN Fichier : main.py, ligne 4045 Impact : Pas d'impact fonctionnel actuel mais semantiquement incorrect et fragile

# Ligne 3943 (debut du try)
is_cloud = tool_id in CLOUD_TOOLS

# Ligne 4045 (dans le try, apres execute_job)
is_cloud = tool_id in CLOUD_TOOLS  # Redefini a l'identique

Et dans le finally (ligne 4117) :

if not is_cloud and tool and tool.gpu_id is not None and tool.gpu_id >= 0:

La variable is_cloud de la ligne 4045 masque celle de la ligne 3943 dans le meme scope. Actuellement pas de divergence de valeur, mais si jamais resolved_tool_id change tool_id entre les deux lignes, le finally utiliserait la mauvaise valeur. C'est un signe de code mal structure.

Reclassification : Apres relecture, les deux evaluations portent sur le meme tool_id (ligne 3931 le fixe une seule fois). Impact nul actuellement. Classification reelle : RISQUE THEORIQUE de regression lors d'un futur refactoring.


F04 — Batch API : httpx client reutilise apres possible timeout

Classification : BUG CERTAIN Fichier : main.py, lignes 2994-3065 Impact : Crash ou requete sur connexion fermee

async with httpx.AsyncClient(timeout=httpx.Timeout(30.0, connect=10.0)) as client:
    # Submit batch (timeout 30s) — OK
    resp = await client.post(...)  # ligne 2996

    # Poll for completion (max 2 hours)
    for i in range(max_polls):
        await asyncio.sleep(30)
        poll_resp = await client.get(...)  # ligne 3020

Le httpx.AsyncClient est cree avec un timeout de 30s. Le polling boucle pendant 2 heures (240 * 30s). Le client est utilise continuellement pendant 2h dans le meme context manager async with. En pratique httpx ne ferme pas le client apres le timeout de 30s (c'est un timeout par requete), mais :

  1. Les connexions TCP sous-jacentes peuvent etre fermees par le serveur ou un intermediaire avant 2h.
  2. Le timeout de 30s est trop court pour les polls individuels si l'API Anthropic est lente.

Plus grave : si poll_resp.status_code != 200 (ligne 3028), le code fait continue et reessaie indefiniment (240 fois). Aucune distinction entre 401 (token invalide, retenter est inutile), 429 (rate limit, devrait attendre plus), et 500 (erreur serveur transitoire).


2. FINDINGS IMPORTANTS (BUG PROBABLE)

F05 — Incoherence Redis/PostgreSQL : job poppe de la queue mais echec DB possible

Classification : BUG PROBABLE Fichier : main.py, lignes 4248-4261 Impact : Job perdu (enleve de Redis, mais jamais traite)

# Ligne 4248 — job RETIRE de Redis (rpop atomique)
job_msg = await queue_pop_job()

# Ligne 4256 — lecture DB
job_data = await db_get_job(job_id)
if not job_data:
    logger.error(f"Job {job_id} not found in database")
    continue  # <-- Job PERDU : retire de Redis, pas remis en queue

if job_data["status"] != "pending":
    continue  # <-- Idem : job retire de Redis, pas remis en queue

Si la DB est temporairement indisponible a la ligne 4256, db_get_job leve une exception qui est attrapee par le except Exception global (ligne 4328), et le job est perdu — il a ete retire de Redis mais jamais traite. L'endpoint /api/queue/sync (lignes 4467-4506) peut resynchroniser, mais il faut une intervention manuelle.

De meme, si le status n'est plus "pending" (ex: cancele entre le push et le pop), le job est simplement ignore sans trace dans les logs.


F06 — Requeue infini sans compteur : un job impossible peut bloquer la queue

Classification : BUG PROBABLE Fichier : main.py, lignes 3951-3955, 3978-3981, 3991-3995, 4313-4316 Impact : Starvation de la queue si un job ne peut jamais etre dispatche

Plusieurs chemins requeue un job en "pending" :

Aucun de ces chemins n'incremente un compteur de retry. Le champ retry_count (defini dans le modele Job ligne 270) n'est jamais incremente nulle part dans le code (search exhaustive : aucun appel retry_count en ecriture).

Un job qui demande un tool casse (ex: ComfyUI ne demarre plus) sera requeue indefiniment : pop -> check -> requeue -> pop -> check -> requeue, a raison d'un cycle toutes les 5 secondes (sleep a la ligne 4316). Cela ne bloque pas les autres jobs (les queues sont high/normal/low et FIFO), mais ce job specifique consomme des cycles et du log.


F07 — tools_state et active_tools sont en memoire pure, pas synchronises avec la realite

Classification : BUG PROBABLE Fichier : main.py, lignes 734-735 Impact : Decisions de scheduling basees sur un etat obsolete

tools_state: Dict[str, ToolStatus] = {tool_id: ToolStatus.STOPPED for tool_id in TOOLS_CONFIG}
active_tools: Dict[int, Optional[str]] = {gpu_id: None for gpu_id in GPU_CONFIG}

Ces dictionnaires sont la "source de verite" pour le scheduling. Ils sont mis a jour quand le code demarre/arrete un outil. Mais :

  1. Si un outil crash sur win11 (ex: OOM, power off), tools_state dit toujours RUNNING.
  2. Le worker utilise active_tools pour decider si un GPU est libre (ligne 3976).
  3. Le seul moyen de resynchroniser est /api/tools/refresh (lignes 1207-1243), qui n'est JAMAIS appele automatiquement.

Le worker fait un health check (check_tool_health) pour les outils non-CLI (ligne 3967), ce qui corrige partiellement le probleme. Mais le scheduling dans determine_preferred_gpu (lignes 3847-3898) et can_dispatch (lignes 4149-4178) se basent sur gpu_slots (interne au worker), pas sur tools_state. Donc l'impact est attenue pour le scheduling, mais les endpoints API (/api/tools, /api/gpu) retournent des infos potentiellement fausses.


F08 — Le GPU cache (30s TTL) peut donner des mauvaises decisions de scheduling

Classification : BUG PROBABLE Fichier : main.py, lignes 930-942, 3949 Impact : Un job est requeue parce que le cache dit GPU unreachable alors qu'il est revenu

# Ligne 3949
gpus = await get_cached_gpu_status()
target_gpu = next((g for g in gpus if g.gpu_id == tool.gpu_id), None)
if target_gpu and target_gpu.vram_free == "N/A":
    # Requeue job

Le cache a un TTL de 30s. Si win11 reboot et revient en 20s, le cache retourne toujours les anciennes donnees pendant 10s. Tous les jobs arrivals pendant ces 10s seront requeues inutilement.

Plus grave : le cache est un dict global mutable (_gpu_cache, ligne 931). En asyncio, deux coroutines peuvent lire/ecrire _gpu_cache simultanément. Techniquement safe en CPython single-threaded asyncio (pas de context switch au milieu d'un dict assignment), mais fragile.


F09 — Webhook callback peut echouer silencieusement, pas de retry

Classification : BUG PROBABLE Fichier : main.py, lignes 1491-1535 Impact : Le consommateur (connectors-api ou autre) ne recoit jamais la notification de completion

async with aiohttp.ClientSession() as session:
    async with session.post(callback_url, json=payload, timeout=...) as resp:
        if resp.status >= 200 and resp.status < 300:
            logger.info(...)
        else:
            logger.warning(...)  # <-- Pas de retry

Le except Exception (ligne 1534) log l'erreur mais ne retente pas. Si le consommateur est temporairement down au moment exact de la completion, le webhook est perdu. Le job est marque completed en DB, mais le consommateur ne le sait jamais. Le consommateur doit alors poll /api/jobs/{id} pour decouvrir le resultat.


F10 — SadTalker cleanup cree une task orpheline

Classification : BUG PROBABLE Fichier : main.py, ligne 3824 Impact : Warning asyncio "coroutine was never awaited" ou task abandonnee

# Ligne 3824
asyncio.create_task(asyncio.create_subprocess_shell(cleanup_cmd))

asyncio.create_subprocess_shell retourne un coroutine, pas un awaitable direct pour create_task. En fait, create_task(coroutine) est correct syntaxiquement, mais la task resultante n'est jamais stockee ni awaitee. Si l'event loop se ferme pendant le cleanup, la task est detruite avec un warning. De plus, le Process retourne par la coroutine n'est jamais wait()ed, ce qui peut laisser un processus zombie.


3. FINDINGS MODEREES (RISQUE THEORIQUE)

F11 — Credentials en dur dans le code source

Classification : RISQUE THEORIQUE (securite) Fichier : main.py, lignes 35-41 Impact : Fuite de credentials si le repo est expose

REDIS_PASSWORD = os.environ.get("REDIS_PASSWORD", "urefsLoXibTZ36mtuygcLyNuqzOmIVkF")
POSTGRES_PASSWORD = os.environ.get("POSTGRES_PASSWORD", "xgvKM65PC3zWspTCXkbxPx2fx5fOHRVK")

Les mots de passe sont en fallback hardcode. En production, les variables d'env sont probablement definies par smart-deploy, mais si elles sont absentes, le code fonctionne quand meme avec les mots de passe en dur. Le docker-compose.yml ne definit PAS ces variables (lignes 14-15 ne contiennent que WIN11_HOST et WIN11_USER), donc en dev local, les mots de passe hardcodes sont utilises.

Le repo est sur un GitLab prive (192.168.1.196), donc le risque est limite au reseau local. Mais c'est une mauvaise pratique.


F12 — CORS allow_origins=["*"] avec allow_credentials=True

Classification : RISQUE THEORIQUE (securite) Fichier : main.py, lignes 176-182 Impact : Potentiel CSRF/cookie stealing si des cookies sont utilises

app.add_middleware(
    CORSMiddleware,
    allow_origins=["*"],
    allow_credentials=True,
    allow_methods=["*"],
    allow_headers=["*"],
)

allow_origins=["*"] avec allow_credentials=True est deconseille par la spec CORS. En pratique, les navigateurs modernes bloquent Access-Control-Allow-Origin: * quand credentials: include est envoye. FastAPI/Starlette gere cela correctement en ne renvoyant pas * avec credentials.

Le service n'utilise pas de cookies ni d'authentification, donc l'impact est nul actuellement. Mais si une authentification est ajoutee plus tard, cette config serait dangereuse.


F13 — Pas de validation de taille sur les uploads

Classification : RISQUE THEORIQUE Fichier : main.py, lignes 4511-4537 Impact : Denial of Service par upload de fichiers tres volumineux

@app.post("/api/upload")
async def upload_input_file(file: UploadFile = File(...)):
    content = await file.read()  # Lit TOUT en memoire
    with open(file_path, "wb") as f:
        f.write(content)

Aucune limite de taille. Un upload de 10 Go serait lu entierement en memoire. Le container a probablement une limite de RAM (non specifiee dans docker-compose.yml), mais cela causerait un OOM kill.


F14 — Pas d'authentification sur aucun endpoint

Classification : RISQUE THEORIQUE (securite) Fichier : main.py, tous les endpoints Impact : N'importe qui sur le reseau peut creer des jobs, uploader des fichiers, arreter des outils

Aucun endpoint n'a de middleware d'authentification. Tous les endpoints sont accessibles sans token ni API key. L'acces est restreint par le reseau (nginx reverse proxy avec split-DNS), mais si quelqu'un accede au reseau local, il a un controle total sur l'orchestrateur.

En particulier :


F15 — Les jobs requeues perdent leur position FIFO

Classification : RISQUE THEORIQUE Fichier : main.py, lignes 3953-3954, 4314 Impact : Inversion de priorite pour les jobs requeues

# Ligne 3954 et 4314
await queue_push_job(job_id, job_data["priority"])

queue_push_job utilise lpush (ligne 1549), qui ajoute en tete de liste (gauche). queue_pop_job utilise rpop (ligne 1557), qui retire de la fin (droite). Donc un lpush ajoute le job au debut de la queue, pas a la fin. Un job requeue sera donc traite en dernier (il est pousse a gauche, et les jobs sont depiles par la droite).

Cela signifie qu'un job requeue perd sa position FIFO et passe derriere tous les jobs en attente dans la meme queue de priorite. Ce n'est pas necessairement un bug (on pourrait argumenter que c'est le comportement souhaite pour eviter la starvation), mais c'est contraire a l'intuition d'un requeue.


F16 — queue_remove_job n'est pas atomique

Classification : RISQUE THEORIQUE Fichier : main.py, lignes 1563-1573 Impact : Race condition lors de l'annulation d'un job pendant que le worker pop

async def queue_remove_job(job_id: str):
    for queue_key in [QUEUE_HIGH, QUEUE_NORMAL, QUEUE_LOW]:
        items = await redis_client.lrange(queue_key, 0, -1)
        for item in items:
            data = json.loads(item)
            if data.get("job_id") == job_id:
                await redis_client.lrem(queue_key, 1, item)
                return True
    return False

Entre le lrange et le lrem, le worker peut avoir rpop le job. Dans ce cas, lrem ne trouve rien (job deja retire par rpop), mais retourne sans erreur. Le job est alors en cours d'execution alors que l'utilisateur croit l'avoir annule. Le db_cancel_job (ligne 4436) ne changera pas le status non plus car process_single_job a deja change le status de "pending" a "loading_model" ou "generating".

Mitigation partielle : cancel_job (ligne 4429) verifie job["status"] != "pending" avant d'essayer. Mais il y a une fenetre entre le check status et le rpop du worker.


F17 — Le watchdog idle peut tuer un job legitime en phase de loading

Classification : RISQUE THEORIQUE Fichier : main.py, lignes 4193-4222 Impact : Faux positif du watchdog

Le watchdog a une grace period de 90 secondes (ligne 4202). Mais certains outils ont un startup_time de 180 secondes (ComfyUI, Fooocus, Bark, MusicGen). Si un tool met 3 minutes a charger son modele en VRAM, le watchdog (qui verifie gpu_util == 0 and vram_used < 500) pourrait declencher apres 90s + 3 cycles (~100s).

En pratique, les exclusive tools en phase de loading auraient deja du VRAM > 500MB avant le timeout du watchdog (le modele se charge incrementalement). Mais pour des outils qui font un pre-processing CPU avant de toucher au GPU, le watchdog pourrait tuer le job.


F18 — run_ssh_command ne nettoie pas le process en cas de timeout

Classification : RISQUE THEORIQUE Fichier : main.py, lignes 810-826 Impact : Processus SSH zombie

async def run_ssh_command(command: str, timeout: int = SSH_TIMEOUT) -> tuple[int, str, str]:
    try:
        proc = await asyncio.create_subprocess_shell(ssh_cmd, ...)
        stdout, stderr = await asyncio.wait_for(proc.communicate(), timeout=timeout)
    except asyncio.TimeoutError:
        return -1, "", "SSH command timeout"
    # proc.kill() n'est JAMAIS appele en cas de timeout

Quand wait_for leve TimeoutError, le process SSH est toujours en cours d'execution. Il n'est pas tue. Il finira par timeout cote SSH (ConnectTimeout=10s), mais pendant ce temps il occupe un descripteur de fichier et un process.


F19 — set_parallel_jobs : le commentaire dit "Clamp between 1 and 5" mais le code clamp entre 1 et 10

Classification : RISQUE THEORIQUE (documentation incorrecte) Fichier : main.py, ligne 1139 Impact : Confusion, pas d'impact fonctionnel

count = max(1, min(10, count))  # Clamp between 1 and 5

Le commentaire dit 5, le code dit 10. L'endpoint docstring (ligne 1135) dit aussi "1-5". Un utilisateur qui lit la doc pourrait croire que 10 est rejete.


F20 — Pas de validation job_type par rapport aux schemas

Classification : RISQUE THEORIQUE Fichier : main.py, lignes 4336-4357 Impact : Un job avec un job_type invalide est accepte, queue, et echoue silencieusement

@app.post("/api/jobs")
async def create_job(job: JobCreate):
    if job.tool_id not in TOOLS_CONFIG:
        raise HTTPException(status_code=400, detail=f"Unknown tool: {job.tool_id}")
    # Pas de validation de job_type !
    job_id = await db_create_job(job)
    await queue_push_job(job_id, job.priority)

Le tool_id est valide, mais job_type ne l'est pas. Un job {"tool_id": "yolo", "job_type": "inexistant"} sera accepte, queue, puis quand execute_job est appele, il passera par le branch YOLO (ligne 1660) qui utilise job_type comme _task dans les params. Le job sera envoye au service YOLO avec task=inexistant, et c'est le service YOLO qui renverra une erreur, pas l'orchestrateur.

Le job sera marque "failed" avec le message d'erreur du service, ce qui est acceptable. Mais il aurait ete mieux de rejeter a la creation.


F21 — Path traversal partiel sur les endpoints de fichiers

Classification : RISQUE THEORIQUE (securite) Fichier : main.py, lignes 4447-4449, 4515 Impact : Acces a des fichiers hors du repertoire prevu

# Ligne 4447-4449 (download_output)
if ".." in filename or "/" in filename:
    raise HTTPException(status_code=400, detail="Invalid filename")

# Ligne 4515 (upload)
filename = os.path.basename(file.filename)
if not filename or ".." in filename:
    raise HTTPException(status_code=400, detail="Invalid filename")

La protection contre le path traversal verifie ".." et "/", mais PAS "\" (backslash Windows). En pratique, le serveur tourne sur Linux dans un container, donc \ n'est pas un separateur de chemin. Et os.path.basename sur Linux ne split pas sur \. Donc pas d'exploitation possible actuellement.

Pour job_id dans l'URL (/api/jobs/{job_id}/output/{filename}), il n'y a aucune validation. Si job_id contient ../../../etc, le path serait {OUTPUTS_PATH}/../../../etc/{filename}. Mais FastAPI valide les segments d'URL et n'accepte pas les / dans {job_id} (path parameter), donc ../ ne peut pas etre dans job_id.

Conclusion : pas exploitable en l'etat, mais la validation est incomplete et fragile.


F22 — Le anthropic tool a port=0 et gpu_id=-1, mais le code assume gpu_id >= 0 a plusieurs endroits

Classification : RISQUE THEORIQUE Fichier : main.py, lignes 425-434, 4003, 4007, 4117 Impact : Aucun actuellement grace aux guards is_cloud

# Ligne 429
"anthropic": Tool(id="anthropic", ..., port=0, gpu_id=-1, vram="0", status=ToolStatus.RUNNING)

# Ligne 4003
if not is_cloud and tool and tool.gpu_id is not None and tool.gpu_id >= 0:
    active_tools[tool.gpu_id] = tool_id

# Ligne 4117
if not is_cloud and tool and tool.gpu_id is not None and tool.gpu_id >= 0:

Les guards not is_cloud protegent correctement. Mais gpu_id=-1 est un magic number non documente. Si un nouveau tool cloud est ajoute sans mettre gpu_id=-1, les guards is_cloud (bases sur le set CLOUD_TOOLS) le protegeraient quand meme. Pas de bug reel, mais code fragile.


4. ANALYSE DES DATA FLOWS

4.1 Flow principal : Creation -> Queue -> Worker -> Execution -> Completion

Client POST /api/jobs
  -> create_job() [L4336]
    -> db_create_job() [L1317] : INSERT PostgreSQL (status=pending)
    -> queue_push_job() [L1540] : LPUSH Redis
  <- {job_id, status: pending}

Worker loop [L4180+]
  -> queue_pop_job() [L1553] : RPOP Redis (priority order: high, normal, low)
  -> db_get_job() [L1333] : SELECT PostgreSQL
  -> determine_preferred_gpu() [L3847] : choix GPU base sur VRAM estimee
  -> can_dispatch() [L4149] : verifie slots GPU
  -> process_single_job() [L3918] via asyncio.create_task
    -> check_tool_health() [L829] : curl vers win11
    -> start_tool() [L980] si necessaire : SSH vers win11
    -> db_update_job_status("loading_model") [L3983]
    -> db_update_job_status("generating") [L4008]
    -> execute_job() [L1610] : dispatch vers la bonne fonction
      -> execute_ollama_job() / execute_comfyui_job() / etc.
        -> HTTP vers win11:PORT (Ollama, ComfyUI, etc.)
        -> ou SSH + CLI (Wan2.1, SadTalker, FaceFusion)
    -> db_update_job_status("completed", output_result=...) [L4056]
    -> send_webhook_callback() [L4080] si callback_url

Points de defaillance identifies :

  1. Entre RPOP Redis et db_get_job, si DB down : job perdu (F05)
  2. Entre check_anthropic_budget et cost_usd write : race condition budget (F02)
  3. Pas de transaction DB pour le cycle status pending->loading->generating->completed
  4. Chaque changement de status est un UPDATE independant, pas d'optimistic locking

4.2 Flow d'erreur : Job failure

process_single_job() try block raises Exception
  -> except [L4091]
    -> db_update_job_status("failed", error_message=...) [L4096]
    -> send_webhook_callback() [L4104] si callback_url
  -> finally [L4113]
    -> cleanup active_tools pour CLI tools

Points de defaillance :

  1. Si db_update_job_status("failed") echoue (DB down), le job reste en status "generating" indefiniment. Au prochain restart, le lifespan le passera en "failed" (stale cleanup, ligne 122-128).
  2. Le finally ne nettoie active_tools que pour les CLI tools (sadtalker, wan21, facefusion). Les server tools restent marques comme actifs.

4.3 Flow upload -> execution

POST /api/upload
  -> file.read() (tout en memoire) [L4525]
  -> write to INPUTS_PATH/{filename} [L4526]

POST /api/jobs (input_params.image_file = "myfile.jpg")
  -> queue -> worker
  -> execute_yolo_job()
    -> image_path = INPUTS_PATH + "/" + image_file [L3530]
    -> os.path.exists(image_path) check
    -> open(image_path, "rb") et envoi HTTP vers win11

Pas de lien formel entre upload et job : le filename est juste une string dans input_params. Deux jobs peuvent utiliser le meme fichier. Un delete supprime le fichier meme si des jobs pending l'utilisent.

4.4 Flow SSH vers win11

run_ssh_command(command) [L810]
  -> subprocess ssh -o ConnectTimeout=10 -o StrictHostKeyChecking=no user@host "command"
  -> wait_for(timeout=SSH_TIMEOUT=30s)
  -> return (returncode, stdout, stderr)

Risques :


5. ANALYSE DE LA GESTION D'ETAT

5.1 Etat PostgreSQL (source de verite pour les jobs)

Le schema n'est pas dans le code audite (pas de migration SQL), mais infere depuis les queries :

Coherence :

5.2 Etat Redis (queue de jobs + metadata ephemere)

Keys utilisees :

Coherence Redis <-> PostgreSQL :

5.3 Etat memoire (in-process, perdu au restart)

Au restart du container :

  1. tools_state est reinitialise a STOPPED pour tous les outils (ligne 734). Correct car on ne sait pas l'etat reel.
  2. active_tools est reinitialise a None pour tous les GPUs (ligne 735). Correct.
  3. Les jobs "generating" ou "loading_model" sont passes en "failed" par le stale cleanup (lignes 121-132). Correct.
  4. La pause Redis est effacee (ligne 117). Intentionnel pour le post-deploy.
  5. Le worker redemarre et pop les jobs pending de Redis. Correct.

Point critique : les jobs cloud fire-and-forget (F01) ne sont pas dans gpu_slots et ne sont pas nettoyees au shutdown.


6. ANALYSE DE LA PROPAGATION D'ERREURS

6.1 Erreurs SSH

Toutes les fonctions SSH retournent (-1, "", "message") en cas d'erreur. Les appelants verifient code != 0 et :

Verdict : Bonne resilience. Les erreurs SSH ne font pas crasher le service.

6.2 Erreurs DB

Les fonctions DB utilisent async with pg_pool.acquire() as conn, qui leve asyncpg.PostgresError ou asyncpg.InterfaceError si la DB est down. Ces exceptions ne sont PAS capturees dans les fonctions DB elles-memes. Elles propagent jusqu'a :

6.3 Erreurs d'execution de job

Toutes les fonctions execute_*_job levent des Exception avec des messages descriptifs. Ces exceptions sont capturees par process_single_job (ligne 4091) qui marque le job comme "failed" avec le message d'erreur.

Risque : si db_update_job_status("failed") echoue lui-meme (DB down), l'exception est propagee au worker qui la log et continue. Le job reste en status "generating" en DB. Il sera nettoye au prochain restart (stale cleanup).

6.4 Erreurs timeout

Les timeouts sont bien geres pour chaque phase :

Mais : quand wait_for timeout sur execute_job, le task est cancelled mais les sous-processes SSH ne sont pas tues (F18). Les fonctions Wan2.1 et FaceFusion ont leur propre cleanup (lignes 2248-2258, 2441-2448), mais les autres non.


7. SECURITE

7.1 Donnees sensibles dans les reponses API

7.2 Injection de commandes SSH

Les fonctions execute_wan21_job, execute_sadtalker_job, execute_facefusion_job utilisent SSH avec des commandes construites par concatenation de strings. Certains parametres viennent de input_params :

Verdict : Risque d'injection de commande faible mais present pour facefusion. Les autres tools utilisent base64 encoding ou des APIs HTTP.

7.3 SSRF via callback_url

send_webhook_callback (ligne 1491) envoie un POST vers callback_url fourni par l'utilisateur. Un attaquant pourrait :

Risque faible car le service est en LAN prive, mais c'est un SSRF classique.


8. SYNTHESE ET RECOMMANDATIONS

Findings par severite

Severite Count References
BUG CERTAIN 3 F01, F02, F04
BUG PROBABLE 6 F05, F06, F07, F08, F09, F10
RISQUE THEORIQUE 12 F03, F11-F22

Top 5 des corrections prioritaires

  1. F01 (CRITIQUE) — Tracker les tasks cloud dans une liste dediee. Ajouter une limite de concurrence (semaphore). Les annuler proprement au shutdown. Effort : 30min.

  2. F02 (CRITIQUE) — Implementer un mecanisme de budget reservation : avant d'envoyer la requete Anthropic, reserver le cout estime (INSERT un row "reserved" ou UPDATE un compteur Redis atomique). Apres la reponse, ajuster au cout reel. Effort : 1h.

  3. F05 (IMPORTANT) — Apres un rpop, wrapper la logique dans un try/except. En cas d'erreur DB, re-push le job dans Redis. Effort : 15min.

  4. F06 (IMPORTANT) — Ajouter un compteur de requeue. Apres N requeues (ex: 10), marquer le job comme failed avec un message explicite. Effort : 30min.

  5. F04 (CRITIQUE) — Pour le batch polling, creer un nouveau httpx client pour chaque poll, ou utiliser un timeout plus long. Gerer les codes HTTP 4xx differemment des 5xx. Effort : 30min.

Architecture : points forts

Architecture : points faibles