Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,12 @@ a2a-grpc = [
a2a-redis = [
"redis>=5.0.0",
]
agui = [
"ag-ui-protocol==0.1.20",
"a2a-sdk[http-server,signing]>=1.0.2,<2",
"starlette>=0.39.0",
"uvicorn[standard]>=0.30.0",
]
diagram = ["termaid>=0.1; python_version >= '3.11'"]

[dependency-groups]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,13 @@ The runner substitutes `{run_id}`, `{stack_name}`, `{vswitch_name}`, and
请使用 alicloud-ros-agent Skill 的 {mode} 模式完成这个真实测试:在 cn-hangzhou 查询已有 VPC,选择其中一个,
只在该 VPC 内通过 ROS Stack 部署一个新 VSwitch。Stack 名称必须是 {stack_name},VSwitch 名称必须是
{vswitch_name},CIDR 不得与已有网段冲突。不要创建或删除 VPC,也不要修改其他资源。
每次调用 Skill bridge 时,命令都必须显式设置 `ALICLOUD_ROS_AGENT_STATE_DIR={state_dir}`,不得使用其他状态目录。
不要在 Qoder 本地自行执行云查询、生成模板或部署;第一条基础设施操作必须通过该 Skill 的 managed `start`
命令交给 ROS Agent,start 命令必须使用精确的小写参数 `--mode {mode_arg}`;之后只按 Skill 的
`follow`、`continue` 或 `respond` 工作流处理同一个 job。
ROS Agent 只使用任务工作区中的文件和已有工具完成任务,不要读取 iac-code 仓库源码或任务工作区外的路径。
在 managed `start` 前可以按 Skill 要求执行一次 readiness `check`;check 完成后必须进入同一次测试的 start。
整个测试只能执行一次 managed `start`,后续不得新建 job。

执行过程中请持续用简短文字解释当前阶段。只读云查询不应申请权限;任何非只读操作都必须等待我明确确认。
部署确认前必须展示部署摘要和 Mermaid 架构图。Pipeline 模式必须生成恰好两个都满足约束且确有差异的候选
Expand All @@ -25,13 +32,15 @@ step 开始/结束和候选选择,并在完成后保留同一 job 的 Normal h
不冲突网段的已有 VPC;只需保留该 VPC 的精简摘要,不要再次返回完整 VPC 列表。继续生成并校验只含一个
VSwitch 的 ROS Stack {stack_name},VSwitch 名称为 {vswitch_name}。在任何部署确认之前,先用简短说明和
Mermaid 架构图展示已有 VPC 与待建 VSwitch 的关系;确认和非只读云操作都必须等待我的明确回答。
每次调用 Skill bridge 时,命令都必须显式设置 `ALICLOUD_ROS_AGENT_STATE_DIR={state_dir}`,不得使用其他状态目录。
```

## Confirm deployment

```text
我已经审阅刚才展示的部署摘要和 Mermaid 架构图,确认仅在所选已有 VPC 中通过 Stack {stack_name} 创建
VSwitch {vswitch_name}。请继续同一个 ROS Agent job;遇到非只读云权限时仍需把权限申请返回给我,不得替我批准。
每次调用 Skill bridge 时,命令都必须显式设置 `ALICLOUD_ROS_AGENT_STATE_DIR={state_dir}`,不得使用其他状态目录。
```

## Cleanup
Expand All @@ -41,13 +50,15 @@ VSwitch {vswitch_name}。请继续同一个 ROS Agent job;遇到非只读云
删除前说明目标并等待我确认;只允许删除这两个本次创建的对象,绝不能删除或修改已有 VPC。Pipeline 场景必须
复用 Pipeline handoff 的 Normal 会话,不得启动新的 Normal job。清理后用只读查询确认 Stack/VSwitch 已不存在,
并确认原有 VPC 仍可用。先展示精简删除摘要并等待我下一条明确确认。
每次调用 Skill bridge 时,命令都必须显式设置 `ALICLOUD_ROS_AGENT_STATE_DIR={state_dir}`,不得使用其他状态目录。
```

## Confirm cleanup

```text
我确认删除本次测试创建的 Stack {stack_name},并让其中的 VSwitch {vswitch_name} 随 Stack 删除。请继续同一个
ROS Agent job;不得删除或修改已有 VPC,遇到非只读云权限时仍需把权限申请返回给我,不得替我批准。
每次调用 Skill bridge 时,命令都必须显式设置 `ALICLOUD_ROS_AGENT_STATE_DIR={state_dir}`,不得使用其他状态目录。
```

## Scripted answers
Expand Down
107 changes: 107 additions & 0 deletions scripts/a2a/e2e/permission_wait/run_start_chat_permission_wait.py
Original file line number Diff line number Diff line change
Expand Up @@ -439,6 +439,14 @@ def _run_qoder(
resume: bool,
run_dir: Path,
) -> dict[str, Any]:
state_dir = str(env.get("ALICLOUD_ROS_AGENT_STATE_DIR") or "")
driver_policy = (
"You are driving one bounded ROS Agent E2E job. Execute at most one ros_agent.py bridge command in each "
"Qoder turn, then stop and return its bounded result. Across this session execute readiness check at most "
"once and managed start exactly once. After a job exists, never start another job; use only follow, continue, "
"or respond for that job. Prefix every ros_agent.py command with ALICLOUD_ROS_AGENT_STATE_DIR={}. Do not "
"replace the remote ROS Agent with local cloud, template, or deployment work."
).format(state_dir)
command = [
str(args.qoder_cli.expanduser().resolve()),
"-p",
Expand All @@ -447,6 +455,8 @@ def _run_qoder(
"--config-dir",
str(args.qoder_config_dir.expanduser().resolve()),
"--dangerously-skip-permissions",
"--append-system-prompt",
driver_policy,
"--cwd",
str(workspace),
]
Expand All @@ -473,6 +483,18 @@ def _run_qoder(
content_block_index = 0
first_mermaid_block_index: int | None = None
first_cloud_permission_block_index: int | None = None
bridge_command_count = 0
bridge_managed_start = False
bridge_managed_start_count = 0
bridge_check_count = 0
bridge_state_dir_bound = False
bridge_state_dir_bound_count = 0
bridge_start_shape_ok = False
bridge_script_path_kinds: set[str] = set()
bridge_tool_use_ids: set[str] = set()
bridge_result_codes: set[str] = set()
bridge_result_states: set[str] = set()
bridge_result_ok_values: set[bool] = set()

def contains_cloud_permission(value: Any) -> bool:
if isinstance(value, dict):
Expand Down Expand Up @@ -511,6 +533,42 @@ def contains_cloud_permission(value: Any) -> bool:
if not isinstance(block, dict):
continue
content_block_index += 1
if block.get("type") == "tool_use":
serialized_block = json.dumps(block, ensure_ascii=False)
if "ros_agent.py" in serialized_block:
tool_use_id = block.get("id")
if isinstance(tool_use_id, str) and tool_use_id:
bridge_tool_use_ids.add(tool_use_id)
bridge_command_count += 1
if " start " in serialized_block:
bridge_managed_start = True
bridge_managed_start_count += 1
bridge_start_shape_ok = bridge_start_shape_ok or all(
token in serialized_block for token in ("--prompt-file", "--mode", "--follow")
)
if " check" in serialized_block:
bridge_check_count += 1
if "ALICLOUD_ROS_AGENT_STATE_DIR" in serialized_block:
bridge_state_dir_bound = True
bridge_state_dir_bound_count += 1
if "<absolute-bridge-path>" in serialized_block:
bridge_script_path_kinds.add("placeholder")
elif "/.qoderwork/skills/alicloud-ros-agent/" in serialized_block:
bridge_script_path_kinds.add("qoderwork")
elif "/.qoder/skills/alicloud-ros-agent/" in serialized_block:
bridge_script_path_kinds.add("qoder")
elif "/skills/alicloud-ros-agent/" in serialized_block:
bridge_script_path_kinds.add("repository")
else:
bridge_script_path_kinds.add("other")
if block.get("type") == "tool_result" and block.get("tool_use_id") in bridge_tool_use_ids:
serialized_result = json.dumps(block.get("content"), ensure_ascii=False)
for code in re.findall(r'"code"\s*:\s*"([A-Za-z0-9_.-]{1,80})"', serialized_result):
bridge_result_codes.add(code)
for state in re.findall(r'"state"\s*:\s*"([A-Za-z0-9_.-]{1,80})"', serialized_result):
bridge_result_states.add(state)
for raw_ok in re.findall(r'"ok"\s*:\s*(true|false)', serialized_result, re.IGNORECASE):
bridge_result_ok_values.add(raw_ok.casefold() == "true")
if item.get("type") == "assistant" and block.get("type") == "text":
text = block.get("text")
if isinstance(text, str) and text.strip():
Expand All @@ -537,6 +595,17 @@ def contains_cloud_permission(value: Any) -> bool:
"firstMermaidBlockIndex": first_mermaid_block_index,
"firstCloudPermissionBlockIndex": first_cloud_permission_block_index,
"mentionsFollow": " follow" in completed.stdout.casefold(),
"bridgeCommandCount": bridge_command_count,
"bridgeManagedStart": bridge_managed_start,
"bridgeManagedStartCount": bridge_managed_start_count,
"bridgeCheckCount": bridge_check_count,
"bridgeStateDirBound": bridge_state_dir_bound,
"bridgeStateDirBoundCount": bridge_state_dir_bound_count,
"bridgeStartShapeOk": bridge_start_shape_ok,
"bridgeScriptPathKinds": sorted(bridge_script_path_kinds),
"bridgeResultCodes": sorted(bridge_result_codes),
"bridgeResultStates": sorted(bridge_result_states),
"bridgeResultOkValues": sorted(bridge_result_ok_values),
}
_append_jsonl(run_dir / "qoder-turns.jsonl", evidence)
if completed.returncode != 0:
Expand Down Expand Up @@ -949,6 +1018,7 @@ def a2a_service(mode: str, config_path: Path, port: int) -> _Service:
deployment_confirmation_attempts = 0
cleanup_confirmation_attempts = 0
architecture_seen = False
readiness_only_turn_seen = False
job_path: Path | None = None
skill_installation_backups: list[_SkillInstallationBackup] = []
read_only_cloud_evidence: list[dict[str, Any]] = []
Expand All @@ -959,6 +1029,8 @@ def a2a_service(mode: str, config_path: Path, port: int) -> _Service:
"stack_name": stack_name,
"vswitch_name": vswitch_name,
"mode": "Normal" if args.mode == "normal" else "Pipeline",
"mode_arg": args.mode,
"state_dir": str(state_root),
},
)
try:
Expand All @@ -985,8 +1057,25 @@ def a2a_service(mode: str, config_path: Path, port: int) -> _Service:
resume=turn > 0,
run_dir=run_dir,
)
bridge_command_count = int(qoder_evidence.get("bridgeCommandCount") or 0)
bridge_start_count = int(qoder_evidence.get("bridgeManagedStartCount") or 0)
if bridge_start_count and qoder_evidence.get("bridgeStateDirBound") is not True:
raise AssertionError("Qoder Skill bridge command did not bind the E2E state directory")
architecture_seen = architecture_seen or bool(qoder_evidence.get("assistantMermaid"))
jobs = _jobs(state_root)
if not jobs:
if bridge_command_count < 1:
raise AssertionError("Qoder did not execute a managed Skill bridge command")
bridge_check_count = int(qoder_evidence.get("bridgeCheckCount") or 0)
if readiness_only_turn_seen or bridge_check_count != bridge_command_count:
raise AssertionError("Qoder did not start the managed ROS Agent job after readiness")
readiness_only_turn_seen = True
next_prompt = (
"readiness check 已完成。现在必须且只执行一次 alicloud-ros-agent Skill 的 managed start,"
"使用精确的小写参数 --mode {} 开始之前给出的部署任务;不得再次 check,也不得在本地替代执行。"
"每条 bridge 命令必须显式设置 ALICLOUD_ROS_AGENT_STATE_DIR={}.".format(args.mode, state_root)
)
continue
if len(jobs) != 1:
raise AssertionError("expected exactly one ROS Agent job, found {}".format(len(jobs)))
job_path, job = jobs[0]
Expand All @@ -1000,6 +1089,14 @@ def a2a_service(mode: str, config_path: Path, port: int) -> _Service:
state = str(job.get("state") or "")
if state in FAILURE_STATES:
raise RuntimeError("ROS Agent job ended in {}".format(state))
text_only_cleanup_summary = (
cleanup_started
and state in TERMINAL_STATES
and int(job.get("turn") or 0) == cleanup_turn
and int(qoder_evidence.get("assistantTextBlocks") or 0) > 0
)
if bridge_command_count < 1 and not text_only_cleanup_summary:
raise AssertionError("Qoder did not execute a managed Skill bridge command")
current = job.get("inputRequired")
if state == "input-required" and isinstance(current, dict):
input_id = current.get("inputId")
Expand Down Expand Up @@ -1066,6 +1163,7 @@ def a2a_service(mode: str, config_path: Path, port: int) -> _Service:
"stack_name": stack_name,
"vswitch_name": vswitch_name,
"mode": args.mode,
"state_dir": str(state_root),
},
)
else:
Expand All @@ -1086,6 +1184,7 @@ def a2a_service(mode: str, config_path: Path, port: int) -> _Service:
"stack_name": stack_name,
"vswitch_name": vswitch_name,
"mode": args.mode,
"state_dir": str(state_root),
},
)
else:
Expand All @@ -1096,6 +1195,7 @@ def a2a_service(mode: str, config_path: Path, port: int) -> _Service:
"stack_name": stack_name,
"vswitch_name": vswitch_name,
"mode": args.mode,
"state_dir": str(state_root),
},
)
continue
Expand All @@ -1112,6 +1212,7 @@ def a2a_service(mode: str, config_path: Path, port: int) -> _Service:
"stack_name": stack_name,
"vswitch_name": vswitch_name,
"mode": args.mode,
"state_dir": str(state_root),
},
)
continue
Expand Down Expand Up @@ -1167,6 +1268,11 @@ def a2a_service(mode: str, config_path: Path, port: int) -> _Service:
metrics_path = run_dir / "relay-metrics.json"
if metrics_path.is_file():
relay_metrics = json.loads(metrics_path.read_text(encoding="utf-8"))
relay_session_ids = {
str(item.get("sessionId"))
for item in relay_metrics.get("requests", [])
if item.get("action") == "StartChat" and item.get("sessionId")
}
qoder_turns = []
qoder_path = run_dir / "qoder-turns.jsonl"
if qoder_path.is_file():
Expand Down Expand Up @@ -1245,6 +1351,7 @@ def a2a_service(mode: str, config_path: Path, port: int) -> _Service:
"native StartChat relay was used": any(
item.get("action") == "StartChat" for item in relay_metrics.get("requests", [])
),
"all StartChat requests stayed in one ROS session": len(relay_session_ids) == 1,
"Qoder emitted explanatory assistant text": sum(
int(item.get("assistantTextBlocks") or 0) for item in qoder_turns
)
Expand Down
40 changes: 40 additions & 0 deletions src/iac_code/a2a/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -39,9 +39,12 @@
resolve_a2a_public_path_roots,
resolve_a2a_public_path_roots_for_data,
)
from iac_code.a2a.types import validate_protocol_id
from iac_code.i18n import _
from iac_code.pipeline.config import get_run_mode
from iac_code.services.configuration_readiness import configuration_readiness
from iac_code.services.session_backup import SessionBackupError
from iac_code.services.session_storage import SessionStorage

logger = logging.getLogger(__name__)
_V03_JSONRPC_METHODS = frozenset(
Expand Down Expand Up @@ -477,6 +480,42 @@ async def get_agent_card(request: Request) -> Response:
async def get_readiness(request: Request) -> JSONResponse:
return JSONResponse(configuration_readiness(model=model))

async def ensure_session_restored(request: Request) -> JSONResponse:
try:
payload = await request.json()
except Exception:
return JSONResponse({"error": "Invalid session restore request."}, status_code=400)
if not isinstance(payload, dict):
return JSONResponse({"error": "Invalid session restore request."}, status_code=400)
raw_cwd = payload.get("cwd")
raw_session_id = payload.get("sessionId")
if not isinstance(raw_cwd, str) or not raw_cwd or not isinstance(raw_session_id, str):
return JSONResponse({"error": "Invalid session restore request."}, status_code=400)

executor = getattr(components.handler, "agent_executor", None)
resolve_cwd = getattr(executor, "_resolve_cwd", None)
backup_service = components.backup_service
if not callable(resolve_cwd) or backup_service is None:
return JSONResponse({"error": "Session restore is unavailable."}, status_code=503)
try:
session_id = validate_protocol_id(raw_session_id)
cwd = resolve_cwd({"iac_code": {"cwd": raw_cwd}})
result = await asyncio.to_thread(backup_service.reconcile_session, cwd, session_id)
exists = SessionStorage().exists(cwd, session_id)
except ValueError:
return JSONResponse({"error": "Invalid session restore request."}, status_code=400)
except SessionBackupError as exc:
logger.warning("A2A session restore failed error_type=%s", type(exc).__name__)
return JSONResponse({"error": "Unable to restore the A2A session."}, status_code=503)
except Exception as exc:
logger.warning("A2A session restore failed error_type=%s", type(exc).__name__)
return JSONResponse({"error": "Unable to restore the A2A session."}, status_code=503)

if not exists:
return JSONResponse({"status": "not_found"}, status_code=404)
status = "restored" if getattr(result, "action", None) == "restored" else "current"
return JSONResponse({"status": status})

recovery_service = A2APipelineRecoveryService(task_store=components.task_store)

async def recovery_path_roots(*, context_id: str | None, task_id: str | None) -> list[dict[str, str]]:
Expand Down Expand Up @@ -523,6 +562,7 @@ async def get_pipeline_state(request: Request) -> JSONResponse:
Route("/health", health, methods=["GET"]),
Route(AGENT_CARD_WELL_KNOWN_PATH, get_agent_card, methods=["GET"]),
Route("/iac-code/readiness", get_readiness, methods=["GET"]),
Route("/iac-code/session/ensure-restored", ensure_session_restored, methods=["POST"]),
Route("/iac-code/pipeline/state", get_pipeline_state, methods=["GET"]),
]
install_jsonrpc_error_data_passthrough()
Expand Down
Loading
Loading