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
132 changes: 132 additions & 0 deletions src/aethermesh_core/runtime_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
)
Expand Down Expand Up @@ -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]:
Expand Down Expand Up @@ -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]:
Expand Down
118 changes: 118 additions & 0 deletions tests/test_capability_advertisement.py
Original file line number Diff line number Diff line change
@@ -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()
Loading