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
5 changes: 5 additions & 0 deletions .github/workflows/test.yml
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,11 @@ concurrency:
group: ${{ github.workflow }}-${{ github.ref }}
cancel-in-progress: true

env:
UV_CONCURRENT_DOWNLOADS: "8"
UV_HTTP_TIMEOUT: "60"
UV_HTTP_RETRIES: "5"

jobs:
skill-bridge-compatibility:
name: Skill bridge (Python ${{ matrix.python-version }})
Expand Down
1 change: 1 addition & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -189,5 +189,6 @@ src/iac_code/pipeline/engine/architecture_rules.json.backup-*
# Super Powers
.superpowers/
.worktrees/
ci-e2e-report/
.agents/
spec/
2 changes: 1 addition & 1 deletion scripts/a2a/e2e/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -460,7 +460,7 @@ the rest of the tests.
| `image-ask-waiting` | `ask_user_question` waits for user input, then the server restarts | Static `ask-first-answer.png` / `ask-second-answer.png` image fixtures without `taskId` | Pending ask input is recovered, image answers hydrate the recovered task, and the pipeline completes with VSwitch evidence. |
| `image-selection-waiting` | Step 4 waits for candidate selection, then the server restarts | Static `selection.png` image fixture without `taskId` | Waiting step4 task is recovered, the image selection is accepted, and VSwitch evidence exists. |
| `image-normal-handoff` | Pipeline completes and hands off to normal chat; the normal follow-up is static `normal-followup.png`, then the server restarts | Normal-chat recovery question without `taskId` | Image follow-up stays in the same `contextId`, uses a new normal-chat task, and completed handoff state survives restart. |
| `image-interrupt` | Step 3 receives static `rollback-interrupt.png` as an image rollback to `intent_parsing`, then the server restarts | `继续`, plus selection when needed | The image interrupt is recognized, the pipeline completes as a security-group task, and final deployment evidence is not VSwitch. |
| `image-interrupt` | Step 3 receives `rollback-interrupt.png` with an image caption explicitly requesting a restart from `intent_parsing`; kill the server only after that new attempt starts | `继续`, plus selection when needed | The image interrupt is recognized, the pipeline completes as a security-group task, and final deployment evidence is not VSwitch. |
| `step1-running` | `intent_parsing` running | `继续` | Running pipeline task is recovered and completes; VSwitch evidence exists. |
| `step2-running` | `architecture_planning` running | `继续` | Running pipeline task is recovered and completes; VSwitch evidence exists. |
| `step3-running` | `evaluate_candidates` candidate/sub-pipeline running | `继续` | Sub-pipeline state is recovered and completes; VSwitch evidence exists. |
Expand Down
2 changes: 1 addition & 1 deletion scripts/a2a/e2e/README.zh-CN.md
Original file line number Diff line number Diff line change
Expand Up @@ -520,7 +520,7 @@ provider、tool、真实云调用场景默认会被保护住。只有确认要
| `image-ask-waiting` | `ask_user_question` 等待用户输入,随后重启 server | 不带 `taskId` 发送静态 `ask-first-answer.png` / `ask-second-answer.png` 图片 fixture | pending ask 输入能恢复,图片回答能 hydrate 到恢复后的 task,最终完成并产生 VSwitch 证据。 |
| `image-selection-waiting` | step4 等待候选方案选择,随后重启 server | 不带 `taskId` 发送静态 `selection.png` 图片 fixture | 能恢复等待中的 step4 task,图片选择被接受,并产生 VSwitch 证据。 |
| `image-normal-handoff` | pipeline 完成并 handoff 到 normal chat;normal follow-up 是静态 `normal-followup.png`,随后重启 server | 不带 `taskId` 发送 normal-chat 恢复问题 | 图片 follow-up 保持同一个 `contextId`,使用新的 normal-chat task;completed handoff 状态重启后仍可恢复。 |
| `image-interrupt` | step3 收到静态 `rollback-interrupt.png` 图片,表示回滚到 `intent_parsing`,随后重启 server | `继续`,必要时再选择方案 | 图片 interrupt 能被识别;pipeline 以安全组任务完成,最终部署证据不是 VSwitch。 |
| `image-interrupt` | step3 收到 `rollback-interrupt.png`,图片附注明确请求从 `intent_parsing` 重新解析;只有该新尝试开始后才强杀并重启 server | `继续`,必要时再选择方案 | 图片 interrupt 能被识别;pipeline 以安全组任务完成,最终部署证据不是 VSwitch。 |
| `step1-running` | `intent_parsing` 运行中 | `继续` | running pipeline task 能恢复并完成;存在 VSwitch 证据。 |
| `step2-running` | `architecture_planning` 运行中 | `继续` | running pipeline task 能恢复并完成;存在 VSwitch 证据。 |
| `step3-running` | `evaluate_candidates` 的 candidate/sub-pipeline 运行中 | `继续` | sub-pipeline 状态能恢复并完成;存在 VSwitch 证据。 |
Expand Down
184 changes: 184 additions & 0 deletions scripts/a2a/e2e/cleanup_owned_stacks.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,184 @@
#!/usr/bin/env python3
"""Delete only ROS stacks with an accepted creation receipt in this E2E case."""

from __future__ import annotations

import argparse
import json
import re
import time
from pathlib import Path
from typing import Any


class CleanupOperationError(RuntimeError):
"""Expose a fixed cleanup stage without leaking the cloud error message."""

def __init__(self, stage: str, cause: Exception) -> None:
self.stage = stage
self.cause_type = type(cause).__name__
code = getattr(cause, "code", None)
self.sdk_code = code if isinstance(code, str) and re.fullmatch(r"[A-Za-z][A-Za-z0-9_.-]{0,79}", code) else ""
super().__init__(f"{stage}: {self.cause_type}")


def _record_cleanup_failure(run_dir: Path, stage: str, exc: Exception) -> None:
known_codes = {
"EntityNotExist.Stack", "NotFound.Stack", "StackNotFound", "ActionInProgress",
"Forbidden", "Forbidden.RAM", "InvalidAccessKeyId.NotFound", "SecurityTokenExpired",
"Throttling", "Throttling.User", "InvalidParameter", "DeleteFailed", "DependencyViolation",
}
code = getattr(exc, "code", None)
diagnostic = {"stage": stage, "errorType": "SDKError",
"code": code if isinstance(code, str) and code in known_codes else "unknown"}
with (run_dir / "cleanup-cloud.log").open("a", encoding="utf-8") as stream:
stream.write(json.dumps({"cleanupDiagnostic": diagnostic}) + "\n")


def _stack_body(client: Any, models: Any, stack_id: str, region: str) -> dict[str, Any] | None:
try:
return client.get_stack(models.GetStackRequest(stack_id=stack_id, region_id=region)).body.to_map()
except Exception as exc:
code = getattr(exc, "code", None)
missing_codes = {"EntityNotExist.Stack", "NotFound.Stack", "StackNotFound"}
if isinstance(code, str):
if code in missing_codes:
return None
elif any(re.search(r"\b" + re.escape(marker) + r"\b", str(exc), re.I) for marker in missing_codes):
return None
raise


def cleanup_owned_stacks(run_dir: Path, *, timeout: float = 840) -> dict[str, Any]:
from iac_code.services.session_storage import SessionStorage
from scripts.ci.stack_ownership import creation_receipts

manifest = json.loads((run_dir / "owned-stacks.json").read_text(encoding="utf-8"))
if not isinstance(manifest.get("configDir"), str) or not isinstance(manifest.get("cwd"), str):
raise ValueError("invalid E2E ownership manifest")
storage = SessionStorage(projects_dir=Path(manifest["configDir"]) / "projects")
directories: list[Path] = []
for context_path in sorted((run_dir / "a2a-persistence" / "contexts").glob("*.json")):
context = json.loads(context_path.read_text(encoding="utf-8"))
session_id = context.get("session_id")
if (not isinstance(session_id, str) or re.fullmatch(r"[A-Za-z0-9_-]{1,128}", session_id) is None
or context.get("cwd") != manifest["cwd"]):
raise ValueError("A2A context does not prove this case's session ownership")
session = storage.session_dir(manifest["cwd"], session_id)
if not session.resolve().is_relative_to((Path(manifest["configDir"]) / "projects").resolve()):
raise ValueError("Stack ownership evidence escaped isolated case")
directories.extend([session / "pipeline", session / "a2a" / "pipeline"])
resources = creation_receipts(directories)
from alibabacloud_ros20190910 import models as ros_models

from iac_code.services.cloud_credentials import CloudCredentials
from iac_code.tools.cloud.aliyun.ros_client import RosClientFactory

try:
credential = CloudCredentials().get_provider("aliyun")
except Exception as exc:
raise CleanupOperationError("credential_lookup", exc) from exc
if credential is None:
raise RuntimeError("Aliyun credential is unavailable for E2E teardown")
region = str(manifest.get("regionId") or credential.region_id)
try:
client = RosClientFactory.create(credential, region)
except Exception as exc:
raise CleanupOperationError("client_create", exc) from exc
deadline = time.monotonic() + timeout
deleted: list[str] = []
remaining: list[str] = []
failures: list[str] = []
for resource in resources:
stack_id, name, region = resource["stackId"], resource["stackName"], resource["regionId"]
client = RosClientFactory.create(credential, region)
delete_submitted = False
while time.monotonic() < deadline:
stage = "get_stack"
try:
body = _stack_body(client, ros_models, stack_id, region)
if body is None or body.get("Status") == "DELETE_COMPLETE":
deleted.append(stack_id)
break
if body.get("StackName") != name or body.get("ParentStackId") or body.get("ServiceManaged"):
failures.append(stack_id + ": Stack identity differs from accepted creation receipt")
break
status = str(body.get("Status") or "")
if status == "DELETE_FAILED" and delete_submitted:
failures.append(stack_id + ": accepted deletion failed")
break
if not status.endswith("_IN_PROGRESS") and not delete_submitted:
stage = "delete_stack"
try:
client.delete_stack(ros_models.DeleteStackRequest(stack_id=stack_id, region_id=region))
delete_submitted = True
except Exception as exc:
if getattr(exc, "code", "") in {"EntityNotExist.Stack", "NotFound.Stack", "StackNotFound"}:
deleted.append(stack_id)
break
if getattr(exc, "code", "") != "ActionInProgress":
raise
time.sleep(min(5, max(0, deadline - time.monotonic())))
except Exception as exc:
_record_cleanup_failure(run_dir, stage, exc)
failures.append(stack_id + ": " + type(exc).__name__)
break
else:
failures.append(stack_id + ": cleanup timeout")
if stack_id not in deleted:
remaining.append(stack_id)
# Real pipeline Stack events can reveal an unproven leak. Read those exact
# IDs without granting deletion authority to observations alone.
from scripts.a2a.debugger import _extract_pipeline_envelopes

observed_ids: set[str] = set()
for path in sorted(run_dir.glob("*.events.jsonl"))[:60]:
if path.stat().st_size > 20_000_000:
raise RuntimeError("observed Stack audit evidence exceeds bounded size")
for line in path.read_text(encoding="utf-8").splitlines():
try:
row = json.loads(line)
except ValueError:
continue
for envelope in _extract_pipeline_envelopes(row):
data = envelope.get("data")
if envelope.get("eventType") != "stack_current_changed" or not isinstance(data, dict):
continue
stack_id = data.get("stackId")
if isinstance(stack_id, str) and stack_id:
observed_ids.add(stack_id)
if len(observed_ids) > 60:
raise RuntimeError("observed Stack audit exceeds bounded count")
try:
for stack_id in sorted(observed_ids - set(deleted) - set(remaining)):
body = _stack_body(client, ros_models, stack_id, region)
if body is not None and body.get("Status") != "DELETE_COMPLETE":
failures.append("observed Stack has no accepted creation receipt or cleanup incomplete")
remaining.append(stack_id)
except Exception as exc:
raise CleanupOperationError("observed_stack_audit", exc) from exc
result = {
"status": "failed" if failures or remaining else "completed",
"deletedStackIds": deleted,
"remainingStackIds": remaining,
"failures": failures,
"resources": resources,
}
(run_dir / "cleanup-result.json").write_text(
json.dumps(result, ensure_ascii=False, indent=2) + "\n", encoding="utf-8"
)
return result


def main() -> int:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--run-dir", type=Path, required=True)
parser.add_argument("--timeout", type=float, default=840)
args = parser.parse_args()
result = cleanup_owned_stacks(args.run_dir, timeout=args.timeout)
print(json.dumps(result, ensure_ascii=False))
return 0 if result["status"] == "completed" else 1


if __name__ == "__main__":
raise SystemExit(main())
Loading
Loading