Spaces:
Running
Running
File size: 19,053 Bytes
24480a0 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 | """
backend/api/execution_fabric.py β Execution Fabric (ARCH-K2.5 + ARCH-K2.6 Health Manager)
CRIT-B fix: circuit breaker Oracle (3 errori β blocco 300s) + fallback Railway con auth completa.
CRIT-C fix: parser 402/429 robusto β body scan + pattern estesi nell'exception handler.
"""
from __future__ import annotations
import asyncio
import logging
import os
import time
import uuid
from enum import Enum
from typing import Any, Optional
import httpx
from fastapi import APIRouter, Depends
from pydantic import BaseModel, Field
from .auth_guard import AuthRole, require_role
from .token_rotator import rotator as _token_rotator
try:
from .telemetry import record_kernel_event as _rke
except Exception:
def _rke(*_a, **_kw): pass
_logger = logging.getLogger("api.execution_fabric")
# ββ Configurazione Timing (ARCH-T1.1) ββββββββββββββββββββββββββββββββββββββββββ
_DEFAULT_TIMEOUT_S = 30.0 # Timeout standard per chiamate API
_ORACLE_TIMEOUT_S = 60.0 # Oracle ha piΓΉ tempo per il calcolo pesante
_HF_WARMUP_S = 10.0 # Tempo di attesa se lo Space Γ¨ in sleep
_RETRY_DELAY_S = 1.5 # Attesa tra i tentativi di rotazione token
_HEALTH_CHECK_INT = 120.0 # Intervallo health check in background
# ββ CRIT-B: Circuit Breaker Oracle βββββββββββββββββββββββββββββββββββββββββββββ
_ORACLE_CB_THRESHOLD = 3 # errori consecutivi prima di aprire il circuit
_ORACLE_CB_TIMEOUT_S = 300.0 # secondi di blocco dopo apertura (5 min)
# ββ CRIT-C: pattern quota/rate-limit (body + eccezioni) βββββββββββββββββββββββ
_QUOTA_PATTERNS = (
"quota", "egress", "rate limit", "rate_limit", "billing",
"402", "429", "limit exceeded", "credits", "insufficient",
"payment", "upgrade", "hours", "quota_or_ratelimit",
)
class ProviderKind(str, Enum):
HF_SPACE = "hf_space"
# RAILWAY rimosso
DOCKER = "docker"
ORACLE = "oracle"
LOCAL = "local"
class ProviderHealth(str, Enum):
OK = "ok"
DEGRADED = "degraded"
DOWN = "down"
UNKNOWN = "unknown"
class AlwaysOn(str, Enum):
YES = "yes"
ON_DEMAND = "on-demand"
NO = "no"
class ProviderSpec(BaseModel):
provider_id: str
name: str
kind: ProviderKind = ProviderKind.LOCAL
base_url: str = ""
capabilities: list[str] = Field(default_factory=list)
gpu: bool = False
always_on: AlwaysOn = AlwaysOn.ON_DEMAND
max_concurrency: int = 10
priority: int = 50
cost_unit: float = 0.0
region: str = "eu"
timeout: float = _DEFAULT_TIMEOUT_S
class ProviderState(BaseModel):
provider_id: str
health: ProviderHealth = ProviderHealth.UNKNOWN
active_tasks: int = 0
last_check_ts: float = 0.0
error_count: int = 0
class DispatchRequest(BaseModel):
capability: str
task_id: str = Field(default_factory=lambda: str(uuid.uuid4()))
payload: dict = Field(default_factory=dict)
require_gpu: bool = False
require_isolation: bool = False
prefer_region: str | None = None
provider_hint: str | None = None
correlation_id: str = Field(default_factory=lambda: str(uuid.uuid4()))
class DispatchResult(BaseModel):
task_id: str
provider_id: str
provider_name: str
provider_kind: str
capability: str
status: str
response: dict = Field(default_factory=dict)
latency_ms: float = 0.0
correlation_id: str = ""
class ExecutionFabric:
def __init__(self):
self._specs: dict[str, ProviderSpec] = {}
self._states: dict[str, ProviderState] = {}
self._initialized = False
# CRIT-B: circuit breaker state per provider (in-memory, per istanza)
# { provider_id: {"errors": int, "blocked_until": float} }
self._circuit: dict[str, dict] = {}
# ββ Inizializzazione provider βββββββββββββββββββββββββββββββββββββββββββββ
async def initialize(self):
if self._initialized:
return
self._initialized = True
providers: list[ProviderSpec] = []
# Local fallback β sempre disponibile, nessun HTTP
providers.append(ProviderSpec(
provider_id="local",
name="Local",
kind=ProviderKind.LOCAL,
base_url="",
capabilities=["exec", "tool", "shell"],
priority=10,
always_on=AlwaysOn.YES,
))
# HF Space principale β no-op se HF_SPACE_URL non configurata
hf_url = os.getenv("HF_SPACE_URL", "").rstrip("/")
if hf_url:
providers.append(ProviderSpec(
provider_id="hf-space-main",
name="HF Space Main",
kind=ProviderKind.HF_SPACE,
base_url=hf_url,
capabilities=["exec", "tool", "llm", "browse"],
priority=50,
always_on=AlwaysOn.ON_DEMAND,
timeout=_DEFAULT_TIMEOUT_S,
))
# Railway core backend rimosso (migrato su HF Spaces)
# Oracle Cloud VM β no-op se ORACLE_CLOUD_VM_URL non impostata
oracle_url = os.getenv("ORACLE_CLOUD_VM_URL", "").rstrip("/")
if oracle_url:
providers.append(ProviderSpec(
provider_id="oracle-cloud-vm-01",
name="Oracle Cloud VM",
kind=ProviderKind.ORACLE,
base_url=oracle_url,
capabilities=["exec", "tool", "llm", "gpu"],
gpu=True,
priority=80,
always_on=AlwaysOn.ON_DEMAND,
timeout=_ORACLE_TIMEOUT_S,
))
for spec in providers:
self._specs[spec.provider_id] = spec
self._states[spec.provider_id] = ProviderState(provider_id=spec.provider_id)
self._circuit[spec.provider_id] = {"errors": 0, "blocked_until": 0.0}
_logger.info(
"[fabric] Providers registrati: %s",
", ".join(f"{p.provider_id}({p.kind.value})" for p in providers)
)
# ββ CRIT-B: Circuit Breaker helpers ββββββββββββββββββββββββββββββββββββββ
def _is_circuit_open(self, provider_id: str) -> bool:
"""True se il provider Γ¨ in blocco circuit breaker."""
cb = self._circuit.get(provider_id, {})
blocked_until = cb.get("blocked_until", 0.0)
if blocked_until > time.time():
return True
# Reset automatico dopo il timeout
if blocked_until > 0.0:
self._circuit[provider_id]["errors"] = 0
self._circuit[provider_id]["blocked_until"] = 0.0
_logger.info("[fabric] Circuit breaker RESET per %s", provider_id)
return False
def _record_oracle_error(self, provider_id: str) -> bool:
"""
Registra un errore Oracle. Se si supera la soglia, apre il circuit.
Ritorna True se il circuit Γ¨ stato appena aperto.
"""
cb = self._circuit.setdefault(provider_id, {"errors": 0, "blocked_until": 0.0})
cb["errors"] += 1
if cb["errors"] >= _ORACLE_CB_THRESHOLD:
cb["blocked_until"] = time.time() + _ORACLE_CB_TIMEOUT_S
_logger.error(
"[fabric] Circuit breaker APERTO per %s (%d errori consecutivi) "
"β blocco per %.0fs fino a %s",
provider_id, cb["errors"], _ORACLE_CB_TIMEOUT_S,
time.strftime("%H:%M:%S", time.localtime(cb["blocked_until"]))
)
return True
return False
def _reset_oracle_errors(self, provider_id: str) -> None:
"""Azzera il contatore errori dopo un successo."""
if provider_id in self._circuit:
self._circuit[provider_id]["errors"] = 0
# ββ Selezione provider ββββββββββββββββββββββββββββββββββββββββββββββββββββ
def _cb_allow(self, provider_id: str) -> bool:
return True
def _select(self, capability: str, exclude: set[str] | None = None) -> Optional[str]:
"""
Seleziona il provider con la prioritΓ piΓΉ alta che:
- supporta la capability richiesta
- non Γ¨ DOWN
- non Γ¨ in circuit breaker aperto
- non Γ¨ nella lista exclude
Ordine: priority DESC (higher = preferred).
"""
exclude = exclude or set()
candidates = [
(spec.priority, pid, spec)
for pid, spec in self._specs.items()
if capability in spec.capabilities
and self._states[pid].health != ProviderHealth.DOWN
and not self._is_circuit_open(pid)
and pid not in exclude
]
if not candidates:
return None
# Ordina per priority decrescente
candidates.sort(key=lambda x: x[0], reverse=True)
return candidates[0][1]
# ββ CRIT-C: rilevamento quota/rate-limit ββββββββββββββββββββββββββββββββββ
@staticmethod
def _is_quota_error(text: str) -> bool:
"""True se il testo (body o eccezione) indica quota/rate-limit."""
t = text.lower()
return any(p in t for p in _QUOTA_PATTERNS)
# ββ Chiamata HTTP al provider βββββββββββββββββββββββββββββββββββββββββββββ
async def _call_provider(self, spec: ProviderSpec, req: DispatchRequest) -> dict:
timeout = _ORACLE_TIMEOUT_S if spec.kind == ProviderKind.ORACLE else spec.timeout
# Provider locale: nessuna chiamata HTTP
if spec.kind == ProviderKind.LOCAL:
return {"status": "ok", "provider": "local", "task_id": req.task_id}
if not spec.base_url:
raise ValueError(f"base_url non configurato per provider {spec.provider_id}")
# ββ Header auth per provider kind ββββββββββββββββββββββββββββββββββββ
headers: dict[str, str] = {"Content-Type": "application/json"}
internal_token = os.getenv("INTERNAL_TOKEN", "")
hf_token = os.getenv("HF_TOKEN", "")
# Railway auth rimosso
if spec.kind in (ProviderKind.ORACLE, ProviderKind.DOCKER):
if internal_token:
headers["X-Internal-Token"] = internal_token
payload = {
"task_id": req.task_id,
"capability": req.capability,
"payload": req.payload,
"correlation_id": req.correlation_id,
}
endpoint = f"{spec.base_url.rstrip('/')}/api/exec"
async with httpx.AsyncClient(timeout=timeout) as client:
resp = await client.post(endpoint, json=payload, headers=headers)
# CRIT-C: rileva quota/rate-limit sia da status code che da body
body_text = resp.text
if resp.status_code in (402, 429) or self._is_quota_error(body_text):
_logger.warning(
"[fabric] %s β quota/rate-limit rilevato (HTTP %s)",
spec.provider_id, resp.status_code
)
return {
"status_code": resp.status_code,
"error": "quota_or_ratelimit",
"body": body_text[:200],
}
resp.raise_for_status()
try:
return resp.json()
except Exception:
return {"status": "ok", "raw": body_text[:500]}
# ββ Health management ββββββββββββββββββββββββββββββββββββββββββββββββββββ
def _update_health(self, pid: str, health: ProviderHealth) -> None:
if pid in self._states:
self._states[pid].health = health
# ββ Dispatch pubblico (con chunking) ββββββββββββββββββββββββββββββββββββ
async def dispatch(self, req: DispatchRequest) -> DispatchResult:
await self.initialize()
# --- LOGICA CHUNKING (S-CHUNK) ---
from tools.payload_chunker import chunk_payload
chunks = chunk_payload(req.payload, max_kb=450)
if len(chunks) > 1:
_logger.info("[fabric] Payload grande β suddivisione in %d pezzi.", len(chunks))
final_responses = []
for chunk in chunks:
chunk_req = req.model_copy(update={"payload": chunk})
res = await self._dispatch_single(chunk_req)
final_responses.append(res.response)
return DispatchResult(
task_id=req.task_id, provider_id="multi", provider_name="Fabric Chunker",
provider_kind="internal", capability=req.capability, status="executed",
response={"chunks": final_responses, "total_chunks": len(chunks), "provider_info": "Fabric Chunker"},
correlation_id=req.correlation_id
)
return await self._dispatch_single(req)
# ββ Dispatch singolo (con retry + circuit breaker) βββββββββββββββββββββββ
async def _dispatch_single(self, req: DispatchRequest) -> DispatchResult:
max_attempts = 4
last_exception: Exception = Exception("nessun tentativo eseguito")
excluded: set[str] = set()
for attempt in range(max_attempts):
provider_id = req.provider_hint if attempt == 0 else None
if provider_id is None:
provider_id = self._select(req.capability, exclude=excluded)
if not provider_id or provider_id not in self._specs:
break # nessun provider disponibile β esci dal loop
spec = self._specs[provider_id]
t0 = time.time()
try:
response = await self._call_provider(spec, req)
# CRIT-C: risposta con quota/rate-limit segnalato nel body
is_quota = (
isinstance(response, dict)
and (
response.get("error") == "quota_or_ratelimit"
or response.get("status_code") in (402, 429)
or self._is_quota_error(str(response))
)
)
if is_quota:
await _token_rotator.rotate()
excluded.add(provider_id)
await asyncio.sleep(_RETRY_DELAY_S)
continue
# Successo β azzera contatore errori Oracle se applicabile
if spec.kind == ProviderKind.ORACLE:
self._reset_oracle_errors(provider_id)
latency = (time.time() - t0) * 1000
_rke("dispatch_ok", provider=provider_id, latency_ms=latency)
return DispatchResult(
task_id=req.task_id, provider_id=spec.provider_id,
provider_name=spec.name, provider_kind=spec.kind.value,
capability=req.capability, status="executed",
response=response, latency_ms=latency,
correlation_id=req.correlation_id
)
except Exception as exc:
last_exception = exc
exc_str = str(exc).lower()
# CRIT-C: pattern quota/rate-limit esteso anche nelle eccezioni
if self._is_quota_error(exc_str):
await _token_rotator.rotate()
excluded.add(provider_id)
await asyncio.sleep(_RETRY_DELAY_S)
continue
# CRIT-B: Oracle error β circuit breaker + fallback sul primo provider disponibile
if spec.kind == ProviderKind.ORACLE:
circuit_opened = self._record_oracle_error(provider_id)
_logger.error(
"[fabric] Oracle error (attempt %d/%d): %s β circuit %s",
attempt + 1, max_attempts, exc,
"APERTO" if circuit_opened else "registrato",
)
# Cerca dinamicamente il primo provider non-Oracle non-excluded (fast-path fallback)
# Evita ricorsione: chiama _call_provider direttamente senza passare da dispatch()
fb_spec = next(
(s for pid, s in self._specs.items()
if pid not in excluded and s.kind != ProviderKind.ORACLE),
None,
)
if fb_spec:
try:
t1 = time.time()
fb_response = await self._call_provider(fb_spec, req)
latency = (time.time() - t1) * 1000
_logger.info("[fabric] Fallback Oracleβ%s riuscito (%.0fms)", fb_spec.provider_id, latency)
return DispatchResult(
task_id=req.task_id, provider_id=fb_spec.provider_id,
provider_name=fb_spec.name, provider_kind=fb_spec.kind.value,
capability=req.capability, status="executed",
response=fb_response, latency_ms=latency,
correlation_id=req.correlation_id
)
except Exception as fb_exc:
_logger.error("[fabric] Fallback Oracleβ%s fallito: %s", fb_spec.provider_id, fb_exc)
last_exception = fb_exc
else:
_logger.warning("[fabric] Nessun provider fallback non-Oracle disponibile")
self._update_health(provider_id, ProviderHealth.DOWN)
excluded.add(provider_id)
return DispatchResult(
task_id=req.task_id, provider_id="failed", provider_name="none",
provider_kind="none", capability=req.capability, status="failed",
response={"error": str(last_exception)},
correlation_id=req.correlation_id
)
|