Spaces:
Running
Running
sync: 142 file da Baida98/AI@ab42bd1e (2026-08-08 08:23 UTC) [deploy-all]
#21
by Baida07 - opened
- api/agent.py +9 -9
api/agent.py
CHANGED
|
@@ -522,12 +522,12 @@ async def agent_kernel_dispatch(body: AgentKernelDispatchIn, role: AuthRole = De
|
|
| 522 |
'goal': goal,
|
| 523 |
'mode': mode,
|
| 524 |
'dispatch_id': _dispatch_id,
|
|
|
|
| 525 |
},
|
| 526 |
priority='HIGH',
|
| 527 |
-
metadata={'workflow': 'agent-kernel.yml'},
|
| 528 |
)).add_done_callback(_log_task_exc)
|
| 529 |
asyncio.create_task(_kernel.publish_event(
|
| 530 |
-
|
| 531 |
payload={'goal': goal[:200], 'mode': mode},
|
| 532 |
)).add_done_callback(_log_task_exc)
|
| 533 |
import httpx as _httpx
|
|
@@ -579,10 +579,10 @@ async def _create_task_internal(task_id: str, goal: str, job: dict) -> dict:
|
|
| 579 |
"goal": goal,
|
| 580 |
"max_steps": job.get("max_steps", 20),
|
| 581 |
"source": "job_queue",
|
|
|
|
| 582 |
},
|
| 583 |
priority="NORMAL",
|
| 584 |
session_id=job.get("session_id", ""),
|
| 585 |
-
metadata={"job_queue": True},
|
| 586 |
)).add_done_callback(_log_task_exc)
|
| 587 |
return {"taskId": task_id, "status": "QUEUED"}
|
| 588 |
|
|
@@ -649,13 +649,13 @@ async def create_agent_task(body: AgentTaskIn, role: AuthRole = Depends(require_
|
|
| 649 |
'max_steps': body.max_steps,
|
| 650 |
'persona': body.persona,
|
| 651 |
'source': 'agent_api',
|
|
|
|
| 652 |
},
|
| 653 |
priority='NORMAL',
|
| 654 |
session_id=body.session_id,
|
| 655 |
-
metadata={'agent_api': True},
|
| 656 |
)).add_done_callback(_log_task_exc)
|
| 657 |
asyncio.create_task(_kernel.publish_event(
|
| 658 |
-
|
| 659 |
payload={'task_id': task_id, 'goal': body.goal[:200], 'status': 'QUEUED'},
|
| 660 |
)).add_done_callback(_log_task_exc)
|
| 661 |
return {'taskId': task_id, 'status': 'QUEUED'}
|
|
@@ -939,7 +939,7 @@ async def stream_agent_task(task_id: str, request: Request, resume: int = 0, rol
|
|
| 939 |
# ARCH-K2.2: pubblica lifecycle event via Kernel
|
| 940 |
if _KERNEL_AVAILABLE and _kernel is not None:
|
| 941 |
asyncio.create_task(_kernel.publish_event(
|
| 942 |
-
|
| 943 |
payload={'task_id': task_id, 'status': 'RUNNING'},
|
| 944 |
)).add_done_callback(_log_task_exc)
|
| 945 |
_prune_agent_tasks()
|
|
@@ -1286,7 +1286,7 @@ async def stream_agent_task(task_id: str, request: Request, resume: int = 0, rol
|
|
| 1286 |
# ARCH-K2.2: pubblica lifecycle event via Kernel
|
| 1287 |
if _KERNEL_AVAILABLE and _kernel is not None:
|
| 1288 |
asyncio.create_task(_kernel.publish_event(
|
| 1289 |
-
|
| 1290 |
payload={'task_id': task_id, 'status': 'SUCCESS'},
|
| 1291 |
)).add_done_callback(_log_task_exc)
|
| 1292 |
_result_text = str(result.get('output', result) if isinstance(result, dict) else result)
|
|
@@ -1309,7 +1309,7 @@ async def stream_agent_task(task_id: str, request: Request, resume: int = 0, rol
|
|
| 1309 |
# ARCH-K2.2: pubblica lifecycle event via Kernel
|
| 1310 |
if _KERNEL_AVAILABLE and _kernel is not None:
|
| 1311 |
asyncio.create_task(_kernel.publish_event(
|
| 1312 |
-
|
| 1313 |
payload={'task_id': task_id, 'status': 'CANCELLED'},
|
| 1314 |
)).add_done_callback(_log_task_exc)
|
| 1315 |
_sse('task_cancelled', {'taskId': task_id})
|
|
@@ -1330,7 +1330,7 @@ async def stream_agent_task(task_id: str, request: Request, resume: int = 0, rol
|
|
| 1330 |
# ARCH-K2.2: pubblica lifecycle event via Kernel
|
| 1331 |
if _KERNEL_AVAILABLE and _kernel is not None:
|
| 1332 |
asyncio.create_task(_kernel.publish_event(
|
| 1333 |
-
|
| 1334 |
payload={'task_id': task_id, 'status': 'ERROR', 'error': str(err)[:500]},
|
| 1335 |
)).add_done_callback(_log_task_exc)
|
| 1336 |
_logger.error('[agent/stream] %s error: %s', task_id, err, exc_info=True)
|
|
|
|
| 522 |
'goal': goal,
|
| 523 |
'mode': mode,
|
| 524 |
'dispatch_id': _dispatch_id,
|
| 525 |
+
'metadata': {'workflow': 'agent-kernel.yml'},
|
| 526 |
},
|
| 527 |
priority='HIGH',
|
|
|
|
| 528 |
)).add_done_callback(_log_task_exc)
|
| 529 |
asyncio.create_task(_kernel.publish_event(
|
| 530 |
+
topic='agent.kernel.dispatched',
|
| 531 |
payload={'goal': goal[:200], 'mode': mode},
|
| 532 |
)).add_done_callback(_log_task_exc)
|
| 533 |
import httpx as _httpx
|
|
|
|
| 579 |
"goal": goal,
|
| 580 |
"max_steps": job.get("max_steps", 20),
|
| 581 |
"source": "job_queue",
|
| 582 |
+
"metadata": {"job_queue": True},
|
| 583 |
},
|
| 584 |
priority="NORMAL",
|
| 585 |
session_id=job.get("session_id", ""),
|
|
|
|
| 586 |
)).add_done_callback(_log_task_exc)
|
| 587 |
return {"taskId": task_id, "status": "QUEUED"}
|
| 588 |
|
|
|
|
| 649 |
'max_steps': body.max_steps,
|
| 650 |
'persona': body.persona,
|
| 651 |
'source': 'agent_api',
|
| 652 |
+
'metadata': {'agent_api': True},
|
| 653 |
},
|
| 654 |
priority='NORMAL',
|
| 655 |
session_id=body.session_id,
|
|
|
|
| 656 |
)).add_done_callback(_log_task_exc)
|
| 657 |
asyncio.create_task(_kernel.publish_event(
|
| 658 |
+
topic='task.created',
|
| 659 |
payload={'task_id': task_id, 'goal': body.goal[:200], 'status': 'QUEUED'},
|
| 660 |
)).add_done_callback(_log_task_exc)
|
| 661 |
return {'taskId': task_id, 'status': 'QUEUED'}
|
|
|
|
| 939 |
# ARCH-K2.2: pubblica lifecycle event via Kernel
|
| 940 |
if _KERNEL_AVAILABLE and _kernel is not None:
|
| 941 |
asyncio.create_task(_kernel.publish_event(
|
| 942 |
+
topic='task.running',
|
| 943 |
payload={'task_id': task_id, 'status': 'RUNNING'},
|
| 944 |
)).add_done_callback(_log_task_exc)
|
| 945 |
_prune_agent_tasks()
|
|
|
|
| 1286 |
# ARCH-K2.2: pubblica lifecycle event via Kernel
|
| 1287 |
if _KERNEL_AVAILABLE and _kernel is not None:
|
| 1288 |
asyncio.create_task(_kernel.publish_event(
|
| 1289 |
+
topic='task.completed',
|
| 1290 |
payload={'task_id': task_id, 'status': 'SUCCESS'},
|
| 1291 |
)).add_done_callback(_log_task_exc)
|
| 1292 |
_result_text = str(result.get('output', result) if isinstance(result, dict) else result)
|
|
|
|
| 1309 |
# ARCH-K2.2: pubblica lifecycle event via Kernel
|
| 1310 |
if _KERNEL_AVAILABLE and _kernel is not None:
|
| 1311 |
asyncio.create_task(_kernel.publish_event(
|
| 1312 |
+
topic='task.cancelled',
|
| 1313 |
payload={'task_id': task_id, 'status': 'CANCELLED'},
|
| 1314 |
)).add_done_callback(_log_task_exc)
|
| 1315 |
_sse('task_cancelled', {'taskId': task_id})
|
|
|
|
| 1330 |
# ARCH-K2.2: pubblica lifecycle event via Kernel
|
| 1331 |
if _KERNEL_AVAILABLE and _kernel is not None:
|
| 1332 |
asyncio.create_task(_kernel.publish_event(
|
| 1333 |
+
topic='task.failed',
|
| 1334 |
payload={'task_id': task_id, 'status': 'ERROR', 'error': str(err)[:500]},
|
| 1335 |
)).add_done_callback(_log_task_exc)
|
| 1336 |
_logger.error('[agent/stream] %s error: %s', task_id, err, exc_info=True)
|