diff --git a/adapters/claude/src/nemo_fabric_adapters/claude/adapter.py b/adapters/claude/src/nemo_fabric_adapters/claude/adapter.py index 0908ca4ae..e004f0485 100644 --- a/adapters/claude/src/nemo_fabric_adapters/claude/adapter.py +++ b/adapters/claude/src/nemo_fabric_adapters/claude/adapter.py @@ -32,6 +32,7 @@ from claude_agent_sdk import HookMatcher from claude_agent_sdk._errors import MessageParseError from nemo_fabric_adapters.common import lifecycle +from nemo_fabric_adapters.common import relay_artifacts from nemo_fabric_adapters.common import relay_gateway from nemo_fabric_adapters.common import relay_hooks from nemo_fabric_adapters.common import utils as common_utils @@ -726,6 +727,8 @@ def child_environment( def _relay_output( output: dict[str, Any], relay: ClaudeRelaySettings, + *, + artifacts: list[dict[str, str]] | None = None, ) -> dict[str, Any]: output["relay_runtime"] = { "enabled": True, @@ -735,8 +738,10 @@ def _relay_output( "gateway_url": relay.gateway.url, "gateway_log_path": str(relay.gateway.log_path), } - output["relay_artifacts"] = common_utils.collect_relay_artifacts( - relay.plugin_config + output["relay_artifacts"] = ( + common_utils.collect_relay_artifacts(relay.plugin_config) + if artifacts is None + else artifacts ) return output @@ -882,7 +887,7 @@ async def invoke(self, invocation: dict[str, Any]) -> dict[str, Any]: if self._unusable: return _failure( "claude_runtime_unavailable", - "Claude runtime cannot accept another invocation after an SDK failure", + "Claude runtime cannot accept another invocation after a runtime failure", ) try: @@ -891,12 +896,42 @@ async def invoke(self, invocation: dict[str, Any]) -> dict[str, Any]: except ClaudeAdapterError as error: output = adapter_failure(error) else: + relay = self._relay + atif_before = ( + relay_artifacts.snapshot_atif_files(relay.plugin_config) + if relay is not None + and relay_artifacts.expects_local_atif(relay.plugin_config) + else None + ) output = await self._run_query( payload, client, prompt, invocation_timeout, ) + if ( + output.get("completed") + and relay is not None + and atif_before is not None + ): + finalized = await relay_artifacts.wait_for_finalized_atif( + relay.plugin_config, atif_before + ) + if finalized is None: + self._unusable = True + return _relay_output( + adapter_failure( + AdapterRelayError( + "claude_relay_atif_timeout", + "NeMo Relay did not finalize an ATIF artifact before the deadline", + metadata={ + "timeout_seconds": relay_artifacts.ATIF_FINALIZATION_TIMEOUT_SECONDS, + }, + ) + ), + relay, + artifacts=[], + ) if self._relay is not None: output = _relay_output(output, self._relay) diff --git a/adapters/codex/src/nemo_fabric_adapters/codex/adapter.py b/adapters/codex/src/nemo_fabric_adapters/codex/adapter.py index 7a13fc9d7..ab8ff1f63 100644 --- a/adapters/codex/src/nemo_fabric_adapters/codex/adapter.py +++ b/adapters/codex/src/nemo_fabric_adapters/codex/adapter.py @@ -32,6 +32,7 @@ import nemo_fabric_adapters.common.relay_gateway as relay_gateway import nemo_fabric_adapters.common.relay_hooks as relay_hooks +import nemo_fabric_adapters.common.relay_artifacts as relay_artifacts import nemo_fabric_adapters.common.utils as common_utils from nemo_fabric_adapters.common import lifecycle @@ -922,7 +923,12 @@ async def _invoke_thread( return sdk_failure(error), False -def _relay_output(output: dict[str, Any], relay: CodexRelaySettings) -> dict[str, Any]: +def _relay_output( + output: dict[str, Any], + relay: CodexRelaySettings, + *, + artifacts: list[dict[str, str]] | None = None, +) -> dict[str, Any]: output["relay_runtime"] = { "enabled": True, "emitter": "codex-sdk/nemo-relay", @@ -931,8 +937,10 @@ def _relay_output(output: dict[str, Any], relay: CodexRelaySettings) -> dict[str "gateway_url": relay.gateway.url, "gateway_log_path": str(relay.gateway.log_path), } - output["relay_artifacts"] = common_utils.collect_relay_artifacts( - relay.plugin_config + output["relay_artifacts"] = ( + common_utils.collect_relay_artifacts(relay.plugin_config) + if artifacts is None + else artifacts ) return output @@ -1064,18 +1072,50 @@ async def invoke(self, invocation: dict[str, Any]) -> dict[str, Any]: "request": invocation.get("request"), } if self._unusable: - output = _failure( + return _failure( "codex_runtime_unavailable", - "Codex runtime cannot accept another invocation after an SDK failure", + "Codex runtime cannot accept another invocation after a runtime failure", ) - return _relay_output(output, self._relay) if self._relay else output try: request_prompt(payload) timeout_seconds(payload) _reasoning_effort(payload) _output_schema(payload) + relay = self._relay + atif_before = ( + relay_artifacts.snapshot_atif_files(relay.plugin_config) + if relay is not None + and relay_artifacts.expects_local_atif(relay.plugin_config) + else None + ) output, usable = await _invoke_thread(payload, self._thread) + if ( + output.get("completed") + and relay is not None + and atif_before is not None + ): + finalized = await relay_artifacts.wait_for_finalized_atif( + relay.plugin_config, atif_before + ) + if finalized is None: + self._unusable = True + return _relay_output( + adapter_failure( + AdapterRelayError( + "codex_relay_atif_timeout", + "NeMo Relay did not finalize an ATIF artifact before the deadline", + metadata={ + "timeout_seconds": relay_artifacts.ATIF_FINALIZATION_TIMEOUT_SECONDS, + }, + ) + ), + relay, + artifacts=[], + ) + except AdapterRelayError as error: + output = adapter_failure(error) + usable = False except CodexAdapterError as error: output = adapter_failure(error) usable = True diff --git a/adapters/common/src/nemo_fabric_adapters/common/relay_artifacts.py b/adapters/common/src/nemo_fabric_adapters/common/relay_artifacts.py new file mode 100644 index 000000000..1b2b624d7 --- /dev/null +++ b/adapters/common/src/nemo_fabric_adapters/common/relay_artifacts.py @@ -0,0 +1,109 @@ +# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +"""Readiness checks for artifacts written asynchronously by NeMo Relay.""" + +from __future__ import annotations + +import asyncio +import json +from pathlib import Path +from typing import Any + +from nemo_fabric_adapters.common import utils as common_utils + +ATIF_FINALIZATION_TIMEOUT_SECONDS = 5.0 +ATIF_POLL_INTERVAL_SECONDS = 0.05 +AtifFileFingerprint = tuple[int, int, int, int] +AtifSnapshot = dict[Path, AtifFileFingerprint] + + +def expects_local_atif(plugin_config: dict[str, Any]) -> bool: + """Return whether Relay is configured to write ATIF to the local runtime. + + Relay treats a non-empty storage list as remote-only, so there is no local + artifact for an adapter to await in that configuration. + """ + + for component in plugin_config.get("components", []): + if ( + not isinstance(component, dict) + or component.get("kind") != "observability" + or component.get("enabled", True) is False + ): + continue + config = component.get("config") + if not isinstance(config, dict): + continue + atif = config.get("atif") + if isinstance(atif, dict) and atif.get("enabled") and not atif.get("storage"): + return True + return False + + +def _atif_fingerprint(path: Path) -> AtifFileFingerprint | None: + try: + status = path.stat() + except (OSError, RuntimeError, ValueError): + return None + return (status.st_dev, status.st_ino, status.st_size, status.st_mtime_ns) + + +def snapshot_atif_files(plugin_config: dict[str, Any]) -> AtifSnapshot: + """Capture ATIF file metadata before an adapter invocation. + + Relay requires ``{session_id}`` in the ATIF filename template and creates a + new session scope for each turn, so a new path is the normal case. Metadata + fingerprints also detect a writer that rewrites an existing path without + reading or hashing artifacts from prior turns. + """ + + snapshot: AtifSnapshot = {} + for artifact in common_utils.collect_relay_artifacts(plugin_config): + if artifact.get("kind") != "atif": + continue + path = Path(artifact["path"]) + fingerprint = _atif_fingerprint(path) + if fingerprint is not None: + snapshot[path] = fingerprint + return snapshot + + +def _finalized_atif_path( + plugin_config: dict[str, Any], before: AtifSnapshot +) -> Path | None: + """Find a new or changed ATIF path containing a complete JSON object.""" + + current = snapshot_atif_files(plugin_config) + for path in sorted(current): + if before.get(path) == current[path]: + continue + try: + # Relay writes directly to the final path, so existence alone does + # not prove that the JSON payload has been written completely. + document = json.loads(path.read_bytes()) + except (OSError, UnicodeDecodeError, json.JSONDecodeError): + continue + if isinstance(document, dict): + return path + return None + + +async def wait_for_finalized_atif( + plugin_config: dict[str, Any], + before: AtifSnapshot, + *, + timeout_seconds: float = ATIF_FINALIZATION_TIMEOUT_SECONDS, + poll_interval_seconds: float = ATIF_POLL_INTERVAL_SECONDS, +) -> Path | None: + """Wait for one new or changed, complete ATIF file until a deadline.""" + + loop = asyncio.get_running_loop() + deadline = loop.time() + timeout_seconds + while True: + if path := _finalized_atif_path(plugin_config, before): + return path + remaining = deadline - loop.time() + if remaining <= 0: + return None + await asyncio.sleep(min(poll_interval_seconds, remaining)) diff --git a/tests/adapters/test_adapters_common_relay_artifacts.py b/tests/adapters/test_adapters_common_relay_artifacts.py new file mode 100644 index 000000000..63865abee --- /dev/null +++ b/tests/adapters/test_adapters_common_relay_artifacts.py @@ -0,0 +1,169 @@ +# SPDX-FileCopyrightText: Copyright (c) 2026, NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# SPDX-License-Identifier: Apache-2.0 + +from __future__ import annotations + +import asyncio +import json +from pathlib import Path +from typing import Any + +import pytest +from nemo_fabric_adapters.common import relay_artifacts + + +def atif_plugin_config( + output_directory: Path, + *, + component_enabled: bool = True, + storage: list[dict[str, Any]] | None = None, +) -> dict[str, Any]: + atif: dict[str, Any] = { + "enabled": True, + "output_directory": str(output_directory), + "filename_template": "trajectory-{session_id}.atif.json", + } + if storage is not None: + atif["storage"] = storage + return { + "version": 1, + "components": [ + { + "kind": "observability", + "enabled": component_enabled, + "config": {"atif": atif}, + } + ], + } + + +def test_expects_local_atif_requires_enabled_local_output(tmp_path): + local = atif_plugin_config(tmp_path / "local") + disabled = atif_plugin_config(tmp_path / "disabled", component_enabled=False) + remote = atif_plugin_config( + tmp_path / "remote", + storage=[{"type": "http", "endpoint": "https://example.test/atif"}], + ) + + assert relay_artifacts.expects_local_atif(local) is True + assert relay_artifacts.expects_local_atif(disabled) is False + assert relay_artifacts.expects_local_atif(remote) is False + assert relay_artifacts.expects_local_atif({"components": []}) is False + + +async def test_wait_for_finalized_atif_ignores_unchanged_and_partial_files( + tmp_path, monkeypatch +): + atif_dir = tmp_path / "atif" + atif_dir.mkdir() + existing = atif_dir / "trajectory-existing.atif.json" + candidate = atif_dir / "trajectory-current.atif.json" + existing.write_text('{"schema_version":"ATIF-v1.7"}', encoding="utf-8") + plugin_config = atif_plugin_config(atif_dir) + before = relay_artifacts.snapshot_atif_files(plugin_config) + reads: list[Path] = [] + read_bytes = Path.read_bytes + + def record_read(path: Path) -> bytes: + reads.append(path) + return read_bytes(path) + + monkeypatch.setattr(Path, "read_bytes", record_read) + + async def finish_candidate(): + candidate.write_text("{", encoding="utf-8") + await asyncio.sleep(0.02) + candidate.write_text( + json.dumps({"schema_version": "ATIF-v1.7", "steps": []}), + encoding="utf-8", + ) + + writer = asyncio.create_task(finish_candidate()) + finalized = await relay_artifacts.wait_for_finalized_atif( + plugin_config, + before, + timeout_seconds=0.5, + poll_interval_seconds=0.001, + ) + await writer + + assert set(before) == {existing.resolve()} + assert finalized == candidate.resolve() + assert set(reads) == {candidate.resolve()} + + +async def test_wait_for_finalized_atif_accepts_modified_existing_path(tmp_path): + atif_dir = tmp_path / "atif" + atif_dir.mkdir() + existing = atif_dir / "trajectory-existing.atif.json" + existing.write_text('{"schema_version":"ATIF-v1.7"}', encoding="utf-8") + plugin_config = atif_plugin_config(atif_dir) + before = relay_artifacts.snapshot_atif_files(plugin_config) + existing.write_text( + json.dumps({"schema_version": "ATIF-v1.7", "steps": [{"step_id": 1}]}), + encoding="utf-8", + ) + + finalized = await relay_artifacts.wait_for_finalized_atif( + plugin_config, + before, + timeout_seconds=0.5, + poll_interval_seconds=0.001, + ) + + assert finalized == existing.resolve() + + +async def test_wait_for_finalized_atif_has_a_hard_deadline(tmp_path): + atif_dir = tmp_path / "atif" + atif_dir.mkdir() + plugin_config = atif_plugin_config(atif_dir) + + finalized = await relay_artifacts.wait_for_finalized_atif( + plugin_config, + relay_artifacts.snapshot_atif_files(plugin_config), + timeout_seconds=0.01, + poll_interval_seconds=0.001, + ) + + assert finalized is None + + +async def test_wait_for_finalized_atif_isolated_by_runtime_directory(tmp_path): + first_dir = tmp_path / "runtime-1" + second_dir = tmp_path / "runtime-2" + first_dir.mkdir() + second_dir.mkdir() + first_config = atif_plugin_config(first_dir) + second_config = atif_plugin_config(second_dir) + + async def write_atif(directory: Path, session_id: str, delay: float): + await asyncio.sleep(delay) + path = directory / f"trajectory-{session_id}.atif.json" + path.write_text( + json.dumps({"schema_version": "ATIF-v1.7", "steps": []}), + encoding="utf-8", + ) + return path.resolve() + + first_writer = asyncio.create_task(write_atif(first_dir, "first", 0.02)) + second_writer = asyncio.create_task(write_atif(second_dir, "second", 0.01)) + first_wait = relay_artifacts.wait_for_finalized_atif( + first_config, + relay_artifacts.snapshot_atif_files(first_config), + timeout_seconds=0.5, + poll_interval_seconds=0.001, + ) + second_wait = relay_artifacts.wait_for_finalized_atif( + second_config, + relay_artifacts.snapshot_atif_files(second_config), + timeout_seconds=0.5, + poll_interval_seconds=0.001, + ) + + first_finalized, second_finalized, first_path, second_path = await asyncio.gather( + first_wait, second_wait, first_writer, second_writer + ) + + assert first_finalized == first_path + assert second_finalized == second_path diff --git a/tests/adapters/test_claude_adapter.py b/tests/adapters/test_claude_adapter.py index f97c53d10..7564aa836 100644 --- a/tests/adapters/test_claude_adapter.py +++ b/tests/adapters/test_claude_adapter.py @@ -752,10 +752,33 @@ async def interrupt(self): assert not relay.plugin_path.exists() -async def test_runtime_reports_relay_artifacts(relay_payload, monkeypatch, tmp_path): +def atif_plugin_config(output_directory: Path) -> dict[str, Any]: + return { + "version": 1, + "components": [ + { + "kind": "observability", + "enabled": True, + "config": { + "atif": { + "enabled": True, + "output_directory": str(output_directory), + "filename_template": "trajectory-{session_id}.atif.json", + } + }, + } + ], + } + + +def relay_settings( + tmp_path: Path, plugin_config: dict[str, Any] +) -> adapter.ClaudeRelaySettings: executable = tmp_path / "nemo-relay" executable.touch() - relay = adapter.ClaudeRelaySettings( + plugin_path = tmp_path / "relay-plugin" + plugin_path.mkdir() + return adapter.ClaudeRelaySettings( gateway=adapter.relay_gateway.RelayGatewayLaunch( executable=executable, config_path=tmp_path / "relay-config" / "config.toml", @@ -763,40 +786,48 @@ async def test_runtime_reports_relay_artifacts(relay_payload, monkeypatch, tmp_p url="http://127.0.0.1:43210", log_path=tmp_path / "relay-config" / "gateway.log", ), - plugin_config={ - "version": 1, - "components": [ - { - "kind": "observability", - "enabled": True, - "config": { - "atif": { - "enabled": True, - "output_directory": str(tmp_path / "atif"), - "filename_template": "trajectory-{session_id}.atif.json", - } - }, - } - ], - }, - plugin_path=tmp_path / "relay-plugin", + plugin_config=plugin_config, + plugin_path=plugin_path, ) - relay.plugin_path.mkdir() + + +def install_mock_relay( + monkeypatch: pytest.MonkeyPatch, relay: adapter.ClaudeRelaySettings +) -> tuple[MagicMock, MagicMock, MagicMock]: relay.gateway.log_path.parent.mkdir() relay.gateway.log_path.write_text("gateway started\n", encoding="utf-8") - atif_path = tmp_path / "atif" / "trajectory-session.atif.json" - atif_path.parent.mkdir() - atif_path.write_text("{}", encoding="utf-8") process = MagicMock() mock_start = MagicMock(return_value=process) mock_stop = MagicMock() monkeypatch.setattr(adapter, "prepare_claude_relay", MagicMock(return_value=relay)) monkeypatch.setattr(adapter.relay_gateway, "start_relay_gateway", mock_start) monkeypatch.setattr(adapter.relay_gateway, "stop_relay_gateway", mock_stop) + return process, mock_start, mock_stop + + +async def test_runtime_waits_for_delayed_relay_artifact( + relay_payload, monkeypatch, tmp_path +): + atif_dir = tmp_path / "atif" + atif_dir.mkdir() + atif_path = atif_dir / "trajectory-session.atif.json" + relay = relay_settings(tmp_path, atif_plugin_config(atif_dir)) + process, mock_start, mock_stop = install_mock_relay(monkeypatch, relay) + write_task = None async def responses(client) -> AsyncIterator[ResultMessage]: + nonlocal write_task assert client.options.env["ANTHROPIC_BASE_URL"] == relay.gateway.url assert Path(client.options.plugins[-1]["path"]) == relay.plugin_path + + async def write_atif(): + await asyncio.sleep(0.05) + atif_path.write_text( + json.dumps({"schema_version": "ATIF-v1.7", "steps": []}), + encoding="utf-8", + ) + + write_task = asyncio.create_task(write_atif()) yield ResultMessage( subtype="success", duration_ms=10, @@ -815,8 +846,12 @@ async def responses(client) -> AsyncIterator[ResultMessage]: key: value for key, value in relay_payload.items() if key != "request" } await runtime.start(start_payload) - output = await runtime.invoke(lifecycle_invocation(relay_payload)) - await runtime.stop() + try: + output = await runtime.invoke(lifecycle_invocation(relay_payload)) + assert write_task is not None + await write_task + finally: + await runtime.stop() assert output["relay_runtime"] == { "enabled": True, @@ -835,6 +870,62 @@ async def responses(client) -> AsyncIterator[ResultMessage]: assert not relay.plugin_path.exists() +async def test_relay_atif_timeout_fails_successful_turn_explicitly( + relay_payload, monkeypatch, tmp_path +): + atif_dir = tmp_path / "atif" + atif_dir.mkdir() + stale_atif = atif_dir / "trajectory-existing.atif.json" + stale_atif.write_text("{}", encoding="utf-8") + relay = relay_settings(tmp_path, atif_plugin_config(atif_dir)) + install_mock_relay(monkeypatch, relay) + wait_for_atif = AsyncMock(return_value=None) + monkeypatch.setattr( + adapter.relay_artifacts, "wait_for_finalized_atif", wait_for_atif + ) + + async def responses(_client) -> AsyncIterator[ResultMessage]: + yield ResultMessage( + subtype="success", + duration_ms=10, + duration_api_ms=8, + is_error=False, + num_turns=1, + session_id="claude-session", + total_cost_usd=0.01, + usage={"input_tokens": 1, "output_tokens": 1}, + result="done", + ) + + install_fake_client(monkeypatch, responses) + runtime = adapter.ClaudeRuntime() + start_payload = { + key: value for key, value in relay_payload.items() if key != "request" + } + await runtime.start(start_payload) + try: + output = await runtime.invoke(lifecycle_invocation(relay_payload)) + unavailable = await runtime.invoke(lifecycle_invocation(relay_payload)) + finally: + await runtime.stop() + + assert output["failed"] is True + assert output["error"] == { + "code": "claude_relay_atif_timeout", + "message": "NeMo Relay did not finalize an ATIF artifact before the deadline", + "retryable": False, + "metadata": { + "timeout_seconds": adapter.relay_artifacts.ATIF_FINALIZATION_TIMEOUT_SECONDS + }, + } + wait_for_atif.assert_awaited_once() + assert output["relay_runtime"]["enabled"] is True + assert output["relay_artifacts"] == [] + assert unavailable["error"]["code"] == "claude_runtime_unavailable" + assert "relay_runtime" not in unavailable + assert "relay_artifacts" not in unavailable + + async def test_runtime_stop_reports_relay_gateway_failure( relay_payload, monkeypatch, tmp_path ): diff --git a/tests/adapters/test_codex_adapter.py b/tests/adapters/test_codex_adapter.py index b17a191c8..f131b98a9 100644 --- a/tests/adapters/test_codex_adapter.py +++ b/tests/adapters/test_codex_adapter.py @@ -136,6 +136,47 @@ def mock_thread(thread_id, result=None): return mock_sdk_thread +def atif_plugin_config(output_directory: Path) -> dict[str, Any]: + return { + "version": 1, + "components": [ + { + "kind": "observability", + "config": { + "atif": { + "enabled": True, + "output_directory": str(output_directory), + "filename_template": "trajectory-{session_id}.atif.json", + } + }, + } + ], + } + + +def relay_settings(tmp_path: Path, plugin_config: dict[str, Any]): + return adapter.CodexRelaySettings( + gateway=adapter.relay_gateway.RelayGatewayLaunch( + executable=tmp_path / "nemo-relay", + config_path=tmp_path / "relay" / "config.toml", + bind="127.0.0.1:43210", + url="http://127.0.0.1:43210", + log_path=tmp_path / "relay" / "gateway.log", + ), + plugin_config=plugin_config, + ) + + +def install_mock_relay(monkeypatch, relay: adapter.CodexRelaySettings): + monkeypatch.setattr(adapter, "prepare_codex_relay", MagicMock(return_value=relay)) + monkeypatch.setattr( + adapter.relay_gateway, + "start_relay_gateway", + MagicMock(return_value=MagicMock()), + ) + monkeypatch.setattr(adapter.relay_gateway, "stop_relay_gateway", MagicMock()) + + @pytest.fixture(name="mock_codex") def mock_codex_fixture(monkeypatch): mock_codex = MagicMock(spec=AsyncCodex) @@ -547,6 +588,97 @@ async def test_persistent_runtime_owns_one_relay_gateway( stop_gateway.assert_called_once_with(process) +async def test_relay_waits_for_delayed_atif_before_collecting_artifacts( + codex_payload, mock_codex, monkeypatch, tmp_path +): + codex_payload["telemetry_plan"] = { + "providers": ["relay"], + "relay_enabled": True, + } + atif_dir = tmp_path / "relay" / "atif" + atif_dir.mkdir(parents=True) + atif_file = atif_dir / "trajectory-session.atif.json" + relay = relay_settings(tmp_path, atif_plugin_config(atif_dir)) + install_mock_relay(monkeypatch, relay) + mock_sdk_thread = mock_thread("thread-123") + write_task = None + + async def finish_turn(): + nonlocal write_task + + async def write_atif(): + await asyncio.sleep(0.05) + atif_file.write_text( + json.dumps({"schema_version": "ATIF-v1.7", "steps": []}), + encoding="utf-8", + ) + + write_task = asyncio.create_task(write_atif()) + return successful_result() + + mock_sdk_thread.handle.run.side_effect = finish_turn + mock_codex.next_thread = mock_sdk_thread + runtime = adapter.CodexRuntime() + + await runtime.start(lifecycle_start_payload(codex_payload)) + try: + output = await runtime.invoke(lifecycle_invocation(codex_payload)) + assert write_task is not None + await write_task + finally: + await runtime.stop() + + assert output["completed"] is True + assert output["relay_artifacts"] == [{"kind": "atif", "path": str(atif_file)}] + + +async def test_relay_atif_timeout_fails_successful_turn_explicitly( + codex_payload, mock_codex, monkeypatch, tmp_path +): + codex_payload["telemetry_plan"] = { + "providers": ["relay"], + "relay_enabled": True, + } + atif_dir = tmp_path / "relay" / "atif" + atif_dir.mkdir(parents=True) + stale_atif = atif_dir / "trajectory-existing.atif.json" + stale_atif.write_text('{"schema_version":"ATIF-v1.7","steps":[]}', encoding="utf-8") + late_atif = atif_dir / "trajectory-late.atif.json" + relay = relay_settings(tmp_path, atif_plugin_config(atif_dir)) + install_mock_relay(monkeypatch, relay) + wait_for_atif = AsyncMock(return_value=None) + monkeypatch.setattr( + adapter.relay_artifacts, "wait_for_finalized_atif", wait_for_atif + ) + runtime = adapter.CodexRuntime() + + await runtime.start(lifecycle_start_payload(codex_payload)) + try: + output = await runtime.invoke(lifecycle_invocation(codex_payload)) + late_atif.write_text( + '{"schema_version":"ATIF-v1.7","steps":[]}', encoding="utf-8" + ) + unavailable = await runtime.invoke(lifecycle_invocation(codex_payload)) + finally: + await runtime.stop() + + assert output["failed"] is True + assert output["error"] == { + "code": "codex_relay_atif_timeout", + "message": "NeMo Relay did not finalize an ATIF artifact before the deadline", + "retryable": False, + "metadata": { + "timeout_seconds": adapter.relay_artifacts.ATIF_FINALIZATION_TIMEOUT_SECONDS + }, + } + wait_for_atif.assert_awaited_once() + assert output["relay_runtime"]["enabled"] is True + assert output["relay_artifacts"] == [] + assert unavailable["error"]["code"] == "codex_runtime_unavailable" + assert "relay_runtime" not in unavailable + assert "relay_artifacts" not in unavailable + + def test_failed_sdk_turn_is_normalized_and_transport_is_closed( codex_payload, mock_codex ): diff --git a/tests/e2e/test_claude.py b/tests/e2e/test_claude.py index ddd681ada..73a0874e8 100644 --- a/tests/e2e/test_claude.py +++ b/tests/e2e/test_claude.py @@ -73,6 +73,7 @@ def fabric_config( *, cli_path=None, relay=False, + atif=True, nemo_relay_command=None, ): tmp_path.mkdir(parents=True, exist_ok=True) @@ -133,7 +134,7 @@ def fabric_config( enabled=True, sinks=[RelayAtofFileSinkConfig()], ), - atif=RelayAtifConfig(enabled=True), + atif=RelayAtifConfig(enabled=atif), ) ) return config @@ -182,10 +183,12 @@ async def test_fabric_claude_relay_supervises_gateway_and_injects_plugin(tmp_pat mock_relay = tmp_path / "nemo-relay" relay_args_path = tmp_path / "relay-args.json" write_mock_relay_gateway(mock_relay, relay_args_path) + # This fake implements gateway lifecycle only, not subscriber artifact writes. config = fabric_config( tmp_path, cli_path=MOCK_CLAUDE_CLI, relay=True, + atif=False, nemo_relay_command=mock_relay, )