Skip to content

Commit b07e85e

Browse files
authored
Merge pull request #325 from aliyun/codex/fix-agui-sequential-permission-resume
fix: fence permission resume on staged backups
2 parents 5c5ddf5 + 4b296d1 commit b07e85e

41 files changed

Lines changed: 2222 additions & 104 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

scripts/a2a/e2e/README.zh-CN.md

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -81,6 +81,18 @@ uv run pytest -q tests/a2a_e2e/test_sub_pipeline_permission_timeout.py
8181
uv run pytest -q tests/a2a_e2e/test_permission_wait_restart.py
8282
```
8383

84+
其中 `staged-backup-generation-fence` 场景使用真实 HTTP A2A 进程、生产 staged backup/worker 和两个连续
85+
权限点:新进程先从共享目录恢复旧 generation;较新的权限 backup 尚未发布时,Resume 必须返回可重试的
86+
`SESSION_BACKUP_NOT_READY`;worker 发布完成后,同一标准权限响应可恢复最新会话并且两个工具各执行一次。
87+
可单独运行:
88+
89+
```bash
90+
uv run python scripts/a2a/e2e/permission_wait/run_permission_wait_restart.py \
91+
--run-dir /tmp/iac-pwait-generation-fence \
92+
--decision allow_once \
93+
--staged-backup-generation-fence
94+
```
95+
8496
进程重启矩阵覆盖 Normal/Pipeline × allow/deny。Sub Pipeline fixture 断言只生成一次拒绝 ToolResult、
8597
Agent loop 继续、父 Pipeline 进入 candidate 选择并完成,且全程没有 grace、持久化 permission checkpoint
8698
或权限关键备份。

scripts/a2a/e2e/permission_wait/permission_wait_fixture_server.py

Lines changed: 30 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -31,10 +31,20 @@ def _parse_args() -> argparse.Namespace:
3131
parser.add_argument("--candidate-first", action="store_true")
3232
parser.add_argument("--pipeline-step-id", choices=SELLING_STAGE_IDS)
3333
parser.add_argument("--handoff-first", action="store_true")
34+
parser.add_argument("--sequential-permissions", action="store_true")
35+
parser.add_argument("--defer-staging-publisher", action="store_true")
36+
parser.add_argument("--staging-dir", type=Path)
37+
parser.add_argument("--backup-dir", type=Path)
3438
return parser.parse_args()
3539

3640

37-
def _create_fixture_runtime(options: Any, *, execution_log: Path, handoff_first: bool = False) -> Any:
41+
def _create_fixture_runtime(
42+
options: Any,
43+
*,
44+
execution_log: Path,
45+
handoff_first: bool = False,
46+
sequential_permissions: bool = False,
47+
) -> Any:
3848
from iac_code.agent.agent_loop import AgentLoop
3949
from iac_code.providers.base import ToolDefinition
4050
from iac_code.services.agent_factory import AgentRuntime
@@ -90,7 +100,7 @@ async def execute(self, *, tool_input: dict[str, Any], context: ToolContext) ->
90100
del context
91101
execution_log.parent.mkdir(parents=True, exist_ok=True)
92102
with execution_log.open("a", encoding="utf-8") as handle:
93-
handle.write("executed\n")
103+
handle.write("{}\n".format(tool_input.get("value", "executed")))
94104
handle.flush()
95105
os.fsync(handle.fileno())
96106
return ToolResult.success("fixture write completed")
@@ -126,6 +136,13 @@ async def stream(
126136
name=tool_name,
127137
input=tool_input,
128138
)
139+
if sequential_permissions:
140+
yield ToolUseStartEvent(tool_use_id="fixture-tool-2", name=tool_name)
141+
yield ToolUseEndEvent(
142+
tool_use_id="fixture-tool-2",
143+
name=tool_name,
144+
input={"value": "executed-2"},
145+
)
129146
yield MessageEndEvent(stop_reason="tool_use", usage=Usage())
130147

131148
provider = FixtureProvider()
@@ -498,12 +515,18 @@ def main() -> int:
498515
os.environ["ALIBABA_CLOUD_ACCESS_KEY_ID"] = "permission-wait-fixture-ak"
499516
os.environ["ALIBABA_CLOUD_ACCESS_KEY_SECRET"] = "permission-wait-fixture-secret"
500517
os.environ["ALIBABA_CLOUD_REGION_ID"] = "cn-hangzhou"
518+
if args.staging_dir is not None or args.backup_dir is not None:
519+
if args.staging_dir is None or args.backup_dir is None:
520+
raise ValueError("--staging-dir and --backup-dir must be provided together")
521+
os.environ["IAC_CODE_CONFIG_BACKUP_TMP_DIR"] = str(args.staging_dir.expanduser().resolve())
522+
os.environ["IAC_CODE_CONFIG_BACKUP_DIR"] = str(args.backup_dir.expanduser().resolve())
501523

502524
import uvicorn
503525

504526
from iac_code.a2a import executor as executor_module
505527
from iac_code.a2a import pipeline_executor as pipeline_executor_module
506528
from iac_code.a2a.app import create_app
529+
from iac_code.services import session_backup_staging as backup_staging_module
507530
from iac_code.services.providers.aliyun_identity import AliyunCallerIdentity, AliyunCallerIdentityResolver
508531

509532
async def resolve_fixture_identity(
@@ -516,10 +539,15 @@ async def resolve_fixture_identity(
516539

517540
AliyunCallerIdentityResolver.resolve = resolve_fixture_identity
518541

542+
if args.defer_staging_publisher:
543+
backup_staging_module.SessionBackupStagingProcess.start = lambda self: None
544+
backup_staging_module.SessionBackupStagingProcess.close = lambda self: None
545+
519546
executor_module.create_agent_runtime = lambda options: _create_fixture_runtime(
520547
options,
521548
execution_log=execution_log,
522549
handoff_first=args.handoff_first,
550+
sequential_permissions=args.sequential_permissions,
523551
)
524552
pipeline_executor_module.create_agent_runtime = executor_module.create_agent_runtime
525553
fixture_pipelines: dict[str, Any] = {}

0 commit comments

Comments
 (0)