Repository navigation
backmerge(OMN-12245): projection runtime dispatch hotfix - #1776
Conversation
|
Warning Review limit reached
More reviews will be available in 42 minutes and 33 seconds. Learn how PR review limits work. Your organization has run out of usage credits. Purchase more in the billing tab. ⌛ How to resolve this issue?After more reviews become available, a review can be triggered using the We recommend that you space out your commits to avoid hitting the rate limit. 🚦 How do rate limits work?CodeRabbit enforces hourly rate limits for each developer per organization. Our paid plans include higher PR review limits than trial, open-source, and free plans. In all cases, reviews become available again over time. During sustained high-volume PR review activity, CodeRabbit may temporarily slow when the next review becomes available. Please see our Fair Usage Limits Policy for further information. ℹ️ Review info⚙️ Run configurationConfiguration used: defaults Review profile: CHILL Plan: Pro Run ID: 📒 Files selected for processing (3)
📝 WalkthroughWalkthroughThis PR adds PR-target enforcement, workspace build provenance and machine-class config seeding, per-database migrations, updated runtime health semantics, dispatch and routing changes, new Linear project-tracker support, multiple boundary model migrations to Pydantic, and broad supporting test and validation updates. ChangesRuntime platform updates
Estimated code review effort🎯 5 (Critical) | ⏱️ ~110 minutes Possibly related PRs
Poem
✨ Finishing Touches🧪 Generate unit tests (beta)
|
There was a problem hiding this comment.
Actionable comments posted: 17
Note
Due to the large number of review comments, Critical, Major severity comments were prioritized as inline comments.
🟡 Minor comments (6)
scripts/runtime_build/stage_workspace.sh-9-10 (1)
9-10:⚠️ Potential issue | 🟡 Minor | ⚡ Quick winCorrect the usage path in the script header.
The example points to
docker/runtime_build/stage_workspace.sh, but this script lives underscripts/runtime_build/stage_workspace.sh.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@scripts/runtime_build/stage_workspace.sh` around lines 9 - 10, Update the script header comment that shows example usage: replace the incorrect path "docker/runtime_build/stage_workspace.sh" with the correct "scripts/runtime_build/stage_workspace.sh" in the commented invocation lines so the usage example points to the actual script location; ensure both commented lines containing OMNI_HOME and the bash invocation are updated accordingly.scripts/runtime_build/compute_workspace_provenance.py-49-52 (1)
49-52:⚠️ Potential issue | 🟡 Minor | ⚡ Quick winFix
*.egg-infoexclusion logic in tree hashing.
"*.egg-info"is matched literally againstpath.parts, so egg-info dirs are not excluded. That makes digests noisy and environment-dependent.Proposed fix
- if any( - part in path.parts - for part in (".git", "__pycache__", ".venv", "*.egg-info") - ): + if any( + part in {".git", "__pycache__", ".venv"} or part.endswith(".egg-info") + for part in path.parts + ): continue🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@scripts/runtime_build/compute_workspace_provenance.py` around lines 49 - 52, The exclusion check in the if any(...) expression is treating "*.egg-info" as a literal part, so egg-info directories aren't being excluded; update the condition to match shell patterns for parts (e.g., use fnmatch.fnmatch(part, "*.egg-info")) or check part.endswith(".egg-info") instead of the literal "*.egg-info", and add an import for fnmatch if you choose fnmatch.fnmatch; keep the other literals (".git", "__pycache__", ".venv") as-is and apply this change where the code inspects path.parts in compute_workspace_provenance (the if any(...) expression).src/omnibase_infra/validation/validation_exemptions.yaml-755-755 (1)
755-755:⚠️ Potential issue | 🟡 Minor | ⚡ Quick winEscape
.infile_patternto avoid over-broad exemptions.At Line 755 and Line 827,
health_checker.pyis interpreted as regex, so.matches any character. This can exempt unintended filenames.Suggested fix
- - file_pattern: 'health_checker.py' + - file_pattern: 'health_checker\.py' ... - - file_pattern: 'health_checker.py' + - file_pattern: 'health_checker\.py'Also applies to: 827-827
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@src/omnibase_infra/validation/validation_exemptions.yaml` at line 755, The exemption uses a regex literal file_pattern 'health_checker.py' where the dot is unescaped and thus matches any character; update each occurrence of file_pattern 'health_checker.py' in validation_exemptions.yaml to escape the dot (e.g. 'health_checker\.py') or use a regex anchor like '^health_checker\.py$' so only the intended filename is exempted. Ensure you update every instance (the two occurrences currently using 'health_checker.py') so the pattern no longer over-broadly matches other filenames.tests/integration/envelope_routing/test_projection_topic_extraction.py-11-17 (1)
11-17:⚠️ Potential issue | 🟡 Minor | ⚡ Quick winAdd a pytest marker to this integration test.
This test currently has no pytest marker, which violates the test-marker guideline for files under
tests/.Suggested fix
from __future__ import annotations +import pytest + from omnibase_core.models.events.model_event_envelope import ModelEventEnvelope from omnibase_infra.runtime.auto_wiring.handler_wiring import _extract_projection_topic +@pytest.mark.integration def test_projection_topic_uses_onex_event_type_when_envelope_topic_is_absent() -> None:As per coding guidelines: "Test files must use pytest markers: mark test functions with
@pytest.mark.unit,@pytest.mark.integration,@pytest.mark.slow,@pytest.mark.chaos, or@pytest.mark.performanceas appropriate".🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@tests/integration/envelope_routing/test_projection_topic_extraction.py` around lines 11 - 17, Add the pytest marker to the test by decorating test_projection_topic_uses_onex_event_type_when_envelope_topic_is_absent with the appropriate marker (e.g., `@pytest.mark.integration`); ensure pytest is imported in the test module if not already. Locate the test function in tests/integration/envelope_routing/test_projection_topic_extraction.py and add the marker decorator immediately above the def to satisfy the test-marker guideline.tests/unit/config/test_machine_class_overlays.py-219-227 (1)
219-227:⚠️ Potential issue | 🟡 Minor | ⚡ Quick winStrengthen the unknown machine-class CLI assertion.
This test can pass on unrelated crashes because it only checks non-zero exit. Assert the argparse failure signature too (e.g., “invalid choice” + the bad value).
Suggested patch
def test_seed_script_cli_machine_class_unknown_exits_nonzero() -> None: """CLI rejects unknown --machine-class value at argparse level.""" result = subprocess.run( [sys.executable, str(SEED_SCRIPT), "--machine-class", "bad-class"], capture_output=True, text=True, check=False, ) assert result.returncode != 0 + combined = f"{result.stdout}\n{result.stderr}".lower() + assert "invalid choice" in combined + assert "bad-class" in combined🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@tests/unit/config/test_machine_class_overlays.py` around lines 219 - 227, The test test_seed_script_cli_machine_class_unknown_exits_nonzero currently only asserts a non-zero return code; update it to also assert the argparse failure message includes the invalid choice signature and the provided bad value by checking result.stderr (from the subprocess.run call invoking SEED_SCRIPT with "--machine-class bad-class") contains both "invalid choice" (or "invalid choice:" depending on platform) and "bad-class", while still asserting non-zero exit; this ensures the failure is from argparse rejecting the value rather than an unrelated crash.tests/unit/runtime/auto_wiring/test_handler_wiring_db_injection.py-230-233 (1)
230-233:⚠️ Potential issue | 🟡 Minor | ⚡ Quick winStabilize standalone-runner skip test by controlling DB env path.
At Line 230, the callback runs without patching DB URL lookup. If DB URL is absent in the test environment, this can pass via the “inactive projection” path instead of validating standalone-runner skip behavior.
✅ Suggested test hardening
- result = asyncio.run(callback(envelope)) + with patch( + _PATCH_ENVIRON_GET, + return_value="postgresql://user:pass@host:5432/omnidash_analytics", + ): + with patch(_PATCH_BUILD_ADAPTER, return_value=MagicMock()): + result = asyncio.run(callback(envelope))🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@tests/unit/runtime/auto_wiring/test_handler_wiring_db_injection.py` around lines 230 - 233, The test runs callback(envelope) without controlling DB URL lookup so the code can take the “inactive projection” path when DATABASE URL is missing; before calling asyncio.run(callback(envelope)) set a deterministic DB URL (e.g., set os.environ['DATABASE_URL'] or patch the function that returns the DB URL such as get_db_url/lookup_db_url) to a known in-memory value and ensure you restore/clear that env/patch after the assertion so the test always exercises the standalone-runner skip behavior; reference the test variables callback, envelope and handler when applying the setup/teardown.
🧹 Nitpick comments (2)
src/omnibase_infra/nodes/node_llm_completion_effect/handlers/handler_llm_completion.py (1)
143-143: ⚡ Quick winConsider explicit validation with a clearer error message.
The direct access to
os.environ["LLM_CODER_URL"]will raiseKeyErrorif the variable is unset. While the broad exception handler on line 103 catches this, the error message will be cryptic (just'LLM_CODER_URL'). For better debuggability in production, consider:url = os.environ.get("LLM_CODER_URL") if not url: raise RuntimeError( "LLM_CODER_URL environment variable is required when no endpoint is specified and token threshold exceeded" ) return url.rstrip("/")This provides a clearer error message while maintaining the contract requirement of no localhost fallback.
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@src/omnibase_infra/nodes/node_llm_completion_effect/handlers/handler_llm_completion.py` at line 143, Replace the direct access os.environ["LLM_CODER_URL"] in src/omnibase_infra/nodes/node_llm_completion_effect/handlers/handler_llm_completion.py with explicit validation: use os.environ.get("LLM_CODER_URL"), check for a truthy value, and if missing raise a RuntimeError with a clear message (e.g., "LLM_CODER_URL environment variable is required when no endpoint is specified and token threshold exceeded"); then return the validated URL with .rstrip("/") to preserve the existing behavior and avoid the cryptic KeyError.src/omnibase_infra/event_bus/topic_constants.py (1)
528-542: ⚡ Quick winRemove duplicated
__all__exports for delegate-skill topics.
TOPIC_DELEGATE_SKILL_COMPLETEDandTOPIC_DELEGATE_SKILL_FAILEDare listed twice in__all__. Keep a single entry per symbol to avoid redundant exports.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@src/omnibase_infra/event_bus/topic_constants.py` around lines 528 - 542, The __all__ list contains duplicate entries for TOPIC_DELEGATE_SKILL_COMPLETED and TOPIC_DELEGATE_SKILL_FAILED; edit the __all__ definition in topic_constants.py to remove the repeated occurrences so each symbol appears only once (ensure the remaining list still includes all unique topic names such as TOPIC_DELEGATION_REQUEST, TOPIC_DELEGATION_ROUTING_DECISION, TOPIC_DELEGATION_INFERENCE_RESPONSE, etc.).
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@docker/migrations/forward/082_swarm_runs.sql`:
- Around line 27-28: The migration defines source_offset as INTEGER which can
overflow; update the column definition in the forward migration (column name
source_offset in the 082_swarm_runs.sql diff) to use BIGINT DEFAULT 0 instead of
INTEGER DEFAULT 0 so Kafka offsets are stored as 64-bit values; ensure any
corresponding backward/rollback migration or references (e.g., indexes,
constraints, inserts) are updated to use BIGINT as well.
In `@docker/migrations/forward/083_create_log_entries.sql`:
- Around line 72-74: Add a new index on the ingested_at column used by the
retention/pruning path: create an index (e.g., idx_log_entries_ingested_at) on
the log_entries table for ingested_at (DESC) with IF NOT EXISTS so retention
queries/scans use ingested_at instead of timestamp; keep the existing
idx_log_entries_ts untouched.
In `@scripts/runtime_build/compute_workspace_provenance.py`:
- Around line 66-69: The _installed_direct_url function currently uses
dist.locate_file("direct_url.json") plus an exists() guard which can falsely
return None; change it to call dist.read_text("direct_url.json"), check that the
returned string is non-empty, and json.loads that content (remove the
locate_file/exists check and variable direct_url_file); additionally update
_hash_tree's exclusion logic which checks for literal `"*.egg-info"` in
path.parts—replace that check with a per-part test using
part.endswith(".egg-info") so entries like "pkg.egg-info" are correctly skipped
(refer to function names _installed_direct_url and _hash_tree and to
dist.read_text and path.parts in the diff).
In `@scripts/seed-infisical.py`:
- Around line 762-769: The machine-class branch returns early to call
_seed_machine_class without loading the operator env, so ensure ~/.omnibase/.env
is loaded first; modify the block that checks args.machine_class to load the
operator environment (e.g., via existing helper or python-dotenv) before
invoking _seed_machine_class so Infisical/database/Kafka credentials are
available when args.machine_class is not None.
In `@scripts/validation/topic_literal_baseline.txt`:
- Line 5: Remove the newly added suppression lines from
topic_literal_baseline.txt and instead fix the hard-coded topic literals in
service_kernel.py: search for raw topic string literals (e.g. any usage like
"orders.created", "user.updated" or similar) inside functions such as
publish_to_topic, subscribe_topic, handle_message or register_handlers and
replace them with the canonical constants/config lookup (e.g. TOPIC_* constants
or a get_topic_name(...) helper) used elsewhere in the codebase; ensure you
import or define the appropriate TOPIC_* symbols and update any unit tests or
callers to use the constants so no new baseline entries are required.
In `@src/omnibase_infra/adapters/llm/adapter_code_review_analysis.py`:
- Line 258: The code directly indexes os.environ for "LLM_CODER_FAST_URL" when
computing resolved_base_url, which raises KeyError before the constructor's
ModelInfraErrorContext.with_correlation(..., operation="validate_config") /
ProtocolConfigurationError logic can run; change the assignment for
resolved_base_url to use os.environ.get("LLM_CODER_FAST_URL") (or os.getenv) and
fall back to the provided base_url (or None) so missing env vars don't raise
immediately and the existing validation/ProtocolConfigurationError path can
handle configuration errors instead; keep the variable name resolved_base_url
and leave the surrounding validation and ModelInfraErrorContext.with_correlation
calls intact.
In `@src/omnibase_infra/adapters/models/model_infisical_batch_result.py`:
- Around line 32-35: The model currently sets model_config =
ConfigDict(frozen=False) which violates repo policy; change to use
ConfigDict(frozen=True, extra="forbid", from_attributes=True) for the Pydantic
model in model_infisical_batch_result.py. If the code was mutating model_config
incrementally, stop mutating it directly—collect any fields or intermediate
state in local variables (e.g., a dict for secrets/errors) and only construct
the final Pydantic model with the policy-compliant model_config
(ConfigDict(frozen=True, extra="forbid", from_attributes=True)) when
instantiating the ModelInfisicalBatchResult (or whatever model class uses
model_config). Ensure references to the symbols model_config,
ModelInfisicalBatchResult, and the secrets/errors Field defaults are updated so
no in-place mutations to model_config remain.
In `@src/omnibase_infra/adapters/models/model_infisical_secret_result.py`:
- Line 25: Update the model_config in model_infisical_secret_result by adding
extra="forbid" and from_attributes=True so it becomes ConfigDict(frozen=True,
extra="forbid", from_attributes=True); locate the existing model_config
assignment (symbol: model_config) in the file and replace the current
ConfigDict(frozen=True) with the full required ConfigDict signature to enforce
immutability, strict validation, and ORM/pytest-xdist compatibility.
In
`@src/omnibase_infra/adapters/project_tracker/adapter_project_tracker_linear.py`:
- Line 40: Replace the non-approved sanitizer `sanitize_error_string` with the
repository-approved `sanitize_error_message`: update the import (replace
sanitize_error_string with sanitize_error_message) and change all call sites
(e.g., the usage around lines ~326 in this module) to call
sanitize_error_message(...) instead; ensure any variable names or error handling
that referenced sanitize_error_string are updated to the new function name so
the module (AdapterProjectTrackerLinear and its exception handling code) uses
the sanctioned util from omnibase_infra.utils.util_error_sanitization.
In `@src/omnibase_infra/runtime/models/model_handshake_check_result.py`:
- Line 30: The Pydantic model ModelHandshakeCheckResult uses an incorrect
model_config; update the class-level model_config to the repository-required
ConfigDict(frozen=True, extra="forbid", from_attributes=True) to enforce
immutability, forbid extra fields, and enable from-attributes behavior; ensure
the existing ConfigDict import (if present) is used/updated accordingly and
replace the current model_config = ConfigDict(frozen=False) with the mandated
settings in the ModelHandshakeCheckResult definition.
In `@src/omnibase_infra/runtime/models/model_handshake_result.py`:
- Line 63: The ModelHandshakeResult Pydantic model currently sets model_config =
ConfigDict(frozen=False); update this to the repository-standard config by
changing model_config on the ModelHandshakeResult class to
ConfigDict(frozen=True, extra="forbid", from_attributes=True) so the model is
immutable, forbids extra fields, and supports from_attributes/ORM compatibility.
In `@src/omnibase_infra/runtime/models/model_plugin_discovery_entry.py`:
- Around line 62-63: The ConfigDict assigned to model_config currently only sets
frozen=True; update the ConfigDict call for the Pydantic model (the model_config
variable in ModelPluginDiscoveryEntry) to include extra="forbid" and
from_attributes=True so it becomes ConfigDict(frozen=True, extra="forbid",
from_attributes=True) to enforce immutability, forbid unknown fields, and enable
ORM/pytest-xdist compatibility.
In `@src/omnibase_infra/runtime/models/model_plugin_discovery_report.py`:
- Around line 43-44: The ConfigDict for the Pydantic model is incomplete: update
the model_config variable (in model_plugin_discovery_report.py) to use
ConfigDict(frozen=True, extra="forbid", from_attributes=True) so the Pydantic
model is immutable, forbids extra fields, and supports from-attributes/ORM
usage; locate the model_config assignment and replace the current
ConfigDict(...) call accordingly.
In `@src/omnibase_infra/runtime/service_delegation_dispatch_port.py`:
- Around line 228-248: The code currently resolves the bridge using the
hard-coded string _BRIDGE_PORT_KEY = "DirectBridgeDelegationDispatchPort";
update _resolve to use the protocol token instead of a concrete name: stop
passing that literal and call registry.resolve_service with the
ProtocolDelegationDispatchPort token (or ProtocolDelegationDispatchPort.__name__
if your DI expects a string). Concretely, change references in _resolve (and the
_BRIDGE_PORT_KEY constant) so resolve_service receives
ProtocolDelegationDispatchPort (or its name) rather than
"DirectBridgeDelegationDispatchPort", keeping the rest of the await
registry.resolve_service(...) logic and the type annotation for port as
ProtocolDelegationDispatchPort.
- Around line 279-289: ContainerBackedDelegationDispatchPort.dispatch currently
accepts an output_schema_key parameter but fails to forward it to the delegated
call; update the return await port.dispatch(...) call inside
ContainerBackedDelegationDispatchPort.dispatch to include
output_schema_key=output_schema_key so the delegated port receives the value
(ensure the parameter name matches the signature expected by port.dispatch).
In `@src/omnibase_infra/runtime/service_kernel.py`:
- Around line 2012-2029: The wiring currently hardcodes Omnimarket terminal
topics for node_delegate_skill_orchestrator using EnumOmnimarketTopic when
constructing DispatchResultApplier; instead, load the terminal publish topics
from the node contract metadata (the contract's event_bus.publish_topics /
subscribe_topics entry for the node_delegate_skill_orchestrator contract) and
pass those values into DispatchResultApplier as output_topic and
allowed_output_topics; update the auto_wiring_result_appliers population to
resolve the contract (node_delegate_skill_orchestrator) metadata at runtime,
extract the configured topic names, and use them rather than EnumOmnimarketTopic
so contract-driven YAML controls the topics.
In `@tests/unit/event_bus/test_kafka_event_bus.py`:
- Around line 2075-2115: The test's fake_start_unlocked does not mark consumers
as started, so EventBusKafka.start_consuming can schedule the same (topic,
group_id) again and make the test flaky; modify fake_start_unlocked (the
replacement for EventBusKafka._start_consumer_for_topic_unlocked used in the
test) to record a started marker into the EventBusKafka instance (e.g., add the
(topic, group_id) to an in-instance set like event_bus._started_consumers or to
the event_bus._consumer_tasks mapping) while holding the same lock used by
start_consuming so subsequent iterations see the consumer as started and do not
append duplicates to call_log; keep the artificial await asyncio.sleep(0) to
preserve concurrency simulation.
---
Minor comments:
In `@scripts/runtime_build/compute_workspace_provenance.py`:
- Around line 49-52: The exclusion check in the if any(...) expression is
treating "*.egg-info" as a literal part, so egg-info directories aren't being
excluded; update the condition to match shell patterns for parts (e.g., use
fnmatch.fnmatch(part, "*.egg-info")) or check part.endswith(".egg-info") instead
of the literal "*.egg-info", and add an import for fnmatch if you choose
fnmatch.fnmatch; keep the other literals (".git", "__pycache__", ".venv") as-is
and apply this change where the code inspects path.parts in
compute_workspace_provenance (the if any(...) expression).
In `@scripts/runtime_build/stage_workspace.sh`:
- Around line 9-10: Update the script header comment that shows example usage:
replace the incorrect path "docker/runtime_build/stage_workspace.sh" with the
correct "scripts/runtime_build/stage_workspace.sh" in the commented invocation
lines so the usage example points to the actual script location; ensure both
commented lines containing OMNI_HOME and the bash invocation are updated
accordingly.
In `@src/omnibase_infra/validation/validation_exemptions.yaml`:
- Line 755: The exemption uses a regex literal file_pattern 'health_checker.py'
where the dot is unescaped and thus matches any character; update each
occurrence of file_pattern 'health_checker.py' in validation_exemptions.yaml to
escape the dot (e.g. 'health_checker\.py') or use a regex anchor like
'^health_checker\.py$' so only the intended filename is exempted. Ensure you
update every instance (the two occurrences currently using 'health_checker.py')
so the pattern no longer over-broadly matches other filenames.
In `@tests/integration/envelope_routing/test_projection_topic_extraction.py`:
- Around line 11-17: Add the pytest marker to the test by decorating
test_projection_topic_uses_onex_event_type_when_envelope_topic_is_absent with
the appropriate marker (e.g., `@pytest.mark.integration`); ensure pytest is
imported in the test module if not already. Locate the test function in
tests/integration/envelope_routing/test_projection_topic_extraction.py and add
the marker decorator immediately above the def to satisfy the test-marker
guideline.
In `@tests/unit/config/test_machine_class_overlays.py`:
- Around line 219-227: The test
test_seed_script_cli_machine_class_unknown_exits_nonzero currently only asserts
a non-zero return code; update it to also assert the argparse failure message
includes the invalid choice signature and the provided bad value by checking
result.stderr (from the subprocess.run call invoking SEED_SCRIPT with
"--machine-class bad-class") contains both "invalid choice" (or "invalid
choice:" depending on platform) and "bad-class", while still asserting non-zero
exit; this ensures the failure is from argparse rejecting the value rather than
an unrelated crash.
In `@tests/unit/runtime/auto_wiring/test_handler_wiring_db_injection.py`:
- Around line 230-233: The test runs callback(envelope) without controlling DB
URL lookup so the code can take the “inactive projection” path when DATABASE URL
is missing; before calling asyncio.run(callback(envelope)) set a deterministic
DB URL (e.g., set os.environ['DATABASE_URL'] or patch the function that returns
the DB URL such as get_db_url/lookup_db_url) to a known in-memory value and
ensure you restore/clear that env/patch after the assertion so the test always
exercises the standalone-runner skip behavior; reference the test variables
callback, envelope and handler when applying the setup/teardown.
---
Nitpick comments:
In `@src/omnibase_infra/event_bus/topic_constants.py`:
- Around line 528-542: The __all__ list contains duplicate entries for
TOPIC_DELEGATE_SKILL_COMPLETED and TOPIC_DELEGATE_SKILL_FAILED; edit the __all__
definition in topic_constants.py to remove the repeated occurrences so each
symbol appears only once (ensure the remaining list still includes all unique
topic names such as TOPIC_DELEGATION_REQUEST, TOPIC_DELEGATION_ROUTING_DECISION,
TOPIC_DELEGATION_INFERENCE_RESPONSE, etc.).
In
`@src/omnibase_infra/nodes/node_llm_completion_effect/handlers/handler_llm_completion.py`:
- Line 143: Replace the direct access os.environ["LLM_CODER_URL"] in
src/omnibase_infra/nodes/node_llm_completion_effect/handlers/handler_llm_completion.py
with explicit validation: use os.environ.get("LLM_CODER_URL"), check for a
truthy value, and if missing raise a RuntimeError with a clear message (e.g.,
"LLM_CODER_URL environment variable is required when no endpoint is specified
and token threshold exceeded"); then return the validated URL with .rstrip("/")
to preserve the existing behavior and avoid the cryptic KeyError.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: defaults
Review profile: CHILL
Plan: Pro
Run ID: 5ad25d91-7e12-4828-a573-b08fc0a327b9
⛔ Files ignored due to path filters (2)
src/omnibase_infra/enums/generated/enum_omnimarket_topic.pyis excluded by!**/generated/**uv.lockis excluded by!**/*.lock
📒 Files selected for processing (255)
.github/workflows/main-target-guard.yml.gitignoreconfig/infisical_projects.yamlconfig/overlays/cloud-k8s.yamlconfig/overlays/linux-server.yamlconfig/overlays/mac-dev.yamlcontracts/OMN-11068.yamlcontracts/OMN-12184.yamlcontracts/OMN-12245.yamlcontracts/runtime/runtime_protocol.lock.jsondocker/Dockerfile.migratedocker/Dockerfile.runtimedocker/catalog/services/consumer-health-projection.yamldocker/docker-compose.infra.ymldocker/migrations/forward/082_swarm_runs.sqldocker/migrations/forward/083_create_log_entries.sqldocker/migrations/intelligence/024_create_llm_delegation_projection_tables.sqldocker/migrations/intelligence/025_fix_llm_delegation_call_log_date_index.sqldocker/migrations/rollback/rollback_083_create_log_entries.sqldocker/migrations/schema_fingerprint.sha256pyproject.tomlscripts/check-env-reads.shscripts/demo_runtime_verification.pyscripts/deploy-agent/deploy_agent/executor.pyscripts/deploy-agent/tests/unit/test_executor_build_source.pyscripts/deploy-agent/tests/unit/test_executor_cache_bust.pyscripts/deploy-agent/tests/unit/test_executor_workspace_provenance.pyscripts/deploy-runtime.shscripts/run-migrations.pyscripts/runtime_build/compute_workspace_provenance.pyscripts/runtime_build/stage_workspace.shscripts/seed-infisical.pyscripts/validation/topic_literal_baseline.txtscripts/validation/validate_clean_root.pysrc/omnibase_infra/adapters/llm/adapter_code_analysis_enrichment.pysrc/omnibase_infra/adapters/llm/adapter_code_review_analysis.pysrc/omnibase_infra/adapters/llm/adapter_documentation_generation.pysrc/omnibase_infra/adapters/llm/adapter_llm_provider_openai.pysrc/omnibase_infra/adapters/llm/adapter_summarization_enrichment.pysrc/omnibase_infra/adapters/llm/adapter_test_boilerplate_generation.pysrc/omnibase_infra/adapters/models/model_infisical_batch_result.pysrc/omnibase_infra/adapters/models/model_infisical_secret_result.pysrc/omnibase_infra/adapters/project_tracker/adapter_project_tracker_linear.pysrc/omnibase_infra/adapters/project_tracker/enum_project_tracker_issue_status_type.pysrc/omnibase_infra/adapters/project_tracker/model_project_tracker_issue_status.pysrc/omnibase_infra/adapters/project_tracker/model_project_tracker_label.pysrc/omnibase_infra/adapters/project_tracker/model_project_tracker_team.pysrc/omnibase_infra/cli/model_demo_reset_config.pysrc/omnibase_infra/cli/model_demo_reset_report.pysrc/omnibase_infra/cli/model_reset_action_result.pysrc/omnibase_infra/configs/pricing_manifest.yamlsrc/omnibase_infra/diagnostics/bus_audit.pysrc/omnibase_infra/diagnostics/models.pysrc/omnibase_infra/docker/catalog/cli.pysrc/omnibase_infra/docker/catalog/manifest_schema.pysrc/omnibase_infra/docker/catalog/resolver.pysrc/omnibase_infra/docker/catalog/validator.pysrc/omnibase_infra/errors/error_catalog.pysrc/omnibase_infra/event_bus/event_bus_kafka.pysrc/omnibase_infra/event_bus/topic_constants.pysrc/omnibase_infra/event_bus/topic_violation_alerter.pysrc/omnibase_infra/gateway/services/service_envelope_validator.pysrc/omnibase_infra/gateway/services/service_policy_engine.pysrc/omnibase_infra/handlers/mcp/adapter_onex_to_mcp.pysrc/omnibase_infra/mixins/mixin_postgres_error_response.pysrc/omnibase_infra/models/dispatch/model_dispatch_context.pysrc/omnibase_infra/nodes/node_decision_store_effect/contract.yamlsrc/omnibase_infra/nodes/node_decision_store_effect/handlers/handler_write_decision.pysrc/omnibase_infra/nodes/node_kafka_replay_compute/contract.yamlsrc/omnibase_infra/nodes/node_kafka_replay_compute/handlers/handler_replay.pysrc/omnibase_infra/nodes/node_llm_completion_effect/contract.yamlsrc/omnibase_infra/nodes/node_llm_completion_effect/handlers/handler_llm_completion.pysrc/omnibase_infra/nodes/node_llm_inference_effect/handlers/bifrost/handler_bifrost_gateway.pysrc/omnibase_infra/nodes/node_node_graph_reducer/reducer.pysrc/omnibase_infra/nodes/node_registration_orchestrator/contract.yamlsrc/omnibase_infra/nodes/node_registration_orchestrator/handlers/handler_node_registration_acked.pysrc/omnibase_infra/nodes/node_registry_api_effect/contract.yamlsrc/omnibase_infra/nodes/node_registry_api_effect/handlers/__init__.pysrc/omnibase_infra/nodes/node_registry_api_effect/handlers/handler_registry_api_get_contract.pysrc/omnibase_infra/nodes/node_registry_api_effect/handlers/handler_registry_api_get_node.pysrc/omnibase_infra/nodes/node_registry_api_effect/handlers/handler_registry_api_get_topic.pysrc/omnibase_infra/nodes/node_registry_api_effect/handlers/handler_registry_api_list_contracts.pysrc/omnibase_infra/nodes/node_registry_api_effect/handlers/handler_registry_api_list_nodes.pysrc/omnibase_infra/nodes/node_registry_api_effect/handlers/handler_registry_api_list_topics.pysrc/omnibase_infra/nodes/node_registry_api_effect/node.pysrc/omnibase_infra/nodes/node_rsd_score_compute/contract.yamlsrc/omnibase_infra/nodes/node_rsd_score_compute/handlers/handler_rsd_score_calculate.pysrc/omnibase_infra/protocols/protocol_dispatch_result_applier.pysrc/omnibase_infra/runtime/__init__.pysrc/omnibase_infra/runtime/_enum_coercion.pysrc/omnibase_infra/runtime/auto_wiring/discovery.pysrc/omnibase_infra/runtime/auto_wiring/handler_wiring.pysrc/omnibase_infra/runtime/chain_aware_dispatch.pysrc/omnibase_infra/runtime/handler_registry.pysrc/omnibase_infra/runtime/message_dispatch_engine.pysrc/omnibase_infra/runtime/models/model_domain_plugin_config.pysrc/omnibase_infra/runtime/models/model_domain_plugin_result.pysrc/omnibase_infra/runtime/models/model_handshake_check_result.pysrc/omnibase_infra/runtime/models/model_handshake_result.pysrc/omnibase_infra/runtime/models/model_plugin_discovery_entry.pysrc/omnibase_infra/runtime/models/model_plugin_discovery_report.pysrc/omnibase_infra/runtime/models/model_runtime_aggregate_health.pysrc/omnibase_infra/runtime/node_invocation_adapter.pysrc/omnibase_infra/runtime/protocols/protocol_delegation_dispatch_port.pysrc/omnibase_infra/runtime/registry_dispatcher.pysrc/omnibase_infra/runtime/render_bifrost_delegation_contract.pysrc/omnibase_infra/runtime/request_response_wiring.pysrc/omnibase_infra/runtime/runtime_contract_config_loader.pysrc/omnibase_infra/runtime/runtime_host_process.pysrc/omnibase_infra/runtime/runtime_local.pysrc/omnibase_infra/runtime/runtime_local_ingress.pysrc/omnibase_infra/runtime/service_delegation_dispatch_port.pysrc/omnibase_infra/runtime/service_dispatch_result_applier.pysrc/omnibase_infra/runtime/service_kernel.pysrc/omnibase_infra/runtime/service_pattern_b_broker.pysrc/omnibase_infra/runtime/strip_runtime_entry_points.pysrc/omnibase_infra/runtime/util_container_wiring.pysrc/omnibase_infra/runtime/util_mcp_auth.pysrc/omnibase_infra/runtime/version_compatibility.pysrc/omnibase_infra/scripts/verify_container_manifest.pysrc/omnibase_infra/services/eval/eval_regression_check.pysrc/omnibase_infra/services/eval/metric_collector.pysrc/omnibase_infra/services/health_checker.pysrc/omnibase_infra/services/observability/llm_cost_aggregation/consumer.pysrc/omnibase_infra/services/observability/savings_estimation/consumer.pysrc/omnibase_infra/services/registry_api/__init__.pysrc/omnibase_infra/services/registry_api/main.pysrc/omnibase_infra/services/registry_api/registry_discovery.pysrc/omnibase_infra/services/registry_api/routes.pysrc/omnibase_infra/services/service_stale_registration_cleanup.pysrc/omnibase_infra/services/session_registry/model_graph_mutation.pysrc/omnibase_infra/topics/contract_topic_extractor.pysrc/omnibase_infra/topics/model_topic_spec.pysrc/omnibase_infra/validation/validation_exemptions.yamlsrc/omnibase_infra/validation/validator_no_direct_adapter.pysrc/omnibase_infra/validation/validator_security.pysrc/omnibase_infra/validators/handler_any_signature.pysrc/omnibase_infra/validators/llm_registry_validator.pysrc/omnibase_infra/validators/llm_topology_validator.pysrc/omnibase_infra/validators/no_plugin_daemon_classes.pysrc/omnibase_infra/validators/type_ignore_budget.pysrc/omnibase_infra/verification/contract.yamlsrc/omnibase_infra/verification/orchestrator.pytests/ci/test_reusable_runtime_boot_schema.pytests/helpers/aiohttp_utils.pytests/integration/envelope_routing/test_projection_topic_extraction.pytests/integration/event_bus/test_kafka_concurrent_subscribe_startup.pytests/integration/gateway/test_runtime_gateway_integration.pytests/integration/infra/test_stability_test_runtime_compose_render.pytests/integration/projectors/test_registry_api_projection_tail_integration.pytests/integration/runtime/test_auto_wiring_deferred_subscription_order.pytests/integration/runtime/test_auto_wiring_duplicate_contract_discovery.pytests/integration/runtime/test_auto_wiring_resolver_end_to_end.pytests/integration/runtime/test_binding_resolution_dispatch_flow.pytests/integration/runtime/test_bootstrap_source_integration.pytests/integration/runtime/test_build_loop_auto_wiring_boot.pytests/integration/runtime/test_contract_routing_introspection.pytests/integration/runtime/test_dispatch_engine_post_freeze_dynamic_registration.pytests/integration/runtime/test_dynamic_contract_registration_e2e.pytests/integration/runtime/test_github_api_poll_golden_chain.pytests/integration/runtime/test_kernel_auto_wiring.pytests/integration/runtime/test_manifest_pool_injection_integration.pytests/integration/runtime/test_pattern_b_broker_omnimarket.pytests/integration/runtime/test_plugin_managed_subscription_integration.pytests/integration/runtime/test_projection_handler_db_injection_integration.pytests/integration/runtime/test_runtime_boot_clean.pytests/integration/runtime/test_runtime_handler_source_mode.pytests/integration/runtime/test_runtime_host_handler_discovery.pytests/integration/runtime/test_runtime_introspection.pytests/integration/runtime/test_runtime_local_ingress.pytests/integration/runtime/test_runtime_startup_readiness.pytests/integration/runtime/test_service_health_skill_endpoint_integration.pytests/integration/runtime/test_shutdown_health_integration.pytests/integration/runtime/test_topic_provisioner_live_materialization.pytests/integration/services/test_introspection_endpoint_integration.pytests/integration/test_handler_event_type_alias_integration.pytests/integration/test_kernel_prefetcher_wiring.pytests/integration/test_omn_9741_health_liveness.pytests/integration/test_omnibase_spi_runtime_protocol_floor.pytests/integration/test_orchestrator_dispatcher_coverage_integration.pytests/integration/test_pattern_b_broker_terminal_waiter.pytests/integration/test_pattern_b_terminal_result_applier.pytests/integration/test_runtime_host_dynamic_registration.pytests/integration/test_runtime_host_plugin_managed_gate.pytests/integration/test_runtime_local_ingress_operation_alias.pytests/unit/adapters/llm/test_adapter_llm_provider_openai.pytests/unit/adapters/test_adapter_project_tracker_linear.pytests/unit/cli/test_service_demo_reset.pytests/unit/config/test_machine_class_overlays.pytests/unit/conftest.pytests/unit/contracts/test_protocol_ownership.pytests/unit/event_bus/test_kafka_event_bus.pytests/unit/event_bus/test_readiness.pytests/unit/infra/test_catalog_validate_runtime.pytests/unit/infra/test_compose_no_silent_fallbacks.pytests/unit/infra/test_stability_test_runtime_lane.pytests/unit/models/runtime/test_model_plugin_discovery_report.pytests/unit/models/test_model_serialization_roundtrip.pytests/unit/nodes/node_registry_api_effect/test_node_registry_api_effect.pytests/unit/runtime/auto_wiring/test_discovery.pytests/unit/runtime/auto_wiring/test_handler_wiring_db_injection.pytests/unit/runtime/auto_wiring/test_handler_wiring_event_type_alias.pytests/unit/runtime/auto_wiring/test_handler_wiring_handle_async_dispatch.pytests/unit/runtime/auto_wiring/test_orchestrator_dispatcher_coverage.pytests/unit/runtime/auto_wiring/test_plugin_managed_subscription.pytests/unit/runtime/auto_wiring/test_raw_event_projection_wiring.pytests/unit/runtime/auto_wiring/test_wiring.pytests/unit/runtime/auto_wiring/test_wiring_subscribe.pytests/unit/runtime/health/test_health_published_events_map_wiring.pytests/unit/runtime/test_batch_response_publisher.pytests/unit/runtime/test_dispatch_context_integration.pytests/unit/runtime/test_dispatch_engine_post_freeze.pytests/unit/runtime/test_enum_coercion.pytests/unit/runtime/test_handler_pool_integration.pytests/unit/runtime/test_handler_registry_stub_removed.pytests/unit/runtime/test_handler_wiring_resolver_integration.pytests/unit/runtime/test_health_config_prefetch_status.pytests/unit/runtime/test_health_detailed.pytests/unit/runtime/test_hook_activations_startup_wiring.pytests/unit/runtime/test_import_order.pytests/unit/runtime/test_introspection_router.pytests/unit/runtime/test_kernel.pytests/unit/runtime/test_kernel_no_hardcoded_topics.pytests/unit/runtime/test_kernel_wiring.pytests/unit/runtime/test_live_contract_materialization.pytests/unit/runtime/test_message_dispatch_engine.pytests/unit/runtime/test_node_invocation_adapter.pytests/unit/runtime/test_node_subscription_wiring.pytests/unit/runtime/test_prefetch_contract_isolation.pytests/unit/runtime/test_render_bifrost_delegation_contract.pytests/unit/runtime/test_runtime_host_architecture_validation.pytests/unit/runtime/test_runtime_host_contract_config.pytests/unit/runtime/test_runtime_host_dynamic_registration.pytests/unit/runtime/test_runtime_host_process.pytests/unit/runtime/test_runtime_host_process_wiring_access.pytests/unit/runtime/test_runtime_host_topic_query.pytests/unit/runtime/test_runtime_idempotency_guard.pytests/unit/runtime/test_runtime_local_ingress.pytests/unit/runtime/test_runtime_parallel_handlers.pytests/unit/runtime/test_service_delegation_dispatch_port.pytests/unit/runtime/test_service_dispatch_result_applier.pytests/unit/runtime/test_service_health.pytests/unit/runtime/test_service_health_skill_endpoint.pytests/unit/runtime/test_service_pattern_b_broker.pytests/unit/runtime/test_topic_provisioner_integration.pytests/unit/scripts/test_run_migrations.pytests/unit/services/eval/test_metric_collector.pytests/unit/services/registry_api/test_feature_flag_routes.pytests/unit/services/registry_api/test_registry_discovery_semantics.pytests/unit/services/registry_api/test_registry_node_identity.pytests/unit/services/registry_api/test_registry_node_projection_fields.pytests/unit/services/registry_api/test_service_update_feature_flag.pytests/unit/test_enum_identity_regression.pytests/unit/validation/test_validation_exemptions.pyworkspace/sibling-repos/.gitkeep
💤 Files with no reviewable changes (2)
- docker/migrations/intelligence/024_create_llm_delegation_projection_tables.sql
- src/omnibase_infra/runtime/handler_registry.py
| source_partition INTEGER DEFAULT 0, | ||
| source_offset INTEGER DEFAULT 0, |
There was a problem hiding this comment.
Store Kafka offsets as BIGINT, not INTEGER.
source_offset INTEGER can overflow on long-lived/high-throughput topics.
Suggested fix
- source_offset INTEGER DEFAULT 0,
+ source_offset BIGINT DEFAULT 0,📝 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.
| source_partition INTEGER DEFAULT 0, | |
| source_offset INTEGER DEFAULT 0, | |
| source_partition INTEGER DEFAULT 0, | |
| source_offset BIGINT DEFAULT 0, |
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@docker/migrations/forward/082_swarm_runs.sql` around lines 27 - 28, The
migration defines source_offset as INTEGER which can overflow; update the column
definition in the forward migration (column name source_offset in the
082_swarm_runs.sql diff) to use BIGINT DEFAULT 0 instead of INTEGER DEFAULT 0 so
Kafka offsets are stored as 64-bit values; ensure any corresponding
backward/rollback migration or references (e.g., indexes, constraints, inserts)
are updated to use BIGINT as well.
| -- Standalone timestamp: full time-range scans and retention sweeps | ||
| CREATE INDEX IF NOT EXISTS idx_log_entries_ts | ||
| ON log_entries (timestamp DESC); |
There was a problem hiding this comment.
Add an index on ingested_at for the documented retention path.
Retention is anchored on ingested_at, but the current index is on timestamp, so pruning can become expensive as data grows.
Suggested fix
+CREATE INDEX IF NOT EXISTS idx_log_entries_ingested_at
+ ON log_entries (ingested_at DESC);📝 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.
| -- Standalone timestamp: full time-range scans and retention sweeps | |
| CREATE INDEX IF NOT EXISTS idx_log_entries_ts | |
| ON log_entries (timestamp DESC); | |
| -- Standalone timestamp: full time-range scans and retention sweeps | |
| CREATE INDEX IF NOT EXISTS idx_log_entries_ts | |
| ON log_entries (timestamp DESC); | |
| CREATE INDEX IF NOT EXISTS idx_log_entries_ingested_at | |
| ON log_entries (ingested_at DESC); |
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@docker/migrations/forward/083_create_log_entries.sql` around lines 72 - 74,
Add a new index on the ingested_at column used by the retention/pruning path:
create an index (e.g., idx_log_entries_ingested_at) on the log_entries table for
ingested_at (DESC) with IF NOT EXISTS so retention queries/scans use ingested_at
instead of timestamp; keep the existing idx_log_entries_ts untouched.
| direct_url_file = dist.locate_file("direct_url.json") | ||
| if not Path(str(direct_url_file)).exists(): | ||
| return None | ||
| return json.loads(Path(str(direct_url_file)).read_text(encoding="utf-8")) |
There was a problem hiding this comment.
🧩 Analysis chain
🏁 Script executed:
#!/usr/bin/env bash
# Verify whether locate_file("direct_url.json") points to an existing file for installed dists.
python - <<'PY'
import importlib.metadata as md
checked = 0
for dist in md.distributions():
name = dist.metadata.get("Name", "<unknown>")
p = dist.locate_file("direct_url.json")
exists = p.exists()
txt = dist.read_text("direct_url.json")
if txt and not exists:
print(f"mismatch: {name} locate_file_missing read_text_present")
checked += 1
if checked >= 10:
break
print("done")
PYRepository: OmniNode-ai/omnibase_infra
Length of output: 262
🏁 Script executed:
#!/usr/bin/env bash
set -euo pipefail
# Locate and print the SPDX header + relevant sections of the file
FILE="scripts/runtime_build/compute_workspace_provenance.py"
echo "=== Top of file (SPDX header check) ==="
nl -ba "$FILE" | head -n 60
echo
echo "=== Lines 60-100 (direct_url lookup check) ==="
nl -ba "$FILE" | sed -n '55,110p'
echo
echo "=== Search for egg-info exclusion logic ==="
rg -n --hidden --no-ignore-vcs "egg-info|egg_info|\\*\\.egg-info|path\\.parts" "$FILE" || true
# Print a bit more around any matches
for m in $(rg -n "egg-info|egg_info|\\*\\.egg-info" "$FILE" | cut -d: -f1); do
echo
echo "=== Around line $m ==="
nl -ba "$FILE" | sed -n "$((m-15)),$((m+15))p"
doneRepository: OmniNode-ai/omnibase_infra
Length of output: 171
🏁 Script executed:
#!/usr/bin/env bash
set -euo pipefail
FILE="scripts/runtime_build/compute_workspace_provenance.py"
echo "=== Top of file (SPDX header check) ==="
cat -n "$FILE" | head -n 80
echo
echo "=== Lines 60-110 (direct_url lookup check) ==="
cat -n "$FILE" | sed -n '55,110p'
echo
echo "=== Search for egg-info exclusion logic ==="
rg -n --hidden --no-ignore-vcs "egg-info|egg_info|\\*\\.egg-info|path\\.parts" "$FILE" || trueRepository: OmniNode-ai/omnibase_infra
Length of output: 6020
🏁 Script executed:
#!/usr/bin/env bash
set -euo pipefail
python - <<'PY'
import importlib.metadata as md
pkgs = ["omnibase-compat", "onex-change-control", "omnimarket"]
for name in pkgs:
try:
dist = md.distribution(name)
except md.PackageNotFoundError:
print(f"{name}: not installed")
continue
p = dist.locate_file("direct_url.json")
exists = p.exists()
txt = dist.read_text("direct_url.json")
print(f"{name}: locate_file={p} exists={exists} read_text_present={bool(txt)}")
PYRepository: OmniNode-ai/omnibase_infra
Length of output: 162
🏁 Script executed:
#!/usr/bin/env bash
set -euo pipefail
python - <<'PY'
import importlib.metadata as md
mismatches = []
with_txt = []
for dist in md.distributions():
name = dist.metadata.get("Name", "<unknown>")
txt = dist.read_text("direct_url.json")
if txt is None:
continue
with_txt.append(name)
p = dist.locate_file("direct_url.json")
exists = p.exists()
if not exists:
mismatches.append((name, str(p), exists, len(txt)))
if len(with_txt) >= 100: # cap to keep output small
break
print(f"distributions_with_direct_url_json_sampled={len(with_txt)}")
print(f"mismatches_sampled={len(mismatches)}")
for name, p, exists, n in mismatches[:20]:
dist = md.distribution(name)
files = dist.files or set()
has_in_files = any(str(f) == "direct_url.json" for f in files)
print(f"- {name}: locate_file={p} exists={exists} read_text_len={n} direct_url_in_dist_files={has_in_files}")
PYRepository: OmniNode-ai/omnibase_infra
Length of output: 578
Fix direct_url.json lookup in scripts/runtime_build/compute_workspace_provenance.py to avoid false missing_install
_installed_direct_urlusesdist.locate_file("direct_url.json")+ anexists()guard; in practice,locate_file()can point to a non-existent path even whendist.read_text("direct_url.json")returns content, so the current code can incorrectly returnNoneand mark valid local installs as missing (lines 66-69). Readdirect_url.jsonviadist.read_text("direct_url.json")and parse when non-empty, dropping thelocate_file/exists()check._hash_treeexclusion tries to skip"*.egg-info"viapart in path.parts(literal match), which won’t exclude real segments likepkg.egg-info—usepart.endswith(".egg-info")(lines 42-53).
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@scripts/runtime_build/compute_workspace_provenance.py` around lines 66 - 69,
The _installed_direct_url function currently uses
dist.locate_file("direct_url.json") plus an exists() guard which can falsely
return None; change it to call dist.read_text("direct_url.json"), check that the
returned string is non-empty, and json.loads that content (remove the
locate_file/exists check and variable direct_url_file); additionally update
_hash_tree's exclusion logic which checks for literal `"*.egg-info"` in
path.parts—replace that check with a per-part test using
part.endswith(".egg-info") so entries like "pkg.egg-info" are correctly skipped
(refer to function names _installed_direct_url and _hash_tree and to
dist.read_text and path.parts in the diff).
| # Handle machine-class overlay seeding (OMN-8904 — Wave B) | ||
| if args.machine_class is not None: | ||
| return _seed_machine_class( | ||
| args.machine_class, | ||
| create_missing=args.create_missing_keys, | ||
| overwrite_existing=args.overwrite_existing, | ||
| dry_run=args.dry_run and not args.execute, | ||
| ) |
There was a problem hiding this comment.
Load ~/.omnibase/.env before invoking machine-class Infisical seeding.
The new --machine-class path calls Infisical write flow without explicitly loading operator env, so behavior depends on caller shell state.
As per coding guidelines **/*.py: “Source ~/.omnibase/.env before any database, Kafka, or Infisical operation”.
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@scripts/seed-infisical.py` around lines 762 - 769, The machine-class branch
returns early to call _seed_machine_class without loading the operator env, so
ensure ~/.omnibase/.env is loaded first; modify the block that checks
args.machine_class to load the operator environment (e.g., via existing helper
or python-dotenv) before invoking _seed_machine_class so
Infisical/database/Kafka credentials are available when args.machine_class is
not None.
| # Format: <repo-relative-path>:<lineno> | ||
| # DO NOT add new entries here — fix violations instead. | ||
| # Total entries: 121 | ||
| # Total entries: 123 |
There was a problem hiding this comment.
Avoid growing the suppression baseline; fix new literals instead.
This change adds new baseline entries despite the file policy forbidding new suppressions. Please remove the added entries and address the underlying topic literals in service_kernel.py.
Also applies to: 123-124
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@scripts/validation/topic_literal_baseline.txt` at line 5, Remove the newly
added suppression lines from topic_literal_baseline.txt and instead fix the
hard-coded topic literals in service_kernel.py: search for raw topic string
literals (e.g. any usage like "orders.created", "user.updated" or similar)
inside functions such as publish_to_topic, subscribe_topic, handle_message or
register_handlers and replace them with the canonical constants/config lookup
(e.g. TOPIC_* constants or a get_topic_name(...) helper) used elsewhere in the
codebase; ensure you import or define the appropriate TOPIC_* symbols and update
any unit tests or callers to use the constants so no new baseline entries are
required.
| model_config = ConfigDict(frozen=True) | ||
|
|
There was a problem hiding this comment.
Align model_config with repository Pydantic model policy.
ConfigDict(frozen=True) is missing extra="forbid" and from_attributes=True.
Suggested fix
- model_config = ConfigDict(frozen=True)
+ model_config = ConfigDict(
+ frozen=True,
+ extra="forbid",
+ from_attributes=True,
+ )As per coding guidelines: "Pydantic models must use ConfigDict(frozen=True, extra="forbid", from_attributes=True) for immutability, strict validation, and ORM/pytest-xdist compatibility".
📝 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.
| model_config = ConfigDict(frozen=True) | |
| model_config = ConfigDict( | |
| frozen=True, | |
| extra="forbid", | |
| from_attributes=True, | |
| ) | |
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@src/omnibase_infra/runtime/models/model_plugin_discovery_report.py` around
lines 43 - 44, The ConfigDict for the Pydantic model is incomplete: update the
model_config variable (in model_plugin_discovery_report.py) to use
ConfigDict(frozen=True, extra="forbid", from_attributes=True) so the Pydantic
model is immutable, forbids extra fields, and supports from-attributes/ORM
usage; locate the model_config assignment and replace the current
ConfigDict(...) call accordingly.
| _BRIDGE_PORT_KEY = "DirectBridgeDelegationDispatchPort" | ||
|
|
||
| def __init__( | ||
| self, | ||
| *, | ||
| container: object, | ||
| fallback_event_bus: ProtocolPatternBBrokerTransport | None = None, | ||
| ) -> None: | ||
| self._container = container | ||
| self._fallback_event_bus = fallback_event_bus | ||
| self._resolved: ProtocolDelegationDispatchPort | None = None | ||
|
|
||
| async def _resolve(self) -> ProtocolDelegationDispatchPort: | ||
| if self._resolved is not None: | ||
| return self._resolved | ||
| registry = getattr(self._container, "service_registry", None) | ||
| if registry is not None: | ||
| try: | ||
| port: ProtocolDelegationDispatchPort = await registry.resolve_service( | ||
| self._BRIDGE_PORT_KEY | ||
| ) |
There was a problem hiding this comment.
🛠️ Refactor suggestion | 🟠 Major | ⚡ Quick win
Resolve delegation port by protocol contract, not concrete class key.
The bridge lookup uses "DirectBridgeDelegationDispatchPort" directly. Please switch to protocol-based DI tokening to avoid brittle concrete-name coupling in container resolution.
As per coding guidelines: "Use protocol names for dependency injection: container.get_service('ProtocolEventBus') instead of concrete class names".
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@src/omnibase_infra/runtime/service_delegation_dispatch_port.py` around lines
228 - 248, The code currently resolves the bridge using the hard-coded string
_BRIDGE_PORT_KEY = "DirectBridgeDelegationDispatchPort"; update _resolve to use
the protocol token instead of a concrete name: stop passing that literal and
call registry.resolve_service with the ProtocolDelegationDispatchPort token (or
ProtocolDelegationDispatchPort.__name__ if your DI expects a string).
Concretely, change references in _resolve (and the _BRIDGE_PORT_KEY constant) so
resolve_service receives ProtocolDelegationDispatchPort (or its name) rather
than "DirectBridgeDelegationDispatchPort", keeping the rest of the await
registry.resolve_service(...) logic and the type annotation for port as
ProtocolDelegationDispatchPort.
| return await port.dispatch( | ||
| prompt=prompt, | ||
| task_type=task_type, | ||
| correlation_id=correlation_id, | ||
| max_tokens=max_tokens, | ||
| source_file_path=source_file_path, | ||
| source_session_id=source_session_id, | ||
| wait=wait, | ||
| quality_contract_mode=quality_contract_mode, | ||
| acceptance_criteria=acceptance_criteria, | ||
| ) |
There was a problem hiding this comment.
Forward output_schema_key in delegated dispatch call.
ContainerBackedDelegationDispatchPort.dispatch() accepts output_schema_key but drops it when calling port.dispatch(...), so callers lose this input.
Proposed fix
return await port.dispatch(
prompt=prompt,
task_type=task_type,
correlation_id=correlation_id,
max_tokens=max_tokens,
source_file_path=source_file_path,
source_session_id=source_session_id,
wait=wait,
+ output_schema_key=output_schema_key,
quality_contract_mode=quality_contract_mode,
acceptance_criteria=acceptance_criteria,
)📝 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.
| return await port.dispatch( | |
| prompt=prompt, | |
| task_type=task_type, | |
| correlation_id=correlation_id, | |
| max_tokens=max_tokens, | |
| source_file_path=source_file_path, | |
| source_session_id=source_session_id, | |
| wait=wait, | |
| quality_contract_mode=quality_contract_mode, | |
| acceptance_criteria=acceptance_criteria, | |
| ) | |
| return await port.dispatch( | |
| prompt=prompt, | |
| task_type=task_type, | |
| correlation_id=correlation_id, | |
| max_tokens=max_tokens, | |
| source_file_path=source_file_path, | |
| source_session_id=source_session_id, | |
| wait=wait, | |
| output_schema_key=output_schema_key, | |
| quality_contract_mode=quality_contract_mode, | |
| acceptance_criteria=acceptance_criteria, | |
| ) |
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@src/omnibase_infra/runtime/service_delegation_dispatch_port.py` around lines
279 - 289, ContainerBackedDelegationDispatchPort.dispatch currently accepts an
output_schema_key parameter but fails to forward it to the delegated call;
update the return await port.dispatch(...) call inside
ContainerBackedDelegationDispatchPort.dispatch to include
output_schema_key=output_schema_key so the delegated port receives the value
(ensure the parameter name matches the signature expected by port.dispatch).
| # node_delegate_skill_orchestrator publishes completed/failed to omnimarket topics. | ||
| # Without this applier the handler result is silently discarded and the CLI adapter | ||
| # times out waiting for onex.evt.omnimarket.delegate-skill-completed.v1 (OMN-11996). | ||
| auto_wiring_result_appliers["node_delegate_skill_orchestrator"] = ( | ||
| DispatchResultApplier( | ||
| event_bus=event_bus, | ||
| output_topic=EnumOmnimarketTopic.EVT_DELEGATE_SKILL_COMPLETED_V1.value, | ||
| allowed_output_topics=[ | ||
| EnumOmnimarketTopic.EVT_DELEGATE_SKILL_COMPLETED_V1.value, | ||
| EnumOmnimarketTopic.EVT_DELEGATE_SKILL_FAILED_V1.value, | ||
| ], | ||
| ) | ||
| ) | ||
| logger.info( | ||
| "Delegate-skill orchestrator terminal result applier registered " | ||
| "(contract=node_delegate_skill_orchestrator, correlation_id=%s)", | ||
| correlation_id, | ||
| ) |
There was a problem hiding this comment.
Avoid hardcoding delegate-skill terminal topics in kernel wiring.
Line 2018–Line 2022 pins terminal topics in Python, which can drift from contract routing. Please resolve these topics from the node contract metadata and feed them into DispatchResultApplier from contract-driven config instead of embedding them here.
As per coding guidelines: "Kafka topics must be declared in node contract YAML under event_bus.publish_topics or subscribe_topics, never hardcoded in Python modules".
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@src/omnibase_infra/runtime/service_kernel.py` around lines 2012 - 2029, The
wiring currently hardcodes Omnimarket terminal topics for
node_delegate_skill_orchestrator using EnumOmnimarketTopic when constructing
DispatchResultApplier; instead, load the terminal publish topics from the node
contract metadata (the contract's event_bus.publish_topics / subscribe_topics
entry for the node_delegate_skill_orchestrator contract) and pass those values
into DispatchResultApplier as output_topic and allowed_output_topics; update the
auto_wiring_result_appliers population to resolve the contract
(node_delegate_skill_orchestrator) metadata at runtime, extract the configured
topic names, and use them rather than EnumOmnimarketTopic so contract-driven
YAML controls the topics.
| async def fake_start_unlocked(topic: str, group_id: str) -> None: | ||
| call_log.append((topic, group_id)) | ||
| # Yield control to simulate concurrent group-join latency | ||
| await asyncio.sleep(0) | ||
|
|
||
| with patch( | ||
| "omnibase_infra.event_bus.event_bus_kafka.AIOKafkaProducer", | ||
| return_value=mock_producer, | ||
| ): | ||
| config = ModelKafkaEventBusConfig(bootstrap_servers=TEST_BOOTSTRAP_SERVERS) | ||
| event_bus = EventBusKafka(config=config) | ||
| await event_bus.start() | ||
|
|
||
| # Pre-register three distinct subscriptions without triggering | ||
| # consumer startup (bus is already started but we bypass it here by | ||
| # directly injecting into _subscribers so we can test start_consuming | ||
| # in isolation). | ||
| group_id_a = "svc-a" | ||
| group_id_b = "svc-b" | ||
|
|
||
| async with event_bus._lock: | ||
| event_bus._subscribers["topic-x"] = [ | ||
| (group_id_a, "sub-1", AsyncMock()), | ||
| (group_id_b, "sub-2", AsyncMock()), | ||
| ] | ||
| event_bus._subscribers["topic-y"] = [ | ||
| (group_id_a, "sub-3", AsyncMock()), | ||
| ] | ||
|
|
||
| event_bus._start_consumer_for_topic_unlocked = fake_start_unlocked # type: ignore[method-assign] | ||
|
|
||
| # Run start_consuming briefly then shut down | ||
| task = asyncio.create_task(event_bus.start_consuming()) | ||
| await asyncio.sleep(0.1) | ||
| await event_bus.shutdown() | ||
| await asyncio.wait_for(task, timeout=2.0) | ||
|
|
||
| # Each distinct (topic, group_id) pair must be started exactly once | ||
| assert len(call_log) == 3 | ||
| assert len(set(call_log)) == 3, "duplicate consumer starts detected" | ||
|
|
There was a problem hiding this comment.
Prevent timing-dependent duplicate starts in the concurrency test.
Line 2108 relies on fixed timing, and fake_start_unlocked never marks consumers as started. If start_consuming() loops again before shutdown, duplicates can be appended and the test becomes flaky.
Suggested patch
call_log: list[tuple[str, str]] = []
async def fake_start_unlocked(topic: str, group_id: str) -> None:
call_log.append((topic, group_id))
+ event_bus._group_consumers[(topic, group_id)] = AsyncMock()
# Yield control to simulate concurrent group-join latency
await asyncio.sleep(0)
@@
- task = asyncio.create_task(event_bus.start_consuming())
- await asyncio.sleep(0.1)
+ task = asyncio.create_task(event_bus.start_consuming())
+
+ async def _wait_for_expected_calls() -> None:
+ while len(call_log) < 3:
+ await asyncio.sleep(0)
+
+ await asyncio.wait_for(_wait_for_expected_calls(), timeout=1.0)
await event_bus.shutdown()
await asyncio.wait_for(task, timeout=2.0)🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@tests/unit/event_bus/test_kafka_event_bus.py` around lines 2075 - 2115, The
test's fake_start_unlocked does not mark consumers as started, so
EventBusKafka.start_consuming can schedule the same (topic, group_id) again and
make the test flaky; modify fake_start_unlocked (the replacement for
EventBusKafka._start_consumer_for_topic_unlocked used in the test) to record a
started marker into the EventBusKafka instance (e.g., add the (topic, group_id)
to an in-instance set like event_bus._started_consumers or to the
event_bus._consumer_tasks mapping) while holding the same lock used by
start_consuming so subsequent iterations see the consumer as started and do not
append duplicates to call_log; keep the artificial await asyncio.sleep(0) to
preserve concurrency simulation.
61e712f to
6eadabf
Compare
6eadabf to
8345029
Compare
Summary
Verification
uv run pytest tests/integration/test_projection_handler_wiring_runtime_dispatch.py tests/unit/runtime/auto_wiring/test_handler_wiring_db_injection.py -q-> 21 passeduv run ruff check src/omnibase_infra/runtime/auto_wiring/handler_wiring.py tests/unit/runtime/auto_wiring/test_handler_wiring_db_injection.py tests/integration/test_projection_handler_wiring_runtime_dispatch.py-> passedEvidence-Ticket: OMN-12245
Evidence-Source: OCC#1790
hotfix-evidence: OCC-1790
backmerge-for: #1775