Repository navigation
feat(OMN-1752): Extract ContractPublisher service to omnibase_infra - #223
Conversation
Implements ARCH-002 (Runtime owns all Kafka plumbing) by extracting the ContractPublisher from omniclaude to omnibase_infra. Key components: - ServiceContractPublisher: Main service with injectable sources - ModelContractPublisherConfig: Configuration model for source selection - Sources: filesystem, package, and composite contract discovery - Result models: ModelPublishResult, ModelPublishStats, error models The service follows the flow: Source → Validate → Normalize → Publish → Report Models split into individual files per ONEX architecture rules: - models/model_contract_error.py - models/model_infra_error.py - models/model_publish_stats.py - models/model_publish_result.py
📝 WalkthroughWalkthroughAdds a Contract Publisher subsystem (config, sources, models, errors, service), extensive unit tests, and exposes the new APIs from the services package; also updates the Changes
Sequence DiagramsequenceDiagram
participant Client as Client
participant Service as ServiceContractPublisher
participant Source as ContractSource
participant Publisher as EventBusPublisher
participant Kafka as KafkaCluster
Client->>Service: publish_all()
activate Service
Service->>Source: discover_contracts()
activate Source
Source-->>Service: [ModelDiscoveredContract...]
deactivate Source
Note over Service: Validate YAML & schema\nCompute content_hash\nDeterministic sort\nAggregate contract errors
Service->>Publisher: publish(event) per contract
activate Publisher
Publisher->>Kafka: send message
activate Kafka
Kafka-->>Publisher: ack / error
deactivate Kafka
Publisher-->>Service: publish result (success/error)
deactivate Publisher
Note over Service: Aggregate infra errors & stats\nBuild ModelPublishResult
Service-->>Client: ModelPublishResult
deactivate Service
Estimated code review effort🎯 4 (Complex) | ⏱️ ~60 minutes Poem
🚥 Pre-merge checks | ✅ 3✅ Passed checks (3 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing touches
🧪 Generate unit tests (beta)
Comment |
There was a problem hiding this comment.
Actionable comments posted: 6
🤖 Fix all issues with AI agents
In `@pyproject.toml`:
- Around line 39-40: The omnibase-core version constraint "^0.10.2" is invalid
because 0.10.2 isn't published; update the dependency entry for "omnibase-core"
in pyproject.toml to a valid released version (e.g., "0.9.10") or to a supported
range within 0.3.x–0.9.x, replacing the "^0.10.2" constraint accordingly;
alternatively revert this change and wait to re-add the release once
omnibase-core 0.10.2 is published.
In `@src/omnibase_infra/services/contract_publisher/config.py`:
- Around line 163-173: The current getter returns the result of
_normalize_environment(self.environment) or _normalize_environment(env_var)
without checking for empty or whitespace-only normalized strings; change the
logic in the method that computes the environment so that after calling
_normalize_environment(...) you verify the returned value is truthy (non-empty,
not just whitespace or dots) before returning it—if the normalized value is
falsy, continue to the next priority (check env_var via env_var =
os.getenv("ONEX_ENV", "") and call _normalize_environment(env_var) similarly)
and finally return the default "dev"; update the branches that call
_normalize_environment(self.environment) and _normalize_environment(env_var) to
perform this truthiness guard.
- Around line 102-129: The validator validate_source_configured currently raises
ValueError; replace those raises with ProtocolConfigurationError and use
explicit error chaining (raise ProtocolConfigurationError("...") from
ValueError("...")) so each branch (filesystem, package, composite) raises
ProtocolConfigurationError with the same message and chains from a ValueError
constructed with that message; update the three raise sites in
validate_source_configured accordingly (and ensure ProtocolConfigurationError is
imported/available).
In `@src/omnibase_infra/services/contract_publisher/service.py`:
- Around line 342-358: The except block is not chaining the original exception
when re-raising, violating the OnexError pattern; update the raise of
ContractPublishingInfraError so it is raised from the caught exception (use
"raise ContractPublishingInfraError(infra_errors) from e") to preserve context,
keeping the existing ModelInfraError creation, logger.exception calls (with
handler_id and correlation_id) and the self._config.fail_fast conditional
intact.
- Around line 418-420: The code currently uses isinstance(self._source,
SourceContractComposite) to decide whether to call get_merge_errors; change this
to duck-typing by checking for and calling the method (e.g., if
getattr(self._source, "get_merge_errors", None) and
callable(self._source.get_merge_errors): errors =
self._source.get_merge_errors()) and remove the isinstance check; also add
get_merge_errors to the appropriate protocol in omnibase_spi so the LSP/typing
reflects the method on the protocol (update protocol interface rather than
concrete class checks).
In `@src/omnibase_infra/services/contract_publisher/sources/source_filesystem.py`:
- Around line 121-140: The read loop that calls
contract_path.read_text(encoding="utf-8") and constructs ModelDiscoveredContract
may raise UnicodeDecodeError (not OSError), so update the exception handling in
the try/except around that read to also catch UnicodeDecodeError (e.g., except
(OSError, UnicodeDecodeError) as e) and log the same warning and skip the file;
locate the block where ModelDiscoveredContract is created and
contract_path.read_text is called and add the additional except for
UnicodeDecodeError to preserve the "skip on read error" behavior.
🧹 Nitpick comments (3)
src/omnibase_infra/services/contract_publisher/sources/source_filesystem.py (1)
90-123: Avoid blocking file IO inside async discovery.
Path.glob/read_textare sync and can block the event loop under load; consider offloading to a thread or using aiofiles.💡 Minimal offload using
asyncio.to_thread+import asyncio @@ - text = contract_path.read_text(encoding="utf-8") + text = await asyncio.to_thread( + contract_path.read_text, encoding="utf-8" + )src/omnibase_infra/services/contract_publisher/sources/source_package.py (1)
91-135: Consider documenting that discovery is synchronous under the hood.The
asyncmethod calls synchronous_discover_recursive. This is acceptable for CPU-bound directory traversal, but the behavior differs from truly async I/O operations. Consider adding a note that large package trees may block the event loop briefly.tests/unit/services/contract_publisher/test_result.py (1)
130-144: Consider using specific exception type for validation test.Using
pytest.raises(Exception)is broad. Consider importing and usingpydantic.ValidationErrorfor more precise assertions.Suggested improvement
+from pydantic import ValidationError + ... def test_non_negative_constraints(self) -> None: """Test stats enforce non-negative values.""" - with pytest.raises(Exception): # Pydantic ValidationError + with pytest.raises(ValidationError): ModelPublishStats(
| # omnibase-core v0.10.2: Adds request_response schema to ModelEventBusSubcontract (OMN-1760) | ||
| omnibase-core = "^0.10.2" |
There was a problem hiding this comment.
🧩 Analysis chain
🏁 Script executed:
#!/bin/bash
# Check PyPI for omnibase-core availability and version info
python <<'PY'
import json
import urllib.request
import urllib.error
try:
url = "https://pypi.org/pypi/omnibase-core/json"
with urllib.request.urlopen(url, timeout=10) as r:
data = json.load(r)
print("Latest version:", data["info"]["version"])
releases = data.get("releases", {})
print("Version 0.10.2 exists:", "0.10.2" in releases)
if "0.10.2" in releases:
print("Release date for 0.10.2:", releases["0.10.2"][0].get("upload_time_iso_8601", "N/A") if releases["0.10.2"] else "N/A")
# List recent versions
recent_versions = sorted([v for v in releases.keys() if v.startswith("0.")], reverse=True)[:10]
print("Recent versions:", recent_versions)
except Exception as e:
print(f"Error: {e}")
PYRepository: OmniNode-ai/omnibase_infra
Length of output: 204
🌐 Web query:
omnibase-core v0.10.2 PyPI release
💡 Result:
There is no omnibase-core / omnibase_core version 0.10.2 published on PyPI. The PyPI project’s release history shows versions in the 0.3.x → 0.9.x range, with the latest release 0.9.10 on Jan 30, 2026. [1]
Source: PyPI project page for omnibase_core. [1]
Fix invalid omnibase-core version constraint—0.10.2 does not exist on PyPI.
The latest published version of omnibase-core is 0.9.10 (Jan 30, 2026). Specifying ^0.10.2 will cause dependency resolution to fail. Either downgrade to a version in the 0.3.x–0.9.x range, or wait for the 0.10.2 release and verify it's published before merging this change.
🤖 Prompt for AI Agents
In `@pyproject.toml` around lines 39 - 40, The omnibase-core version constraint
"^0.10.2" is invalid because 0.10.2 isn't published; update the dependency entry
for "omnibase-core" in pyproject.toml to a valid released version (e.g.,
"0.9.10") or to a supported range within 0.3.x–0.9.x, replacing the "^0.10.2"
constraint accordingly; alternatively revert this change and wait to re-add the
release once omnibase-core 0.10.2 is published.
| @model_validator(mode="after") | ||
| def validate_source_configured(self) -> Self: | ||
| """Validate that required source fields are configured for the mode. | ||
|
|
||
| Rules: | ||
| - filesystem mode requires filesystem_root | ||
| - package mode requires package_module | ||
| - composite mode requires at least one source | ||
|
|
||
| Returns: | ||
| Self if validation passes | ||
|
|
||
| Raises: | ||
| ValueError: If required source configuration is missing | ||
| """ | ||
| match self.mode: | ||
| case "filesystem": | ||
| if not self.filesystem_root: | ||
| raise ValueError("filesystem mode requires filesystem_root") | ||
| case "package": | ||
| if not self.package_module: | ||
| raise ValueError("package mode requires package_module") | ||
| case "composite": | ||
| if not self.filesystem_root and not self.package_module: | ||
| raise ValueError( | ||
| "composite mode requires at least one source " | ||
| "(filesystem_root or package_module)" | ||
| ) |
There was a problem hiding this comment.
🧩 Analysis chain
🏁 Script executed:
# First, check if the file exists and read it
cat -n src/omnibase_infra/services/contract_publisher/config.py | head -150Repository: OmniNode-ai/omnibase_infra
Length of output: 6489
🏁 Script executed:
# Search for OnexError usage in the codebase to understand the pattern
rg "OnexError" --type py -A 3 | head -100Repository: OmniNode-ai/omnibase_infra
Length of output: 8011
🏁 Script executed:
# Check imports in the config.py file
head -50 src/omnibase_infra/services/contract_publisher/config.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 1724
🏁 Script executed:
# Search for config-related error classes or OnexError subclasses
rg "class.*Error.*OnexError" --type pyRepository: OmniNode-ai/omnibase_infra
Length of output: 543
🏁 Script executed:
# Check how config validators in similar files handle errors
rg "model_validator|ValidationError|OnexError" src/omnibase_infra --type py -B 2 -A 2 | head -80Repository: OmniNode-ai/omnibase_infra
Length of output: 7127
🏁 Script executed:
# Check what's in omnibase_core.errors module to understand available error classes
rg "from omnibase_core.errors import" src/omnibase_infra --type py | head -20Repository: OmniNode-ai/omnibase_infra
Length of output: 2237
🏁 Script executed:
# Check if there are other config files with validators and how they handle errors
rg "model_validator|validator" src/omnibase_infra/services --type py -lRepository: OmniNode-ai/omnibase_infra
Length of output: 288
🏁 Script executed:
# Look at config.py more completely to see if there's resolve_environment and other validators
sed -n '130,200p' src/omnibase_infra/services/contract_publisher/config.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 1975
🏁 Script executed:
# Check if there are any config-specific error classes in the errors module
ls -la src/omnibase_infra/errors/Repository: OmniNode-ai/omnibase_infra
Length of output: 1079
🏁 Script executed:
# Check error_infra.py for config-related error classes
rg "class.*Error" src/omnibase_infra/errors/error_infra.py --type py | head -20Repository: OmniNode-ai/omnibase_infra
Length of output: 561
🏁 Script executed:
# Also check the __init__.py in errors to see what's exported
head -100 src/omnibase_infra/errors/__init__.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 4413
🏁 Script executed:
# Check if there are validators in other config files
rg "model_validator" src/omnibase_infra/services --type py -B 2 -A 10Repository: OmniNode-ai/omnibase_infra
Length of output: 7909
🏁 Script executed:
# Check how emit_daemon/config.py uses OnexError with validators or config
head -100 src/omnibase_infra/runtime/emit_daemon/config.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 3772
🏁 Script executed:
# Search for PydanticErrorMixin or ValidationError handling patterns
rg "ValidationError|PydanticError" src/omnibase_infra --type py | head -20Repository: OmniNode-ai/omnibase_infra
Length of output: 2482
🏁 Script executed:
# Check if there's a configuration-specific error class
rg "ConfigurationError|Config.*Error" src/omnibase_infra/errors --type pyRepository: OmniNode-ai/omnibase_infra
Length of output: 1054
🏁 Script executed:
# Look at config_consumer.py to see ProtocolConfigurationError usage pattern
grep -A 5 "ProtocolConfigurationError" src/omnibase_infra/services/session/config_consumer.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 482
🏁 Script executed:
# Check what errors are imported in config_consumer.py
head -40 src/omnibase_infra/services/session/config_consumer.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 1291
🏁 Script executed:
# Check ProtocolConfigurationError definition to understand if it accepts error chaining
grep -A 15 "class ProtocolConfigurationError" src/omnibase_infra/errors/error_infra.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 648
🏁 Script executed:
# See how other validators in the codebase handle error chaining
rg "from omnibase_infra.errors import|raise.*Error.*from" src/omnibase_infra/services --type py -A 2 | head -40Repository: OmniNode-ai/omnibase_infra
Length of output: 3086
🏁 Script executed:
# Check if config_store.py also raises ValueError to see consistency
grep -A 10 "validate_pool_sizes" src/omnibase_infra/services/session/config_store.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 411
🏁 Script executed:
# Check what happens in emit_daemon config validators
grep -A 10 "model_validator" src/omnibase_infra/runtime/emit_daemon/config.py | head -50Repository: OmniNode-ai/omnibase_infra
Length of output: 926
🏁 Script executed:
# Get more details on the emit_daemon config validators
sed -n '200,350p' src/omnibase_infra/runtime/emit_daemon/config.pyRepository: OmniNode-ai/omnibase_infra
Length of output: 5073
🏁 Script executed:
# Also check service_timeout_emitter.py to see how ProtocolConfigurationError is raised in validators
grep -B 5 -A 10 "raise ProtocolConfigurationError" src/omnibase_infra/services/service_timeout_emitter.py | head -30Repository: OmniNode-ai/omnibase_infra
Length of output: 1491
Replace ValueError with ProtocolConfigurationError and include error chaining.
Current code violates the mandatory error policy by raising ValueError. Use ProtocolConfigurationError (the config-specific subclass) with the required error chaining pattern.
Corrected pattern
+from omnibase_infra.errors import ProtocolConfigurationError
@@
case "filesystem":
if not self.filesystem_root:
- raise ValueError("filesystem mode requires filesystem_root")
+ err = ValueError("filesystem mode requires filesystem_root")
+ raise ProtocolConfigurationError(str(err)) from err
case "package":
if not self.package_module:
- raise ValueError("package mode requires package_module")
+ err = ValueError("package mode requires package_module")
+ raise ProtocolConfigurationError(str(err)) from err
case "composite":
if not self.filesystem_root and not self.package_module:
- raise ValueError(
- "composite mode requires at least one source "
- "(filesystem_root or package_module)"
- )
+ err = ValueError(
+ "composite mode requires at least one source "
+ "(filesystem_root or package_module)"
+ )
+ raise ProtocolConfigurationError(str(err)) from errPer coding guidelines: Only raise OnexError or its subclasses with error chaining pattern raise ... from e.
🤖 Prompt for AI Agents
In `@src/omnibase_infra/services/contract_publisher/config.py` around lines 102 -
129, The validator validate_source_configured currently raises ValueError;
replace those raises with ProtocolConfigurationError and use explicit error
chaining (raise ProtocolConfigurationError("...") from ValueError("...")) so
each branch (filesystem, package, composite) raises ProtocolConfigurationError
with the same message and chains from a ValueError constructed with that
message; update the three raise sites in validate_source_configured accordingly
(and ensure ProtocolConfigurationError is imported/available).
| errors: list[ModelContractError] = [] | ||
| if isinstance(self._source, SourceContractComposite): | ||
| errors = self._source.get_merge_errors() |
There was a problem hiding this comment.
🛠️ Refactor suggestion | 🟠 Major
Avoid isinstance for protocol resolution.
Switch to duck-typing (e.g., getattr/callable) or promote get_merge_errors() onto the protocol to keep LSP-friendly behavior.
♻️ Suggested change
- errors: list[ModelContractError] = []
- if isinstance(self._source, SourceContractComposite):
- errors = self._source.get_merge_errors()
+ errors: list[ModelContractError] = []
+ get_merge_errors = getattr(self._source, "get_merge_errors", None)
+ if callable(get_merge_errors):
+ errors = get_merge_errors()As per coding guidelines: Protocol resolution: use duck typing through protocols, never use isinstance checks; protocols defined in omnibase_spi.
🤖 Prompt for AI Agents
In `@src/omnibase_infra/services/contract_publisher/service.py` around lines 418 -
420, The code currently uses isinstance(self._source, SourceContractComposite)
to decide whether to call get_merge_errors; change this to duck-typing by
checking for and calling the method (e.g., if getattr(self._source,
"get_merge_errors", None) and callable(self._source.get_merge_errors): errors =
self._source.get_merge_errors()) and remove the isinstance check; also add
get_merge_errors to the appropriate protocol in omnibase_spi so the LSP/typing
reflects the method on the protocol (update protocol interface rather than
concrete class checks).
…sive tests Fixes identified during code review: - Extract handler_id from YAML in composite source BEFORE merge (was always None, breaking conflict detection) - Track and return dedup_count from composite source to service stats - Refactor _validate_contract to parse YAML once (removed double parsing) - Remove redundant _create_validation_error method - Fix node_version type (was str, should be ModelSemVer) Adds comprehensive test coverage: - test_service.py: 20 tests for ServiceContractPublisher - test_sources.py: 30 tests for sources (filesystem, package, composite) - Total: 80 tests in contract_publisher suite (was 30)
- Fix double-hashing: skip content_hash computation for contracts that already have a hash (composite source computes hashes internally) - Fix sort ordering: extract handler_id for ALL contracts before sorting to ensure deterministic ordering by (handler_id, origin, ref) - Add comprehensive documentation explaining dedup tracking behavior (only tracked for composite sources as single sources cannot have duplicates) - Add 3 new tests verifying handler_id extraction and sort consistency
…tion - Replace bare Exception catch with specific handlers in publish_all(): - ValidationError, TypeError, ValueError → serialization_failed - InfraTimeoutError → kafka_timeout - InfraUnavailableError → publisher_unavailable - InfraConnectionError → broker_down - Keep Exception fallback for unexpected errors (logged as warning) - Add docstring notes to source classes explaining async-for-protocol consistency pattern (sync I/O is acceptable for startup operations) - Remove unused imports (json, ModelSemVer)
…t failures - SourceContractFilesystem now raises ContractSourceNotConfiguredError when: - Root directory doesn't exist - Root path is not a directory - SourceContractPackage now raises ContractSourceNotConfiguredError when: - Package module not found - Package is invalid for resource discovery - DRY refactor: Moved duplicate _extract_handler_id logic to ModelDiscoveredContract.extract_handler_id() method - Updated tests to expect new error-throwing behavior
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Fix all issues with AI agents
In `@src/omnibase_infra/services/contract_publisher/sources/source_package.py`:
- Around line 201-206: The package resource read currently only catches OSError
in the except block that logs "Failed to read package resource %s: %s" for
child_path; add an additional except for UnicodeDecodeError around the
read_text(encoding="utf-8") call (the same place handling child_path) so that
invalid UTF-8 content is skipped and logged the same way as other read errors
(use logger.warning with the child_path and exception details). Ensure the new
except handles UnicodeDecodeError separately or as a multi-except with OSError
to avoid letting decoding errors propagate and break discovery.
| except OSError as e: | ||
| logger.warning( | ||
| "Failed to read package resource %s: %s", | ||
| child_path, | ||
| e, | ||
| ) |
There was a problem hiding this comment.
Catch UnicodeDecodeError for package resource reads.
Similar to the filesystem source, read_text(encoding="utf-8") at line 186 can raise UnicodeDecodeError for invalid UTF-8 content. This won't be caught by the OSError handler, causing the discovery to fail instead of skipping the problematic resource.
🛠️ Proposed fix
- except OSError as e:
+ except (OSError, UnicodeDecodeError) as e:
logger.warning(
"Failed to read package resource %s: %s",
child_path,
e,
)📝 Committable suggestion
‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.
| except OSError as e: | |
| logger.warning( | |
| "Failed to read package resource %s: %s", | |
| child_path, | |
| e, | |
| ) | |
| except (OSError, UnicodeDecodeError) as e: | |
| logger.warning( | |
| "Failed to read package resource %s: %s", | |
| child_path, | |
| e, | |
| ) |
🤖 Prompt for AI Agents
In `@src/omnibase_infra/services/contract_publisher/sources/source_package.py`
around lines 201 - 206, The package resource read currently only catches OSError
in the except block that logs "Failed to read package resource %s: %s" for
child_path; add an additional except for UnicodeDecodeError around the
read_text(encoding="utf-8") call (the same place handling child_path) so that
invalid UTF-8 content is skipped and logged the same way as other read errors
(use logger.warning with the child_path and exception details). Ensure the new
except handles UnicodeDecodeError separately or as a multi-except with OSError
to avoid letting decoding errors propagate and break discovery.
Resolves merge conflict in pyproject.toml by combining changelog comments for omnibase-core v0.10.2 from both branches.
Fixes identified during code review: 1. service.py - Exception chaining (ONEX compliance) - Added `from e` to all 6 ContractPublishingInfraError raises in fail_fast blocks - Preserves exception context for debugging 2. service.py - Duck typing over isinstance (ONEX compliance) - Replaced isinstance(self._source, SourceContractComposite) with getattr/callable - Uses duck typing pattern: getattr(obj, "method", None) + callable() 3. config.py - Whitespace edge case in resolve_environment() - Fixed whitespace-only strings falling through to default - " " now correctly returns "dev" instead of empty string 4. source_filesystem.py - UnicodeDecodeError handling - Added UnicodeDecodeError to exception handling for file reads - Binary files no longer crash discovery 5. test_result.py - Specific exception type - Changed pytest.raises(Exception) to pytest.raises(ValidationError) - More precise test assertions
Implements ARCH-002 (Runtime owns all Kafka plumbing) by extracting the ContractPublisher from omniclaude to omnibase_infra.
Key components:
The service follows the flow: Source → Validate → Normalize → Publish → Report
Models split into individual files per ONEX architecture rules:
Summary by CodeRabbit
New Features
Tests
Chores
✏️ Tip: You can customize this high-level summary in your review settings.