diff --git a/src/aethermesh_core/runtime_service.py b/src/aethermesh_core/runtime_service.py index 54b5f3f..ce9a11f 100644 --- a/src/aethermesh_core/runtime_service.py +++ b/src/aethermesh_core/runtime_service.py @@ -46,6 +46,7 @@ from aethermesh_core.runner import LocalRunner, run_local_job from aethermesh_core.validation import validate_job_result from aethermesh_core.validation_receipt_schema import ( + VALIDATION_RECEIPT_SCHEMA_VERSION, capture_validator_software_metadata, validate_validator_software_metadata, ) @@ -534,6 +535,67 @@ def inspect_capability_records(self) -> dict[str, Any]: "note": "Local capability records only; validation is not network consensus.", } + def emit_capability_advertisement(self) -> dict[str, Any]: + """Write an honest, deterministic local capability advertisement artifact.""" + + config = self.load_config() + creator_node_id = _config_node_id(config) + if creator_node_id is None: + raise RuntimeServiceError( + "capability advertisement requires initialized local node identity" + ) + enabled_work_types = _config_enabled_work_types(config) + current_workers = len(self.list_jobs()["current"]) + capability_manifest_dir = self.paths.data_dir / "capability-manifests" + advertised_capabilities: list[dict[str, Any]] = [] + for identifier, description, work_type in LOCAL_CAPABILITY_DEFINITIONS: + if work_type not in enabled_work_types: + continue + availability = self._work_capability_availability( + work_type=work_type, + enabled=True, + creator_node_id=creator_node_id, + current_workers=current_workers, + ) + if availability["status"] != "available": + continue + manifest_ref = f"data/capability-manifests/{identifier}.json" + manifest = _capability_advertisement_manifest( + identifier=identifier, + description=description, + work_type=work_type, + creator_node_id=creator_node_id, + ) + atomic_write_json(capability_manifest_dir / f"{identifier}.json", manifest) + advertised_capabilities.append( + _advertised_capability( + identifier=identifier, + description=description, + work_type=work_type, + creator_node_id=creator_node_id, + manifest_ref=manifest_ref, + ) + ) + advertisement = { + "schema_version": 1, + "scope": "local-only-no-p2p", + "prototype_status": "prototype-local-validation-required", + "creator_node_id": creator_node_id, + "node_manifest_ref": "config.json", + "validation_receipt_format": ( + f"local-validation-receipt-schema-v{VALIDATION_RECEIPT_SCHEMA_VERSION}" + ), + "capabilities": advertised_capabilities, + "note": ( + "Local prototype advertisement only; it is not peer discovery, " + "network consensus, or reward eligibility." + ), + } + atomic_write_json( + self.paths.data_dir / "capability-advertisement.json", advertisement + ) + return advertisement + def _capability_record_summary( self, path: Path, local_node_id: str | None ) -> dict[str, Any]: @@ -3254,6 +3316,76 @@ def _capability_provenance( } +def _capability_advertisement_manifest( + *, + identifier: str, + description: str, + work_type: str, + creator_node_id: str, +) -> dict[str, Any]: + """Build the local manifest that backs one emitted capability claim.""" + + return { + "schema_version": 1, + "manifest_id": _capability_manifest_id(identifier), + "creator_node_id": creator_node_id, + "capability_id": identifier, + "task_type": work_type, + "description": description, + "required_inputs": _capability_check_payload(work_type), + "expected_outputs": f"local-job-result-schema-v{JOB_RESULT_SCHEMA_VERSION}", + "validation": { + "method": "local-deterministic-job-result-validation", + "receipt_format": ( + f"local-validation-receipt-schema-v{VALIDATION_RECEIPT_SCHEMA_VERSION}" + ), + }, + "lineage": { + "node_manifest_ref": "config.json", + "result_manifest_template": "data/job-results/{job_id}.json", + }, + "contribution_attribution": {"creator_node_id": creator_node_id}, + } + + +def _advertised_capability( + *, + identifier: str, + description: str, + work_type: str, + creator_node_id: str, + manifest_ref: str, +) -> dict[str, Any]: + """Return one validation-gated local advertisement entry.""" + + return { + "capability_id": identifier, + "task_types": [work_type], + "description": description, + "scope": "local-only-no-p2p", + "prototype_status": "prototype-local-validation-required", + "runtime_limits": {"max_concurrent_jobs": LOCAL_CAPABILITY_WORKER_CAPACITY}, + "required_inputs": _capability_check_payload(work_type), + "expected_outputs": f"local-job-result-schema-v{JOB_RESULT_SCHEMA_VERSION}", + "validation_requirements": { + "method": "local-deterministic-job-result-validation", + "receipt_format": ( + f"local-validation-receipt-schema-v{VALIDATION_RECEIPT_SCHEMA_VERSION}" + ), + "required": True, + }, + "lineage": { + "node_manifest_ref": "config.json", + "capability_manifest_ref": manifest_ref, + "validation_receipt_format": ( + f"local-validation-receipt-schema-v{VALIDATION_RECEIPT_SCHEMA_VERSION}" + ), + "result_manifest_template": "data/job-results/{job_id}.json", + }, + "contribution_attribution": {"creator_node_id": creator_node_id}, + } + + def _availability( status: str, reason: str | None, worker_capacity: dict[str, int] ) -> dict[str, Any]: diff --git a/tests/test_capability_advertisement.py b/tests/test_capability_advertisement.py new file mode 100644 index 0000000..277acdd --- /dev/null +++ b/tests/test_capability_advertisement.py @@ -0,0 +1,118 @@ +import json +import tempfile +import unittest +from pathlib import Path +from unittest.mock import patch + +from aethermesh_core.job_result_schema import JOB_RESULT_SCHEMA_VERSION +from aethermesh_core.runtime_service import NodeRuntimeService, RuntimeServiceError + + +class CapabilityAdvertisementTests(unittest.TestCase): + def test_emits_deterministic_manifest_linked_local_advertisement(self) -> None: + with tempfile.TemporaryDirectory() as temp_dir: + service = NodeRuntimeService.from_home(temp_dir) + initialized = service.initialize_local_node_data() + config = service.load_config() + config["capabilities"]["enabled_work_types"] = ["echo"] + Path(temp_dir, "config.json").write_text( + json.dumps(config), encoding="utf-8" + ) + + advertisement = service.emit_capability_advertisement() + artifact_path = Path(temp_dir, "data", "capability-advertisement.json") + first_bytes = artifact_path.read_bytes() + repeated = service.emit_capability_advertisement() + + self.assertEqual(advertisement, repeated) + self.assertEqual(first_bytes, artifact_path.read_bytes()) + self.assertEqual(advertisement["creator_node_id"], initialized["node_id"]) + self.assertEqual(advertisement["scope"], "local-only-no-p2p") + self.assertEqual( + advertisement["prototype_status"], "prototype-local-validation-required" + ) + self.assertEqual(advertisement["node_manifest_ref"], "config.json") + self.assertEqual(len(advertisement["capabilities"]), 1) + capability = advertisement["capabilities"][0] + self.assertEqual(capability["capability_id"], "work.echo") + self.assertEqual(capability["task_types"], ["echo"]) + self.assertEqual( + capability["contribution_attribution"]["creator_node_id"], + initialized["node_id"], + ) + self.assertTrue(capability["validation_requirements"]["required"]) + self.assertEqual( + capability["expected_outputs"], + f"local-job-result-schema-v{JOB_RESULT_SCHEMA_VERSION}", + ) + self.assertEqual( + capability["lineage"]["result_manifest_template"], + "data/job-results/{job_id}.json", + ) + manifest_path = Path( + temp_dir, capability["lineage"]["capability_manifest_ref"] + ) + manifest = json.loads(manifest_path.read_text(encoding="utf-8")) + self.assertEqual(manifest["creator_node_id"], initialized["node_id"]) + self.assertEqual(manifest["task_type"], "echo") + self.assertEqual( + manifest["expected_outputs"], + f"local-job-result-schema-v{JOB_RESULT_SCHEMA_VERSION}", + ) + self.assertEqual( + manifest["lineage"]["result_manifest_template"], + "data/job-results/{job_id}.json", + ) + self.assertEqual( + manifest["validation"]["receipt_format"], + capability["validation_requirements"]["receipt_format"], + ) + + def test_omits_disabled_and_unsupported_configured_capabilities(self) -> None: + with tempfile.TemporaryDirectory() as temp_dir: + service = NodeRuntimeService.from_home(temp_dir) + service.initialize_local_node_data() + config = service.load_config() + config["capabilities"]["enabled_work_types"] = ["not-an-executor"] + Path(temp_dir, "config.json").write_text( + json.dumps(config), encoding="utf-8" + ) + + advertisement = service.emit_capability_advertisement() + + self.assertEqual(advertisement["capabilities"], []) + self.assertFalse(Path(temp_dir, "data", "capability-manifests").exists()) + + def test_omits_a_configured_capability_when_local_validation_is_degraded( + self, + ) -> None: + with tempfile.TemporaryDirectory() as temp_dir: + service = NodeRuntimeService.from_home(temp_dir) + service.initialize_local_node_data() + config = service.load_config() + config["capabilities"]["enabled_work_types"] = ["echo"] + Path(temp_dir, "config.json").write_text( + json.dumps(config), encoding="utf-8" + ) + + with patch.object( + service, + "_work_capability_availability", + return_value={"status": "degraded"}, + ): + advertisement = service.emit_capability_advertisement() + + self.assertEqual(advertisement["capabilities"], []) + + def test_requires_a_local_node_identity(self) -> None: + with tempfile.TemporaryDirectory() as temp_dir: + service = NodeRuntimeService.from_home(temp_dir) + + with self.assertRaisesRegex( + RuntimeServiceError, "initialized local node identity" + ): + service.emit_capability_advertisement() + + +if __name__ == "__main__": + unittest.main()