diff --git a/frontend/server/intelligent_development.py b/frontend/server/intelligent_development.py index a42886ece..bf24dc8bd 100644 --- a/frontend/server/intelligent_development.py +++ b/frontend/server/intelligent_development.py @@ -223,7 +223,7 @@ def main(): fail("Delivery request is invalid") request_path = Path(sys.argv[1]) request = json.loads(read_regular(request_path).decode("utf-8")) - if set(request) != {"projectRoot", "report", "secretPath", "agentName", "entryPoint", "manifestSha256"}: + if set(request) != {"projectRoot", "report", "secretPath", "agentName", "entryPoint", "fallbackEntryPoint", "manifestSha256"}: fail("Delivery request fields are invalid") project = project_path(request["projectRoot"]) secret_path = Path(request["secretPath"]) @@ -258,9 +258,15 @@ def main(): fail("Delivery agentkit.yaml changed during packaging") agent_name = request["agentName"] entry_point = request["entryPoint"] + fallback_entry_point = request["fallbackEntryPoint"] if not isinstance(agent_name, str) or not agent_name.strip(): fail("Delivery agent name is invalid") - if not isinstance(entry_point, str) or entry_point not in {name for name, _ in files}: + if not isinstance(entry_point, str) or not isinstance(fallback_entry_point, str): + fail("Delivery entry point is invalid") + source_names = {name for name, _ in files} + if entry_point not in source_names and fallback_entry_point in source_names: + entry_point = fallback_entry_point + if entry_point not in source_names: fail("Delivery entry point is invalid") report["agentName"] = agent_name.strip() report["entryPoint"] = entry_point diff --git a/frontend/server/intelligent_development_routes.py b/frontend/server/intelligent_development_routes.py index 4635b3fb6..cd66917b1 100644 --- a/frontend/server/intelligent_development_routes.py +++ b/frontend/server/intelligent_development_routes.py @@ -1255,6 +1255,7 @@ async def _interrupt(session_id: str, request: Request) -> dict[str, bool]: async def _message(session_id: str, request: Request) -> StreamingResponse: owner = owner_resolver(request) project_context = "" + trusted_manifest_metadata: tuple[str, str] | None = None try: data = await _request_object(request, 128 * 1024) if set(data) != {"message"}: @@ -1280,6 +1281,7 @@ async def _message(session_id: str, request: Request) -> StreamingResponse: ) base = await project_service.base_metadata(owner, session_id) if base is not None: + trusted_manifest_metadata = (base.agent_name, base.entry_point) project_context = json.dumps( { "intentSummary": base.intent_summary, @@ -1468,6 +1470,7 @@ async def cleanup_task_files() -> None: completion=completion, exact_secrets=lease.exact_secrets, acceptance_criteria=decision.acceptance_criteria, + trusted_manifest_metadata=trusted_manifest_metadata, ) stored_version = None persistence_error: dict[str, object] | None = None @@ -1596,7 +1599,13 @@ async def cleanup_task_files() -> None: payload = _stream_error_payload(failure) yield f"event: error\ndata: {json.dumps(payload, ensure_ascii=False)}\n\n" yield 'event: done\ndata: {"reason":"failed"}\n\n' - except Exception: # noqa: BLE001 + except Exception as error: # noqa: BLE001 + logger.error( + "Unexpected intelligent development turn failure stage=%s error_type=%s session_id=%s", + failure_stage, + type(error).__name__, + session_id, + ) try: await cleanup_task_files() except SandboxError as cleanup_error: diff --git a/frontend/server/intelligent_development_task.py b/frontend/server/intelligent_development_task.py index 8d91a6db8..b6ba4504e 100644 --- a/frontend/server/intelligent_development_task.py +++ b/frontend/server/intelligent_development_task.py @@ -726,18 +726,9 @@ async def read_completion_contract( return parse_completion_contract(content) -def _delivery_manifest_metadata(content: bytes) -> tuple[str, str]: - if len(content) > _MAX_MANIFEST_BYTES: - raise ValueError("Delivery agentkit.yaml is too large") - try: - manifest = yaml.safe_load(content) - except (UnicodeDecodeError, yaml.YAMLError) as error: - raise ValueError("Delivery agentkit.yaml is invalid") from error - common = manifest.get("common") if isinstance(manifest, dict) else None - if not isinstance(common, dict): - raise ValueError("Delivery agentkit.yaml common is invalid") - agent_name = common.get("agent_name") or common.get("name") - entry_point = common.get("entry_point") +def _validate_delivery_metadata( + agent_name: object, entry_point: object +) -> tuple[str, str]: if ( not isinstance(agent_name, str) or not agent_name.strip() @@ -759,6 +750,48 @@ def _delivery_manifest_metadata(content: bytes) -> tuple[str, str]: return agent_name.strip(), entry_point +def _delivery_manifest_metadata( + content: bytes, + *, + trusted_fallback: tuple[str, str] | None = None, +) -> tuple[str, str]: + if len(content) > _MAX_MANIFEST_BYTES: + raise ValueError("Delivery agentkit.yaml is too large") + if not content: + if trusted_fallback is None: + raise ValueError("Delivery agentkit.yaml is missing") + return _validate_delivery_metadata(*trusted_fallback) + try: + manifest = yaml.safe_load(content) + except (UnicodeDecodeError, yaml.YAMLError) as error: + raise ValueError("Delivery agentkit.yaml is invalid") from error + if not isinstance(manifest, dict): + raise ValueError("Delivery agentkit.yaml is invalid") + common = manifest.get("common") + if common is None and trusted_fallback is not None: + return _validate_delivery_metadata(*trusted_fallback) + if not isinstance(common, dict): + raise ValueError("Delivery agentkit.yaml common is invalid") + if common.get("agent_name"): + agent_name = common["agent_name"] + elif "name" in common: + agent_name = common["name"] + elif "agent_name" in common: + agent_name = common["agent_name"] + elif trusted_fallback is not None: + agent_name = trusted_fallback[0] + else: + agent_name = None + entry_point = ( + common["entry_point"] + if "entry_point" in common + else trusted_fallback[1] + if trusted_fallback is not None + else None + ) + return _validate_delivery_metadata(agent_name, entry_point) + + class DeliveryPublisher: """Package an immutable source snapshot without re-running validation.""" @@ -774,7 +807,13 @@ async def publish( completion: CompletionContract | None, exact_secrets: tuple[str, ...], acceptance_criteria: tuple[str, ...] = (), + trusted_manifest_metadata: tuple[str, str] | None = None, ) -> DeliveryReference: + trusted_metadata = ( + _validate_delivery_metadata(*trusted_manifest_metadata) + if trusted_manifest_metadata is not None + else None + ) token = uuid4().hex worker_path = f"{task_root}/delivery-{token}.py" request_path = f"{task_root}/delivery-{token}.json" @@ -813,16 +852,21 @@ async def publish( ), "steps": steps, } - manifest_bytes = await self._transport.download( - f"{project_root}/agentkit.yaml", max_bytes=_MAX_MANIFEST_BYTES + manifest_bytes = await self._manifest_bytes( + project_root, + allow_missing=trusted_metadata is not None, + ) + agent_name, entry_point = _delivery_manifest_metadata( + manifest_bytes, + trusted_fallback=trusted_metadata, ) - agent_name, entry_point = _delivery_manifest_metadata(manifest_bytes) request = { "projectRoot": project_root, "report": report, "secretPath": secret_path, "agentName": agent_name, "entryPoint": entry_point, + "fallbackEntryPoint": trusted_metadata[1] if trusted_metadata else "", "manifestSha256": hashlib.sha256(manifest_bytes).hexdigest(), } await self._transport.upload( @@ -857,6 +901,40 @@ async def publish( finally: await self._unlink_many(secret_path, request_path, worker_path) + async def _manifest_bytes(self, project_root: str, *, allow_missing: bool) -> bytes: + manifest_path = f"{project_root}/agentkit.yaml" + if not allow_missing: + return await self._transport.download( + manifest_path, max_bytes=_MAX_MANIFEST_BYTES + ) + source = ( + "import json,os,stat\n" + f"path={manifest_path!r}\n" + "try: metadata=os.lstat(path)\n" + "except FileNotFoundError: value={'state':'missing'}\n" + "else:\n" + " value={'state':'regular','size':metadata.st_size} if stat.S_ISREG(metadata.st_mode) else {'state':'unsafe'}\n" + "print(json.dumps(value,separators=(',',':')))\n" + ) + status = await self._transport.exec_json( + f"python3 -c {shlex.quote(source)}", timeout=12 + ) + if status == {"state": "missing"}: + return b"" + size = status.get("size") + if ( + set(status) != {"state", "size"} + or status.get("state") != "regular" + or isinstance(size, bool) + or not isinstance(size, int) + or size < 0 + or size > _MAX_MANIFEST_BYTES + ): + raise ValueError("Delivery agentkit.yaml is unsafe") + return await self._transport.download( + manifest_path, max_bytes=_MAX_MANIFEST_BYTES + ) + async def _unlink_many(self, *paths: str) -> None: source = ( "import os\n" diff --git a/frontend/service/studio_scheduler/deploy.py b/frontend/service/studio_scheduler/deploy.py index 78ba41332..ad5c19c93 100644 --- a/frontend/service/studio_scheduler/deploy.py +++ b/frontend/service/studio_scheduler/deploy.py @@ -29,6 +29,11 @@ from pathlib import Path from typing import Any +from frontend.service.studio_release_server.offline_runtime import ( + STUDIO_RUNTIME_LOCK, + STUDIO_RUNTIME_WHEELHOUSE, +) + from .diagnostics import sanitize_diagnostic _SCAN_TIMER_NAME = "veadk-studio-cronjobs-minute" @@ -256,6 +261,17 @@ def _stage_package(package_root: Path, destination: Path) -> None: if not requirements.is_file(): raise ValueError("Studio scheduler package is missing requirements.txt") shutil.copy2(requirements, destination / requirements.name) + + runtime_lock = package_root / STUDIO_RUNTIME_LOCK + if runtime_lock.is_file(): + shutil.copy2(runtime_lock, destination / runtime_lock.name) + + wheelhouse = package_root / STUDIO_RUNTIME_WHEELHOUSE + if wheelhouse.is_dir(): + shutil.copytree(wheelhouse, destination / wheelhouse.name) + + # Preserve compatibility with packages produced before the offline + # wheelhouse layout was introduced. for wheel in package_root.glob("*.whl"): shutil.copy2(wheel, destination / wheel.name) run_script = destination / "run.sh" diff --git a/tests/frontend/server/test_intelligent_development_routes.py b/tests/frontend/server/test_intelligent_development_routes.py index 060feb274..9de1f14ca 100644 --- a/tests/frontend/server/test_intelligent_development_routes.py +++ b/tests/frontend/server/test_intelligent_development_routes.py @@ -1610,9 +1610,8 @@ def test_restored_project_context_selects_incremental_builder_mode( routes, "read_completion_contract", AsyncMock(return_value=_partial()) ) monkeypatch.setattr(routes, "remove_completion_file", AsyncMock()) - monkeypatch.setattr( - routes, "DeliveryPublisher", lambda _transport: _publisher_mock() - ) + publisher = _publisher_mock() + monkeypatch.setattr(routes, "DeliveryPublisher", lambda _transport: publisher) with TestClient(_app(gateway, project_service=project_service)) as client: _connect(client) @@ -1630,6 +1629,10 @@ def test_restored_project_context_selects_incremental_builder_mode( assert '"agentName":"weather_agent"' in builder assert "Do not run `ak init`" in builder assert "use `ak init --template agent_server` by default" not in builder + assert publisher.publish.await_args.kwargs["trusted_manifest_metadata"] == ( + "weather_agent", + "app.py", + ) def test_builder_uses_preinstalled_skill_without_discovery_or_injection( @@ -2387,6 +2390,43 @@ def test_delivery_persistence_failures_keep_distinct_sse_semantics( assert "event: development.succeeded" not in response.text +def test_unexpected_snapshot_failure_logs_stage_and_type_without_error_detail( + monkeypatch: pytest.MonkeyPatch, + caplog: pytest.LogCaptureFixture, +) -> None: + gateway = _FakeGateway() + gateway.sessions["dev-session"] = _cloud() + gateway.codex.turns = [ + [CodexAppServerEvent(kind="text", text="源码已生成")], + ] + lease = _Lease(_Remote(gateway.sessions["dev-session"].endpoint)) + monkeypatch.setattr( + routes, "create_credential_lease", AsyncMock(return_value=lease) + ) + monkeypatch.setattr(routes, "invalidate_current_delivery", AsyncMock()) + monkeypatch.setattr( + routes, "read_completion_contract", AsyncMock(return_value=_partial()) + ) + monkeypatch.setattr(routes, "remove_completion_file", AsyncMock()) + publisher = _publisher_mock() + publisher.publish.side_effect = RuntimeError("private upstream detail") + monkeypatch.setattr(routes, "DeliveryPublisher", lambda _transport: publisher) + + with caplog.at_level("ERROR", logger=routes.__name__): + with TestClient(_app(gateway)) as client: + _connect(client) + response = client.post( + "/web/intelligent-development/sessions/dev-session/messages", + headers={"X-Test-User": "alice"}, + json={"message": "做一个天气 Agent"}, + ) + + assert '"code": "INTELLIGENT_DEVELOPMENT_FAILED"' in response.text + assert "stage=delivery_publish" in caplog.text + assert "error_type=RuntimeError" in caplog.text + assert "private upstream detail" not in caplog.text + + def test_builder_response_cannot_replace_a_missing_completion_file( monkeypatch: pytest.MonkeyPatch, ) -> None: diff --git a/tests/frontend/server/test_intelligent_development_task.py b/tests/frontend/server/test_intelligent_development_task.py index f80e99acf..1eb4147d8 100644 --- a/tests/frontend/server/test_intelligent_development_task.py +++ b/tests/frontend/server/test_intelligent_development_task.py @@ -583,6 +583,41 @@ def test_delivery_manifest_rejects_non_ascii_agent_names(agent_name: str) -> Non task_module._delivery_manifest_metadata(manifest) +def test_delivery_manifest_uses_trusted_metadata_for_legacy_fields() -> None: + manifest = b"common:\n agent_name: optimized_agent\n" + + assert task_module._delivery_manifest_metadata( + manifest, + trusted_fallback=("migrated_agent", "main.py"), + ) == ("optimized_agent", "main.py") + + +def test_delivery_manifest_keeps_complete_manifest_authoritative() -> None: + manifest = b"common:\n agent_name: optimized_agent\n entry_point: optimized.py\n" + + assert task_module._delivery_manifest_metadata( + manifest, + trusted_fallback=("migrated_agent", "main.py"), + ) == ("optimized_agent", "optimized.py") + + +@pytest.mark.parametrize( + "manifest", + [ + b"common:\n agent_name: ../unsafe\n", + b"common:\n agent_name: optimized_agent\n entry_point: ../unsafe.py\n", + ], +) +def test_delivery_manifest_rejects_unsafe_explicit_metadata_with_fallback( + manifest: bytes, +) -> None: + with pytest.raises(ValueError): + task_module._delivery_manifest_metadata( + manifest, + trusted_fallback=("migrated_agent", "main.py"), + ) + + @pytest.mark.asyncio async def test_credentials_are_uploaded_once_outside_workspace_and_cleaned( monkeypatch: pytest.MonkeyPatch, @@ -765,6 +800,7 @@ async def exec_text(self, command: str, *, timeout: int) -> str: assert request["projectRoot"] == "/home/gem/workspace/session" assert request["agentName"] == "weather" assert request["entryPoint"] == "weather.py" + assert request["fallbackEntryPoint"] == "" assert request["manifestSha256"] == hashlib.sha256(manifest).hexdigest() assert set(request) == { "projectRoot", @@ -772,6 +808,7 @@ async def exec_text(self, command: str, *, timeout: int) -> str: "secretPath", "agentName", "entryPoint", + "fallbackEntryPoint", "manifestSha256", } secret_path = request["secretPath"] @@ -806,12 +843,95 @@ async def exec_text(self, command: str, *, timeout: int) -> str: assert source_only.gate_summary == () +@pytest.mark.asyncio +async def test_delivery_publisher_packages_migrated_source_without_root_manifest() -> ( + None +): + artifact_digest = "a" * 64 + report_digest = "b" * 64 + release = ( + f"/home/gem/.intelligent-development/releases/{artifact_digest}-{report_digest}" + ) + + class Remote: + def __init__(self) -> None: + self.downloads: list[str] = [] + self.uploads: dict[str, tuple[bytes, int | None]] = {} + self.exec_json_calls = 0 + + async def download(self, path: str, *, max_bytes: int) -> bytes: + del max_bytes + self.downloads.append(path) + raise AssertionError("missing manifest must not be downloaded") + + async def upload( + self, + path: str, + content: bytes, + *, + media_type: str = "application/octet-stream", + max_bytes: int = 20 * 1024 * 1024, + mode: int | None = None, + ) -> None: + del media_type, max_bytes + self.uploads[path] = (content, mode) + + async def exec_json(self, command: str, *, timeout: int) -> dict[str, object]: + del command, timeout + self.exec_json_calls += 1 + if self.exec_json_calls == 1: + return {"state": "missing"} + return { + "sessionId": "session", + "artifactSha256": artifact_digest, + "artifactSize": 128, + "agentName": "travel_planner", + "entryPoint": "main.py", + "fileCount": 2, + "artifactPath": f"{release}/artifact.zip", + "descriptorPath": f"{release}/descriptor.json", + "validationReportPath": f"{release}/validation/{report_digest}.json", + "validationReportSha256": report_digest, + "releasePath": release, + } + + async def exec_text(self, command: str, *, timeout: int) -> str: + del command, timeout + return "" + + remote = Remote() + delivery = await DeliveryPublisher(remote).publish( # type: ignore[arg-type] + session_id="session", + project_root="/home/gem/workspace/session", + task_root="/home/gem/.intelligent-development/tasks/task", + completion=None, + exact_secrets=(), + trusted_manifest_metadata=("travel_planner", "main.py"), + ) + + request_path = next( + path + for path in remote.uploads + if path.endswith(".json") and "secrets" not in path + ) + request = json.loads(remote.uploads[request_path][0]) + assert remote.downloads == [] + assert request["agentName"] == "travel_planner" + assert request["entryPoint"] == "main.py" + assert request["fallbackEntryPoint"] == "main.py" + assert request["manifestSha256"] == hashlib.sha256(b"").hexdigest() + assert delivery.agent_name == "travel_planner" + assert delivery.entry_point == "main.py" + + def _run_delivery_worker( tmp_path: Path, *, files: dict[str, bytes], secrets: tuple[str, ...] = (), report: dict[str, object] | None = None, + trusted_metadata: tuple[str, str] | None = None, + fallback_entry_point: str = "", ) -> subprocess.CompletedProcess[str]: workspace_root = tmp_path / "workspace" project = workspace_root / "session" @@ -836,17 +956,20 @@ def _run_delivery_worker( secret_path.write_text(json.dumps(list(secrets)), encoding="utf-8") os.chmod(secret_path, 0o600) request = tmp_path / "request.json" - manifest_bytes = files["agentkit.yaml"] - manifest = yaml.safe_load(manifest_bytes) - common = manifest["common"] + manifest_bytes = files.get("agentkit.yaml", b"") + if trusted_metadata is None: + manifest = yaml.safe_load(manifest_bytes) + common = manifest["common"] + trusted_metadata = (common["agent_name"], common["entry_point"]) request.write_text( json.dumps( { "projectRoot": str(project), "report": report or {"sessionId": "session"}, "secretPath": str(secret_path), - "agentName": common["agent_name"], - "entryPoint": common["entry_point"], + "agentName": trusted_metadata[0], + "entryPoint": trusted_metadata[1], + "fallbackEntryPoint": fallback_entry_point, "manifestSha256": hashlib.sha256(manifest_bytes).hexdigest(), } ), @@ -887,6 +1010,46 @@ def test_delivery_worker_packages_final_project_and_excludes_local_state( ] +def test_delivery_worker_packages_migrated_project_without_root_manifest( + tmp_path: Path, +) -> None: + result = _run_delivery_worker( + tmp_path, + files={ + "main.py": b"root_agent = object()\n", + ".agentkit/agentkit.yaml": b"name: legacy\n", + }, + trusted_metadata=("travel_planner", "main.py"), + ) + + assert result.returncode == 0, result.stderr + descriptor = json.loads(result.stdout) + assert descriptor["agentName"] == "travel_planner" + assert descriptor["entryPoint"] == "main.py" + with zipfile.ZipFile(descriptor["artifactPath"]) as archive: + assert archive.namelist() == ["main.py"] + + +def test_delivery_worker_uses_trusted_entry_point_when_manifest_target_is_absent( + tmp_path: Path, +) -> None: + result = _run_delivery_worker( + tmp_path, + files={ + "agentkit.yaml": ( + b"common:\n agent_name: travel_planner\n entry_point: agent.py\n" + ), + "main.py": b"root_agent = object()\n", + }, + trusted_metadata=("travel_planner", "agent.py"), + fallback_entry_point="main.py", + ) + + assert result.returncode == 0, result.stderr + descriptor = json.loads(result.stdout) + assert descriptor["entryPoint"] == "main.py" + + def test_delivery_worker_rejects_supplied_credentials(tmp_path: Path) -> None: result = _run_delivery_worker( tmp_path, diff --git a/tests/frontend/service/studio_scheduler/test_scheduler_deploy.py b/tests/frontend/service/studio_scheduler/test_scheduler_deploy.py index 87625d628..ddc2da3fa 100644 --- a/tests/frontend/service/studio_scheduler/test_scheduler_deploy.py +++ b/tests/frontend/service/studio_scheduler/test_scheduler_deploy.py @@ -21,6 +21,7 @@ import pytest from frontend.service.studio_scheduler.deploy import ( + _stage_package, deploy_scheduler, deploy_scheduler_for_studio_update, scheduler_function_name, @@ -183,6 +184,36 @@ def test_deploy_extends_vefaas_sdk_request_timeout(tmp_path: Path) -> None: assert service.client.list_functions == original_list_functions +def test_stage_package_preserves_offline_runtime_dependencies( + tmp_path: Path, +) -> None: + package_root = tmp_path / "package" + destination = tmp_path / "scheduler" + package_root.mkdir() + destination.mkdir() + (package_root / "requirements.txt").write_text( + "--no-index\n--find-links ./wheelhouse\n-r ./studio-runtime.lock\n", + encoding="utf-8", + ) + (package_root / "studio-runtime.lock").write_text( + "fastapi==1.0 --hash=sha256:test\n", + encoding="utf-8", + ) + wheelhouse = package_root / "wheelhouse" + wheelhouse.mkdir() + (wheelhouse / "fastapi-1.0-py3-none-any.whl").write_bytes(b"dependency") + (package_root / "veadk_python-1.0-py3-none-any.whl").write_bytes(b"legacy") + + _stage_package(package_root, destination) + + assert (destination / "requirements.txt").is_file() + assert (destination / "studio-runtime.lock").is_file() + assert ( + destination / "wheelhouse" / "fastapi-1.0-py3-none-any.whl" + ).read_bytes() == b"dependency" + assert (destination / "veadk_python-1.0-py3-none-any.whl").read_bytes() == b"legacy" + + def test_scheduler_function_name_is_safe_and_bounded() -> None: name = scheduler_function_name("studio_" + "a" * 100) worker_name = scheduler_worker_function_name("studio_" + "a" * 100)