Files
Agentswarm/orchestrator/swarm_runtime.py
T
FastheiandClaude Opus 4.8 1edf73aba4 fix(swarm/#70): 任务超时透传+提默认 / 事件时间线降噪 / 失败 termination_reason 准确
真机端到端实测(#70)暴露三问题,本 PR 全部修复(仅 agent + orchestrator,不跨仓):

问题1【阻断】单任务超时只有 60s,生成类任务必挂
- agent/main.py: TASK_TIMEOUT_SECONDS 默认 60→300(仅对外部/独立启动 agent 生效)。
- agent_launcher.py: 新增 DEFAULT_TASK_TIMEOUT_SECONDS=300、_budget_duration_seconds、
  resolve_task_timeout(base=env 默认 300,与 run budget.duration_seconds 取较小);
  plan_launch_specs 把 TASK_TIMEOUT_SECONDS 透传进每个 agent env(非敏感,inline,
  k8s 不进 Secret)。

问题2【体验】事件时间线全是内部噪音(纯附加,未碰冻结契约)
- swarm_runtime.py: is_client_visible(=event_type∈FROZEN_CLIENT_EVENT_TYPES,单一真源);
  emit_event 给 envelope 加 metadata.client_visible 布尔 + 关键客户端事件回填可选
  payload.message(人话进度,仅取已有字段,不伪造)。task.heartbeat/retried/
  deployment.status_changed/timeline/budget 标 client_visible=false,仍持久化+回调
  但客户端据此过滤出时间线。冻结事件集/类型/sequence/artifact 形状一字未动。
- event-schema.md: 文档化两个附加字段 + 新增 §6.1,明确未解冻。

问题3【正确性】失败/超时 termination_reason 仍报 "tasks_completed"
- convergence.py: 新增 TIMEOUT/MAX_RETRIES_EXCEEDED/TASK_FAILED;classify_failure_reason
  按 timeout→max_retries→task_failed 取最具体(仅凭真实 per-task 信号);FAILED 分支
  再不会返回 tasks_completed(该 reason 仅用于成功),budget/rounds 仅在通用失败时才覆盖。
- task_queue.py: fail_task 永久失败时把 reason 落到 task.result({"success":false,"error":reason}),
  不覆盖已有结果,供 convergence 读取。
- main.py: compute_convergence_report 快照补 retry_count/max_retries。

测试:新增 test_resolve_task_timeout、扩 test-convergence(failed_timeout/max_retries/
generic + "FAILED 永不报 tasks_completed"不变量)。本地全过:test-agent-launcher /
test-convergence / test-runtime-contract / test-contract-freeze / test-merge-smoke /
test-workflow-e2e / test-security-boundary。

影响:agent + orchestrator + 文档;不动 Manager↔Swarm 冻结契约字段(问题2 纯附加)。
栈在 #64(agent_swarm git 注入)之上,#64 合并后本 PR base 自动转 main。

Closes #70

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-16 18:29:39 +08:00

1066 lines
47 KiB
Python

"""Agent Manager compatible swarm runtime bridge."""
import asyncio
from copy import deepcopy
import hashlib
import hmac
import json
import logging
import os
import time
import uuid
from datetime import datetime, timezone
from typing import Any, Dict, List, Optional
import httpx
from pydantic import BaseModel, Field
from .redis_client import redis_client
logger = logging.getLogger(__name__)
class CallbackConfig(BaseModel):
"""Callback settings provided by Agent Manager."""
url: Optional[str] = None
signing_secret_ref: Optional[str] = None
subscribed_events: List[str] = Field(default_factory=list)
class SwarmRun(BaseModel):
"""Runtime-owned representation of a swarm run."""
deployment_id: str
swarm_id: str
mode: str = "swarm"
status: str
objective: str
manager_deployment_id: Optional[str] = None
correlation_id: Optional[str] = None
callback: CallbackConfig = Field(default_factory=CallbackConfig)
task_ids: List[str] = Field(default_factory=list)
approvals: Dict[str, Dict[str, Any]] = Field(default_factory=dict)
# Internal peer-communication telemetry (NOT a Manager event — kept off the HM event
# registry on purpose). Feeds benchmark S_communication. correlation_ids are de-duplicated
# so a retried/echoed message is not double counted.
collaboration: Dict[str, Any] = Field(
default_factory=lambda: {
"request_correlations": [], # correlation_ids of peer requests routed
"reply_correlations": [], # correlation_ids of peer replies routed
"delivery_failures": 0, # peer messages dropped (target not connected)
}
)
# Benchmark Group B: quality grade of the run's generated code (fixture TestPassRate →
# Q_quality). Populated only when ENABLE_QUALITY_EVAL is set and a fixture is bound. Feeds
# benchmark reward(); empty dict means not graded (collector leaves reward NaN).
quality: Dict[str, Any] = Field(default_factory=dict)
# Benchmark Group A: per-assignment ACO decision telemetry ({task_id, agent_id, tau, eta,
# p_norm, p_score, explored}). Internal state (NOT a Manager event). Appended only when
# ENABLE_ACO_DISPATCH selected the assignment; empty → tau/eta/p_decision stay NaN.
decisions: List[Dict[str, Any]] = Field(default_factory=list)
# Issue #9: explainable dispatch-scoring decisions (`dispatch.decision_made` payloads:
# candidate agents + per-dimension score breakdown + exclusion reasons). Internal state
# (NOT a Manager event — a Manager-facing dispatch event needs event-schema registration,
# cross-team). Appended only when ENABLE_DISPATCH_SCORE selected the assignment. Audit/
# benchmark-replayable from here.
dispatch_decisions: List[Dict[str, Any]] = Field(default_factory=list)
metadata: Dict[str, Any] = Field(default_factory=dict)
request_body: Dict[str, Any] = Field(default_factory=dict)
created_at: float = Field(default_factory=time.time)
updated_at: float = Field(default_factory=time.time)
class RuntimeValidationError(ValueError):
"""Validation error returned in the Agent Manager runtime envelope."""
def __init__(self, message: str, code: str = "VALIDATION_ERROR"):
super().__init__(message)
self.message = message
self.code = code
# Frozen client-facing event-type set (agent_swarm#15.2 / HM #28). These 13 types are the
# contract the desktop "task cockpit" renders against; HM (agent_callback.go) registers exactly
# these. The runtime emits other types too (deployment.status_changed, timeline.updated,
# budget.alert, task.heartbeat/blocked/retried/released) and many INTERNAL-only telemetry events
# that are NOT emitted to Manager (see assess_swarm_health / proposals / bids / reviews).
FROZEN_CLIENT_EVENT_TYPES = (
"task.created",
"task.claimed",
"task.running",
"task.completed",
"task.failed",
"handoff.created",
"approval.requested",
"approval.approved",
"approval.rejected",
"artifact.created",
"swarm.completed",
"swarm.failed",
"swarm.stopped",
# Review/rework timeline (agent_swarm#34) — emitted as REDACTED client projections from
# run_cross_review (cross_review.*_client_payload): reviewer id/verdict/criteria/targets/
# root_cause only, no evidence/summary/reason free text.
"review.started",
"review.decision_made",
"rework.requested",
"rework.completed",
)
# Terminal client events (agent_swarm#40): once one of these is streamed, the run has reached a
# final state and the SSE stream can close — the client cockpit keys its terminal UI on these.
TERMINAL_CLIENT_EVENT_TYPES = frozenset({"swarm.completed", "swarm.failed", "swarm.stopped"})
_FROZEN_CLIENT_EVENT_SET = frozenset(FROZEN_CLIENT_EVENT_TYPES)
def is_client_visible(event_type: str) -> bool:
"""Whether an event type is part of the frozen client cockpit contract (vs internal noise).
The runtime stores/forwards many non-client event types (task.heartbeat, task.retried,
deployment.status_changed, timeline.updated, budget.alert, …). Only FROZEN_CLIENT_EVENT_TYPES
are rendered on the client timeline. This is the single source of truth both the runtime
(envelope `metadata.client_visible`, agent_swarm#70) and HM/clients filter on — no event is
removed, and the frozen type set is untouched.
"""
return event_type in _FROZEN_CLIENT_EVENT_SET
#: Human-readable message templates for key client-visible events (agent_swarm#70). Purely
#: additive: an optional `message` is set on the payload only when the emitter did not already
#: provide one. Templates pull from already-present payload fields, so no new data is fabricated.
def _client_message(event_type: str, payload: Dict[str, Any]) -> Optional[str]:
role = payload.get("agent_role")
title = payload.get("title") or payload.get("summary")
if event_type == "task.created":
return f"New task created: {title}" if title else "New task created"
if event_type == "task.claimed":
return f"{role} picked up a task" if role else "An agent picked up a task"
if event_type == "task.running":
return f"{role} is working on a task" if role else "An agent is working on a task"
if event_type == "task.completed":
return "Task completed"
if event_type == "task.failed":
reason = payload.get("reason")
return f"Task failed: {reason}" if reason else "Task failed"
if event_type == "swarm.completed":
return "Swarm run completed"
if event_type == "swarm.failed":
tr = payload.get("termination_reason")
return f"Swarm run failed ({tr})" if tr else "Swarm run failed"
if event_type == "swarm.stopped":
return "Swarm run stopped"
return None
class SwarmRuntime:
"""Stores swarm runs and emits Agent Manager callback events."""
RUN_KEY_PREFIX = "swarm:"
TASK_RUN_KEY_PREFIX = "swarm_task:"
IDEMPOTENCY_KEY_PREFIX = "swarm_idempotency:"
EVENT_KEY_PREFIX = "swarm_events:"
EVENT_SEQ_KEY_PREFIX = "swarm_event_seq:"
# Per-user "retained" run set (agent_swarm#16/#60 — model-key revocation handshake). Holds the
# swarm_ids whose per-user `sk-` must stay valid: a run is retained from create until it is
# *stopped*. completed/failed runs stay retained because they are reopenable via POST …/input
# (their relaunched agents reuse the same per-user key). Only when ALL of a user's runs are
# stopped does the runtime emit `swarm.pool_terminated` so HM can revoke the `sk-` + clear KV.
USER_RETAINED_KEY_PREFIX = "swarm_user_retained:"
POOL_TERMINATED_FLAG_PREFIX = "swarm_user_pool_terminated:"
def __init__(self):
self.callback_service_token = (
os.getenv("AGENT_CALLBACK_SERVICE_TOKEN")
or os.getenv("AGNET_CALLBACK_SERVICE_TOKEN")
)
self.callback_signing_secret = (
os.getenv("AGENT_CALLBACK_SIGNING_SECRET")
or os.getenv("AGNET_CALLBACK_SIGNING_SECRET")
)
self.runtime_source = os.getenv("SWARM_RUNTIME_SOURCE", "heicode-swarm-runtime")
def multi_agent_workflow_enabled(self, body: Optional[Dict[str, Any]] = None) -> bool:
"""Return whether the DAG-style runtime workflow is enabled."""
if os.getenv("ENABLE_SUBTASK_HANDOFF", "false").lower() not in {"1", "true", "yes"}:
return False
plan = (body or {}).get("orchestration_plan") or {}
agents = plan.get("agents") or (body or {}).get("agents") or []
return len(agents) > 0
async def health(self) -> Dict[str, Any]:
"""Return Agent Manager compatible runtime health."""
await redis_client.client.ping()
return {
"success": True,
"data": {
"status": "healthy",
"service": "heicode-swarm-runtime",
"version": "1.0.0",
"runtime": os.getenv("SWARM_RUNTIME_PLATFORM", "aks"),
"time": self._now_iso(),
"capabilities": [
"swarm.create",
"task.flow",
"handoff.events",
"artifact.events",
"approval.pause_resume",
"deployment.stop",
"runtime.tasks.query",
"runtime.logs.query",
"runtime.events.query",
"runtime.metrics.query",
"runtime.workflow.query",
"runtime.diagnostics.query",
],
},
}
def normalize_create_request(self, body: Dict[str, Any]) -> Dict[str, Any]:
"""Normalize the Manager swarm request into the runtime request shape."""
normalized = deepcopy(body or {})
metadata = dict(normalized.get("metadata") or {})
callback = dict(normalized.get("callback") or {})
requirement = dict(normalized.get("requirement") or {})
model_selection = dict(normalized.get("model_selection") or {})
existing_plan = dict(normalized.get("orchestration_plan") or {})
uses_new_shape = normalized.get("mode") == "swarm" or bool(requirement)
if uses_new_shape:
primary_model = (
model_selection.get("primary_model")
or existing_plan.get("model_id")
or normalized.get("model_id")
)
budget = deepcopy(
normalized.get("budget")
or existing_plan.get("budget")
or {}
)
if not existing_plan.get("objective") and requirement.get("objective"):
existing_plan["objective"] = requirement.get("objective")
existing_plan.setdefault("sub_mode", "swarm")
existing_plan.setdefault("risk_level", normalized.get("risk_level") or "medium")
existing_plan.setdefault("budget", budget)
existing_plan["model_id"] = primary_model
existing_plan["requirement"] = self._redact_sensitive(requirement)
existing_plan["attachments"] = self._redact_sensitive(
requirement.get("attachments") or []
)
existing_plan["constraints"] = self._redact_sensitive(
requirement.get("constraints") or []
)
existing_plan["acceptance_criteria"] = self._redact_sensitive(
requirement.get("acceptance_criteria") or []
)
if not existing_plan.get("agents"):
existing_plan["agents"] = normalized.get("agents") or []
normalized["mode"] = "swarm"
normalized["conversation_id"] = (
normalized.get("conversation_id")
or metadata.get("conversation_id")
)
normalized["requirement"] = requirement
normalized["model_selection"] = model_selection
normalized["model_id"] = primary_model
normalized["orchestration_plan"] = existing_plan
if normalized.get("conversation_id"):
metadata.setdefault("conversation_id", normalized.get("conversation_id"))
if requirement.get("objective"):
metadata.setdefault("objective", requirement.get("objective"))
if model_selection.get("type"):
metadata.setdefault("model_selection_type", model_selection.get("type"))
else:
normalized["orchestration_plan"] = existing_plan
normalized["metadata"] = metadata
normalized["callback"] = callback
return normalized
def validate_create_request(self, body: Dict[str, Any]):
"""Validate the Manager create request before a runtime run is persisted."""
normalized = self.normalize_create_request(body)
plan = normalized.get("orchestration_plan")
metadata = normalized.get("metadata")
callback = normalized.get("callback")
required_paths = {
"orchestration_plan.objective": plan.get("objective") if isinstance(plan, dict) else None,
"callback.url": callback.get("url") if isinstance(callback, dict) else None,
"metadata.manager_deployment_id": metadata.get("manager_deployment_id") if isinstance(metadata, dict) else None,
}
missing = [path for path, value in required_paths.items() if value in (None, "", [], {})]
if missing:
raise RuntimeValidationError(f"Missing required field(s): {', '.join(missing)}")
if normalized.get("mode") not in (None, "swarm"):
raise RuntimeValidationError("mode must be 'swarm'")
if normalized.get("model_selection"):
selection_type = (normalized.get("model_selection") or {}).get("type")
if selection_type not in (None, "primary"):
raise RuntimeValidationError("model_selection.type must be 'primary' for swarm")
billing_secret = (normalized.get("billing_context") or {}).get("secret_ref")
if billing_secret and not self._is_azkv_ref(billing_secret):
raise RuntimeValidationError("billing_context.secret_ref must use azkv://")
self._validate_resource_grants(normalized.get("resource_grants") or [], "resource_grants")
for index, agent in enumerate(plan.get("agents") or []):
self._validate_resource_grants(
agent.get("resource_grants") or [],
f"orchestration_plan.agents[{index}].resource_grants",
)
for index, agent in enumerate(normalized.get("agents") or []):
self._validate_resource_grants(
agent.get("resource_grants") or [],
f"agents[{index}].resource_grants",
)
for path in ("metadata", "resource_grants", "callback"):
value = normalized.get(path)
if value is not None:
self._reject_plaintext_secrets(value, path)
async def get_or_create_run(
self,
body: Dict[str, Any],
idempotency_key: Optional[str],
correlation_id: Optional[str],
) -> tuple[SwarmRun, bool]:
"""Create a swarm run or return the previous run for an idempotency key."""
body = self.normalize_create_request(body)
if idempotency_key:
existing_swarm_id = await redis_client.get(
f"{self.IDEMPOTENCY_KEY_PREFIX}{idempotency_key}"
)
if existing_swarm_id:
existing_run = await self.get_run(existing_swarm_id)
if existing_run:
return existing_run, False
metadata = body.get("metadata") or {}
plan = body.get("orchestration_plan") or {}
callback = CallbackConfig.model_validate(body.get("callback") or {})
swarm_id = f"swarm-{uuid.uuid4().hex[:12]}"
deployment_id = f"runtime-dep-{uuid.uuid4().hex[:12]}"
manager_deployment_id = (
metadata.get("manager_deployment_id")
or metadata.get("heicode_deployment_id")
)
objective = (
plan.get("objective")
or metadata.get("objective")
or "Swarm runtime task"
)
run = SwarmRun(
deployment_id=deployment_id,
swarm_id=swarm_id,
mode="swarm",
status="running",
objective=objective,
manager_deployment_id=manager_deployment_id,
correlation_id=correlation_id or metadata.get("correlation_id"),
callback=callback,
metadata=self._redact_sensitive(metadata),
request_body=self._redact_sensitive(body),
)
if self._requires_approval(body):
approval = self._build_approval(body, run)
run.status = "waiting_approval"
run.approvals[approval["approval_id"]] = approval
await self.save_run(run)
if idempotency_key:
await redis_client.set(
f"{self.IDEMPOTENCY_KEY_PREFIX}{idempotency_key}",
run.swarm_id,
ex=86400,
)
await self.emit_event(
run,
"deployment.status_changed",
payload=self.status_payload(run, phase="Plan"),
)
for approval in run.approvals.values():
await self.emit_event(run, "approval.requested", payload=approval)
# Mark the per-user `sk-` as in-use for the model-key revocation handshake (#16/#60). Held
# until every run of this user is stopped (completed/failed stay retained — reopenable).
await self.retain_run_for_user(run)
return run, True
async def save_run(self, run: SwarmRun):
"""Persist a swarm run."""
run.updated_at = time.time()
await redis_client.set(
f"{self.RUN_KEY_PREFIX}{run.swarm_id}",
run.model_dump_json(),
)
async def get_run(self, swarm_id: str) -> Optional[SwarmRun]:
"""Get a swarm run by id."""
data = await redis_client.get(f"{self.RUN_KEY_PREFIX}{swarm_id}")
if not data:
return None
return SwarmRun.model_validate_json(data)
async def get_run_by_identifier(self, identifier: str) -> Optional[SwarmRun]:
"""Find a run by swarm id, runtime deployment id, or manager deployment id."""
direct = await self.get_run(identifier)
if direct:
return direct
for key in await redis_client.keys(f"{self.RUN_KEY_PREFIX}*"):
data = await redis_client.get(key)
if not data:
continue
run = SwarmRun.model_validate_json(data)
if identifier in {run.deployment_id, run.manager_deployment_id, run.swarm_id}:
return run
return None
async def get_run_for_task(self, task_id: str) -> Optional[SwarmRun]:
"""Find the swarm run that owns a task."""
swarm_id = await redis_client.get(f"{self.TASK_RUN_KEY_PREFIX}{task_id}")
if not swarm_id:
return None
return await self.get_run(swarm_id)
async def record_peer_message(
self,
run: SwarmRun,
*,
correlation_id: Optional[str],
is_reply: bool,
delivered: bool,
) -> None:
"""Record one routed peer-collaboration message for S_communication.
Internal telemetry only: counters live on the run, not on the Manager event stream.
Requests and replies are tracked by correlation_id so the collector can compute the
share of requests that received an answer (a real communication-success rate).
"""
collab = run.collaboration or {}
requests = collab.setdefault("request_correlations", [])
replies = collab.setdefault("reply_correlations", [])
if correlation_id:
bucket = replies if is_reply else requests
if correlation_id not in bucket:
bucket.append(correlation_id)
if not delivered:
collab["delivery_failures"] = int(collab.get("delivery_failures", 0)) + 1
run.collaboration = collab
await self.save_run(run)
async def record_quality(self, run: SwarmRun, quality: Dict[str, Any]) -> None:
"""Store the benchmark quality grade for a completed run (Group B)."""
run.quality = quality or {}
await self.save_run(run)
async def record_benchmark_metrics(self, run: SwarmRun, payload: Dict[str, Any]) -> None:
"""Persist a terminal run's benchmark SwarmMetrics snapshot (阶段2 数据采集).
Stored on run.metadata['benchmark'] so real-user runs accrue data for later empirical
calibration (O / coefficients). Read-only over run state; never fabricates (NaN → null).
"""
run.metadata["benchmark"] = payload or {}
await self.save_run(run)
async def record_decision(self, run: SwarmRun, decision: Dict[str, Any]) -> None:
"""Append one ACO assignment decision (Group A telemetry). Bounded list."""
run.decisions = (run.decisions or [])[-999:] + [decision]
await self.save_run(run)
async def record_dispatch_decision(self, run: SwarmRun, decision: Dict[str, Any]) -> None:
"""Append one explainable dispatch-scoring decision (#9). Internal, bounded list."""
run.dispatch_decisions = (run.dispatch_decisions or [])[-999:] + [decision]
await self.save_run(run)
async def attach_task(self, run: SwarmRun, task_id: str):
"""Associate a queue task with a swarm run."""
if task_id not in run.task_ids:
run.task_ids.append(task_id)
await redis_client.set(f"{self.TASK_RUN_KEY_PREFIX}{task_id}", run.swarm_id)
await self.save_run(run)
# ------------------------------------------------------------------------------------------
# Per-user model-key revocation handshake (agent_swarm#16 / HM #60)
# ------------------------------------------------------------------------------------------
@staticmethod
def _run_user_id(run: SwarmRun) -> Optional[str]:
"""The initiating user for a run (from the create `x-user-id`, passed via runtime_headers)."""
headers = (run.metadata or {}).get("runtime_headers") or {}
uid = headers.get("x_user_id") or headers.get("x-user-id")
return str(uid) if uid else None
@staticmethod
def _run_secret_ref(run: SwarmRun) -> Optional[str]:
"""The azkv:// secret_ref the per-user `sk-` was resolved from (for HM to target the KV
secret on revocation). This is a reference, never a plaintext key (see _redact_sensitive)."""
ref = ((run.request_body or {}).get("billing_context") or {}).get("secret_ref")
return ref or None
async def _read_retained(self, user_id: str) -> List[str]:
raw = await redis_client.get(f"{self.USER_RETAINED_KEY_PREFIX}{user_id}")
if not raw:
return []
try:
return list(json.loads(raw))
except (ValueError, TypeError):
return []
async def _write_retained(self, user_id: str, swarm_ids: List[str]) -> None:
key = f"{self.USER_RETAINED_KEY_PREFIX}{user_id}"
if swarm_ids:
await redis_client.set(key, json.dumps(sorted(set(swarm_ids))))
else:
await redis_client.delete(key)
async def retain_run_for_user(self, run: SwarmRun) -> None:
"""Mark a run's per-user `sk-` as in-use. Idempotent. Clears any prior pool-terminated flag
so a future drain can signal again (the user has live work after a previous full teardown)."""
user_id = self._run_user_id(run)
if not user_id:
return
retained = await self._read_retained(user_id)
if run.swarm_id not in retained:
retained.append(run.swarm_id)
await self._write_retained(user_id, retained)
await redis_client.delete(f"{self.POOL_TERMINATED_FLAG_PREFIX}{user_id}")
async def release_run_and_maybe_terminate_pool(self, run: SwarmRun, reason: str) -> None:
"""Drop a *stopped* run from its user's retained set; if that user now has no retained runs
(all stopped — completed/failed stay retained, being reopenable), emit `swarm.pool_terminated`
once so HM revokes the per-user `sk-` and clears the KV secret. Best-effort; never raises."""
user_id = self._run_user_id(run)
if not user_id:
return
try:
retained = [sid for sid in await self._read_retained(user_id) if sid != run.swarm_id]
await self._write_retained(user_id, retained)
if retained:
return # user still has live/reopenable runs — keep the key valid
flag_key = f"{self.POOL_TERMINATED_FLAG_PREFIX}{user_id}"
if await redis_client.get(flag_key):
return # already signalled since the last active run — don't double-revoke
await redis_client.set(flag_key, "1", ex=7 * 24 * 3600)
payload = {"user_id": user_id, "reason": reason or "user swarm pool torn down"}
secret_ref = self._run_secret_ref(run)
if secret_ref:
payload["secret_ref"] = secret_ref # azkv ref (not the key) — HM clears this KV secret
await self.emit_event(run, "swarm.pool_terminated", payload=payload)
except Exception as exc: # never break stop on the revocation signal
logger.warning("pool_terminated signal failed for user %s: %s", user_id, exc)
async def stop_run(self, deployment_id: str, reason: str = "") -> Optional[SwarmRun]:
"""Stop a run by runtime deployment id or swarm id."""
run = await self.find_run_by_deployment_id(deployment_id)
if not run:
return None
run.status = "stopped"
await self.save_run(run)
await self.emit_event(
run,
"deployment.status_changed",
payload=self.status_payload(
run,
phase="Deliver",
reason=reason or "Heicode Manager requested stop",
),
)
# Frozen terminal swarm event (agent_swarm#15.2): the client cockpit keys its terminal
# banner off swarm.{completed|failed|stopped}, not deployment.status_changed.
await self.emit_event(run, "swarm.stopped", payload={
"status": "stopped",
"reason": reason or "Heicode Manager requested stop",
})
# Model-key revocation handshake (#16/#60): a stop is the only *final* state (completed/
# failed are reopenable via …/input). If this drains the user's retained set, signal HM.
await self.release_run_and_maybe_terminate_pool(run, reason or "Heicode Manager requested stop")
return run
async def find_run_by_deployment_id(self, deployment_id: str) -> Optional[SwarmRun]:
"""Find a run by runtime deployment id or swarm id."""
return await self.get_run_by_identifier(deployment_id)
async def record_approval_decision(
self,
swarm_id: str,
approval_id: str,
decision: Dict[str, Any],
) -> Optional[SwarmRun]:
"""Persist an approval decision and update run state."""
run = await self.get_run_by_identifier(swarm_id)
if not run:
return None
approval = run.approvals.get(approval_id, {"approval_id": approval_id})
approval["decision"] = decision.get("decision")
credential_lease = decision.get("credential_lease") or {}
approval["credential_ref"] = (
decision.get("credential_ref") or credential_lease.get("credential_ref")
)
approval["lease_id"] = decision.get("lease_id") or credential_lease.get("lease_id")
approval["lease_expires_at"] = (
decision.get("lease_expires_at") or credential_lease.get("expires_at")
)
approval["decided_at"] = self._now_iso()
run.approvals[approval_id] = self._redact_sensitive(approval)
if decision.get("decision") == "approved":
run.status = "running"
elif decision.get("decision") == "rejected":
run.status = "blocked"
await self.save_run(run)
await self.emit_event(
run,
"deployment.status_changed",
payload=self.status_payload(
run,
phase="Review",
approval_id=approval_id,
decision=decision.get("decision"),
reason=decision.get("reason"),
),
)
# Frozen approval-outcome events (agent_swarm#15.2): the client renders an approval
# timeline from approval.requested -> approval.approved|rejected.
decision_value = decision.get("decision")
if decision_value in ("approved", "rejected"):
await self.emit_event(run, f"approval.{decision_value}", payload={
"approval_id": approval_id,
"decision": decision_value,
"reason": decision.get("reason"),
})
if run.status == "blocked":
await self.emit_event(run, "task.blocked", payload={
"approval_id": approval_id,
"reason": decision.get("reason") or "Approval rejected",
"runtime_deployment_id": run.deployment_id,
})
await self.emit_event(run, "timeline.updated", payload={
"summary": "Swarm blocked by approval rejection",
"approval_id": approval_id,
})
return run
async def emit_event(
self,
run: SwarmRun,
event_type: str,
task_id: Optional[str] = None,
agent_instance_id: Optional[str] = None,
payload: Optional[Dict[str, Any]] = None,
artifact: Optional[Dict[str, Any]] = None,
):
"""Record and optionally emit a Manager callback event."""
callback = run.callback
occurred_at = self._now_iso()
redacted_payload = self._redact_sensitive(payload or {})
redacted_artifact = self._redact_sensitive(artifact) if artifact else None
if redacted_artifact:
# Freeze (agent_swarm#15.4): the client artifact view needs a stable flat shape
# {uri, checksum, task_id, size_bytes, created_at}. Default created_at to the event time;
# leave size_bytes absent when unknown (rule #9: do not fabricate a size).
redacted_artifact.setdefault("created_at", occurred_at)
if task_id and "task_id" not in redacted_artifact:
redacted_artifact["task_id"] = task_id
if redacted_artifact and event_type == "artifact.created":
redacted_payload = {
**redacted_artifact,
**redacted_payload,
}
# Conform to HM event contract required fields (heicode agent_callback.go):
# timeline.updated requires `title`; budget.alert requires `threshold_pct`.
if event_type == "timeline.updated" and "title" not in redacted_payload:
redacted_payload["title"] = redacted_payload.get("summary") or "timeline"
if event_type == "budget.alert" and "threshold_pct" not in redacted_payload:
threshold = redacted_payload.get("threshold")
if isinstance(threshold, (int, float)):
redacted_payload["threshold_pct"] = threshold * 100 if threshold <= 1 else threshold
# Natural-language progress (agent_swarm#70): give client-visible events a human `message`
# so the cockpit timeline reads as progress instead of raw event types. Additive + optional;
# only set when the emitter did not already provide one. Internal/noise events get none.
client_visible = is_client_visible(event_type)
if client_visible and "message" not in redacted_payload:
message = _client_message(event_type, redacted_payload)
if message:
redacted_payload["message"] = message
event_id = f"evt_{uuid.uuid4().hex}"
# Per-swarm strictly-increasing sequence (agent_swarm#15.1): the client polls
# `events?after=<sequence>` and dedups/orders by it. INCR is atomic so concurrent emits on
# the same swarm get distinct, gap-free numbers starting at 1.
sequence = await redis_client.incr(f"{self.EVENT_SEQ_KEY_PREFIX}{run.swarm_id}")
body: Dict[str, Any] = {
"event_id": event_id,
"idempotency_key": event_id,
"sequence": sequence,
"event_type": event_type,
"deployment_id": run.manager_deployment_id or run.deployment_id,
"runtime_deployment_id": run.deployment_id,
"swarm_id": run.swarm_id,
"agent_instance_id": agent_instance_id,
"task_id": task_id,
"occurred_at": occurred_at,
"correlation_id": run.correlation_id,
"source": self.runtime_source,
# client_visible (agent_swarm#70): explicit, type-derived flag so HM/clients can filter
# internal noise (heartbeat/retried/deployment.status_changed/timeline/budget) off the
# cockpit timeline without hard-coding the frozen set. Additive metadata field — does not
# touch the frozen envelope keys, event types, sequence, or artifact shape.
"metadata": {"client_visible": client_visible},
"payload": redacted_payload,
}
if redacted_artifact:
body["artifact"] = redacted_artifact
raw_body = json.dumps(body, ensure_ascii=False, separators=(",", ":"))
await self._store_event(run.swarm_id, raw_body)
if not callback.url:
return
if callback.subscribed_events and event_type not in callback.subscribed_events:
return
headers = {
"Content-Type": "application/json",
"X-Agent-Event-Id": event_id,
"X-Agnet-Event-Id": event_id,
}
if run.correlation_id:
headers["X-Correlation-ID"] = run.correlation_id
if self.callback_service_token:
headers["X-Agent-Service-Token"] = self.callback_service_token
headers["X-Agnet-Service-Token"] = self.callback_service_token
if self.callback_signing_secret:
timestamp = str(int(time.time() * 1000))
signature_payload = f"{timestamp}.{event_id}.{raw_body}"
signature = hmac.new(
self.callback_signing_secret.encode("utf-8"),
signature_payload.encode("utf-8"),
hashlib.sha256,
).hexdigest()
headers["X-Agent-Timestamp"] = timestamp
headers["X-Agent-Signature"] = f"sha256={signature}"
headers["X-Agnet-Timestamp"] = timestamp
headers["X-Agnet-Signature"] = f"sha256={signature}"
if callback.url and (not callback.subscribed_events or event_type in callback.subscribed_events):
await self._upsert_callback_attempt(
run.swarm_id,
event_id,
{
"event_id": event_id,
"event_type": event_type,
"url": callback.url,
"status": "pending",
"attempted_at": self._now_iso(),
},
)
asyncio.create_task(
self._post_callback(run.swarm_id, callback.url, raw_body, headers, event_type, event_id)
)
async def list_events(
self,
swarm_id: str,
limit: int = 100,
cursor: Optional[str] = None,
) -> Dict[str, Any]:
"""Return stored runtime events for a swarm."""
start = int(cursor or 0)
safe_limit = max(1, min(limit, 500))
items = await redis_client.lrange(
f"{self.EVENT_KEY_PREFIX}{swarm_id}",
start,
start + safe_limit - 1,
)
events = [json.loads(item) for item in items]
next_cursor = str(start + safe_limit) if len(events) == safe_limit else None
return {"events": events, "next_cursor": next_cursor}
async def _store_event(self, swarm_id: str, raw_body: str):
await redis_client.rpush(f"{self.EVENT_KEY_PREFIX}{swarm_id}", raw_body)
async def _post_callback(
self,
swarm_id: str,
url: str,
raw_body: str,
headers: Dict[str, str],
event_type: str,
event_id: str,
):
"""Send callback without blocking the agent WebSocket loop."""
if not url:
return
try:
async with httpx.AsyncClient(timeout=10.0) as client:
response = await client.post(url, content=raw_body, headers=headers)
await self._upsert_callback_attempt(
swarm_id,
event_id,
{
"event_id": event_id,
"event_type": event_type,
"url": url,
"status": "delivered" if response.status_code < 400 else "failed",
"response_status": response.status_code,
"attempted_at": self._now_iso(),
},
)
if response.status_code >= 400:
logger.warning(
"Manager callback %s failed with status %s: %s",
event_type,
response.status_code,
response.text[:500],
)
except Exception as exc:
await self._upsert_callback_attempt(
swarm_id,
event_id,
{
"event_id": event_id,
"event_type": event_type,
"url": url,
"status": "failed",
"error": str(exc),
"attempted_at": self._now_iso(),
},
)
logger.warning("Manager callback %s failed: %s", event_type, exc)
def build_task_descriptions(self, body: Dict[str, Any]) -> List[Dict[str, Any]]:
"""Create a simple task graph from a Manager orchestration request."""
body = self.normalize_create_request(body)
plan = body.get("orchestration_plan") or {}
agents = plan.get("agents") or body.get("agents") or []
requirement = body.get("requirement") or {}
objective = (
requirement.get("objective")
or plan.get("objective")
or "Complete swarm objective"
)
multi_agent_enabled = self.multi_agent_workflow_enabled(body)
context = {
"orchestration_plan": self._redact_sensitive(plan),
"resource_grants": self._redact_sensitive(body.get("resource_grants") or []),
"sub_mode": body.get("sub_mode") or plan.get("sub_mode"),
"conversation_id": body.get("conversation_id"),
"mode": body.get("mode") or "swarm",
"requirement": self._redact_sensitive(requirement),
}
if not agents or not multi_agent_enabled:
return [{
"task_id": "task-1",
"title": "Swarm objective",
"description": objective,
"agent_role": "general",
"required_capabilities": ["general"],
"depends_on": [],
"parent_task_id": None,
"root_task_id": "task-1",
"source": "runtime_bridge",
"workflow_mode": "single_agent",
"allow_handoff": False,
"context": context,
}]
tasks = []
for index, agent in enumerate(agents, start=1):
role = agent.get("role") or f"agent-{index}"
title = agent.get("title") or f"{role} task"
description = agent.get("description") or f"[{role}] {objective}"
task_id = agent.get("task_id") or f"task-{index}"
required_capabilities = agent.get("required_capabilities") or [role]
tasks.append({
"task_id": task_id,
"title": title,
"description": description,
"agent_role": role,
"required_capabilities": required_capabilities,
"depends_on": agent.get("depends_on") or [],
"parent_task_id": None,
"root_task_id": task_id,
"source": "runtime_bridge",
"workflow_mode": "multi_agent",
"allow_handoff": True,
"context": {
**context,
"agent_role": role,
"workflow_mode": "multi_agent",
"resource_grants": self._redact_sensitive(
agent.get("resource_grants") or []
),
},
})
return tasks
def task_event_payload(self, task: Any, task_spec: Dict[str, Any]) -> Dict[str, Any]:
"""Return the minimum task graph payload required by Manager."""
return {
"task_id": task.task_id,
"title": task.title or task_spec.get("title", task.description[:80]),
"description": task.description,
"agent_role": task.agent_role or task_spec.get("agent_role", "general"),
"status": task.status.value if hasattr(task.status, "value") else task.status,
"depends_on": task.depends_on or task_spec.get("depends_on") or [],
"parent_task_id": task.parent_task_id or task_spec.get("parent_task_id"),
"root_task_id": task.root_task_id or task_spec.get("root_task_id"),
"required_capabilities": task.required_capabilities or task_spec.get("required_capabilities") or [],
"source": task.source or task_spec.get("source", "runtime_bridge"),
"attempt": getattr(task, "retry_count", 0),
}
def _requires_approval(self, body: Dict[str, Any]) -> bool:
plan = body.get("orchestration_plan") or {}
agile = plan.get("agile_context") or {}
risk = str(plan.get("risk_level") or body.get("risk_level") or "").lower()
return bool(agile.get("requires_user_approval")) or risk == "high"
def status_payload(
self,
run: SwarmRun,
phase: str,
reason: Optional[str] = None,
approval_id: Optional[str] = None,
decision: Optional[str] = None,
deliverable: Optional[Dict[str, Any]] = None,
) -> Dict[str, Any]:
"""Build the standard runtime status payload."""
payload: Dict[str, Any] = {
"status": run.status,
"runtime_execution_status": run.status,
"runtime_deployment_id": run.deployment_id,
"manager_deployment_id": run.manager_deployment_id,
"phase": phase,
}
if reason:
payload["reason"] = reason
if approval_id:
payload["approval_id"] = approval_id
if decision:
payload["decision"] = decision
if deliverable is not None:
payload["deliverable"] = deliverable
return payload
async def _upsert_callback_attempt(
self,
swarm_id: str,
event_id: str,
attempt: Dict[str, Any],
):
run = await self.get_run(swarm_id)
if not run:
return
callback_attempts = list(run.metadata.get("callback_attempts") or [])
filtered = [item for item in callback_attempts if item.get("event_id") != event_id]
filtered.append(self._redact_sensitive(attempt))
run.metadata["callback_attempts"] = filtered[-50:]
await self.save_run(run)
def _build_approval(self, body: Dict[str, Any], run: SwarmRun) -> Dict[str, Any]:
grants = body.get("resource_grants") or []
grant = grants[0] if grants else {}
return {
"approval_id": f"appr_{uuid.uuid4().hex[:12]}",
"operation": "git.write",
"resource_id": grant.get("resource_id", "repo-main"),
"resource_type": grant.get("resource_type", "git"),
"resource_scope": ",".join(grant.get("permission_scope") or []),
"target_role": grant.get("target_role", "general"),
"risk_level": "high",
"requires_credential": True,
"secret_ref": grant.get("secret_ref") or grant.get("ref"),
"ttl_seconds": 900,
"reason": "High-risk swarm run requires Manager approval",
"runtime_deployment_id": run.deployment_id,
}
def _validate_resource_grants(self, grants: List[Dict[str, Any]], path: str):
for index, grant in enumerate(grants):
for key in ("secret_ref", "ref"):
value = grant.get(key)
if value and not self._is_azkv_ref(value):
raise RuntimeValidationError(f"{path}[{index}].{key} must use azkv://")
def _reject_plaintext_secrets(self, value: Any, path: str):
if isinstance(value, dict):
for key, item in value.items():
child_path = f"{path}.{key}"
key_lower = key.lower()
is_ref_key = key_lower.endswith("_ref") or key_lower == "ref"
if self._looks_sensitive_key(key_lower) and not is_ref_key:
raise RuntimeValidationError(
f"Plaintext secret-like field is not allowed: {child_path}"
)
if is_ref_key and isinstance(item, str) and item and not self._is_azkv_ref(item):
raise RuntimeValidationError(f"{child_path} must use azkv://")
self._reject_plaintext_secrets(item, child_path)
elif isinstance(value, list):
for index, item in enumerate(value):
self._reject_plaintext_secrets(item, f"{path}[{index}]")
def _looks_sensitive_key(self, key: str) -> bool:
return (
any(part in key for part in ("token", "password", "passwd", "secret", "private_key", "api_key"))
or key.endswith("_key")
)
def _is_azkv_ref(self, value: Any) -> bool:
return isinstance(value, str) and value.startswith("azkv://")
def _redact_sensitive(self, value: Any) -> Any:
"""Remove obvious plaintext secrets from callback-safe payloads."""
sensitive_keys = {
"password",
"passwd",
"token",
"api_token",
"access_token",
"refresh_token",
"private_key",
"access_key",
"secret",
"client_secret",
"connection_string",
}
if isinstance(value, dict):
redacted = {}
for key, item in value.items():
key_lower = key.lower()
if key_lower in {"secret_ref", "credential_ref", "signing_secret_ref"}:
redacted[key] = item
elif key_lower in sensitive_keys or key_lower.endswith(("_token", "_secret", "_password", "_key")):
redacted[key] = "[redacted]"
else:
redacted[key] = self._redact_sensitive(item)
return redacted
if isinstance(value, list):
return [self._redact_sensitive(item) for item in value]
return value
def _now_iso(self) -> str:
return datetime.now(timezone.utc).isoformat().replace("+00:00", "Z")
swarm_runtime = SwarmRuntime()