sync: 142 file da Baida98/AI@ab42bd1e (2026-08-08 08:23 UTC) [deploy-all]

#21
by Baida07 - opened
Files changed (1) hide show
  1. 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
- event_type='agent.kernel.dispatched',
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
- event_type='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,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
- event_type='task.running',
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
- event_type='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,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
- event_type='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,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
- event_type='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)
 
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)