Spaces:
Running
Running
sync: 181 file da Baida98/AI@0f8ce228 (2026-08-25 11:12 UTC) [deploy-all]
Browse files- api/agent.py +14 -0
- api/vfs_sync.py +24 -0
- tests/test_vfs_sync.py +24 -0
api/agent.py
CHANGED
|
@@ -51,6 +51,7 @@ from .state import (
|
|
| 51 |
ReasonLoopIn, AgentTaskIn,
|
| 52 |
)
|
| 53 |
from .speculative import fire_speculative_tools
|
|
|
|
| 54 |
try:
|
| 55 |
from .quality_guardian import run_quality_check as _run_quality_check
|
| 56 |
except Exception:
|
|
@@ -1092,6 +1093,10 @@ async def stream_agent_task(task_id: str, request: Request, resume: int = 0, rol
|
|
| 1092 |
)
|
| 1093 |
step_idx = [0]
|
| 1094 |
_backend_steps: list[dict] = [] # GAP-SYNC-FIX: log step per resume preciso
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1095 |
|
| 1096 |
async def step_cb(step_data: dict) -> None:
|
| 1097 |
step_idx[0] += 1
|
|
@@ -1221,6 +1226,7 @@ async def stream_agent_task(task_id: str, request: Request, resume: int = 0, rol
|
|
| 1221 |
# Frontend scrive direttamente nel VFS locale senza fetch aggiuntivo
|
| 1222 |
if _action == 'file_written' and step_data.get('content'):
|
| 1223 |
_vfs_evt['content'] = str(step_data['content'])[:60_000]
|
|
|
|
| 1224 |
_sse('vfs_update', _vfs_evt)
|
| 1225 |
|
| 1226 |
# S363-UI: thought event — emitted when planner completes
|
|
@@ -1332,6 +1338,14 @@ async def stream_agent_task(task_id: str, request: Request, resume: int = 0, rol
|
|
| 1332 |
payload={'task_id': task_id, 'status': 'SUCCESS'},
|
| 1333 |
)).add_done_callback(_log_task_exc)
|
| 1334 |
_result_text = str(result.get('output', result) if isinstance(result, dict) else result)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1335 |
_sse('task_done', {'taskId': task_id, 'result': _result_text[:8000]})
|
| 1336 |
asyncio.create_task(_tg_done(task_id, task.get('goal', ''), _result_text[:500], _task_started_ms)).add_done_callback(_log_task_exc)
|
| 1337 |
|
|
|
|
| 51 |
ReasonLoopIn, AgentTaskIn,
|
| 52 |
)
|
| 53 |
from .speculative import fire_speculative_tools
|
| 54 |
+
from .vfs_sync import build_vfs_sync_complete
|
| 55 |
try:
|
| 56 |
from .quality_guardian import run_quality_check as _run_quality_check
|
| 57 |
except Exception:
|
|
|
|
| 1093 |
)
|
| 1094 |
step_idx = [0]
|
| 1095 |
_backend_steps: list[dict] = [] # GAP-SYNC-FIX: log step per resume preciso
|
| 1096 |
+
# P34/SYNC-1: i file arrivano nel frontend in staging; raccogliamo
|
| 1097 |
+
# soltanto quelli con contenuto così il commit atomico può avvenire
|
| 1098 |
+
# una sola volta dopo un task riuscito.
|
| 1099 |
+
_vfs_written_paths: set[str] = set()
|
| 1100 |
|
| 1101 |
async def step_cb(step_data: dict) -> None:
|
| 1102 |
step_idx[0] += 1
|
|
|
|
| 1226 |
# Frontend scrive direttamente nel VFS locale senza fetch aggiuntivo
|
| 1227 |
if _action == 'file_written' and step_data.get('content'):
|
| 1228 |
_vfs_evt['content'] = str(step_data['content'])[:60_000]
|
| 1229 |
+
_vfs_written_paths.add(str(_vfs_file)[:500])
|
| 1230 |
_sse('vfs_update', _vfs_evt)
|
| 1231 |
|
| 1232 |
# S363-UI: thought event — emitted when planner completes
|
|
|
|
| 1338 |
payload={'task_id': task_id, 'status': 'SUCCESS'},
|
| 1339 |
)).add_done_callback(_log_task_exc)
|
| 1340 |
_result_text = str(result.get('output', result) if isinstance(result, dict) else result)
|
| 1341 |
+
_vfs_commit = build_vfs_sync_complete(
|
| 1342 |
+
task_id, _vfs_written_paths, result,
|
| 1343 |
+
)
|
| 1344 |
+
if _vfs_commit is not None:
|
| 1345 |
+
# P34: completa l’atomic swap frontend solo dopo che il loop ha
|
| 1346 |
+
# confermato il task. Su errore/cancellazione lo staging rimane
|
| 1347 |
+
# intenzionalmente non committato.
|
| 1348 |
+
_sse('vfs_sync_complete', _vfs_commit)
|
| 1349 |
_sse('task_done', {'taskId': task_id, 'result': _result_text[:8000]})
|
| 1350 |
asyncio.create_task(_tg_done(task_id, task.get('goal', ''), _result_text[:500], _task_started_ms)).add_done_callback(_log_task_exc)
|
| 1351 |
|
api/vfs_sync.py
ADDED
|
@@ -0,0 +1,24 @@
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1 |
+
"""Contratto SSE per il commit atomico degli artefatti VFS lato frontend."""
|
| 2 |
+
from __future__ import annotations
|
| 3 |
+
|
| 4 |
+
from collections.abc import Iterable, Mapping
|
| 5 |
+
from typing import Any
|
| 6 |
+
|
| 7 |
+
|
| 8 |
+
def build_vfs_sync_complete(
|
| 9 |
+
task_id: str,
|
| 10 |
+
written_paths: Iterable[str],
|
| 11 |
+
result: Any,
|
| 12 |
+
) -> dict[str, object] | None:
|
| 13 |
+
"""Restituisce il payload di commit solo dopo un loop riuscito.
|
| 14 |
+
|
| 15 |
+
I contenuti viaggiano in eventi ``vfs_update`` separati e restano nello
|
| 16 |
+
staging del frontend. Questo segnale consente l'atomic swap esclusivamente
|
| 17 |
+
per i file effettivamente inoltrati dal backend.
|
| 18 |
+
"""
|
| 19 |
+
files = sorted({path for path in written_paths if path})
|
| 20 |
+
if not files:
|
| 21 |
+
return None
|
| 22 |
+
if isinstance(result, Mapping) and not bool(result.get("success", True)):
|
| 23 |
+
return None
|
| 24 |
+
return {"taskId": task_id, "files": files}
|
tests/test_vfs_sync.py
ADDED
|
@@ -0,0 +1,24 @@
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1 |
+
from api.vfs_sync import build_vfs_sync_complete
|
| 2 |
+
|
| 3 |
+
|
| 4 |
+
def test_vfs_sync_commit_includes_sorted_unique_files_after_success():
|
| 5 |
+
event = build_vfs_sync_complete(
|
| 6 |
+
"task-123",
|
| 7 |
+
{"reports/data.json", "e2e_metrics.json", "reports/data.json"},
|
| 8 |
+
{"success": True, "output": "ok"},
|
| 9 |
+
)
|
| 10 |
+
|
| 11 |
+
assert event == {
|
| 12 |
+
"taskId": "task-123",
|
| 13 |
+
"files": ["e2e_metrics.json", "reports/data.json"],
|
| 14 |
+
}
|
| 15 |
+
|
| 16 |
+
|
| 17 |
+
def test_vfs_sync_commit_is_omitted_without_files():
|
| 18 |
+
assert build_vfs_sync_complete("task-123", set(), {"success": True}) is None
|
| 19 |
+
|
| 20 |
+
|
| 21 |
+
def test_vfs_sync_commit_is_omitted_after_failed_loop():
|
| 22 |
+
assert build_vfs_sync_complete(
|
| 23 |
+
"task-123", {"e2e_metrics.json"}, {"success": False, "error": "write failed"}
|
| 24 |
+
) is None
|