Date : 15/02/2026 04:00
Status : IMPLEMENTEE
Auditeur : Claude Opus 4.6 (Pass 3 independant)
Fichier principal : /stock_8to/33800-stack/projects/ai-orchestrator/app/main.py (4675 lignes)
Fichiers secondaires : Dockerfile, requirements.txt, docker-compose.yml, conf.prod.gouroubleu.yml, .env.deploy, .gitlab-ci.yml
| Severite | Nombre |
|---|---|
| BUG CERTAIN | 6 |
| BUG PROBABLE | 11 |
| RISQUE THEORIQUE | 12 |
| TOTAL | 29 |
Les problemes les plus graves concernent :
app/main.py:4422-4441 et app/main.py:4248-4261DELETE /api/jobs/{job_id} verifie job["status"] == "pending" puis fait queue_remove_job() puis db_cancel_job(). Mais entre le moment ou le worker fait queue_pop_job() (ligne 4248) et le moment ou il verifie job_data["status"] != "pending" (ligne 4261), le cancel peut avoir retire le job de la queue et mis le statut a "cancelled" en DB. Le worker a deja poppe le message Redis, donc queue_remove_job ne le trouvera pas. Le worker procedera avec un job dont le statut DB est "cancelled", car il a lu le status AVANT le cancel.
# Worker (ligne 4248):
job_msg = await queue_pop_job() # Job removed from Redis
# ... entre ces deux lignes, cancel peut executer
job_data = await db_get_job(job_id) # Lit le status
if job_data["status"] != "pending": # Pourrait etre "cancelled" ici
continue # OK dans ce cas
- **Impact** : Un job "cancel" peut continuer a s'executer. L'API retourne "failed" pour le cancel alors que le job a deja commence, sans indication a l'utilisateur.
- **Fix suggere** : Dans `process_single_job`, re-verifier le status DB juste avant de passer a "loading_model/generating". Ajouter un mecanisme de cancellation token (Redis key `ai:cancel:{job_id}`) que le worker verifie periodiquement.
---
### F-02 : Jobs requeue cree des doublons en Redis si le worker crash entre pop et requeue
- **Severite** : BUG CERTAIN
- **Fichier** : `app/main.py:3952-3954` et `app/main.py:3978-3980` et `app/main.py:3991-3994`
- **Description** : Le pattern requeue fait `db_update_job_status(job_id, "pending")` puis `queue_push_job(job_id, ...)`. Si le worker crash (OOM, container restart) apres le db_update mais avant le queue_push, le job reste "pending" en DB mais absent de Redis. Inversement, si le crash arrive apres le push, le job est "pending" en DB et present en Redis. L'endpoint `/api/queue/sync` resout le cas "absent de Redis" mais ne deduplique pas les entrees Redis existantes avant de pousser les jobs.
- **Code concerne** :
```python
# Requeue pattern (3 occurrences):
await db_update_job_status(job_id, "pending") # Etape 1: DB
await queue_push_job(job_id, job_data["priority"]) # Etape 2: Redis
# Si crash entre 1 et 2 -> job perdu dans les limbes
/api/queue/sync n'est pas appele manuellement. Apres un sync, des doublons possibles (un job traite deux fois).queue_sync devrait TOUJOURS etre execute au demarrage (il ne l'est pas dans lifespan). 2) Ajouter une deduplication dans le worker : avant de traiter un job poppe, verifier qu'il est toujours "pending" en DB (deja fait ligne 4261, mais un meme job_id pourrait etre poppe par deux cycles worker differents si duplique en Redis).app/main.py:35-40REDIS_PASSWORD = os.environ.get("REDIS_PASSWORD", "urefsLoXibTZ36mtuygcLyNuqzOmIVkF")
POSTGRES_PASSWORD = os.environ.get("POSTGRES_PASSWORD", "xgvKM65PC3zWspTCXkbxPx2fx5fOHRVK")
docker/secrets/prod.env uniquement.app/main.py:1326-1329 et app/main.py:1401-1478ai_jobs (cost_usd, provider, tokens_input, tokens_output, cache_read_tokens, cache_creation_tokens) dans db_update_job_status sans qu'aucun fichier de migration ou schema SQL n'existe dans le repo. Si la table ne contient pas ces colonnes, les UPDATE echoueront silencieusement (asyncpg leve une exception, mais elle est capturee par le try/except global dans process_single_job).# db_update_job_status references ces colonnes (ajoutees pour le tracking Anthropic):
elif key == "cost_usd":
updates.append(f"cost_usd = ${idx}")
elif key == "provider":
updates.append(f"provider = ${idx}")
elif key == "tokens_input":
updates.append(f"tokens_input = ${idx}")
# ... etc
# Aucun schema SQL dans le repo pour verifier que ces colonnes existent
ai_jobs a partir du repo. Un deploiement sur une nouvelle instance DB echouera. Les colonnes de tracking de cout pourraient ne pas exister et les jobs echoueraient tous au moment du completed update.schema.sql ou un systeme de migrations (Alembic) au repo. Ajouter une verification des colonnes au demarrage dans lifespan.app/main.py:4336-4357POST /api/jobs fait db_create_job() puis queue_push_job() en deux operations independantes. Si Redis est down ou la connexion echoue entre les deux, le job existe en DB avec status "pending" mais n'est pas dans la queue Redis. Il ne sera jamais traite.@app.post("/api/jobs")
async def create_job(job: JobCreate):
job_id = await db_create_job(job) # DB OK
await queue_push_job(job_id, job.priority) # Redis fail -> job perdu
stats = await queue_get_stats()
return {"job_id": job_id, "status": "pending", "queue_position": stats.total_pending}
queue_push_job echoue, soit annuler le job en DB, soit retourner une erreur 503. Ou ajouter queue_sync au demarrage (voir F-02).app/main.py:1059-1062/health fait await redis_client.get(WORKER_PAUSED_KEY) sans try/except. Si Redis est down, cette requete leve une exception non capturee, et le health check retourne un 500. Comme le health check est utilise par Docker (healthcheck dans docker-compose.yml), le container sera redemarrer apres 3 echecs. Si Redis est temporairement indisponible, cela cree un cycle de redemarrage.@app.get("/health")
async def health():
is_paused = await redis_client.get(WORKER_PAUSED_KEY) # Crash si Redis down
return {"status": "healthy", "active_tools": active_tools, "worker_paused": is_paused == "true"}
"worker_paused": "unknown" si Redis indisponible, et ne pas retourner un 500 pour le healthcheck.app/main.py:4267-4271asyncio.create_task() sans etre ajoute a gpu_slots ni a aucune structure de tracking. Le task retourne est ignore (variable locale, pas stockee). Si 100 jobs Anthropic sont en queue, ils seront TOUS dispatches immediatement dans la meme iteration de boucle (le worker fait continue et poppe le suivant).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 # Pas de tracking, pas de limite, pas d'await
app/main.py:2895-2906check_anthropic_budget() lit le cout mensuel en DB et compare au budget. Mais si N jobs Anthropic sont dispatches en parallele (voir F-07), ils verifient tous le budget AVANT que l'un d'eux ait ecrit son cout en DB (le cout n'est ecrit qu'a la completion, pas au demarrage). Avec N jobs simultanes, chacun voit le meme budget restant.async def check_anthropic_budget() -> None:
monthly_cost = await db_get_monthly_cost("anthropic") # Lit uniquement les jobs 'completed'
if monthly_cost >= ANTHROPIC_BUDGET_USD:
raise Exception(...)
# Le cost_usd n'est ecrit qu'a db_update_job_status(..., cost_usd=job_cost) apres completion
run_ssh_command ne kill pas le process en cas de timeoutapp/main.py:810-826TimeoutError dans run_ssh_command, le code retourne (-1, "", "SSH command timeout") mais ne kill() pas le sous-process. Le process SSH reste orphelin en arriere-plan.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)
return proc.returncode, stdout.decode(), stderr.decode()
except asyncio.TimeoutError:
return -1, "", "SSH command timeout" # proc n'est PAS kill
except Exception as e:
return -1, "", str(e) # proc n'est PAS kill
proc.kill() et await proc.wait() dans le bloc except TimeoutError.check_tool_health bare except avale toutes les exceptionsapp/main.py:829-842check_tool_health a un except: sans type qui avale TOUTES les exceptions y compris KeyboardInterrupt, SystemExit, et asyncio.CancelledError. Cela empeche le shutdown propre du container.async def check_tool_health(tool: Tool) -> bool:
try:
proc = await asyncio.create_subprocess_shell(...)
stdout, _ = await asyncio.wait_for(proc.communicate(), timeout=10)
return proc.returncode == 0 and len(stdout) > 0
except: # Attrape TOUT, y compris CancelledError
return False
CancelledError est avalee, ce qui casse le protocole de cancellation asyncio.except: par except Exception: pour laisser passer CancelledError, KeyboardInterrupt, SystemExit.start_tool escape des backslashes casse les chemins Windowsapp/main.py:1023-1024start_tool fait escaped_cmd = final_cmd.replace('\\', '\\\\'). Le final_cmd contient deja tool.start_cmd qui utilise des raw strings (r"C:\Users\gouro\start_ollama_gpu0.bat"). Le replace double tous les backslashes. Ensuite, le tout est envoye via SSH dans des double-quotes, ce qui peut casser l'interpretation du chemin Windows.final_cmd = f'set CUDA_VISIBLE_DEVICES={gpu_id} && {tool.start_cmd}'
# final_cmd = "set CUDA_VISIBLE_DEVICES=0 && C:\Users\gouro\start_ollama_gpu0.bat"
escaped_cmd = final_cmd.replace('\\', '\\\\')
# escaped_cmd = "set CUDA_VISIBLE_DEVICES=0 && C:\\Users\\gouro\\start_ollama_gpu0.bat"
ssh_cmd = f'ssh -f ... "{escaped_cmd}"'
# Le SSH passe "C:\\Users\\gouro\\start_ollama_gpu0.bat" au shell Windows
# cmd.exe interprete \\ comme un seul \ -> OK dans ce cas
cmd /c, le double backslash peut fonctionner dans la plupart des cas. Mais si le path contient deja des doubles backslashes (par erreur dans la config), ils seraient quadruples. Le vrai risque est avec run_ssh_command (ligne 813) qui fait son propre escaping de double-quotes, et start_tool qui bypass run_ssh_command avec son propre SSH call - incoherence d'escaping entre les deux.run_ssh_command, soit documenter clairement les conventions d'escaping.stop_tool ne gere pas les outils cloud (gpu_id=-1)app/main.py:945-977stop_tool appelle check_port_listening(tool.port) pour l'outil "anthropic" qui a port=0 et gpu_id=-1. Cela lance une commande SSH vers win11 pour tester le port 0, qui echouera. De plus, active_tools.get(tool.gpu_id) avec gpu_id=-1 ne correspond a aucune cle dans active_tools (qui n'a que 0 et 1).# L'outil anthropic a: port=0, gpu_id=-1
async def stop_tool(tool_id: str) -> bool:
tool = TOOLS_CONFIG.get(tool_id)
# Pour anthropic:
kill_by_port_cmd = f"... Get-NetTCPConnection -LocalPort 0 ..." # Port 0 -> erreur
await run_ssh_command(kill_by_port_cmd) # SSH inutile
# ...
if not await check_port_listening(tool.port): # Port 0 -> False
tools_state[tool_id] = ToolStatus.STOPPED
if active_tools.get(tool.gpu_id) == tool_id: # gpu_id=-1 -> None
active_tools[tool.gpu_id] = None # Cree une nouvelle cle -1 dans active_tools
POST /api/tools/anthropic/stop, des commandes SSH inutiles sont envoyees a win11. L'entree active_tools[-1] est creee, polluant le dictionnaire. L'endpoint POST /api/tools/stop-all (ligne 1246-1254) itere tools_state et peut tenter de stopper "anthropic".stop_tool et start_tool pour les outils cloud : if tool_id in CLOUD_TOOLS: return True. Filtrer les cloud tools dans stop-all.refresh_tools_status inclut anthropic dans les health checks SSHapp/main.py:1207-1243refresh_tools_status itere TOUS les TOOLS_CONFIG y compris "anthropic" (qui n'a pas de health_url). check_tool_health retourne False pour health_url=None, puis check_port_listening(0) est appele (SSH inutile). Puis le status est mis a STOPPED pour anthropic.async def check_single_tool(tool_id: str, tool: Tool) -> tuple:
is_healthy = await check_tool_health(tool) # False pour anthropic (no health_url)
port_open = await check_port_listening(tool.port) # SSH pour port 0
if port_open:
return (tool_id, ToolStatus.STARTING)
return (tool_id, ToolStatus.STOPPED) # anthropic marqué STOPPED
refresh_tools_status. Ou ajouter if not tool.health_url and tool.port == 0: return (tool_id, tools_state.get(tool_id, ToolStatus.STOPPED)).app/main.py:4206-4207int(gpu.vram_used.replace(" Mo", "")) et int(gpu.gpu_util.replace("%", "")). Si le GPU est unreachable, get_gpu_status retourne un fallback avec vram_used="N/A" et gpu_util="N/A". Le replace(" Mo", "") produit "N/A", puis int("N/A") leve un ValueError. Ce ValueError est capture par le except Exception: pass global (ligne 4221), donc le watchdog est silencieusement desactive.gpu = await asyncio.wait_for(get_gpu_status(gpu_id), timeout=5)
vram_used = int(gpu.vram_used.replace(" Mo", "")) # "N/A" -> int("N/A") -> ValueError
gpu_util = int(gpu.gpu_util.replace("%", "")) # idem
except le gere), mais c'est une logique morte qui pollue les logs en mode debug et montre un manque de robustesse dans le parsing.if "N/A" in gpu.vram_used: continue avant le parsing.Severite : BUG PROBABLE
Fichier : app/main.py:2994-3065
Description : Le httpx.AsyncClient pour le batch est cree avec timeout=httpx.Timeout(30.0, connect=10.0). Le poll (client.get(...)) a l'interieur de la boucle for utilise le meme client avec ce timeout de 30s par requete, ce qui est correct. MAIS : le async with httpx.AsyncClient(...) garde la connexion ouverte pendant toute la duree du polling (jusqu'a 2 heures = 240 polls * 30s). Certains proxies ou load balancers ferment les connexions idle apres un certain temps.
De plus, le premier appel client.post(...) soumet le batch, mais si un des polls intermediaires retourne une erreur non-200, le code fait continue sans logger ni verifier s'il s'agit d'une erreur fatale (ex: batch_id invalide, auth revokee).
Code concerne :
async with httpx.AsyncClient(timeout=httpx.Timeout(30.0, connect=10.0)) as client:
resp = await client.post(...) # Submit
batch_id = batch_data.get("id")
for i in range(240):
await asyncio.sleep(30)
poll_resp = await client.get(...)
if poll_resp.status_code != 200:
continue # Ignore TOUTES les erreurs de poll
Impact : Si Anthropic retourne un 401 (API key revokee) ou 404 (batch expire), le code poll en boucle pendant 2 heures avant de timeout. Gaspillage de ressources.
Fix suggere : Distinguer les erreurs fatales (401, 404) des erreurs transitoires (500, 429). Logger chaque erreur. Considerer un nouveau client par poll pour eviter les problemes de connexion longue.
app/main.py:176-182Access-Control-Allow-Origin: * combine avec Access-Control-Allow-Credentials: true. Les navigateurs modernes refusent les reponses dans cette configuration.app.add_middleware(
CORSMiddleware,
allow_origins=["*"],
allow_credentials=True, # Incompatible avec origins=["*"]
allow_methods=["*"],
allow_headers=["*"],
)
*, mais le comportement est non-standard et peut varier selon les versions.allow_origins=["https://dashboard.33800.nowhere84.com"]), soit mettre allow_credentials=False.is_cloud re-declaree dans le meme scope, masque la valeur initialeapp/main.py:3943 et app/main.py:4045is_cloud est calculee a la ligne 3943 (is_cloud = tool_id in CLOUD_TOOLS). Puis elle est recalculee a la ligne 4045 (is_cloud = tool_id in CLOUD_TOOLS) avec la meme valeur. Bien que le resultat soit identique, la re-declaration au milieu du code est confuse et suggere un refactoring incomplet.is_cloud = tool_id in CLOUD_TOOLS # Ligne 3943
# ... 100 lignes plus tard ...
is_cloud = tool_id in CLOUD_TOOLS # Ligne 4045 - idem, re-declaration inutile
active_tools est un dict mutable partage entre coroutines sans lockapp/main.py:735 + multipleactive_tools est un dict global mutable accede en lecture/ecriture par le worker, les endpoints API (refresh_tools_status, start_tool, stop_tool, health), et les tasks de processing. Python asyncio est single-threaded, donc il n'y a PAS de race condition stricto sensu. Mais certaines operations (iterate + modify dans refresh_tools_status) pourraient poser probleme si une coroutine cede le controle au milieu d'une iteration._gpu_cache n'est pas thread-safe et pas async-safe entre worker et APIapp/main.py:931-942_gpu_cache) est un dict global avec un TTL de 30s. Le worker et les endpoints API utilisent get_cached_gpu_status(). Si deux coroutines appellent simultanement (worker + API request), elles peuvent toutes les deux voir le cache expire et faire deux SSH calls.queue_sync efface les queues Redis SANS verifier si le worker est en train de traiterapp/main.py:4467-4506queue_sync fait redis_client.delete(QUEUE_HIGH, QUEUE_NORMAL, QUEUE_LOW) puis re-pousse les jobs "pending". Si le worker a JUSTE poppe un job et est entre le pop et le premier db_update, ce job sera re-ajoute a Redis par le sync. Le worker le traitera en parallele avec le meme job poppe precedemment.WORKER_PAUSED_KEY) avant le sync, attendre la fin des jobs en cours, syncer, puis resumer.docker-compose.ymlREDIS_HOST, REDIS_PASSWORD, POSTGRES_HOST, POSTGRES_PASSWORD. Le container se connecte aux IPs hardcodees dans le code (192.168.1.12). Mais le conf.prod.gouroubleu.yml ne les definit pas non plus. Le container depend de smart-deploy pour injecter ces variables, mais il n'y a aucune trace de ces variables dans la config.REDIS_HOST, REDIS_PASSWORD, POSTGRES_HOST, POSTGRES_PASSWORD dans conf.prod.gouroubleu.yml section env:.conf.prod.gouroubleu.yml definit ANTHROPIC_API_KEY videconf.prod.gouroubleu.yml:18ANTHROPIC_API_KEY: "" est defini vide dans la config. Le code fait os.environ.get("ANTHROPIC_API_KEY", ""). Si la cle n'est pas injectee par un autre mecanisme (secret, env file), tous les jobs Anthropic echoueront avec "ANTHROPIC_API_KEY not configured".app/main.py:4511-4537POST /api/upload fait content = await file.read() sans verifier la taille du fichier. Un fichier de plusieurs Go remplira la RAM du container.db_update_job_status construit du SQL dynamique sans parameterisation des noms de colonnesapp/main.py:1401-1478$1, $2...), les NOMS de colonnes viennent des kwargs keys qui sont controlees par le code interne. Ce n'est pas une injection SQL car les cles sont hardcodees dans les if/elif. Mais le pattern est fragile : toute nouvelle colonne ajoutee necessite un nouveau elif.execute_anthropic_batch_job ne verifie pas le type de reponse "succeeded" vs "errored"app/main.py:3046-3058result_data.get("result", {}).get("message", {}). Mais selon l'API Anthropic Batch, chaque result a un result.type qui peut etre "succeeded" ou "errored". Le code ne verifie pas ce type et tente de convertir un resultat d'erreur comme s'il etait une reponse reussie.first_line = results_resp.text.strip().split("\n")[0]
result_data = json.loads(first_line)
response_data = result_data.get("result", {}).get("message", {})
# Si result.type == "errored", result.message n'existe pas -> response_data = {}
result = _convert_anthropic_to_ollama(response_data, model)
# -> retourne un message vide avec 0 tokens
result_data["result"]["type"]. Si "errored", lever une exception avec le message d'erreur.sadtalker cleanup async fire-and-forget sans awaitapp/main.py:3823-3824asyncio.create_task(asyncio.create_subprocess_shell(cleanup_cmd)) cree une coroutine qui cree un sous-processus. Mais le task n'est ni stocke ni awaite. Si le programme s'arrete, le cleanup ne se fait pas. De plus, create_subprocess_shell retourne un Process, pas un resultat — la coroutine se termine immediatement apres avoir spawne le process, sans attendre sa completion.cleanup_cmd = f'ssh ... "del /q ..."'
asyncio.create_task(asyncio.create_subprocess_shell(cleanup_cmd))
# Le task spawne le process SSH puis se termine
# Mais personne n'attend la fin du process SSH
I:\SadTalker\inputs\.proc = await create_subprocess_shell(...) puis await proc.wait(), et passer cette fonction a create_task.python:3.12-slim mais requirements.txt n'a pas de version pinned pour redis et asyncpgrequirements.txt:6-7 et Dockerfile:1redis>=5.0.0 et asyncpg>=0.29.0 et aiohttp>=3.9.0 utilisent des bornes minimales, pas des versions exactes. Un rebuild du Docker image pourra installer des versions differentes a des moments differents, causant des comportements imprevisibles.redis==5.2.1, asyncpg==0.30.0, aiohttp==3.11.0 etc.) ou generer un requirements.lock.set_parallel_jobs clamp "between 1 and 5" mais le code clamp a 10app/main.py:1133-1139max(1, min(10, count)) ce qui autorise jusqu'a 10.@app.post("/api/worker/parallel")
async def set_parallel_jobs(count: int = 2):
"""Set number of parallel Ollama jobs (1-5) # <-- doc dit 5
...
"""
count = max(1, min(10, count)) # <-- code dit 10
costs_summary a des parametres par defaut hardcodes (year=2026, month=2)app/main.py:4595async def costs_summary(period: str = "month", year: int = 2026, month: int = 2) a des defaults pour year et month qui correspondent a la date de developpement, pas a "le mois courant". Si l'endpoint est appele sans parametres en mars 2026, il retourne les couts de fevrier.@app.get("/api/costs/summary")
async def costs_summary(period: str = "month", year: int = 2026, month: int = 2):
year: int = None, month: int = None puis calculer la date courante si None.aiohttp.ClientConnectorError dans execute_ollama_jobactive_tools[gpu_id] reste sur "ollama" (pas de cleanup pour les server tools)active_tools est desynchronisepg_pool.acquire() leve une exception/api/queue/sync/health crash (F-06) -> Docker restart le containerPOST /api/jobsrun_ssh_command timeout apres 30s -> process zombie (F-09)check_tool_health retourne False -> tools marques STOPPEDschema.sql avec la definition de ai_jobs et une verification au demarragequeue_sync au demarrage dans la fonction lifespanrun_ssh_commandstop_tool, refresh_tools_status, et stop-all