Skip to content

feat(OMN-1932): Wire snapshot publisher into registration handlers - #254

Merged
jonahgabriel merged 11 commits into
mainfrom
jonah/omn-1932-phase-3-wire-snapshot-publisher-p1
Feb 7, 2026
Merged

jonahgabriel merged 11 commits into
mainfrom
jonah/omn-1932-phase-3-wire-snapshot-publisher-p1

Conversation

@jonahgabriel

@jonahgabriel jonahgabriel commented Feb 6, 2026 •

Copy link
Copy Markdown
Collaborator

Summary

  • Wire SnapshotPublisherRegistration into the registration orchestrator pipeline so handler state transitions produce compacted Kafka snapshots for read optimization
  • Add schema_version field (semver string) to ModelRegistrationSnapshot and simplify Kafka key to entity_id UUID only
  • Add 500ms trailing-edge debounce per entity, kernel bootstrap wiring, and handler integration for publish-on-transition and tombstone-on-expiry
  • Add integration test verifying end-to-end state transition → snapshot publishing with mock Kafka producer

Changes by Phase

P3.1 — Snapshot Schema + Key Format

  • ModelRegistrationSnapshot: Added schema_version: str = "1.0.0" with semver validation via ModelSemVer.parse()
  • to_kafka_key() simplified from "{domain}:{entity_id}" to str(entity_id) (UUID only)
  • ModelSnapshotTopicConfig.get_snapshot_key() updated to match

P3.2 — Ops Documentation

  • docs/ops/snapshot_topic_configuration.md: rpk commands, dev/prod config, consumer debugging

P3.3 — Kernel Bootstrap Wiring

  • service_kernel.py: Creates AIOKafkaProducer + SnapshotPublisherRegistration, passes to handler wiring, manages lifecycle
  • util_container_wiring.py and wiring.py: Accept and forward snapshot_publisher parameter
  • All best-effort: if Kafka unavailable, logs warning and continues

P3.4 — Debounce (500ms)

  • Trailing-edge debounce per entity_id using asyncio.call_later()
  • stop() flushes pending publishes; tombstones bypass debounce
  • debounce_ms=0 disables (used in test fixtures)

P3.5 — Handler Integration

Handler Trigger Action
HandlerNodeIntrospected After projection persist publish_from_projection()
HandlerNodeRegistrationAcked After ACTIVE transition publish_from_projection()
HandlerRuntimeTick On LIVENESS_EXPIRED delete_snapshot() (tombstone)

P3.6 — Integration Test

  • tests/integration/projectors/test_snapshot_state_transitions.py
  • 4 tests: introspection→publish, ack→publish, expiry→tombstone, debounce coalescing

Test plan

  • All 4 new integration tests pass (test_snapshot_state_transitions.py)
  • All 6 new debounce unit tests pass
  • Full test suite: 12,705 passed, 0 regressions (9 pre-existing failures from external service deps)
  • Pre-commit hooks pass (ruff, mypy, architecture validation)
  • CI green

Summary by CodeRabbit

  • New Features

    • Optional snapshot publishing to Kafka with per-entity debounce (coalescing), best-effort non-blocking sends, and tombstone publishing on liveness expiry.
    • Snapshot messages now use the entity UUID as the message key and include a semver-validated schema_version field.
    • Runtime wiring and banner report Snapshot Publisher status; publisher started/stopped when Kafka is configured.
  • Documentation

    • Added a comprehensive snapshot topic configuration and operations guide.
  • Tests

    • Added/updated integration and unit tests for publishing, tombstones, debounce behavior, and key/schema semantics.

…hase 3)

Wire SnapshotPublisherRegistration into the registration orchestrator
pipeline so state transitions produce compacted Kafka snapshots for
read optimization.

P3.1: Add schema_version (semver string) to ModelRegistrationSnapshot,
      simplify Kafka key to entity_id UUID only.
P3.2: Add ops docs for compacted topic configuration.
P3.3: Wire SnapshotPublisherRegistration in kernel bootstrap with
      best-effort Kafka producer lifecycle.
P3.4: Add 500ms trailing-edge debounce per entity_id with
      flush-on-stop and tombstone bypass.
P3.5: Handlers publish snapshots on introspection/ack transitions
      and tombstones on LIVENESS_EXPIRED.
P3.6: Integration test verifying state transitions produce snapshots
      end-to-end with mock Kafka producer.
@coderabbitai

coderabbitai Bot commented Feb 6, 2026 •

Copy link
Copy Markdown
📝 Walkthrough

Walkthrough

Adds snapshot publishing wired through runtime and handlers, changes Kafka snapshot keying to entity_id-only, adds schema_version semver validation on snapshots, implements per-entity debounce/coalescing in SnapshotPublisherRegistration, updates wiring/bootstrap lifecycle, documentation, and tests for keying and debounce behavior.

Changes

Cohort / File(s) Summary
Docs
docs/ops/snapshot_topic_configuration.md
New operational guide for snapshot topic configuration, creation, verification, overrides, compaction semantics, and commands.
Models: snapshot schema & topic config
src/omnibase_infra/models/projection/model_registration_snapshot.py, src/omnibase_infra/models/projection/model_snapshot_topic_config.py
Add schema_version: str with semver validation and change snapshot key semantics to return raw entity_id (UUID string). ModelSnapshotTopicConfig.get_snapshot_key simplified to accept entity_id only.
Publisher: debounce, keying & lifecycle
src/omnibase_infra/projectors/snapshot_publisher_registration.py, src/omnibase_infra/protocols/protocol_snapshot_publisher.py
Switch compaction/keying to entity_id-only; add per-entity debounce (debounce_ms, timers, pending snapshots, lock); add publish_from_projection API (correlation_id); flush pending publishes on stop; tombstone and cancel semantics updated.
Handlers: best-effort publish integration
src/omnibase_infra/nodes/node_registration_orchestrator/handlers/handler_node_introspected.py, .../handler_node_registration_acked.py, .../handler_runtime_tick.py
Handlers accept optional snapshot_publisher and perform non-blocking, best-effort snapshot publishes (and tombstones) after projection persistence or on liveness expiry; errors are logged and do not block flow.
Wiring & runtime bootstrap
src/omnibase_infra/nodes/node_registration_orchestrator/wiring.py, src/omnibase_infra/runtime/service_kernel.py, src/omnibase_infra/runtime/util_container_wiring.py
Thread snapshot_publisher through wiring; attempt to create/start publisher when Kafka configured at bootstrap; stop publisher on shutdown; new parameter propagated through wiring functions.
Tests: unit & integration
tests/unit/.../test_snapshot_*, tests/integration/projectors/test_snapshot_state_transitions.py, tests/replay/test_snapshot_plus_tail.py, tests/integration/handlers/test_handler_no_publish_constraint.py
New integration tests for publish flows and debounce; unit tests updated for entity_id-only keys and debounce_ms constructor param; handler no-publish tests switched to type-based bus-detection logic.
Other (protocols, typing)
src/omnibase_infra/protocols/protocol_snapshot_publisher.py
Protocol updated to entity_id-only keys, get_latest_snapshot now returns ModelRegistrationSnapshot

Sequence Diagram(s)

sequenceDiagram
  participant Handler
  participant Publisher as "SnapshotPublisherRegistration"
  participant Debounce as "Debounce Scheduler"
  participant Kafka as "Kafka Producer"
  participant Cache

  Handler->>Publisher: publish_from_projection(projection, correlation_id?)
  alt debounce_ms > 0
    Publisher->>Debounce: schedule(entity_id, latest_snapshot)
    Debounce-->>Publisher: trigger execute_debounced_publish(entity_id)
    Publisher->>Kafka: send(key=entity_id, value=snapshot)
  else immediate
    Publisher->>Kafka: send(key=entity_id, value=snapshot)
  end
  Kafka-->>Publisher: ack
  Publisher->>Cache: update latest snapshot/version for entity_id
Loading

Estimated code review effort

🎯 4 (Complex) | ⏱️ ~45 minutes

Possibly related PRs

"a rabbit wrote this note in the log:
I hop and I coalesce, latest wins the race 🐇
Keys now plain UUIDs in flight,
Debounce keeps quick flutters light,
Docs, wiring, tests in tidy rows,
A tiny hop for cleaner flows."

🚥 Pre-merge checks | ✅ 3
✅ Passed checks (3 passed)
Check name Status Explanation
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Title check ✅ Passed The PR title 'feat(OMN-1932): Wire snapshot publisher into registration handlers' clearly summarizes the main change: wiring SnapshotPublisherRegistration into the registration handlers for snapshot publishing on state transitions.
Docstring Coverage ✅ Passed Docstring coverage is 100.00% which is sufficient. The required threshold is 80.00%.

✏️ Tip: You can configure your own custom pre-merge checks in the settings.

✨ Finishing touches
  • 📝 Generate docstrings
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Post copyable unit tests in a comment
  • Commit unit tests in branch jonah/omn-1932-phase-3-wire-snapshot-publisher-p1

Comment @coderabbitai help to get the list of available commands and usage tips.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 9

Caution

Some comments are outside the diff and can’t be posted inline due to platform limitations.

⚠️ Outside diff range comments (1)
src/omnibase_infra/projectors/snapshot_publisher_registration.py (1)

195-279: ⚠️ Potential issue | 🟠 Major

Propagate correlation_id through debounced publishes.
The debounce queue currently drops the caller correlation_id, so errors/logs can’t be traced back to the triggering request. Carry it from the projection (or generate one if missing) into the pending queue and into _publish_snapshot_model.

🛠️ Proposed fix (propagate correlation_id end-to-end)
@@
-        self._pending_snapshots: dict[str, ModelRegistrationSnapshot] = {}
+        self._pending_snapshots: dict[str, tuple[ModelRegistrationSnapshot, UUID]] = {}
@@
-    async def _publish_snapshot_model(
-        self,
-        snapshot: ModelRegistrationSnapshot,
-    ) -> None:
+    async def _publish_snapshot_model(
+        self,
+        snapshot: ModelRegistrationSnapshot,
+        *,
+        correlation_id: UUID | None = None,
+    ) -> None:
@@
-        correlation_id = uuid4()
+        correlation_id = correlation_id or uuid4()
@@
-        version = await self._get_next_version(entity_id_str, projection.domain)
+        version = await self._get_next_version(entity_id_str, projection.domain)
+        correlation_id = projection.correlation_id or uuid4()
@@
-        if self._debounce_ms > 0:
-            await self._schedule_debounced_publish(entity_id_str, snapshot)
+        if self._debounce_ms > 0:
+            await self._schedule_debounced_publish(
+                entity_id_str,
+                snapshot,
+                correlation_id,
+            )
         else:
-            await self._publish_snapshot_model(snapshot)
+            await self._publish_snapshot_model(snapshot, correlation_id=correlation_id)
@@
-    async def _schedule_debounced_publish(
-        self,
-        entity_id: str,
-        snapshot: ModelRegistrationSnapshot,
-    ) -> None:
+    async def _schedule_debounced_publish(
+        self,
+        entity_id: str,
+        snapshot: ModelRegistrationSnapshot,
+        correlation_id: UUID,
+    ) -> None:
@@
-            self._pending_snapshots[entity_id] = snapshot
+            self._pending_snapshots[entity_id] = (snapshot, correlation_id)
@@
-            snapshot = self._pending_snapshots.pop(entity_id, None)
+            pending = self._pending_snapshots.pop(entity_id, None)
             self._debounce_timers.pop(entity_id, None)
 
-        if snapshot is not None:
+        if pending is not None:
+            snapshot, correlation_id = pending
             try:
-                await self._publish_snapshot_model(snapshot)
+                await self._publish_snapshot_model(
+                    snapshot,
+                    correlation_id=correlation_id,
+                )
@@
-        for entity_id, snapshot in pending.items():
+        for entity_id, pending_item in pending.items():
             try:
-                await self._publish_snapshot_model(snapshot)
+                snapshot, correlation_id = pending_item
+                await self._publish_snapshot_model(
+                    snapshot,
+                    correlation_id=correlation_id,
+                )

As per coding guidelines, Propagate correlation_id from incoming requests; auto-generate with uuid4() if missing; include in all error context.

Also applies to: 543-623, 1281-1445

🤖 Fix all issues with AI agents
In `@docs/ops/snapshot_topic_configuration.md`:
- Around line 39-41: Two fenced code blocks in snapshot_topic_configuration.md
are missing language identifiers (markdownlint MD040): add a neutral language
tag like "text" to the fenced blocks that contain the UUID line
"550e8400-e29b-41d4-a716-446655440000" and the ASCII table starting with "|---
min_compaction_lag_ms (60s) ---|--- max_compaction_lag_ms (300s) ---|". Edit
those triple-backtick fences (the ones wrapping the UUID and the table) to be
"```text" so linting passes.

In `@src/omnibase_infra/models/projection/model_registration_snapshot.py`:
- Around line 261-282: Update the ModelRegistrationSnapshot documentation and
error messaging to reflect that to_kafka_key() now returns only the entity_id
(UUID string) for compaction instead of a domain:entity_id pair, and change
is_newer_than() so its mismatch error message includes both snapshot domains
(e.g., f"domain mismatch: {self.domain} != {other.domain}") to avoid ambiguity;
edit the class docstring to remove or replace references to
"{domain}:{entity_id}" and clarify that Kafka keys are entity_id-only, and
modify the is_newer_than() exception text to print the domains and entity_ids
involved for clear diagnostics.

In `@src/omnibase_infra/models/projection/model_snapshot_topic_config.py`:
- Around line 568-586: Update the module and class-level docstrings that
describe snapshot key semantics to match the new behavior of
ModelSnapshotTopicConfig.get_snapshot_key (which now returns only the entity_id
UUID string); find references to key format in the ModelSnapshotTopicConfig
class and top-of-module docstring and change any `{domain}:{entity_id}` or
domain-prefixed examples to state and show that the Kafka compaction key is the
entity_id UUID only, and ensure any documentation examples and guidance for
operators mention the entity_id-only key format.

In
`@src/omnibase_infra/nodes/node_registration_orchestrator/handlers/handler_node_introspected.py`:
- Around line 461-497: The snapshot publish exception is logged raw (snap_err)
and may leak credentials; import and use the sanitizer from
omnibase_infra.utils.util_error_sanitization (e.g., sanitize_error or
sanitize_exception) and pass the sanitized string into logger.warning instead of
snap_err, e.g., call the sanitizer on snap_err and log that sanitized message
(keep existing context keys like "node_id"/"correlation_id" and use
type(sanitized) or type(snap_err).__name__ as needed); apply this change around
the _snapshot_publisher.publish_from_projection try/except and replace any
direct usage of snap_err in the log call with the sanitized output.

In
`@src/omnibase_infra/nodes/node_registration_orchestrator/handlers/handler_node_registration_acked.py`:
- Around line 315-320: The snapshot is being published from the local variable
projection while it still represents ACCEPTED/AWAITING_ACK, causing stale
snapshots; update the flow in the handler (the block that calls
_snapshot_publisher.publish_from_projection(projection, node_name=None)) to
publish only after the projection has been updated/persisted to ACTIVE (or
re-read the projection from the authoritative store) — e.g., move the publish
call to occur after the projection save/confirmation step or fetch the fresh
projection by its id/state before calling
_snapshot_publisher.publish_from_projection; keep the call non-blocking and
preserve the try/except/error logging around _snapshot_publisher to avoid
regressions.
- Around line 321-329: The logger call in the except block in
handler_node_registration_acked (the snap_err handler and logger.warning call)
must not log the raw exception; instead import and use the sanitization helpers
from omnibase_infra.utils.util_error_sanitization to produce a sanitized message
or safe fields and log only those (for example obtain a sanitized_message or
safe_error_fields from the util and pass that into logger.warning along with
node_id and correlation_id), replacing the direct snap_err usage so no sensitive
data (passwords/API keys/PII/connection strings) appears in logs.

In
`@src/omnibase_infra/nodes/node_registration_orchestrator/handlers/handler_runtime_tick.py`:
- Around line 375-391: The snapshot tombstone exception is logged raw in the
handler_runtime_tick handler; wrap/sanitize the error before logging by calling
sanitize_error_message(snap_err) (from
omnibase_infra.utils.util_error_sanitization) and use that sanitized string in
the logger.warning call for the _snapshot_publisher.delete_snapshot error path,
ensuring the log extra still includes node_id/correlation_id/error_type.

In `@src/omnibase_infra/runtime/service_kernel.py`:
- Around line 924-975: The snapshot publisher may be started
(SnapshotPublisherRegistration and its snapshot_producer) but then set to None
on a later error leaving background tasks running; update the error/cleanup
paths to explicitly stop/close the started publisher and underlying producer:
when catching downstream exceptions (or when setting snapshot_publisher = None
on partial failure), call and await snapshot_publisher.stop() (or shutdown/close
the snapshot_producer) and handle exceptions from that shutdown, and ensure the
global/final finally block checks the actual started
SnapshotPublisherRegistration instance (not just a None placeholder) and
properly awaits its stop/close; apply the same change to the other occurrences
referenced (the blocks around the other startup sites).

In `@tests/unit/models/projection/test_model_snapshot_topic_config.py`:
- Around line 196-215: Add a pytest marker to classify these tests as unit
tests: import pytest if not already imported in the test module and add a
module-level marker assignment pytestmark = [pytest.mark.unit]; alternatively,
apply `@pytest.mark.unit` to the Test class or the test functions (e.g.,
test_get_snapshot_key_format, test_get_snapshot_key_with_uuid,
test_get_snapshot_key_different_entities) in
tests/unit/models/projection/test_model_snapshot_topic_config.py so the test
module is properly marked as unit.
🧹 Nitpick comments (1)
src/omnibase_infra/runtime/util_container_wiring.py (1)

838-921: Use ProtocolSnapshotPublisher type for the snapshot_publisher parameter.

Change the parameter type from SnapshotPublisherRegistration to ProtocolSnapshotPublisher and import the protocol. This decouples the wiring function from the concrete implementation, adhering to the dependency injection guideline: use protocol names instead of concrete class names.

from omnibase_infra.protocols import ProtocolSnapshotPublisher

async def wire_registration_handlers(
    container: ModelONEXContainer,
    pool: asyncpg.Pool,
    liveness_interval_seconds: int | None = None,
    projector: ProjectorShell | None = None,
    consul_handler: HandlerConsul | None = None,
    snapshot_publisher: ProtocolSnapshotPublisher | None = None,
) -> WiringResult:

Comment thread docs/ops/snapshot_topic_configuration.md Outdated
Comment on lines 261 to +282
def to_kafka_key(self) -> str:
"""Generate Kafka compaction key for this snapshot.

Returns a key suitable for Kafka topic compaction. The key format
is "{domain}:{entity_id}" which ensures:
- Per-entity compaction (only latest snapshot retained per entity)
- Multi-domain support (entities in different domains are distinct)
Returns a key suitable for Kafka topic compaction. The key is the
node_id (entity_id) as a UUID string, which ensures:
- Per-node compaction (only latest snapshot retained per node)
- Simple cross-language consumer compatibility
- Aligns with "node_id is the partition key" principle

Returns:
Compaction key in format "domain:entity_id"
Node UUID as string

Example:
>>> from uuid import UUID
>>> snapshot = ModelRegistrationSnapshot(
... entity_id=UUID("550e8400-e29b-41d4-a716-446655440000"),
... domain="registration",
... ...
... )
>>> snapshot.to_kafka_key()
'registration:550e8400-e29b-41d4-a716-446655440000'
'550e8400-e29b-41d4-a716-446655440000'
"""
return f"{self.domain}:{self.entity_id!s}"
return str(self.entity_id)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🟡 Minor

Align snapshot key documentation (and error messaging) with entity_id-only compaction keys.

Now that to_kafka_key() drops the domain, the class docstring still references {domain}:{entity_id} keys and is_newer_than() error messages become ambiguous when domains differ. Consider updating the docstring and including domain in the mismatch error to avoid confusion.

🤖 Prompt for AI Agents
In `@src/omnibase_infra/models/projection/model_registration_snapshot.py` around
lines 261 - 282, Update the ModelRegistrationSnapshot documentation and error
messaging to reflect that to_kafka_key() now returns only the entity_id (UUID
string) for compaction instead of a domain:entity_id pair, and change
is_newer_than() so its mismatch error message includes both snapshot domains
(e.g., f"domain mismatch: {self.domain} != {other.domain}") to avoid ambiguity;
edit the class docstring to remove or replace references to
"{domain}:{entity_id}" and clarify that Kafka keys are entity_id-only, and
modify the is_newer_than() exception text to print the domains and entity_ids
involved for clear diagnostics.

Comment thread src/omnibase_infra/models/projection/model_snapshot_topic_config.py
Comment thread src/omnibase_infra/runtime/service_kernel.py
Comment on lines 196 to +215
def test_get_snapshot_key_format(self) -> None:
"""Test snapshot key follows domain:entity_id format."""
"""Test snapshot key returns entity_id directly."""
config = ModelSnapshotTopicConfig.default()
key = config.get_snapshot_key("registration", "node-123")
assert key == "registration:node-123"
key = config.get_snapshot_key("node-123")
assert key == "node-123"

def test_get_snapshot_key_with_uuid(self) -> None:
"""Test snapshot key with UUID entity_id."""
config = ModelSnapshotTopicConfig.default()
key = config.get_snapshot_key(
"registration", "550e8400-e29b-41d4-a716-446655440000"
)
assert key == "registration:550e8400-e29b-41d4-a716-446655440000"
key = config.get_snapshot_key("550e8400-e29b-41d4-a716-446655440000")
assert key == "550e8400-e29b-41d4-a716-446655440000"

def test_get_snapshot_key_different_domains(self) -> None:
"""Test snapshot keys for different domains."""
def test_get_snapshot_key_different_entities(self) -> None:
"""Test snapshot keys for different entities are distinct."""
config = ModelSnapshotTopicConfig.default()
reg_key = config.get_snapshot_key("registration", "node-1")
disc_key = config.get_snapshot_key("discovery", "node-1")
assert reg_key == "registration:node-1"
assert disc_key == "discovery:node-1"
assert reg_key != disc_key
key_1 = config.get_snapshot_key("node-1")
key_2 = config.get_snapshot_key("node-2")
assert key_1 == "node-1"
assert key_2 == "node-2"
assert key_1 != key_2

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🟡 Minor

Add a unit marker for this test module/class.

These tests are unmarked; per guidelines, they should carry a pytest classification marker. Consider adding a module-level pytestmark = [pytest.mark.unit] (or class-level decorators).
As per coding guidelines, Use pytest markers: @pytest.mark.unit, @pytest.mark.integration, @pytest.mark.slow, @pytest.mark.chaos, @pytest.mark.serial for test classification.

Proposed fix
@@
 import pytest
 from pydantic import ValidationError
 
 from omnibase_infra.errors import ProtocolConfigurationError
@@
 from omnibase_infra.topics import SUFFIX_REGISTRATION_SNAPSHOTS
+
+pytestmark = [pytest.mark.unit]
🤖 Prompt for AI Agents
In `@tests/unit/models/projection/test_model_snapshot_topic_config.py` around
lines 196 - 215, Add a pytest marker to classify these tests as unit tests:
import pytest if not already imported in the test module and add a module-level
marker assignment pytestmark = [pytest.mark.unit]; alternatively, apply
`@pytest.mark.unit` to the Test class or the test functions (e.g.,
test_get_snapshot_key_format, test_get_snapshot_key_with_uuid,
test_get_snapshot_key_different_entities) in
tests/unit/models/projection/test_model_snapshot_topic_config.py so the test
module is properly marked as unit.

…ion in no-publish constraint tests

The no-publish constraint test used substring matching ("_publisher" in
attr_name) to detect bus infrastructure, which false-positived on
_snapshot_publisher—a domain dependency, not an EventBus publisher.

Replace with isinstance(value, ProtocolEventBusLike) for precise type-based
detection. Add regression tests proving real bus types fail while snapshot
publishers pass.
- Fix MAJOR: publish post-transition ACTIVE state instead of pre-transition
  AWAITING_ACK in ack handler snapshot (model_copy with target state)
- Fix MAJOR: stop snapshot publisher before nulling on error path in
  service_kernel to prevent Kafka producer leak
- Fix MINOR: update stale domain:entity_id references to entity_id-only
  key format in model and publisher docstrings
- Fix MINOR: remove unused domain parameter from _get_next_version()
- Fix MINOR: assert actual ACTIVE state value in snapshot state test

Review iteration: 1/10
… publisher docs

- Fix MINOR: update module and class docstrings in model_snapshot_topic_config
  from {domain}:{entity_id} to entity_id-only key format
- Fix NIT: update publish_snapshot docstring in snapshot_publisher_registration

Review iteration: 2/10

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 0

Caution

Some comments are outside the diff and can’t be posted inline due to platform limitations.

⚠️ Outside diff range comments (1)
src/omnibase_infra/projectors/snapshot_publisher_registration.py (1)

1192-1215: ⚠️ Potential issue | 🟠 Major

Cancel pending debounced publishes before tombstoning.

If a debounced publish is pending for the same entity, a later timer can re‑publish the snapshot after the tombstone, effectively resurrecting deleted state. Clear the pending snapshot/timer when deleting.

🧹 Proposed fix
@@
         try:
+            # Cancel any pending debounced publish for this entity
+            async with self._debounce_lock:
+                timer = self._debounce_timers.pop(entity_id, None)
+                if timer is not None:
+                    timer.cancel()
+                self._pending_snapshots.pop(entity_id, None)
+
             # Build key for tombstone (node_id only, matches to_kafka_key())
             key = entity_id.encode()

…tion, types

- Propagate correlation_id through debounced snapshot publishes
- Cancel pending debounced publishes before tombstoning to prevent resurrection
- Sanitize snapshot publish errors in all 3 registration handlers
- Stop snapshot publisher in finally block on partial startup failures
- Use ProtocolSnapshotPublisher type annotations instead of concrete class
- Align key format docs to entity_id-only (not domain:entity_id)
- Add language identifiers to markdown fenced blocks (MD040)
- Add publish_from_projection to ProtocolSnapshotPublisher protocol
- Move snapshot_publisher init to function-level cleanup guard block
  to prevent UnboundLocalError on early startup failures
…shot-projector coupling

- Updated protocol_snapshot_publisher.py docstrings: replaced 12 stale
  references to old domain:entity_id composite key format with entity_id-only
  format matching the actual to_kafka_key() implementation
- Removed _build_key helper from example code (no longer needed)
- Added comment in handler_node_introspected.py explaining why snapshot
  publishing is intentionally gated on projector availability

Review iteration: 1/10

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 1

Caution

Some comments are outside the diff and can’t be posted inline due to platform limitations.

⚠️ Outside diff range comments (1)
src/omnibase_infra/protocols/protocol_snapshot_publisher.py (1)

34-38: ⚠️ Potential issue | 🟡 Minor

Protocol docs still describe domain‑prefixed keys.
Compaction semantics and examples mention domain:entity_id; update them to entity_id‑only to avoid implementers reintroducing the old key format.

🤖 Fix all issues with AI agents
In
`@src/omnibase_infra/nodes/node_registration_orchestrator/handlers/handler_node_registration_acked.py`:
- Around line 318-338: The snapshot publish call in
handler_node_registration_acked is missing propagation of the request
correlation_id; locate the code that derives or should derive correlation_id
(the incoming request handling in the same method) and pass that correlation_id
into self._snapshot_publisher.publish_from_projection(active_projection,
node_name=None, correlation_id=correlation_id). Ensure correlation_id is set
earlier in the method (use uuid.uuid4() when missing) and include that
correlation_id in any error/log context around publish_from_projection so the
trace chain is preserved.
🧹 Nitpick comments (1)
src/omnibase_infra/projectors/snapshot_publisher_registration.py (1)

1394-1447: Fire-and-forget task in debounce callback is acceptable but consider error handling.

The loop.create_task() call creates a fire-and-forget task. While this is the correct pattern for call_later callbacks, unhandled exceptions in fire-and-forget tasks can be silently lost or only logged to stderr by asyncio's default handler.

The current implementation logs exceptions with exc_info=True (lines 1441-1447), which is good. However, consider adding the task to a set to prevent garbage collection warnings and enable graceful cancellation during stop().

♻️ Optional: Track fire-and-forget tasks for cleaner lifecycle
 def __init__(self, ...):
     ...
     self._debounce_lock = asyncio.Lock()
+    self._pending_tasks: set[asyncio.Task[None]] = set()
 
 def _debounce_timer_callback(
     self,
     entity_id: str,
     loop: asyncio.AbstractEventLoop,
 ) -> None:
-    loop.create_task(self._execute_debounced_publish(entity_id))
+    task = loop.create_task(self._execute_debounced_publish(entity_id))
+    self._pending_tasks.add(task)
+    task.add_done_callback(self._pending_tasks.discard)

- Changed util_container_wiring.py parameter type from concrete
  SnapshotPublisherRegistration to ProtocolSnapshotPublisher (matches
  inner wiring function signature)
- Fixed get_latest_snapshot return type in ProtocolSnapshotPublisher from
  ModelRegistrationProjection to ModelRegistrationSnapshot (matches
  concrete implementation)

Review iteration: 2/10
…ection

- Added missing correlation_id kwarg in handler_node_introspected.py
  publish_from_projection call, preserving trace chain from introspection
  event through debounce pipeline to Kafka
- Added missing correlation_id kwarg in handler_node_registration_acked.py
  publish_from_projection call for ACTIVE transition snapshots

Review iteration: 3/10
…docstring

- Fixed Args description in SnapshotPublisherRegistration.publish_snapshot
  to match entity_id-only key format (domain isolation at topic level)

Review iteration: 4/10

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 2

🤖 Fix all issues with AI agents
In `@src/omnibase_infra/projectors/snapshot_publisher_registration.py`:
- Around line 1478-1491: The logger.warning call inside the pending loop (where
you call self._publish_snapshot_model for each entity_id) currently passes
exc_info=True and may leak sensitive connection data; change the except block to
capture the exception (except Exception as e:), use the sanitization utility
from omnibase_infra.utils.util_error_sanitization to produce a sanitized message
or sanitized exception details, and log that sanitized output (and set
exc_info=False) in the logger.warning for the "Failed to flush debounced
snapshot for %s during stop" message so no raw credentials/PII are logged.
- Around line 1410-1447: In _execute_debounced_publish, remove exc_info=True
from the logger.warning call and instead log a sanitized error string using
omnibase_infra.utils.util_error_sanitization.sanitize_error_message (import it
if not present); catch the exception as e in the except block and call
sanitize_error_message(e) and include that sanitized message (and optionally
str(type(e))) in the warning text for "Failed to publish debounced snapshot for
%s version %d" to avoid leaking credentials or sensitive details.
🧹 Nitpick comments (1)
src/omnibase_infra/nodes/node_registration_orchestrator/handlers/handler_node_introspected.py (1)

466-468: Consider moving the import outside the try block.

The inline import of ModelRegistrationProjection inside the try block is unusual. While it works, it adds overhead on each call and may obscure import errors. Since this file already imports from omnibase_infra.models.projection at the module level (via the projection reader), this could be a top-level import.

♻️ Suggested refactor

Add to imports near line 84:

from omnibase_infra.models.projection import ModelRegistrationProjection

Then remove lines 466-468.

Comment thread src/omnibase_infra/projectors/snapshot_publisher_registration.py
Comment thread src/omnibase_infra/projectors/snapshot_publisher_registration.py
… module level

Address remaining CodeRabbit PR #254 review feedback:
- Use sanitize_error_message() in debounce publish/flush error paths
  to prevent leaking sensitive data in logs
- Move ModelRegistrationProjection import from inline to module level
  in handler_node_introspected.py

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 0

Caution

Some comments are outside the diff and can’t be posted inline due to platform limitations.

⚠️ Outside diff range comments (1)
src/omnibase_infra/projectors/snapshot_publisher_registration.py (1)

1272-1336: ⚠️ Potential issue | 🟡 Minor

Auto-generate correlation_id at the API boundary.

For consistency, create a correlation_id in publish_from_projection when missing so the same ID flows through debounce and publishing.

🔧 Suggested fix
-        entity_id_str = str(projection.entity_id)
+        if correlation_id is None:
+            correlation_id = uuid4()
+        entity_id_str = str(projection.entity_id)
         version = await self._get_next_version(entity_id_str)

As per coding guidelines, "Propagate correlation_id from incoming requests; auto-generate with uuid4() if missing; include in all error context".

@jonahgabriel
jonahgabriel merged commit 6868d92 into main Feb 7, 2026
19 checks passed
@jonahgabriel
jonahgabriel deleted the jonah/omn-1932-phase-3-wire-snapshot-publisher-p1 branch February 7, 2026 14:44
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant