diff --git a/docs/contracts/registry-v2.md b/docs/contracts/registry-v2.md new file mode 100644 index 0000000..f7d08fc --- /dev/null +++ b/docs/contracts/registry-v2.md @@ -0,0 +1,116 @@ +# Canonical registry v2 contract + +Registry v2 is the single catalog identity consumed by future CLI, plugin, and supervisor adapters. This module defines the catalog and identity service only. It does not wire those adapters, alter the legacy registry loader, write a registry, move a database, or mutate a live profile. + +## Canonical schema + +A v2 document has exactly these root fields: + +```json +{ + "schema_version": 2, + "state_root": "state", + "dbs": { + "palmer": {"path": "workflows.sqlite"} + }, + "workflows": { + "palmer-trip-planning": { + "workflow_ref": "palmer_workflows.trip_planning:palmer_trip_planning_workflow", + "db": "palmer", + "defaults_overlay": "local" + } + }, + "runner": {"dbs": ["palmer"], "lease_seconds": 30} +} +``` + +Rules: + +- `schema_version` is the integer `2`. +- `state_root` and every `dbs..path` are normalized relative POSIX paths. Absolute, home-relative, drive-prefixed, backslash, empty-segment, `.`, `..`, NUL, and overlong values are rejected by the FND-REGLOC contract. +- A DB entry has exactly `path`; the legacy string spelling is not valid v2. +- DB, workflow, tag, runner-DB, and public consumer IDs are bounded lowercase canonical IDs. Duplicate JSON keys, tags, or runner DB aliases fail closed. +- A workflow has `workflow_ref` and `db`. `workflow_ref` has one importable `module:symbol` spelling; filesystem refs are rejected. Optional catalog metadata is `title`, `description`, `tags`, `default_input`, `trusted_resume`, `kanban_policy`, `dashboard_policy`, and `defaults_overlay` (`local` only). +- Workflow and runner DB references must name declared aliases. `runner.dbs` is nonempty and `lease_seconds` is an integer from 1 through 3600. +- Canonical output omits default-valued optional workflow fields, sorts object keys, workflow/DB aliases, tags, and runner DB aliases, normalizes strings to Unicode NFC, rejects non-finite JSON numbers, and emits compact UTF-8 JSON. + +The complete fixture is `tests/fixtures/registry_v2_valid.json`. + +## Resolution and identity + +Given registry `/copy/.hermes/workflows.registry.json`: + +1. Resolve `state_root` once from the registry directory, producing `/copy/.hermes/state`. +2. Resolve the selected DB's relative `path` beneath that state root. +3. Recheck registry, state-root, and DB containment through FND-REGLOC. Registry symlinks, intermediate symlinks, DB symlink escapes, noncanonical file paths, unstable files, and root escapes fail closed. +4. Accept a configured DB alias only. Raw public paths are never interpreted as aliases. + +`RegistryCatalogV2.fingerprint` is `sha256:` plus SHA-256 of canonical normalized registry JSON. It is copy-independent: identical relocated catalogs have the same catalog fingerprint. + +`RegistryIdentityV1` is the public consumer identity: + +```json +{ + "schema_version": 1, + "registry_fingerprint": "sha256:...", + "registry_identity": "sha256:...", + "db_alias": "palmer", + "resolved_db_identity": "sha256:..." +} +``` + +The registry and resolved-DB identities are domain-separated hashes of the canonical resolved paths. They distinguish two copy-local databases with the same alias and filename without exposing private paths. Public identity contains no path field, registry filename, database filename, workflow defaults, or secret-bearing value. + +Future adapters register one `RegistryIdentityServiceV1` under FND-OP service ID `registry.identity`, contract version `1`. `require_consumer_parity()` compares the exact identity returned to each named consumer. Catalog, registry-path, alias, or resolved-DB drift raises `registry_drift`; it never chooses a winner or silently falls back. + +The service rereads the bounded canonical registry before each resolution. A post-load source-version or fingerprint change is drift rather than an implicit reload. + +## Registry-v1 compatibility window + +For one release, the parser accepts the existing read-only v1 object spelling: + +- root `dbs` and `workflows`, with optional integer `schema_version: 1`; +- DB entries as a relative string or `{ "path": ... }`; +- workflow entries using the existing canonical `workflow_ref` key. + +A migration-safe v1 catalog must place every DB below one shared first path component, such as `state/palmer.sqlite`. The dry-run migrator lifts that component into v2 `state_root`, converts DB entries to `{ "path": ... }`, normalizes workflow metadata, and adds the explicit runner policy over all aliases. + +`dry_run_migrate_registry_file()` returns canonical target JSON, target fingerprint, and `would_write: false`. There is deliberately no write/migrate/apply function. V1 absolute paths, direct registry-directory DB files, divergent roots, traversal, aliases, or alternate workflow-ref spellings fail closed rather than changing DB identity. + +Rollback order is consumers first, parser second. Drift refusal stays enabled while the v1 read window exists. + +## Errors and diagnostics + +All registry input, path, alias, migration, and drift failures raise `RegistryContractError` with `exit_code = 2`. Its FND-OP-style envelope is exactly: + +```json +{"code":"registry_invalid","message":"...","fields":{},"conflict_id":null} +``` + +Messages are at most 256 UTF-8 bytes; complete envelopes are at most 4096 bytes. Parser details, rejected JSON values, registry paths, DB paths, secrets, and private defaults are never copied into diagnostics. Drift fields contain only bounded canonical consumer IDs. + +Primary codes are: + +- `registry_invalid` +- `registry_path_invalid` +- `registry_alias_required` +- `registry_unknown_alias` +- `registry_drift` +- `registry_invalid_consumer` +- `registry_migration_not_required` + +## Verification + +Focused contract: + +```console +uv run pytest -q tests/test_registry_v2_contract.py tests/contracts/test_registry_consumer_contract.py +``` + +Location compatibility: + +```console +uv run pytest -q tests/test_registry_location_contract.py tests/test_registry.py +``` + +The focused tests cover canonical fixture round-trips, normalized fingerprints, v1 dry-run migration without file mutation, registry-relative resolution, alias-only public identity, redaction, two similar copy-local DBs, consumer parity, exit-2 envelopes, malformed/duplicate IDs, traversal, symlink/root escape refusal, bounds, and secret/private-path non-disclosure. All tests use temporary roots or immutable fixtures; no live registry or database is read or written. diff --git a/src/hermes_workflows/registry_v2.py b/src/hermes_workflows/registry_v2.py new file mode 100644 index 0000000..eca8d35 --- /dev/null +++ b/src/hermes_workflows/registry_v2.py @@ -0,0 +1,821 @@ +from __future__ import annotations + +import hashlib +import json +import math +import os +import re +import stat +import unicodedata +from collections.abc import Mapping +from dataclasses import dataclass, field +from pathlib import Path +from types import MappingProxyType +from typing import Any + +from .registry_location import ( + RegistryLocationV1, + RelativeDbPathV1, + resolve_registry_location, + resolve_relative_db_path, +) + + +REGISTRY_IDENTITY_SERVICE_ID = "registry.identity" +REGISTRY_IDENTITY_CONTRACT_VERSION = 1 + +_ALIAS_PATTERN = re.compile(r"^[a-z][a-z0-9_-]{0,63}$") +_ID_PATTERN = re.compile(r"^[a-z][a-z0-9_.-]{0,63}$") +_WORKFLOW_REF_PATTERN = re.compile( + r"^[A-Za-z_][A-Za-z0-9_]*(?:\.[A-Za-z_][A-Za-z0-9_]*)*:[A-Za-z_][A-Za-z0-9_]*$" +) +_ROOT_FIELDS = frozenset({"schema_version", "state_root", "dbs", "workflows", "runner"}) +_LEGACY_ROOT_FIELDS = frozenset({"schema_version", "dbs", "workflows"}) +_WORKFLOW_REQUIRED_FIELDS = frozenset({"workflow_ref", "db"}) +_WORKFLOW_OPTIONAL_FIELDS = frozenset( + { + "title", + "description", + "tags", + "default_input", + "trusted_resume", + "kanban_policy", + "dashboard_policy", + "defaults_overlay", + } +) +_MAX_REGISTRY_BYTES = 1_048_576 +_MAX_DEFAULT_INPUT_BYTES = 65_536 +_MAX_ERROR_BYTES = 4096 +_MAX_ERROR_MESSAGE_BYTES = 256 +_MAX_CONSUMERS = 32 +_MAX_COLLECTION_ITEMS = 10_000 +_MAX_JSON_DEPTH = 32 + + +class RegistryContractError(Exception): + """Bounded, redacted configuration error suitable for doctor-style exit 2.""" + + exit_code = 2 + + def __init__( + self, + code: str, + message: str, + *, + fields: Mapping[str, object] | None = None, + conflict_id: str | None = None, + ) -> None: + if _ID_PATTERN.fullmatch(code) is None: + raise ValueError("registry error code must be a canonical id") + bounded_message = _bounded_text(message, _MAX_ERROR_MESSAGE_BYTES) + normalized_fields = _bounded_error_fields(fields or {}) + if conflict_id is not None and _ID_PATTERN.fullmatch(conflict_id) is None: + raise ValueError("registry conflict_id must be a canonical id or None") + self.code = code + self.message = bounded_message + self.fields = normalized_fields + self.conflict_id = conflict_id + super().__init__(bounded_message) + + def to_dict(self) -> dict[str, object]: + payload: dict[str, object] = { + "code": self.code, + "message": self.message, + "fields": dict(self.fields), + "conflict_id": self.conflict_id, + } + if len(_canonical_json(payload).encode("utf-8")) > _MAX_ERROR_BYTES: + payload["fields"] = {} + return payload + + +@dataclass(frozen=True) +class RegistryDbV2: + alias: str + path: str + + def to_dict(self) -> dict[str, str]: + return {"path": self.path} + + +@dataclass(frozen=True) +class RegistryWorkflowV2: + alias: str + workflow_ref: str + db: str + title: str | None = None + description: str | None = None + tags: tuple[str, ...] = () + default_input: Mapping[str, Any] = field(default_factory=lambda: MappingProxyType({})) + trusted_resume: bool = False + kanban_policy: str = "comment" + dashboard_policy: str = "receipt" + defaults_overlay: str | None = None + + def to_dict(self) -> dict[str, object]: + payload: dict[str, object] = { + "workflow_ref": self.workflow_ref, + "db": self.db, + } + if self.title is not None: + payload["title"] = self.title + if self.description is not None: + payload["description"] = self.description + if self.tags: + payload["tags"] = list(self.tags) + if self.default_input: + payload["default_input"] = _thaw_json(self.default_input) + if self.trusted_resume: + payload["trusted_resume"] = True + if self.kanban_policy != "comment": + payload["kanban_policy"] = self.kanban_policy + if self.dashboard_policy != "receipt": + payload["dashboard_policy"] = self.dashboard_policy + if self.defaults_overlay is not None: + payload["defaults_overlay"] = self.defaults_overlay + return payload + + +@dataclass(frozen=True) +class RegistryRunnerV2: + dbs: tuple[str, ...] + lease_seconds: int + + def to_dict(self) -> dict[str, object]: + return {"dbs": list(self.dbs), "lease_seconds": self.lease_seconds} + + +@dataclass(frozen=True) +class RegistryCatalogV2: + state_root: str + dbs: Mapping[str, RegistryDbV2] + workflows: Mapping[str, RegistryWorkflowV2] + runner: RegistryRunnerV2 + schema_version: int = 2 + + def to_dict(self) -> dict[str, object]: + return { + "schema_version": self.schema_version, + "state_root": self.state_root, + "dbs": {alias: db.to_dict() for alias, db in sorted(self.dbs.items())}, + "workflows": { + alias: workflow.to_dict() for alias, workflow in sorted(self.workflows.items()) + }, + "runner": self.runner.to_dict(), + } + + @property + def fingerprint(self) -> str: + digest = hashlib.sha256(encode_registry_v2(self).encode("utf-8")).hexdigest() + return f"sha256:{digest}" + + def resolve_db(self, registry_path: str | Path, db_alias: str) -> "ResolvedRegistryDbV2": + alias = _require_public_alias(db_alias) + db = self.dbs.get(alias) + if db is None: + raise RegistryContractError( + "registry_unknown_alias", + "the requested DB alias is not present in the registry", + fields={"db_alias": alias}, + ) + try: + canonical_registry = _canonical_registry_path(registry_path) + location = resolve_registry_location( + canonical_registry.parent, + RegistryLocationV1( + registry_file=canonical_registry.name, + state_root=self.state_root, + ), + ) + if Path(location.registry_path) != canonical_registry: + raise ValueError("registry location identity mismatch") + resolved_db = Path( + resolve_relative_db_path( + location, + RelativeDbPathV1(alias=alias, path=db.path), + ) + ) + except RegistryContractError: + raise + except (OSError, RuntimeError, TypeError, ValueError) as exc: + raise RegistryContractError( + "registry_path_invalid", + "registry state resolves outside its canonical registry-relative root", + ) from exc + return ResolvedRegistryDbV2( + registry_path=canonical_registry, + state_root=Path(location.state_root_path), + db_path=resolved_db, + db_alias=alias, + ) + + +@dataclass(frozen=True) +class LoadedRegistry: + catalog: RegistryCatalogV2 + source_schema_version: int + + +@dataclass(frozen=True) +class ResolvedRegistryDbV2: + """Private internal receipt. Public consumers use RegistryIdentityV1 instead.""" + + registry_path: Path + state_root: Path + db_path: Path + db_alias: str + + +@dataclass(frozen=True) +class RegistryIdentityV1: + schema_version: int + registry_fingerprint: str + registry_identity: str + db_alias: str + resolved_db_identity: str + + def to_dict(self) -> dict[str, object]: + return { + "schema_version": self.schema_version, + "registry_fingerprint": self.registry_fingerprint, + "registry_identity": self.registry_identity, + "db_alias": self.db_alias, + "resolved_db_identity": self.resolved_db_identity, + } + + +@dataclass(frozen=True) +class RegistryMigrationPlanV1: + catalog: RegistryCatalogV2 + source_schema_version: int = 1 + target_schema_version: int = 2 + schema_version: int = 1 + would_write: bool = False + + @property + def canonical_target_json(self) -> str: + return encode_registry_v2(self.catalog) + + def to_dict(self) -> dict[str, object]: + return { + "schema_version": self.schema_version, + "source_schema_version": self.source_schema_version, + "target_schema_version": self.target_schema_version, + "would_write": self.would_write, + "target_fingerprint": self.catalog.fingerprint, + "target_registry": self.catalog.to_dict(), + } + + +class RegistryIdentityServiceV1: + """One catalog-backed identity service for CLI, plugin, and supervisor adapters.""" + + service_id = REGISTRY_IDENTITY_SERVICE_ID + contract_version = REGISTRY_IDENTITY_CONTRACT_VERSION + + def __init__( + self, + registry_path: Path, + catalog: RegistryCatalogV2, + *, + source_schema_version: int, + ) -> None: + self.registry_path = registry_path + self.catalog = catalog + self.source_schema_version = source_schema_version + + @classmethod + def from_file(cls, registry_path: str | Path) -> "RegistryIdentityServiceV1": + canonical_path, payload = _read_registry_file(registry_path) + loaded = decode_registry(payload) + return cls( + canonical_path, + loaded.catalog, + source_schema_version=loaded.source_schema_version, + ) + + def resolve_db(self, db_alias: str) -> ResolvedRegistryDbV2: + self._require_current_catalog() + return self.catalog.resolve_db(self.registry_path, db_alias) + + def identity(self, db_alias: str) -> RegistryIdentityV1: + resolved = self.resolve_db(db_alias) + return RegistryIdentityV1( + schema_version=1, + registry_fingerprint=self.catalog.fingerprint, + registry_identity=_identity_digest("registry", str(resolved.registry_path)), + db_alias=resolved.db_alias, + resolved_db_identity=_identity_digest( + "db", + self.catalog.fingerprint, + resolved.db_alias, + str(resolved.db_path), + ), + ) + + def _require_current_catalog(self) -> None: + _, payload = _read_registry_file(self.registry_path) + current = decode_registry(payload) + if ( + current.source_schema_version != self.source_schema_version + or current.catalog.fingerprint != self.catalog.fingerprint + ): + raise RegistryContractError( + "registry_drift", + "the registry changed after the consumer identity service loaded it", + ) + + +def decode_registry(value: str | bytes | bytearray) -> LoadedRegistry: + try: + payload = _decode_json_object(value) + schema_version = payload.get("schema_version") + if schema_version is None or (type(schema_version) is int and schema_version == 1): + return LoadedRegistry(catalog=_parse_legacy_registry(payload), source_schema_version=1) + if type(schema_version) is int and schema_version == 2: + return LoadedRegistry(catalog=_parse_registry_v2(payload), source_schema_version=2) + raise ValueError("unsupported schema version") + except RegistryContractError: + raise + except (json.JSONDecodeError, UnicodeDecodeError, RecursionError, TypeError, ValueError) as exc: + raise RegistryContractError( + "registry_invalid", + "registry does not satisfy the canonical registry-v2 contract", + ) from exc + + +def load_registry_file(registry_path: str | Path) -> LoadedRegistry: + _, payload = _read_registry_file(registry_path) + return decode_registry(payload) + + +def encode_registry_v2(catalog: RegistryCatalogV2) -> str: + if not isinstance(catalog, RegistryCatalogV2): + raise TypeError("catalog must be RegistryCatalogV2") + return _canonical_json(catalog.to_dict()) + + +def dry_run_migrate_registry_file(registry_path: str | Path) -> RegistryMigrationPlanV1: + loaded = load_registry_file(registry_path) + if loaded.source_schema_version != 1: + raise RegistryContractError( + "registry_migration_not_required", + "dry-run migration accepts a read-only registry-v1 source", + ) + return RegistryMigrationPlanV1(catalog=loaded.catalog) + + +def require_consumer_parity( + identities: Mapping[str, RegistryIdentityV1], +) -> RegistryIdentityV1: + if not isinstance(identities, Mapping) or not identities: + raise RegistryContractError( + "registry_invalid_consumer", + "consumer identities must be a nonempty mapping", + ) + if len(identities) > _MAX_CONSUMERS: + raise RegistryContractError( + "registry_invalid_consumer", + "consumer identity comparison exceeds its fixed bound", + ) + validated: list[tuple[str, RegistryIdentityV1]] = [] + for consumer, identity in identities.items(): + if not isinstance(consumer, str) or _ID_PATTERN.fullmatch(consumer) is None: + raise RegistryContractError( + "registry_invalid_consumer", + "consumer names must be bounded canonical ids", + ) + if not isinstance(identity, RegistryIdentityV1): + raise RegistryContractError( + "registry_invalid_consumer", + "consumer values must be registry identities", + ) + validated.append((consumer, identity)) + validated.sort(key=lambda item: item[0]) + first = validated[0][1] + if any(identity != first for _, identity in validated[1:]): + raise RegistryContractError( + "registry_drift", + "registry consumers do not share one registry identity", + fields={"consumers": [consumer for consumer, _ in validated]}, + ) + return first + + +def _parse_registry_v2(payload: Mapping[str, Any]) -> RegistryCatalogV2: + _require_exact_fields(payload, required=_ROOT_FIELDS, optional=frozenset()) + if type(payload["schema_version"]) is not int or payload["schema_version"] != 2: + raise ValueError("schema_version must equal 2") + state_root = RegistryLocationV1( + registry_file="workflows.registry.json", + state_root=_normalize_path_string(payload["state_root"]), + ).state_root + dbs = _parse_v2_dbs(payload["dbs"]) + workflows = _parse_workflows(payload["workflows"], dbs=dbs) + runner = _parse_runner(payload["runner"], dbs=dbs) + return RegistryCatalogV2( + state_root=state_root, + dbs=MappingProxyType(dbs), + workflows=MappingProxyType(workflows), + runner=runner, + ) + + +def _parse_legacy_registry(payload: Mapping[str, Any]) -> RegistryCatalogV2: + _require_exact_fields( + payload, + required=frozenset({"dbs", "workflows"}), + optional=frozenset({"schema_version"}), + ) + if "schema_version" in payload and ( + type(payload["schema_version"]) is not int or payload["schema_version"] != 1 + ): + raise ValueError("legacy schema_version must equal 1") + raw_dbs = payload["dbs"] + if not isinstance(raw_dbs, Mapping) or not raw_dbs: + raise ValueError("legacy dbs must be a nonempty object") + + legacy_paths: dict[str, str] = {} + for alias, raw_db in raw_dbs.items(): + _validate_alias(alias) + if isinstance(raw_db, str): + path = raw_db + elif isinstance(raw_db, Mapping): + _require_exact_fields(raw_db, required=frozenset({"path"}), optional=frozenset()) + path = raw_db["path"] + else: + raise TypeError("legacy DB entries must be strings or path objects") + validated = RelativeDbPathV1(alias=alias, path=_normalize_path_string(path)).path + legacy_paths[alias] = validated + + state_roots = {path.split("/", 1)[0] for path in legacy_paths.values() if "/" in path} + if len(state_roots) != 1 or any("/" not in path for path in legacy_paths.values()): + raise ValueError("legacy DB paths do not share one migration-safe state root") + state_root = next(iter(state_roots)) + dbs = { + alias: RegistryDbV2(alias=alias, path=path.split("/", 1)[1]) + for alias, path in sorted(legacy_paths.items()) + } + workflows = _parse_workflows(payload["workflows"], dbs=dbs, legacy=True) + runner = RegistryRunnerV2(dbs=tuple(sorted(dbs)), lease_seconds=30) + return RegistryCatalogV2( + state_root=state_root, + dbs=MappingProxyType(dbs), + workflows=MappingProxyType(workflows), + runner=runner, + ) + + +def _parse_v2_dbs(value: object) -> dict[str, RegistryDbV2]: + if not isinstance(value, Mapping) or not value: + raise ValueError("dbs must be a nonempty object") + dbs: dict[str, RegistryDbV2] = {} + for alias, raw_db in value.items(): + _validate_alias(alias) + if not isinstance(raw_db, Mapping): + raise TypeError("registry-v2 DB entries must be objects") + _require_exact_fields(raw_db, required=frozenset({"path"}), optional=frozenset()) + path = RelativeDbPathV1(alias=alias, path=_normalize_path_string(raw_db["path"])).path + dbs[alias] = RegistryDbV2(alias=alias, path=path) + return dict(sorted(dbs.items())) + + +def _parse_workflows( + value: object, + *, + dbs: Mapping[str, RegistryDbV2], + legacy: bool = False, +) -> dict[str, RegistryWorkflowV2]: + if not isinstance(value, Mapping): + raise TypeError("workflows must be an object") + workflows: dict[str, RegistryWorkflowV2] = {} + for alias, raw_workflow in value.items(): + _validate_alias(alias) + if not isinstance(raw_workflow, Mapping): + raise TypeError("workflow entries must be objects") + required = frozenset({"workflow_ref"}) if legacy and len(dbs) == 1 else _WORKFLOW_REQUIRED_FIELDS + _require_exact_fields( + raw_workflow, + required=required, + optional=_WORKFLOW_OPTIONAL_FIELDS | ({"db"} if legacy else frozenset()), + ) + workflow_ref = raw_workflow["workflow_ref"] + if not isinstance(workflow_ref, str) or _WORKFLOW_REF_PATTERN.fullmatch(workflow_ref) is None: + raise ValueError("workflow_ref must use one importable module:symbol spelling") + if "db" in raw_workflow: + db_alias = raw_workflow["db"] + else: + db_alias = next(iter(dbs)) + if not isinstance(db_alias, str) or db_alias not in dbs: + raise ValueError("workflow db must reference a configured alias") + + title = _optional_text(raw_workflow.get("title"), label="title", max_bytes=512) + description = _optional_text(raw_workflow.get("description"), label="description", max_bytes=4096) + tags = _parse_tags(raw_workflow.get("tags", [])) + default_input_raw = raw_workflow.get("default_input", {}) + if not isinstance(default_input_raw, Mapping): + raise TypeError("default_input must be an object") + normalized_default = _normalize_json(default_input_raw) + if not isinstance(normalized_default, dict): + raise TypeError("default_input must normalize to an object") + if len(_canonical_json(normalized_default).encode("utf-8")) > _MAX_DEFAULT_INPUT_BYTES: + raise ValueError("default_input exceeds its fixed bound") + trusted_resume = raw_workflow.get("trusted_resume", False) + if not isinstance(trusted_resume, bool): + raise TypeError("trusted_resume must be a boolean") + kanban_policy = _policy(raw_workflow.get("kanban_policy", "comment"), label="kanban_policy") + dashboard_policy = _policy(raw_workflow.get("dashboard_policy", "receipt"), label="dashboard_policy") + defaults_overlay = raw_workflow.get("defaults_overlay") + if defaults_overlay is not None and defaults_overlay != "local": + raise ValueError("defaults_overlay must equal local when present") + + workflows[alias] = RegistryWorkflowV2( + alias=alias, + workflow_ref=workflow_ref, + db=db_alias, + title=title, + description=description, + tags=tags, + default_input=_freeze_json(normalized_default), + trusted_resume=trusted_resume, + kanban_policy=kanban_policy, + dashboard_policy=dashboard_policy, + defaults_overlay=defaults_overlay, + ) + return dict(sorted(workflows.items())) + + +def _parse_runner(value: object, *, dbs: Mapping[str, RegistryDbV2]) -> RegistryRunnerV2: + if not isinstance(value, Mapping): + raise TypeError("runner must be an object") + _require_exact_fields( + value, + required=frozenset({"dbs", "lease_seconds"}), + optional=frozenset(), + ) + raw_dbs = value["dbs"] + if not isinstance(raw_dbs, list) or not raw_dbs: + raise TypeError("runner dbs must be a nonempty list") + if len(raw_dbs) != len(set(raw_dbs)): + raise ValueError("runner db aliases must be unique") + for alias in raw_dbs: + if not isinstance(alias, str) or alias not in dbs: + raise ValueError("runner dbs must reference configured aliases") + lease_seconds = value["lease_seconds"] + if type(lease_seconds) is not int or not 1 <= lease_seconds <= 3600: + raise ValueError("runner lease_seconds must be an integer from 1 through 3600") + return RegistryRunnerV2(dbs=tuple(sorted(raw_dbs)), lease_seconds=lease_seconds) + + +def _parse_tags(value: object) -> tuple[str, ...]: + if not isinstance(value, list): + raise TypeError("tags must be a list") + if len(value) > 32: + raise ValueError("tags exceed their fixed bound") + tags: list[str] = [] + for tag in value: + if not isinstance(tag, str) or _ALIAS_PATTERN.fullmatch(tag) is None: + raise ValueError("tags must be canonical ids") + tags.append(tag) + if len(tags) != len(set(tags)): + raise ValueError("tags must be unique") + return tuple(sorted(tags)) + + +def _decode_json_object(value: str | bytes | bytearray) -> Mapping[str, Any]: + if isinstance(value, str): + encoded = value.encode("utf-8") + text = value + elif isinstance(value, (bytes, bytearray)): + encoded = bytes(value) + text = encoded.decode("utf-8", errors="strict") + else: + raise TypeError("registry JSON must be text or bytes") + if len(encoded) > _MAX_REGISTRY_BYTES: + raise ValueError("registry JSON exceeds its fixed bound") + payload = json.loads( + text, + object_pairs_hook=_reject_duplicate_keys, + parse_constant=_reject_nonfinite_constant, + ) + if not isinstance(payload, Mapping): + raise TypeError("registry JSON must be an object") + return payload + + +def _read_registry_file(registry_path: str | Path) -> tuple[Path, bytes]: + try: + canonical = _canonical_registry_path(registry_path) + required_flags = ("O_CLOEXEC", "O_NOFOLLOW", "O_NONBLOCK") + if not all(hasattr(os, name) for name in required_flags): + raise OSError("safe registry file primitives unavailable") + before = canonical.lstat() + flags = os.O_RDONLY | os.O_CLOEXEC | os.O_NOFOLLOW | os.O_NONBLOCK + descriptor = os.open(canonical, flags) + try: + descriptor_before = os.fstat(descriptor) + if not stat.S_ISREG(descriptor_before.st_mode): + raise OSError("registry is not a regular file") + chunks: list[bytes] = [] + remaining = _MAX_REGISTRY_BYTES + 1 + while remaining: + chunk = os.read(descriptor, min(remaining, 64 * 1024)) + if not chunk: + break + chunks.append(chunk) + remaining -= len(chunk) + payload = b"".join(chunks) + descriptor_after = os.fstat(descriptor) + finally: + os.close(descriptor) + after = canonical.lstat() + if not ( + _stable_stat_identity(before) + == _stable_stat_identity(descriptor_before) + == _stable_stat_identity(descriptor_after) + == _stable_stat_identity(after) + ): + raise OSError("registry file changed while being read") + if len(payload) > _MAX_REGISTRY_BYTES: + raise ValueError("registry exceeds its fixed bound") + return canonical, payload + except RegistryContractError: + raise + except (OSError, RuntimeError, TypeError, ValueError) as exc: + raise RegistryContractError( + "registry_path_invalid", + "registry_path must identify one stable canonical regular file", + ) from exc + + +def _canonical_registry_path(registry_path: str | Path) -> Path: + if not isinstance(registry_path, (str, Path)): + raise TypeError("registry_path must be a path") + raw = Path(registry_path) + if ".." in raw.parts: + raise ValueError("registry_path must not contain parent traversal") + absolute = raw if raw.is_absolute() else Path.cwd() / raw + if not absolute.exists() or absolute.is_symlink() or not absolute.is_file(): + raise ValueError("registry_path must identify a regular file") + resolved = absolute.resolve(strict=True) + if str(absolute) != str(resolved): + raise ValueError("registry_path must be canonical and symlink-free") + return resolved + + +def _stable_stat_identity(value: os.stat_result) -> tuple[int, int, int, int, int, int]: + return ( + value.st_dev, + value.st_ino, + stat.S_IFMT(value.st_mode), + value.st_size, + value.st_mtime_ns, + value.st_ctime_ns, + ) + + +def _require_exact_fields( + value: Mapping[str, Any], + *, + required: frozenset[str], + optional: frozenset[str], +) -> None: + actual = set(value) + if not required <= actual or actual - required - optional: + raise ValueError("object fields do not match the canonical registry schema") + + +def _validate_alias(value: object) -> str: + if not isinstance(value, str) or _ALIAS_PATTERN.fullmatch(value) is None: + raise ValueError("alias must be a bounded canonical id") + return value + + +def _require_public_alias(value: object) -> str: + if not isinstance(value, str) or _ALIAS_PATTERN.fullmatch(value) is None: + raise RegistryContractError( + "registry_alias_required", + "public registry consumers require a configured DB alias", + ) + return value + + +def _normalize_path_string(value: object) -> str: + if not isinstance(value, str): + raise TypeError("registry paths must be strings") + return unicodedata.normalize("NFC", value) + + +def _optional_text(value: object, *, label: str, max_bytes: int) -> str | None: + if value is None: + return None + if not isinstance(value, str) or not value.strip(): + raise ValueError(f"{label} must be a nonblank string or absent") + normalized = unicodedata.normalize("NFC", value) + if len(normalized.encode("utf-8")) > max_bytes: + raise ValueError(f"{label} exceeds its fixed bound") + return normalized + + +def _policy(value: object, *, label: str) -> str: + if not isinstance(value, str) or _ID_PATTERN.fullmatch(value) is None: + raise ValueError(f"{label} must be a canonical id") + return value + + +def _normalize_json(value: object, *, depth: int = 0) -> object: + if depth > _MAX_JSON_DEPTH: + raise ValueError("JSON value exceeds its depth bound") + if value is None or isinstance(value, bool) or type(value) is int: + return value + if isinstance(value, str): + return unicodedata.normalize("NFC", value) + if isinstance(value, float): + if not math.isfinite(value): + raise ValueError("JSON numbers must be finite") + return value + if isinstance(value, Mapping): + if len(value) > _MAX_COLLECTION_ITEMS: + raise ValueError("JSON object exceeds its item bound") + normalized: dict[str, object] = {} + for key, item in value.items(): + if not isinstance(key, str): + raise TypeError("JSON object keys must be strings") + normalized_key = unicodedata.normalize("NFC", key) + if normalized_key in normalized: + raise ValueError("JSON keys collide after Unicode normalization") + normalized[normalized_key] = _normalize_json(item, depth=depth + 1) + return normalized + if isinstance(value, list): + if len(value) > _MAX_COLLECTION_ITEMS: + raise ValueError("JSON list exceeds its item bound") + return [_normalize_json(item, depth=depth + 1) for item in value] + raise TypeError("value is not JSON-compatible") + + +def _freeze_json(value: object) -> Any: + if isinstance(value, Mapping): + return MappingProxyType({key: _freeze_json(item) for key, item in value.items()}) + if isinstance(value, list): + return tuple(_freeze_json(item) for item in value) + return value + + +def _thaw_json(value: object) -> object: + if isinstance(value, Mapping): + return {key: _thaw_json(item) for key, item in value.items()} + if isinstance(value, tuple): + return [_thaw_json(item) for item in value] + return value + + +def _reject_duplicate_keys(pairs: list[tuple[str, Any]]) -> dict[str, Any]: + value: dict[str, Any] = {} + for key, item in pairs: + if key in value: + raise ValueError("registry JSON contains a duplicate object key") + value[key] = item + return value + + +def _reject_nonfinite_constant(value: str) -> None: + raise ValueError("registry JSON numbers must be finite") + + +def _identity_digest(kind: str, *parts: str) -> str: + digest = hashlib.sha256() + digest.update(f"hermes-workflows:{kind}:v1\0".encode("utf-8")) + for part in parts: + digest.update(part.encode("utf-8")) + digest.update(b"\0") + return f"sha256:{digest.hexdigest()}" + + +def _bounded_text(value: str, max_bytes: int) -> str: + if not isinstance(value, str): + raise TypeError("bounded text must be a string") + encoded = value.encode("utf-8") + if len(encoded) <= max_bytes: + return value + return encoded[: max_bytes - 3].decode("utf-8", errors="ignore") + "..." + + +def _bounded_error_fields(value: Mapping[str, object]) -> Mapping[str, object]: + normalized = _normalize_json(value) + if not isinstance(normalized, dict): + return MappingProxyType({}) + encoded = _canonical_json(normalized).encode("utf-8") + if len(encoded) > _MAX_ERROR_BYTES // 2: + return MappingProxyType({}) + return MappingProxyType(normalized) + + +def _canonical_json(value: object) -> str: + return json.dumps( + value, + sort_keys=True, + separators=(",", ":"), + ensure_ascii=False, + allow_nan=False, + ) diff --git a/tests/contracts/test_registry_consumer_contract.py b/tests/contracts/test_registry_consumer_contract.py new file mode 100644 index 0000000..f2cefe6 --- /dev/null +++ b/tests/contracts/test_registry_consumer_contract.py @@ -0,0 +1,103 @@ +from __future__ import annotations + +import json +import shutil +from pathlib import Path + +import pytest + +from hermes_workflows.operator_services import OperatorServicesV1 +from hermes_workflows.registry_v2 import ( + REGISTRY_IDENTITY_CONTRACT_VERSION, + REGISTRY_IDENTITY_SERVICE_ID, + RegistryContractError, + RegistryIdentityServiceV1, + require_consumer_parity, +) + + +FIXTURES = Path(__file__).parents[1] / "fixtures" + + +def _service(root: Path, fixture: str = "registry_v2_valid.json") -> RegistryIdentityServiceV1: + registry = root / ".hermes" / "workflows.registry.json" + registry.parent.mkdir(parents=True) + shutil.copy2(FIXTURES / fixture, registry) + return RegistryIdentityServiceV1.from_file(registry) + + +def _consume(service: RegistryIdentityServiceV1, db_alias: str) -> dict[str, object]: + return service.identity(db_alias).to_dict() + + +def test_cli_plugin_and_supervisor_contract_consumers_share_one_service_identity(tmp_path: Path) -> None: + service = _service(tmp_path / "workspace") + services = OperatorServicesV1(services={REGISTRY_IDENTITY_SERVICE_ID: service}) + + resolved = services.resolve(REGISTRY_IDENTITY_SERVICE_ID, REGISTRY_IDENTITY_CONTRACT_VERSION) + + assert resolved is service + cli = _consume(service, "primary") + plugin = _consume(service, "primary") + supervisor = _consume(service, "primary") + identity = require_consumer_parity( + { + "cli": service.identity("primary"), + "plugin": service.identity("primary"), + "supervisor": service.identity("primary"), + } + ) + assert cli == plugin == supervisor == identity.to_dict() + assert set(cli) == { + "schema_version", + "registry_fingerprint", + "registry_identity", + "db_alias", + "resolved_db_identity", + } + assert not any("path" in key for key in cli) + + +def test_every_public_consumer_requires_an_alias_not_a_raw_path(tmp_path: Path) -> None: + service = _service(tmp_path / "workspace") + + for consumer in (_consume, _consume, _consume): + with pytest.raises(RegistryContractError) as raised: + consumer(service, "state/primary/workflows.sqlite") + assert raised.value.code == "registry_alias_required" + assert raised.value.exit_code == 2 + + +def test_consumer_drift_is_a_redacted_doctor_style_exit_2_error(tmp_path: Path) -> None: + canonical = _service(tmp_path / "canonical") + drifted = _service(tmp_path / "drifted", "registry_v2_drift.json") + + with pytest.raises(RegistryContractError) as raised: + require_consumer_parity( + { + "cli": canonical.identity("primary"), + "plugin": canonical.identity("primary"), + "supervisor": drifted.identity("primary"), + } + ) + + error = raised.value + payload = error.to_dict() + encoded = json.dumps(payload, sort_keys=True) + assert error.code == "registry_drift" + assert error.exit_code == 2 + assert payload["fields"] == {"consumers": ["cli", "plugin", "supervisor"]} + assert "canonical" not in encoded + assert "drifted" not in encoded + assert "workflows.sqlite" not in encoded + assert len(encoded.encode("utf-8")) <= 4096 + + +def test_consumer_names_are_bounded_canonical_ids(tmp_path: Path) -> None: + identity = _service(tmp_path / "workspace").identity("primary") + + for name in ("", "CLI", "has/slash", "a" * 65): + with pytest.raises(RegistryContractError) as raised: + require_consumer_parity({name: identity}) + assert raised.value.code == "registry_invalid_consumer" + assert raised.value.exit_code == 2 diff --git a/tests/fixtures/registry_v1_legacy.json b/tests/fixtures/registry_v1_legacy.json new file mode 100644 index 0000000..f1b0292 --- /dev/null +++ b/tests/fixtures/registry_v1_legacy.json @@ -0,0 +1,37 @@ +{ + "dbs": { + "audit": { + "path": "state/audit/workflows.sqlite" + }, + "primary": "state/primary/workflows.sqlite" + }, + "workflows": { + "audit-workflow": { + "workflow_ref": "example_workflows.audit:audit_workflow", + "db": "audit", + "title": "Audit workflow", + "tags": [ + "review", + "audit" + ], + "trusted_resume": false + }, + "review-workflow": { + "workflow_ref": "example_workflows.review:review_workflow", + "db": "primary", + "title": "Review workflow", + "description": "A deterministic registry-v2 contract fixture.", + "tags": [ + "workflow", + "review" + ], + "default_input": { + "mode": "review" + }, + "trusted_resume": true, + "kanban_policy": "comment", + "dashboard_policy": "receipt", + "defaults_overlay": "local" + } + } +} diff --git a/tests/fixtures/registry_v2_drift.json b/tests/fixtures/registry_v2_drift.json new file mode 100644 index 0000000..18a3129 --- /dev/null +++ b/tests/fixtures/registry_v2_drift.json @@ -0,0 +1,45 @@ +{ + "schema_version": 2, + "state_root": "state", + "dbs": { + "audit": { + "path": "audit/workflows.sqlite" + }, + "primary": { + "path": "primary-archive/workflows.sqlite" + } + }, + "workflows": { + "audit-workflow": { + "workflow_ref": "example_workflows.audit:audit_workflow", + "db": "audit", + "title": "Audit workflow", + "tags": [ + "audit", + "review" + ] + }, + "review-workflow": { + "workflow_ref": "example_workflows.review:review_workflow", + "db": "primary", + "title": "Review workflow", + "description": "A deterministic registry-v2 contract fixture.", + "tags": [ + "review", + "workflow" + ], + "default_input": { + "mode": "review" + }, + "trusted_resume": true, + "defaults_overlay": "local" + } + }, + "runner": { + "dbs": [ + "audit", + "primary" + ], + "lease_seconds": 30 + } +} diff --git a/tests/fixtures/registry_v2_valid.json b/tests/fixtures/registry_v2_valid.json new file mode 100644 index 0000000..7bf249a --- /dev/null +++ b/tests/fixtures/registry_v2_valid.json @@ -0,0 +1,45 @@ +{ + "schema_version": 2, + "state_root": "state", + "dbs": { + "audit": { + "path": "audit/workflows.sqlite" + }, + "primary": { + "path": "primary/workflows.sqlite" + } + }, + "workflows": { + "audit-workflow": { + "workflow_ref": "example_workflows.audit:audit_workflow", + "db": "audit", + "title": "Audit workflow", + "tags": [ + "audit", + "review" + ] + }, + "review-workflow": { + "workflow_ref": "example_workflows.review:review_workflow", + "db": "primary", + "title": "Review workflow", + "description": "A deterministic registry-v2 contract fixture.", + "tags": [ + "review", + "workflow" + ], + "default_input": { + "mode": "review" + }, + "trusted_resume": true, + "defaults_overlay": "local" + } + }, + "runner": { + "dbs": [ + "audit", + "primary" + ], + "lease_seconds": 30 + } +} diff --git a/tests/test_registry_v2_contract.py b/tests/test_registry_v2_contract.py new file mode 100644 index 0000000..5123868 --- /dev/null +++ b/tests/test_registry_v2_contract.py @@ -0,0 +1,382 @@ +from __future__ import annotations + +import copy +import json +import shutil +from pathlib import Path + +import pytest + +from hermes_workflows.registry_v2 import ( + RegistryContractError, + RegistryIdentityServiceV1, + decode_registry, + dry_run_migrate_registry_file, + encode_registry_v2, + load_registry_file, + require_consumer_parity, +) + + +FIXTURES = Path(__file__).parent / "fixtures" +VALID = FIXTURES / "registry_v2_valid.json" +DRIFT = FIXTURES / "registry_v2_drift.json" +LEGACY = FIXTURES / "registry_v1_legacy.json" + + +def _write_registry(root: Path, fixture: Path = VALID) -> Path: + registry = root / ".hermes" / "workflows.registry.json" + registry.parent.mkdir(parents=True) + shutil.copy2(fixture, registry) + return registry + + +def test_target_schema_fixture_round_trip_and_normalized_fingerprint() -> None: + raw = json.loads(VALID.read_text(encoding="utf-8")) + loaded = decode_registry(VALID.read_bytes()) + + assert loaded.source_schema_version == 2 + assert json.loads(encode_registry_v2(loaded.catalog)) == raw + assert loaded.catalog.fingerprint.startswith("sha256:") + assert len(loaded.catalog.fingerprint) == 71 + + reordered = copy.deepcopy(raw) + reordered["dbs"] = dict(reversed(tuple(reordered["dbs"].items()))) + reordered["workflows"] = dict(reversed(tuple(reordered["workflows"].items()))) + reordered["runner"]["dbs"].reverse() + reordered["workflows"]["review-workflow"]["tags"].reverse() + reordered_loaded = decode_registry(json.dumps(reordered, indent=7)) + + assert reordered_loaded.catalog.fingerprint == loaded.catalog.fingerprint + assert encode_registry_v2(reordered_loaded.catalog) == encode_registry_v2(loaded.catalog) + + +def test_v2_state_root_uses_one_nfc_canonical_spelling() -> None: + composed = json.loads(VALID.read_text(encoding="utf-8")) + decomposed = copy.deepcopy(composed) + composed["state_root"] = "café" + decomposed["state_root"] = "cafe\u0301" + + composed_catalog = decode_registry(json.dumps(composed)).catalog + decomposed_catalog = decode_registry(json.dumps(decomposed)).catalog + + assert composed_catalog.state_root == decomposed_catalog.state_root == "café" + assert encode_registry_v2(composed_catalog) == encode_registry_v2(decomposed_catalog) + assert composed_catalog.fingerprint == decomposed_catalog.fingerprint + + +def test_v2_db_paths_use_one_nfc_canonical_spelling() -> None: + composed = json.loads(VALID.read_text(encoding="utf-8")) + decomposed = copy.deepcopy(composed) + composed["dbs"]["primary"]["path"] = "café/workflows.sqlite" + decomposed["dbs"]["primary"]["path"] = "cafe\u0301/workflows.sqlite" + + composed_catalog = decode_registry(json.dumps(composed)).catalog + decomposed_catalog = decode_registry(json.dumps(decomposed)).catalog + + assert composed_catalog.dbs["primary"].path == decomposed_catalog.dbs["primary"].path == ( + "café/workflows.sqlite" + ) + assert encode_registry_v2(composed_catalog) == encode_registry_v2(decomposed_catalog) + assert composed_catalog.fingerprint == decomposed_catalog.fingerprint + + +def test_v1_is_read_only_compatible_and_migration_is_a_deterministic_dry_run() -> None: + before = LEGACY.read_bytes() + + loaded = load_registry_file(LEGACY) + migration = dry_run_migrate_registry_file(LEGACY) + + assert LEGACY.read_bytes() == before + assert loaded.source_schema_version == 1 + assert json.loads(encode_registry_v2(loaded.catalog)) == json.loads(VALID.read_text(encoding="utf-8")) + assert migration.to_dict() == { + "schema_version": 1, + "source_schema_version": 1, + "target_schema_version": 2, + "would_write": False, + "target_fingerprint": loaded.catalog.fingerprint, + "target_registry": json.loads(encode_registry_v2(loaded.catalog)), + } + assert migration.canonical_target_json == encode_registry_v2(loaded.catalog) + + +def test_v1_dry_run_normalizes_paths_to_one_nfc_canonical_target(tmp_path: Path) -> None: + composed = json.loads(LEGACY.read_text(encoding="utf-8")) + decomposed = copy.deepcopy(composed) + composed["dbs"]["audit"]["path"] = "café/audit/workflows.sqlite" + composed["dbs"]["primary"] = "café/primary/workflows.sqlite" + decomposed["dbs"]["audit"]["path"] = "cafe\u0301/audit/workflows.sqlite" + decomposed["dbs"]["primary"] = "cafe\u0301/primary/workflows.sqlite" + composed_path = tmp_path / "composed.json" + decomposed_path = tmp_path / "decomposed.json" + composed_path.write_text(json.dumps(composed), encoding="utf-8") + decomposed_path.write_text(json.dumps(decomposed), encoding="utf-8") + + composed_plan = dry_run_migrate_registry_file(composed_path) + decomposed_plan = dry_run_migrate_registry_file(decomposed_path) + + assert composed_plan.catalog.state_root == decomposed_plan.catalog.state_root == "café" + assert composed_plan.canonical_target_json == decomposed_plan.canonical_target_json + assert composed_plan.catalog.fingerprint == decomposed_plan.catalog.fingerprint + assert composed_plan.would_write is decomposed_plan.would_write is False + + +def test_v1_marker_is_accepted_but_v2_never_emits_it() -> None: + payload = json.loads(LEGACY.read_text(encoding="utf-8")) + payload["schema_version"] = 1 + + loaded = decode_registry(json.dumps(payload)) + + assert loaded.source_schema_version == 1 + assert json.loads(encode_registry_v2(loaded.catalog))["schema_version"] == 2 + + +def test_relative_resolution_is_registry_directory_then_state_root(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> None: + root = tmp_path / "copy" + elsewhere = tmp_path / "elsewhere" + elsewhere.mkdir() + registry = _write_registry(root) + monkeypatch.chdir(elsewhere) + + service = RegistryIdentityServiceV1.from_file(registry) + resolved = service.resolve_db("primary") + + assert resolved.registry_path == registry + assert resolved.state_root == root / ".hermes" / "state" + assert resolved.db_path == root / ".hermes" / "state" / "primary" / "workflows.sqlite" + assert resolved.db_alias == "primary" + + +def test_identity_is_alias_only_and_redacts_registry_and_db_paths(tmp_path: Path) -> None: + root = tmp_path / "private-user" / "secret-project" + registry = _write_registry(root) + service = RegistryIdentityServiceV1.from_file(registry) + + identity = service.identity("primary") + encoded = json.dumps(identity.to_dict(), sort_keys=True) + + assert identity.db_alias == "primary" + assert identity.registry_fingerprint == service.catalog.fingerprint + assert identity.resolved_db_identity.startswith("sha256:") + assert identity.registry_identity.startswith("sha256:") + assert str(root) not in encoded + assert "workflows.sqlite" not in encoded + for path_like in (str(registry), "state/primary/workflows.sqlite", "primary/workflows.sqlite"): + with pytest.raises(RegistryContractError) as raised: + service.identity(path_like) + assert raised.value.code == "registry_alias_required" + assert raised.value.exit_code == 2 + + +def test_two_copy_local_databases_with_the_same_alias_and_catalog_are_not_confused(tmp_path: Path) -> None: + left = RegistryIdentityServiceV1.from_file(_write_registry(tmp_path / "left")) + right = RegistryIdentityServiceV1.from_file(_write_registry(tmp_path / "right")) + + left_identity = left.identity("primary") + right_identity = right.identity("primary") + + assert left_identity.registry_fingerprint == right_identity.registry_fingerprint + assert left_identity.db_alias == right_identity.db_alias == "primary" + assert left_identity.registry_identity != right_identity.registry_identity + assert left_identity.resolved_db_identity != right_identity.resolved_db_identity + with pytest.raises(RegistryContractError) as raised: + require_consumer_parity({"cli": left_identity, "plugin": right_identity}) + assert raised.value.code == "registry_drift" + assert raised.value.exit_code == 2 + + +def test_changed_catalog_fingerprint_detects_similar_database_drift(tmp_path: Path) -> None: + valid = RegistryIdentityServiceV1.from_file(_write_registry(tmp_path / "valid", VALID)) + drift = RegistryIdentityServiceV1.from_file(_write_registry(tmp_path / "drift", DRIFT)) + + assert valid.catalog.fingerprint != drift.catalog.fingerprint + with pytest.raises(RegistryContractError, match="do not share one registry identity"): + require_consumer_parity({"cli": valid.identity("primary"), "supervisor": drift.identity("primary")}) + + +def test_v2_uses_one_canonical_spelling_and_rejects_unknown_or_legacy_shapes() -> None: + valid = json.loads(VALID.read_text(encoding="utf-8")) + invalid_payloads = [] + + extra_root = copy.deepcopy(valid) + extra_root["version"] = 2 + invalid_payloads.append(extra_root) + + string_db = copy.deepcopy(valid) + string_db["dbs"]["primary"] = "primary/workflows.sqlite" + invalid_payloads.append(string_db) + + alternate_ref = copy.deepcopy(valid) + alternate_ref["workflows"]["review-workflow"]["ref"] = alternate_ref["workflows"]["review-workflow"].pop( + "workflow_ref" + ) + invalid_payloads.append(alternate_ref) + + unknown_workflow_field = copy.deepcopy(valid) + unknown_workflow_field["workflows"]["review-workflow"]["database"] = "primary" + invalid_payloads.append(unknown_workflow_field) + + for payload in invalid_payloads: + with pytest.raises(RegistryContractError) as raised: + decode_registry(json.dumps(payload)) + assert raised.value.code == "registry_invalid" + assert raised.value.exit_code == 2 + + +def test_malformed_duplicate_and_unknown_ids_fail_closed() -> None: + valid = json.loads(VALID.read_text(encoding="utf-8")) + + malformed_alias = copy.deepcopy(valid) + malformed_alias["dbs"]["Primary"] = malformed_alias["dbs"].pop("primary") + + duplicate_runner = copy.deepcopy(valid) + duplicate_runner["runner"]["dbs"] = ["primary", "primary"] + + unknown_db = copy.deepcopy(valid) + unknown_db["workflows"]["review-workflow"]["db"] = "missing" + + filesystem_ref = copy.deepcopy(valid) + filesystem_ref["workflows"]["review-workflow"]["workflow_ref"] = "../workflow.py:run" + + for payload in (malformed_alias, duplicate_runner, unknown_db, filesystem_ref): + with pytest.raises(RegistryContractError) as raised: + decode_registry(json.dumps(payload)) + assert raised.value.code == "registry_invalid" + + duplicate_key = VALID.read_text(encoding="utf-8").replace( + '"schema_version": 2,', + '"schema_version": 2, "schema_version": 2,', + 1, + ) + with pytest.raises(RegistryContractError) as raised: + decode_registry(duplicate_key) + assert raised.value.code == "registry_invalid" + + +def test_runner_and_json_values_are_strict_and_deterministic() -> None: + valid = json.loads(VALID.read_text(encoding="utf-8")) + invalid_payloads = [] + for lease in (True, 0, 3601, 1.5, "30"): + payload = copy.deepcopy(valid) + payload["runner"]["lease_seconds"] = lease + invalid_payloads.append(payload) + + nan_payload = VALID.read_text(encoding="utf-8").replace('"mode": "review"', '"mode": NaN') + for payload in invalid_payloads: + with pytest.raises(RegistryContractError): + decode_registry(json.dumps(payload)) + with pytest.raises(RegistryContractError): + decode_registry(nan_payload) + + +def test_path_traversal_absolute_home_drive_and_backslash_values_fail_closed() -> None: + valid = json.loads(VALID.read_text(encoding="utf-8")) + bad_paths = [ + "../escape", + "nested/../escape", + "/private/state", + "~/state", + "C:/state", + "nested\\state", + "nested//state", + "nested/./state", + ] + for bad in bad_paths: + bad_state = copy.deepcopy(valid) + bad_state["state_root"] = bad + bad_db = copy.deepcopy(valid) + bad_db["dbs"]["primary"]["path"] = bad + for payload in (bad_state, bad_db): + with pytest.raises(RegistryContractError) as raised: + decode_registry(json.dumps(payload)) + assert raised.value.code == "registry_invalid" + assert raised.value.exit_code == 2 + + +def test_registry_state_root_and_db_symlinks_are_refused(tmp_path: Path) -> None: + outside = tmp_path / "outside" + outside.mkdir() + + real_registry = _write_registry(tmp_path / "registry-link-source") + linked_registry = tmp_path / "registry-link.json" + linked_registry.symlink_to(real_registry) + with pytest.raises(RegistryContractError) as raised: + load_registry_file(linked_registry) + assert raised.value.code == "registry_path_invalid" + + state_root = tmp_path / "state-link-root" + registry = _write_registry(state_root) + (state_root / ".hermes" / "state").symlink_to(outside, target_is_directory=True) + with pytest.raises(RegistryContractError) as raised: + RegistryIdentityServiceV1.from_file(registry).identity("primary") + assert raised.value.code == "registry_path_invalid" + + db_link_root = tmp_path / "db-link-root" + registry = _write_registry(db_link_root) + state = db_link_root / ".hermes" / "state" + state.mkdir() + (state / "primary").symlink_to(outside, target_is_directory=True) + with pytest.raises(RegistryContractError) as raised: + RegistryIdentityServiceV1.from_file(registry).identity("primary") + assert raised.value.code == "registry_path_invalid" + + +def test_noncanonical_registry_file_path_is_refused(tmp_path: Path) -> None: + registry = _write_registry(tmp_path / "root") + noncanonical = registry.parent / ".." / ".hermes" / registry.name + + with pytest.raises(RegistryContractError) as raised: + load_registry_file(noncanonical) + assert raised.value.code == "registry_path_invalid" + + +def test_errors_use_bounded_doctor_exit_2_envelopes_without_private_values(tmp_path: Path) -> None: + private_root = tmp_path / ("private-" + "x" * 100) + registry = _write_registry(private_root) + service = RegistryIdentityServiceV1.from_file(registry) + + with pytest.raises(RegistryContractError) as raised: + service.identity("missing") + + error = raised.value + envelope = error.to_dict() + encoded = json.dumps(envelope, sort_keys=True) + assert error.exit_code == 2 + assert set(envelope) == {"code", "message", "fields", "conflict_id"} + assert len(encoded.encode("utf-8")) <= 4096 + assert len(envelope["message"].encode("utf-8")) <= 256 + assert str(private_root) not in encoded + assert "workflows.sqlite" not in encoded + + oversized = '{"schema_version":2,"private":"' + "secret-value-" * 100_000 + '"}' + with pytest.raises(RegistryContractError) as oversized_error: + decode_registry(oversized) + oversized_encoded = json.dumps(oversized_error.value.to_dict(), sort_keys=True) + assert oversized_error.value.exit_code == 2 + assert len(oversized_encoded.encode("utf-8")) <= 4096 + assert "secret-value" not in oversized_encoded + + +def test_secret_like_defaults_never_appear_in_public_identity_or_drift_error(tmp_path: Path) -> None: + payload = json.loads(VALID.read_text(encoding="utf-8")) + payload["workflows"]["review-workflow"]["default_input"] = { + "api_token": "super-secret-token", + "private_path": "/Users/example/private", + } + left_path = tmp_path / "left" / ".hermes" / "workflows.registry.json" + left_path.parent.mkdir(parents=True) + left_path.write_text(json.dumps(payload), encoding="utf-8") + right_path = _write_registry(tmp_path / "right", DRIFT) + + left = RegistryIdentityServiceV1.from_file(left_path).identity("primary") + right = RegistryIdentityServiceV1.from_file(right_path).identity("primary") + public = json.dumps(left.to_dict(), sort_keys=True) + assert "super-secret-token" not in public + assert "/Users/example/private" not in public + + with pytest.raises(RegistryContractError) as raised: + require_consumer_parity({"cli": left, "plugin": right}) + error = json.dumps(raised.value.to_dict(), sort_keys=True) + assert "super-secret-token" not in error + assert "/Users/example/private" not in error