Repository navigation
fix(OMN-13149): savings consumer decodes typed ModelEventMessage instead of calling .get() - #1984
Conversation
…ead of calling .get()
The savings-estimation consumer is wired in service_kernel via
event_bus.subscribe; both the Kafka and in-memory buses deliver a typed
ModelEventMessage to the on_message callback. The pre-fix _savings_on_message
json.loads-ed only str/bytes and otherwise passed the value through unchanged,
so a ModelEventMessage reached ServiceSavingsEstimator.ingest_event ->
_session_id_for_topic, which calls payload.get('session_id'). ModelEventMessage
has no .get(), raising AttributeError at runtime — the consumer wrote zero rows.
Fix: add module-level decode_event_message(message) that reads the JSON payload
off the typed .value field directly (strongly typed, fail-fast: TypeError on a
non-object payload — no getattr-default, no dict fallback). The kernel callback
decodes then calls ingest_event(topic, payload). Removed the now-redundant local
ModelEventMessage import in the runtime-error-triage consumer.
Regression tests through the REAL consumer path (EventBusInmemory subscribe/
publish -> on_message(ModelEventMessage) -> decode_event_message -> ingest_event):
typed message correlates a session; pre-fix path reproduces AttributeError.
Same dict/typed-model family as OMN-13141 / golden-chain BUG A.
|
Warning Review limit reached
More reviews will be available in 2 hours, 5 minutes, and 32 seconds. Learn how PR review limits work. Your organization has used up its prepaid credits, and credit purchases are no longer available. Enable the review add-on in the billing tab to keep reviews running — you're only billed for reviews past your plan's rate limits ($0.25/file). ⌛ 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 (4)
✨ Finishing Touches🧪 Generate unit tests (beta)
Comment |
OMN-13149 — savings_estimation consumer AttributeError (ModelEventMessage has no .get())
Problem
The savings-estimation consumer is wired in
service_kernel.pyviaevent_bus.subscribe. Both the Kafka bus (EventBusKafka) and the in-memory bus (EventBusInmemory) deliver a typedModelEventMessageto theon_messagecallback — never a rawdictorstr(callback signature isCallable[[ModelEventMessage], Awaitable[None]]).The pre-fix
_savings_on_messagedid:A
ModelEventMessageis neitherstrnorbytes, so it fell into theelse msgbranch and was passed straight toServiceSavingsEstimator.ingest_event(topic, payload). The correlation path_session_id_for_topicthen callspayload.get("session_id").ModelEventMessageis a frozen Pydantic model with no.get(), so this raised at runtime:The subscriber callback crashed on every consumed event → the consumer correlated nothing and wrote zero rows. Same dict/typed-model family as OMN-13141 / golden-chain BUG A.
The repo's own runtime-error-triage consumer in the same kernel already showed the correct pattern: it reads
message.valueandjson.loadsit. The savings callback never decoded.value.Fix
decode_event_message(message: ModelEventMessage) -> tuple[str, dict[str, object]]inconsumer.pyreads the JSON payload off the typed.valuefield directly — strongly typed, fail-fast: raisesTypeErrorif the decoded payload is not a JSON object. Nogetattr-default, no dict-style fallback, no silent guard.service_kernel._savings_on_messagenow takesmessage: ModelEventMessage, callsdecode_event_message(message), theningest_event(topic, payload).ModelEventMessageimport in the runtime-error-triage consumer (module-level import added).decode_event_messageis a module-level function (not a class method) soServiceSavingsEstimatorstays at 10 methods — keeping the ONEX Pattern Validation gate green rather than tripping the per-class method-count threshold.Regression tests (REAL consumer path, not handler isolation)
tests/unit/services/observability/savings_estimation/test_consumer.py:test_ingest_event_rejects_typed_message_documents_bug— pins the pre-fix crash: a typedModelEventMessagestraight intoingest_eventraisesAttributeError: ...has no attribute 'get'.test_decode_event_message_correlates_typed_message— decode +ingest_eventcorrelates a session.test_decode_event_message_rejects_non_object_payload— non-object JSON fails fast (TypeError).test_real_bus_delivers_typed_message_to_consumer— end-to-end through the realEventBusInmemory(subscribe/publish→on_message(ModelEventMessage)→decode_event_message→ingest_event), mirroring the kernel callback. Reproduced failing-first: with the pre-fix decode the callback raises the exactAttributeErrorinsideEventBusInmemory.publish; with the fix the session is correlated.These run under
@pytest.mark.unit(no infra markers) → wired as a pre-merge CI gate.Verification (local, in worktree)
uv run ruff format src/ tests/ && uv run ruff check --fix src/ tests/— cleanuv run mypy src/omnibase_infra/services/observability/savings_estimation/consumer.py src/omnibase_infra/runtime/service_kernel.py --strict— Success, no issuesuv run pytest tests/unit/services/observability/savings_estimation/test_consumer.py -q— 16 passeduv run pytest tests/unit/runtime/ tests/unit/services/observability/savings_estimation/ -q— 4902 passed, 5 skipped (pre-existing)uv run pytest tests/unit/ -m "not slow and not kafka and not postgres and not consul"— 20472 passed, 1 failure (tests/unit/verification/test_cli.py::test_registration_only, pre-existing — fails identically on cleanorigin/dev, unrelated to this change)pre-commit runon changed files — all hooks pass (incl. ONEX Pattern Validation)dod_evidence
Per
contracts/OMN-13149.yaml(this PR):dod-contract-artifact,dod-typed-message-regression,dod-runtime-suite,dod-deploy.OCC pairing
Paired OCC receipt PR: OmniNode-ai/onex_change_control#2629 (contract
contracts/OMN-13149.yaml+ PASS receipts underdrift/dod_receipts/OMN-13149/; verifierreceipt-gate-local≠ runner; honesty gate 0 violations).Receipt Gate / DoD Gate evidence for OMN-13149:
Evidence-Source: OCC#2629
Evidence-Ticket: OMN-13149
Note on repo
The OMN-13149 ticket text guessed
omnimarketfor the consumer. The actual savings_estimation consumer lives in omnibase_infra (src/omnibase_infra/services/observability/savings_estimation/consumer.py), packaged into omnimarket's venv as a dependency. The fix and the runtime wiring (service_kernel.py) are both in omnibase_infra.Closes OMN-13149. Parent: OMN-13119.