Skip to content

feat(projection): track last_heartbeat_at in registration projection [OMN-1006] - #94

Merged
jonahgabriel merged 6 commits into
mainfrom
jonah/omn-1006-track-last_heartbeat_at-in-registration-projection
Dec 25, 2025
Merged

jonahgabriel merged 6 commits into
mainfrom
jonah/omn-1006-track-last_heartbeat_at-in-registration-projection

Conversation

@jonahgabriel

@jonahgabriel jonahgabriel commented Dec 25, 2025 •

Copy link
Copy Markdown
Collaborator

Summary

Add last_heartbeat_at field to track when the last heartbeat was received from a node, enabling accurate reporting in ModelNodeLivenessExpired events.

Linear Issue: OMN-1006

Changes

File Change
model_registration_projection.py Added last_heartbeat_at: datetime | None field
schema_registration_projection.sql Added last_heartbeat_at TIMESTAMPTZ column
projector_registration.py Updated upsert SQL + added update_heartbeat() method
projection_reader_registration.py Populate last_heartbeat_at from DB rows
timeout_emitter.py Use projection.last_heartbeat_at instead of None
handler_node_heartbeat.py New - Draft heartbeat handler for future wiring
test_projection_reader_registration.py Fixed test fixtures to include new field
infra_validators.py Updated INFRA_MAX_UNIONS threshold (589 → 600)

Acceptance Criteria

  • ModelRegistrationProjection has last_heartbeat_at field
  • SQL schema includes last_heartbeat_at column
  • Projector supports updating last_heartbeat_at via update_heartbeat()
  • ModelNodeLivenessExpired includes accurate last_heartbeat_at
  • All existing tests pass

Test plan

  • Unit tests for projection model pass
  • Unit tests for projector pass
  • Unit tests for timeout emitter pass
  • Pre-commit hooks pass (including union validation)
  • Integration test with actual heartbeat events (follow-up)

Follow-up

The heartbeat handler (handler_node_heartbeat.py) is created but not yet wired into event routing. This can be done in a follow-up ticket when heartbeat event consumption is implemented.

Summary by CodeRabbit

  • New Features

    • Heartbeat handling: nodes can send heartbeats to update last-seen timestamps and extend liveness deadlines.
    • Orchestrator now accepts and delegates direct heartbeat events for liveness processing.
    • Public projector API to atomically update heartbeat/liveness.
  • Infrastructure

    • Database schema adds a heartbeat timestamp column for registration projections.
    • Timeout emitter now includes heartbeat timestamps when emitting expiration events.
  • Tests

    • New integration and unit tests covering heartbeat flows, concurrency, and error scenarios.
  • Documentation

    • Docstrings and contract updated to document heartbeat handling and timestamp accuracy.
  • Chores

    • Internal validator threshold increased to accommodate growth.

✏️ Tip: You can customize this high-level summary in your review settings.

…[OMN-1006]

Add last_heartbeat_at field to track when the last heartbeat was received
from a node, enabling accurate reporting in NodeLivenessExpired events.

Changes:
- Add last_heartbeat_at field to ModelRegistrationProjection
- Add last_heartbeat_at TIMESTAMPTZ column to SQL schema
- Update projector upsert to include last_heartbeat_at
- Add update_heartbeat() method to projector for heartbeat processing
- Update projection reader to populate last_heartbeat_at from DB rows
- Update timeout emitter to use projection.last_heartbeat_at
- Add heartbeat handler (draft, wiring in follow-up)
- Fix test fixtures to include new field
- Update INFRA_MAX_UNIONS threshold (589 -> 600)

Acceptance Criteria:
- [x] ModelRegistrationProjection has last_heartbeat_at field
- [x] SQL schema includes last_heartbeat_at column
- [x] Projector supports updating last_heartbeat_at
- [x] ModelNodeLivenessExpired includes accurate last_heartbeat_at
- [x] All existing tests pass
@linear

linear Bot commented Dec 25, 2025

Copy link
Copy Markdown

OMN-1006

@coderabbitai

coderabbitai Bot commented Dec 25, 2025 •

Copy link
Copy Markdown

Warning

Rate limit exceeded

@jonahgabriel has exceeded the limit for the number of commits that can be reviewed per hour. Please wait 7 minutes and 35 seconds before requesting another review.

⌛ How to resolve this issue?

After the wait time has elapsed, a review can be triggered using the @coderabbitai review command as a PR comment. Alternatively, push new commits to this PR.

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 have higher rate limits than the trial, open-source and free plans. In all cases, we re-allow further reviews after a brief timeout.

Please see our FAQ for further information.

📥 Commits

Reviewing files that changed from the base of the PR and between b9f1cf4 and a08d980.

📒 Files selected for processing (9)
  • src/omnibase_infra/orchestrators/__init__.py
  • src/omnibase_infra/orchestrators/registration/handlers/handler_node_heartbeat.py
  • src/omnibase_infra/schemas/schema_registration_projection.sql
  • src/omnibase_infra/services/timeout_emitter.py
  • src/omnibase_infra/validation/infra_validators.py
  • tests/integration/registration/handlers/conftest.py
  • tests/integration/registration/handlers/test_handler_node_heartbeat_integration.py
  • tests/unit/services/test_timeout_emitter.py
  • tests/unit/validation/test_validator_defaults.py
📝 Walkthrough

Walkthrough

Adds node heartbeat support: new projection field last_heartbeat_at, heartbeat handler and result model, projector update method for heartbeats, schema column, orchestrator hooks and contract entry, timeout emitter payload change, validator bump, and accompanying unit and integration tests.

Changes

Cohort / File(s) Summary
Projection model & DB schema
src/omnibase_infra/models/projection/model_registration_projection.py, src/omnibase_infra/schemas/schema_registration_projection.sql
Added optional `last_heartbeat_at: datetime
Projector & Reader
src/omnibase_infra/projectors/projection_reader_registration.py, src/omnibase_infra/projectors/projector_registration.py
Reader maps last_heartbeat_at from DB rows. Projector: include last_heartbeat_at in upserts and persist payloads; added update_heartbeat(...) to atomically set last_heartbeat_at and liveness_deadline with error mapping and circuit-breaker handling.
Heartbeat handler & orchestrator exports
src/omnibase_infra/orchestrators/registration/handlers/handler_node_heartbeat.py, src/omnibase_infra/orchestrators/registration/__init__.py, src/omnibase_infra/orchestrators/registration/handlers/__init__.py, src/omnibase_infra/orchestrators/__init__.py
New HandlerNodeHeartbeat, ModelHeartbeatHandlerResult, and DEFAULT_LIVENESS_WINDOW_SECONDS constant; package inits re-export public names.
Node orchestrator integration & contract
src/omnibase_infra/nodes/node_registration_orchestrator/node.py, src/omnibase_infra/nodes/node_registration_orchestrator/__init__.py, src/omnibase_infra/nodes/node_registration_orchestrator/contract.yaml
Orchestrator: _heartbeat_handler attribute, set_heartbeat_handler, has_heartbeat_handler, handle_heartbeat(...) added. Contract: consumed event for node heartbeat added with direct_handler: true.
Timeout emitter
src/omnibase_infra/services/timeout_emitter.py
Liveness-expiration events now include last_heartbeat_at from the projection payload.
Validation threshold
src/omnibase_infra/validation/infra_validators.py, tests/unit/validation/test_validator_defaults.py
Bumped INFRA_MAX_UNIONS from 589 to 600; updated test expectation and changelog entry.
Tests — unit & integration
tests/unit/projectors/test_projection_reader_registration.py, tests/unit/nodes/*, tests/unit/*, tests/integration/registration/handlers/*, tests/integration/registration/handlers/test_handler_node_heartbeat_integration.py, tests/integration/registration/handlers/conftest.py
Unit tests: include last_heartbeat_at in mock rows and adapt orchestrator minimality/whitelist checks for heartbeat methods. Integration: new fixtures and comprehensive integration tests for HandlerNodeHeartbeat.

Sequence Diagram

sequenceDiagram
    autonumber
    participant Event as Heartbeat Event
    participant Handler as HandlerNodeHeartbeat
    participant Reader as ProjectionReader
    participant Projector as ProjectorRegistration
    participant DB as Database
    Note over Handler,Projector: New heartbeat flow (OMN-1006)

    Event->>Handler: receive ModelNodeHeartbeatEvent
    activate Handler
    Handler->>Reader: fetch projection by node_id
    activate Reader
    Reader->>DB: SELECT registration_projections
    DB-->>Reader: row (includes last_heartbeat_at, liveness_deadline, state)
    Reader-->>Handler: ModelRegistrationProjection
    deactivate Reader

    alt projection not found
        Handler-->>Event: ModelHeartbeatHandlerResult(node_not_found=true)
    else projection found
        Handler->>Handler: compute new_liveness_deadline = event.timestamp + window
        Handler->>Projector: update_heartbeat(node_id, domain, last_heartbeat_at, liveness_deadline, correlation_id)
        activate Projector
        Projector->>DB: UPDATE registration_projections SET last_heartbeat_at, liveness_deadline ...
        DB-->>Projector: rows affected
        Projector-->>Handler: bool (updated?)
        deactivate Projector

        alt updated
            Handler-->>Event: ModelHeartbeatHandlerResult(success=true, previous_state, last_heartbeat_at, liveness_deadline)
        else not updated
            Handler-->>Event: ModelHeartbeatHandlerResult(success=false, node_not_found=true)
        end
    end
    deactivate Handler
Loading

Estimated code review effort

🎯 4 (Complex) | ⏱️ ~40 minutes

Poem

🐰 A soft thump echoes, timestamps set in place,
Deadlines stretch forward, keeping nodes in pace,
Projections now whisper when hearts last beat,
Handlers update, projector and DB meet—
Hooray, the cluster hums with steady grace!


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: 1

🧹 Nitpick comments (1)
src/omnibase_infra/orchestrators/registration/handlers/__init__.py (1)

5-13: Consider re-exporting DEFAULT_LIVENESS_WINDOW_SECONDS for API consistency.

The handler_node_heartbeat.py module includes DEFAULT_LIVENESS_WINDOW_SECONDS in its __all__ (line 299), but the package __init__.py does not re-export it. This creates an inconsistency where users can access from omnibase_infra.orchestrators.registration.handlers.handler_node_heartbeat import DEFAULT_LIVENESS_WINDOW_SECONDS but not from the package level.

Either add the constant to this file's exports, or remove it from the module's __all__ if it's meant to be internal.

🔎 Suggested fix to include the constant
 from omnibase_infra.orchestrators.registration.handlers.handler_node_heartbeat import (
+    DEFAULT_LIVENESS_WINDOW_SECONDS,
     HandlerNodeHeartbeat,
     ModelHeartbeatHandlerResult,
 )
 
 __all__ = [
+    "DEFAULT_LIVENESS_WINDOW_SECONDS",
     "HandlerNodeHeartbeat",
     "ModelHeartbeatHandlerResult",
 ]
📜 Review details

Configuration used: defaults

Review profile: CHILL

Plan: Lite

📥 Commits

Reviewing files that changed from the base of the PR and between 6e3d158 and 9276251.

📒 Files selected for processing (12)
  • src/omnibase_infra/models/projection/model_registration_projection.py
  • src/omnibase_infra/orchestrators/__init__.py
  • src/omnibase_infra/orchestrators/registration/__init__.py
  • src/omnibase_infra/orchestrators/registration/handlers/__init__.py
  • src/omnibase_infra/orchestrators/registration/handlers/handler_node_heartbeat.py
  • src/omnibase_infra/projectors/projection_reader_registration.py
  • src/omnibase_infra/projectors/projector_registration.py
  • src/omnibase_infra/schemas/schema_registration_projection.sql
  • src/omnibase_infra/services/timeout_emitter.py
  • src/omnibase_infra/validation/infra_validators.py
  • tests/unit/projectors/test_projection_reader_registration.py
  • tests/unit/validation/test_validator_defaults.py
🧰 Additional context used
📓 Path-based instructions (1)
**/*.py

📄 CodeRabbit inference engine (CLAUDE.md)

**/*.py: NEVER use Any types in Python - Always use specific types, use object for generic dispatchers instead
All data structures must be proper Pydantic models - one model per file
Use nullable type annotation X | None (PEP 604 union syntax) over Optional[X] for null types in Python
All services MUST use ModelONEXContainer for container-based dependency injection
Error classes must raise OnexError not base Exception - use raise OnexError(...) from e pattern
Never use isinstance for protocol resolution - use duck typing through protocols instead
NEVER include passwords, API keys, tokens, secrets, full connection strings with credentials, PII, or private keys in error messages or context

Files:

  • src/omnibase_infra/orchestrators/__init__.py
  • src/omnibase_infra/orchestrators/registration/handlers/__init__.py
  • src/omnibase_infra/orchestrators/registration/handlers/handler_node_heartbeat.py
  • src/omnibase_infra/projectors/projector_registration.py
  • src/omnibase_infra/orchestrators/registration/__init__.py
  • tests/unit/projectors/test_projection_reader_registration.py
  • src/omnibase_infra/models/projection/model_registration_projection.py
  • src/omnibase_infra/projectors/projection_reader_registration.py
  • src/omnibase_infra/validation/infra_validators.py
  • tests/unit/validation/test_validator_defaults.py
  • src/omnibase_infra/services/timeout_emitter.py
🧠 Learnings (5)
📚 Learning: 2025-12-03T16:55:49.755Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-03T16:55:49.755Z
Learning: Applies to agents/**/*.py : Implement ONEX-compliant agent architecture with four node types: Effect (External I/O), Compute (Pure transforms), Reducer (State/persistence), and Orchestrator (Workflow coordination)

Applied to files:

  • src/omnibase_infra/orchestrators/__init__.py
📚 Learning: 2025-11-24T16:32:55.606Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: .cursor/rules/node_standards.mdc:0-0
Timestamp: 2025-11-24T16:32:55.606Z
Learning: ORCHESTRATOR Nodes must inherit from `NodeOrchestrator` or use `NodeOrchestratorService` and must coordinate workflows, manage node interactions, and handle process/event orchestration

Applied to files:

  • src/omnibase_infra/orchestrators/__init__.py
📚 Learning: 2025-11-24T17:23:49.777Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: .cursor/rules/node_standards.mdc:0-0
Timestamp: 2025-11-24T17:23:49.777Z
Learning: All ONEX nodes must conform to the canonical structure, code generation, and interface patterns established in the `node_cli` node, using it as the primary source of truth for directory structure, contract schema patterns, linked document architecture, base state patterns, shared schema references, extensibility patterns, CLI interface declarations, code generation, dependency injection, error handling, testing, and documentation

Applied to files:

  • src/omnibase_infra/orchestrators/__init__.py
📚 Learning: 2025-12-07T17:50:13.678Z
Learnt from: CR
Repo: OmniNode-ai/omniintelligence PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-07T17:50:13.678Z
Learning: Decompose intelligence operations into specialized ONEX nodes following a four-node pattern: Orchestrator (coordinate workflows), Reducer (manage state, FSM transitions), Compute (pure data processing), and Effect (external I/O)

Applied to files:

  • src/omnibase_infra/orchestrators/__init__.py
📚 Learning: 2025-12-06T22:21:32.649Z
Learnt from: CR
Repo: OmniNode-ai/omniagent PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-06T22:21:32.649Z
Learning: Applies to nodes/**/*.py : Use `omnibase_infra` handlers for OmniIntelligence queries via HttpRestAdapter envelope pattern

Applied to files:

  • src/omnibase_infra/orchestrators/registration/handlers/handler_node_heartbeat.py
🧬 Code graph analysis (3)
src/omnibase_infra/orchestrators/registration/handlers/__init__.py (1)
src/omnibase_infra/orchestrators/registration/handlers/handler_node_heartbeat.py (2)
  • HandlerNodeHeartbeat (98-293)
  • ModelHeartbeatHandlerResult (45-95)
src/omnibase_infra/orchestrators/registration/handlers/handler_node_heartbeat.py (6)
src/omnibase_infra/enums/enum_infra_transport_type.py (1)
  • EnumInfraTransportType (28-52)
src/omnibase_infra/enums/enum_registration_state.py (1)
  • EnumRegistrationState (24-258)
src/omnibase_infra/errors/infra_errors.py (1)
  • RuntimeHostError (35-100)
src/omnibase_infra/models/registration/model_node_heartbeat_event.py (1)
  • ModelNodeHeartbeatEvent (19-95)
src/omnibase_infra/projectors/projection_reader_registration.py (2)
  • ProjectionReaderRegistration (45-656)
  • get_entity_state (146-226)
src/omnibase_infra/projectors/projector_registration.py (2)
  • ProjectorRegistration (52-840)
  • update_heartbeat (717-840)
src/omnibase_infra/orchestrators/registration/__init__.py (1)
src/omnibase_infra/orchestrators/registration/handlers/handler_node_heartbeat.py (2)
  • HandlerNodeHeartbeat (98-293)
  • ModelHeartbeatHandlerResult (45-95)
🔇 Additional comments (17)
src/omnibase_infra/models/projection/model_registration_projection.py (1)

137-140: LGTM! Field addition is clean and well-documented.

The last_heartbeat_at field is properly defined with:

  • Correct type annotation using PEP 604 union syntax (datetime | None)
  • Appropriate default value (None)
  • Clear description for liveness reporting
  • Logical placement between deadline and emission marker fields
src/omnibase_infra/projectors/projection_reader_registration.py (1)

134-134: LGTM! Database row mapping is correct.

The last_heartbeat_at field is properly extracted from the database row and passed to the projection constructor, following the same pattern as other timestamp fields.

src/omnibase_infra/validation/infra_validators.py (1)

373-377: LGTM! Threshold adjustment accounts for new heartbeat functionality.

The INFRA_MAX_UNIONS threshold is updated from 589 to 600 to accommodate the new heartbeat handler and projection update logic. The 11-unit increase (for ~6 unions mentioned) provides a reasonable buffer above the current baseline for near-term codebase growth, which aligns with the stated strategy.

src/omnibase_infra/services/timeout_emitter.py (1)

626-626: LGTM! Liveness event now includes actual heartbeat timestamp.

The event now correctly uses projection.last_heartbeat_at instead of None, providing accurate reporting of when the last heartbeat was received. The field is properly nullable (datetime | None), so the event can handle cases where no heartbeat was ever received.

tests/unit/projectors/test_projection_reader_registration.py (1)

88-88: LGTM! Test mock data updated correctly.

The last_heartbeat_at field is properly added to the mock row with a None default, maintaining compatibility with the updated projection schema while allowing tests to override when needed.

src/omnibase_infra/orchestrators/__init__.py (1)

1-3: LGTM! Package initializer follows conventions.

The orchestrator package __init__.py is properly structured with license headers and a clear docstring describing its purpose.

src/omnibase_infra/schemas/schema_registration_projection.sql (2)

67-67: LGTM! Schema column added correctly.

The last_heartbeat_at column is properly defined as TIMESTAMPTZ (timezone-aware) and nullable, matching the model definition. The placement between liveness_deadline and the timeout emission markers is logical.

Note: Ensure that any database migration strategy accounts for adding this column to existing registration_projections tables. The schema is idempotent for new installations, but existing deployments may need explicit migration steps.


175-177: LGTM! Column documentation is clear.

The comment accurately describes the purpose of last_heartbeat_at and when it's updated, providing helpful context for database administrators and future developers.

src/omnibase_infra/orchestrators/registration/__init__.py (1)

1-13: LGTM! Registration orchestrator package exports are clean.

The registration orchestrator package properly re-exports HandlerNodeHeartbeat and ModelHeartbeatHandlerResult, making them available for integration. The structure follows Python conventions with clear __all__ definition.

Note: The PR description mentions the heartbeat handler is "draft" and "not yet wired into event routing." This export structure supports future integration when the handler is ready to be activated.

tests/unit/validation/test_validator_defaults.py (1)

58-65: LGTM!

The threshold adjustment is well-documented with ticket reference (OMN-1006), follows the established history pattern, and the assertion is consistent with the updated value.

src/omnibase_infra/projectors/projector_registration.py (2)

259-278: LGTM!

The upsert SQL correctly incorporates last_heartbeat_at:

  • Added to INSERT column list (line 262)
  • Placeholder updated with correct position $9 (line 268)
  • ON CONFLICT SET clause properly references EXCLUDED.last_heartbeat_at (line 278)

717-840: LGTM!

The new update_heartbeat method is well-implemented:

  • Follows the established pattern from update_ack_timeout_marker and update_liveness_timeout_marker
  • Atomically updates both last_heartbeat_at and liveness_deadline
  • Proper circuit breaker integration
  • Consistent error handling with typed exceptions
  • Comprehensive docstring with example and ticket references
src/omnibase_infra/orchestrators/registration/handlers/handler_node_heartbeat.py (5)

1-42: LGTM!

Module structure is well-organized with proper license header, docstring with ticket references, and clean imports. The TYPE_CHECKING usage correctly avoids circular import issues with projector/reader dependencies.


45-95: LGTM!

The result model follows project conventions:

  • Immutable with frozen=True and extra="forbid"
  • Uses PEP 604 union syntax (X | None) per coding guidelines
  • All fields have descriptive Field definitions

133-155: LGTM!

Clean initialization with sensible default for liveness_window_seconds. The 90-second default (3x the 30-second heartbeat interval) is well-documented and allows for 2 missed heartbeats before liveness expiry.


193-280: LGTM!

The happy path implementation is well-structured:

  • Gracefully handles missing projections with descriptive error messages
  • Appropriately logs warnings for non-active nodes while still processing (handles race conditions)
  • Correctly calculates new liveness deadline from event timestamp
  • Returns comprehensive result with previous state and updated timestamps

296-300: LGTM!

Exports are appropriate for the module's public API.

Merge:
- Integrate OMN-816 handler naming changes from main

PR Review Fixes:
- Fix exception handling to preserve InfraConnectionError/InfraTimeoutError types
- Re-export DEFAULT_LIVENESS_WINDOW_SECONDS from handlers __init__.py
@claude

claude Bot commented Dec 25, 2025

Copy link
Copy Markdown

PR Review: feat(projection): track last_heartbeat_at in registration projection [OMN-1006]

Summary

This PR adds heartbeat tracking to the registration projection system, enabling accurate reporting of when nodes last sent heartbeats in liveness expiration events. The implementation is well-architected and follows ONEX infrastructure patterns closely.


✅ Strengths

1. Excellent ONEX Compliance

  • Type annotations: Perfect use of X | None (PEP 604) syntax throughout ✅
  • Error handling: Proper error hierarchy with InfraConnectionError, InfraTimeoutError, RuntimeHostError ✅
  • Correlation ID tracking: Consistent propagation across all operations ✅
  • Circuit breaker integration: Correct usage of MixinAsyncCircuitBreaker pattern in projector ✅
  • Error context sanitization: No sensitive data exposed in error messages ✅

2. Clean Architecture

  • Model-first design: ModelHeartbeatHandlerResult properly typed with Pydantic ✅
  • Separation of concerns: Handler reads via projection reader, writes via projector ✅
  • Idempotency: SQL RETURNING clause verifies update success ✅
  • Thread safety: Proper circuit breaker lock usage in projector ✅

3. Database Design

  • Schema migration: last_heartbeat_at TIMESTAMPTZ column added correctly ✅
  • SQL comments: Clear documentation of field purpose ✅
  • Atomic updates: Single UPDATE statement with proper WHERE clause ✅
  • Index support: Uses existing (entity_id, domain) primary key ✅

4. Documentation Quality

  • Comprehensive docstrings: All classes and methods well-documented ✅
  • Ticket tracing: Related tickets (OMN-1006, OMN-932, OMN-881) referenced ✅
  • Examples: Handler includes usage examples in docstrings ✅
  • Error documentation: Raises clauses clearly document error conditions ✅

🔍 Code Quality Observations

handler_node_heartbeat.py:284-289

Issue: Exception handling preserves InfraConnectionError and InfraTimeoutError correctly ✅

The PR description mentions "Fix exception handling to preserve InfraConnectionError/InfraTimeoutError types" - this is correctly implemented:

except InfraConnectionError:
    # Re-raise infrastructure connection errors directly (preserve error type)
    raise
except InfraTimeoutError:
    # Re-raise infrastructure timeout errors directly (preserve error type)
    raise

This follows ONEX error handling patterns perfectly.

projector_registration.py:717-840

Strong implementation of update_heartbeat() method:

  • ✅ Circuit breaker check before operation
  • ✅ Proper lock management (acquire before check, after reset/failure)
  • ✅ SQL uses RETURNING clause to verify row update
  • ✅ Distinguished PostgresConnectionError vs QueryCanceledError exceptions
  • ✅ All error paths record circuit breaker failure
  • ✅ Returns bool for idempotent handling (False if entity not found)

timeout_emitter.py:626

Excellent fix - Now passes projection.last_heartbeat_at instead of None:

last_heartbeat_at=projection.last_heartbeat_at,  # Was: None

This directly fulfills the ticket's acceptance criteria.


🎯 Performance Considerations

✅ Optimized Operations

  1. Single UPDATE query: No unnecessary SELECTs in heartbeat update path
  2. Connection pooling: Uses self._pool.acquire() for efficient connection management
  3. Index usage: UPDATE uses primary key (entity_id, domain) - very efficient
  4. Minimal logging: Debug-level success logs, warning-level anomalies

💡 Potential Optimization (Future)

The handler performs a read-then-write pattern:

projection = await self._projection_reader.get_entity_state(...)  # SELECT
# ...
updated = await self._projector.update_heartbeat(...)  # UPDATE

Consideration: For high-throughput heartbeat processing, this could be optimized to a single UPDATE with RETURNING * to get the previous state. However, the current pattern is:

  • ✅ Correct for the current use case
  • ✅ More maintainable (separation of read/write concerns)
  • ✅ Validates node existence before update
  • ✅ Allows state checking (is_active() validation)

Recommendation: Keep current implementation unless profiling shows heartbeat processing is a bottleneck.


🔒 Security Review

✅ No Security Concerns

  1. SQL injection: All queries use parameterized statements ($1, $2, etc.) ✅
  2. Error sanitization: No sensitive data in error messages ✅
  3. UUID validation: entity_id is strongly typed as UUID ✅
  4. Correlation ID tracking: Proper distributed tracing without leaking data ✅
  5. No hardcoded secrets: All configuration via parameters ✅

🧪 Test Coverage

✅ Test Updates

  • test_projection_reader_registration.py: Fixed to include last_heartbeat_at: None in mock rows ✅
  • test_validator_defaults.py: Updated INFRA_MAX_UNIONS threshold with rationale ✅

⚠️ Missing Test Coverage

The PR notes "Integration test with actual heartbeat events (follow-up)" is still pending. This is acceptable for this PR because:

  1. Unit tests for the projection model pass ✅
  2. Unit tests for the projector pass ✅
  3. The handler is created but not yet wired into event routing
  4. Integration testing makes more sense when the handler is fully connected

Recommendation: Create a follow-up ticket (if not already exists) to:

  • Wire HandlerNodeHeartbeat into the registration orchestrator
  • Add integration tests for end-to-end heartbeat flow
  • Test ModelNodeLivenessExpired event emission with accurate last_heartbeat_at

📐 ONEX Pattern Compliance

Circuit Breaker Pattern ✅

# Check circuit breaker (caller-held lock pattern)
async with self._circuit_breaker_lock:
    await self._check_circuit_breaker("update_heartbeat", corr_id)

# I/O operation outside lock
async with self._pool.acquire() as conn:
    result = await conn.fetchrow(...)

# Record success/failure (caller-held lock pattern)
async with self._circuit_breaker_lock:
    await self._reset_circuit_breaker()

This perfectly follows the caller-held lock pattern documented in CLAUDE.md.

Error Recovery Pattern ✅

The handler implements graceful degradation:

  • Returns ModelHeartbeatHandlerResult with success=False for missing nodes
  • Sets node_not_found=True instead of raising exceptions
  • Logs warnings for non-active node heartbeats but still processes them
  • Only raises exceptions for true infrastructure failures

This aligns with ONEX error recovery patterns.


🚀 Migration Considerations

✅ Backwards Compatibility

The last_heartbeat_at field is nullable (TIMESTAMPTZ without NOT NULL), which means:

  1. Existing rows can have NULL for last_heartbeat_at ✅
  2. No data migration required for existing nodes ✅
  3. First heartbeat will populate the field ✅
  4. ModelNodeLivenessExpired accepts last_heartbeat_at: datetime | None ✅

No breaking changes - excellent design decision.

SQL Migration Path

ALTER TABLE registration_projections 
ADD COLUMN last_heartbeat_at TIMESTAMPTZ;

This is a non-blocking DDL operation in PostgreSQL - safe for production deployment.


📊 Union Type Validation Update

infra_validators.py:378

INFRA_MAX_UNIONS = 600  # Was: 589

Rationale: +~6 unions from:

  • ModelHeartbeatHandlerResult fields (multiple X | None types)
  • Handler method signatures
  • Projector method signatures

Assessment: ✅ Acceptable increase with clear documentation

  • Threshold updated with ticket reference (OMN-1006)
  • Historical context preserved in comments
  • Target still documented: "Reduce to <200 through dict[str, object] -> JsonValue migration"

This is transparent technical debt tracking - well done.


🎯 Acceptance Criteria Review

Let's verify against the PR's stated acceptance criteria:

  • ModelRegistrationProjection has last_heartbeat_at field ✅

    • src/omnibase_infra/models/projection/model_registration_projection.py:137-140
  • SQL schema includes last_heartbeat_at column ✅

    • src/omnibase_infra/schemas/schema_registration_projection.sql:67
  • Projector supports updating last_heartbeat_at via update_heartbeat() ✅

    • src/omnibase_infra/projectors/projector_registration.py:717-840
  • ModelNodeLivenessExpired includes accurate last_heartbeat_at ✅

    • src/omnibase_infra/services/timeout_emitter.py:626
  • All existing tests pass ✅

    • Test fixtures updated in test_projection_reader_registration.py

All acceptance criteria met ✅


🔧 Minor Suggestions (Optional)

1. Handler Export Consistency

The __init__.py files export DEFAULT_LIVENESS_WINDOW_SECONDS - consider whether this constant needs to be part of the public API or if it should remain internal to the handler.

Current:

__all__ = [
    "DEFAULT_LIVENESS_WINDOW_SECONDS",  # Public constant
    "HandlerNodeHeartbeat",
    "ModelHeartbeatHandlerResult",
]

Rationale: If external code needs to configure liveness windows, they pass it via the constructor. The default is already available as handler.liveness_window_seconds property.

Impact: Low - keeping it exported is fine for transparency.

2. Test Coverage Note

Consider adding a unit test for HandlerNodeHeartbeat itself (not just the projector method) to validate:

  • node_not_found=True when projection is None
  • Warning logged for non-active nodes
  • Success result when update succeeds
  • Error result when update returns False

This would provide better coverage before integration testing.


✅ Final Recommendation

APPROVE - This PR is ready to merge.

Highlights:

  • ✅ Excellent ONEX pattern compliance
  • ✅ Proper error handling and circuit breaker integration
  • ✅ Clean separation of concerns (reader/projector/handler)
  • ✅ Backwards-compatible schema change
  • ✅ Well-documented code with ticket tracing
  • ✅ All acceptance criteria met
  • ✅ No security concerns
  • ✅ Transparent technical debt tracking (union validation)

Follow-up Work (As Noted in PR):

  • Wire HandlerNodeHeartbeat into registration orchestrator event routing
  • Add integration tests for end-to-end heartbeat processing
  • Verify ModelNodeLivenessExpired events include accurate timestamps in production

Great work on this implementation! The code quality is high and follows ONEX infrastructure standards closely.


Reviewed by: Claude Sonnet 4.5 (ONEX Infrastructure Specialist)
Review Date: 2025-12-25
Ticket: OMN-1006

- Fix exception handling to preserve all RuntimeHostError subtypes
  (InfraConnectionError, InfraTimeoutError, etc.) instead of wrapping
- Wire HandlerNodeHeartbeat into registration orchestrator with
  set_heartbeat_handler(), has_heartbeat_handler, handle_heartbeat()
- Add direct_handler flag to contract.yaml for non-workflow events
- Add 27 integration tests for heartbeat handler (8 test classes)
- Verify ModelNodeLivenessExpired timestamps are accurate (documented)
- Re-export DEFAULT_LIVENESS_WINDOW_SECONDS for API consistency
- Update orchestrator tests to whitelist heartbeat handler methods
@claude

claude Bot commented Dec 25, 2025

Copy link
Copy Markdown

PR Review: Track last_heartbeat_at in Registration Projection [OMN-1006]

Summary

This PR adds heartbeat timestamp tracking to the registration projection, enabling accurate liveness reporting. The implementation follows ONEX infrastructure patterns and includes comprehensive test coverage (877 lines of integration tests).


✅ Strengths

1. Excellent ONEX Compliance

  • ✅ Strong typing: No Any types, uses X | None (PEP 604) consistently
  • ✅ Error handling: Proper infrastructure error hierarchy with RuntimeHostError preservation
  • ✅ Circuit breaker: update_heartbeat() method uses MixinAsyncCircuitBreaker correctly
  • ✅ Correlation ID: Proper propagation throughout the call chain
  • ✅ Naming conventions: HandlerNodeHeartbeat, ModelHeartbeatHandlerResult follow patterns

2. Database Design

  • ✅ Schema evolution: Clean addition of last_heartbeat_at TIMESTAMPTZ column
  • ✅ SQL correctness: Atomic UPDATE with RETURNING for update verification
  • ✅ Idempotent: Schema uses IF NOT EXISTS, safely re-runnable
  • ✅ Documentation: Comprehensive SQL comments explaining the field purpose

3. Error Handling Excellence

except RuntimeHostError:
    # Re-raise all infrastructure errors directly (preserves error type)
    raise
except Exception as e:
    # Wrap only non-infrastructure errors
    raise RuntimeHostError(...) from e
  • ✅ Preserves specific error types (InfraConnectionError, InfraTimeoutError)
  • ✅ Follows CLAUDE.md error recovery patterns
  • ✅ Proper error sanitization (no credentials in error messages)

4. Test Coverage

  • ✅ 27 integration tests across 8 test classes (comprehensive)
  • ✅ Real PostgreSQL: Uses testcontainers, not mocks
  • ✅ Graceful CI skip: Tests skip if Docker unavailable (good for CI/CD)
  • ✅ Concurrent testing: Validates concurrent heartbeat processing
  • ✅ Error scenarios: Tests connection failures, timeouts, missing nodes

5. Documentation Quality

  • ✅ Extensive docstrings with examples
  • ✅ Clear design rationale in model_node_liveness_expired.py (timestamp verification)
  • ✅ Thread safety notes in handler_node_heartbeat.py
  • ✅ Proper __all__ exports for API clarity

🔍 Issues & Recommendations

1. 🟡 Medium: Handler Wiring Not Demonstrated

Issue: The heartbeat handler is created but wiring to event routing is deferred to "follow-up". This creates integration risk.

Location: node.py:257 - handle_heartbeat() method exists but no caller demonstrated

Recommendation:

# Consider adding integration test showing full event flow:
# 1. Heartbeat event received from Kafka
# 2. Dispatcher routes to orchestrator.handle_heartbeat()
# 3. Projection updated
# 4. Liveness deadline extended

Risk: Medium - functionality exists but integration path unverified


2. 🟡 Medium: Liveness Window Calculation Inconsistency

Issue: Deadline calculated from heartbeat_timestamp (event time), not now (processing time).

Location: handler_node_heartbeat.py:234

new_liveness_deadline = heartbeat_timestamp + timedelta(
    seconds=self._liveness_window_seconds
)

Problem: If event processing is delayed (e.g., Kafka lag), deadline may already be expired when calculated.

Example:

  • Heartbeat timestamp: 2025-12-25T10:00:00Z
  • Processing time: 2025-12-25T10:02:00Z (2 min lag)
  • Window: 90 seconds
  • Calculated deadline: 10:01:30Z (already expired!)

Recommendation:

# Option 1: Use processing time (safer for delayed events)
now = datetime.now(UTC)
new_liveness_deadline = now + timedelta(seconds=self._liveness_window_seconds)

# Option 2: Use max(event_time, now) to handle both cases
reference_time = max(heartbeat_timestamp, datetime.now(UTC))
new_liveness_deadline = reference_time + timedelta(seconds=self._liveness_window_seconds)

Impact: Could cause false liveness expirations in high-latency scenarios


3. 🟢 Low: Unused Variable in Error Context

Issue: ctx variable created but only used in exception handlers.

Location: handler_node_heartbeat.py:189-194

Recommendation: Consider lazy construction or clarify intent with comment:

# Error context prepared for potential exception handling
ctx = ModelInfraErrorContext(...)

4. 🟢 Low: Test Helper Function Naming

Location: test_handler_node_heartbeat_integration.py:73

def make_projection(

ONEX Convention: Factory functions should be create_* or build_*. make_* is acceptable but less common in ONEX codebase.

Recommendation: Consider renaming to create_test_projection() for consistency.


5. 🟢 Low: Union Validator Threshold Bump

Location: infra_validators.py

INFRA_MAX_UNIONS = 600  # Was 589

Question: Is this increase due to new heartbeat models, or unrelated growth?

Recommendation: Add comment explaining what pushed the limit:

# OMN-1006: Heartbeat handler models increased union count by ~11
INFRA_MAX_UNIONS = 600

🔒 Security Review

✅ No security concerns identified

  • Proper error sanitization (no credentials exposed)
  • UUID correlation IDs used (not user-controlled)
  • SQL uses parameterized queries (no injection risk)
  • No PII in heartbeat events

🚀 Performance Considerations

Positive:

  • ✅ Atomic single-row UPDATE (O(1) with index on entity_id)
  • ✅ Circuit breaker prevents cascade failures
  • ✅ Proper index coverage (idx_registration_liveness_timeout_scan)

Watch:

  • ⚠️ High heartbeat frequency (every 30s) × many nodes could generate DB load
  • ⚠️ Consider batch heartbeat processing if scaling beyond 1000 nodes

📊 Code Quality Metrics

Metric Value Assessment
Test Coverage 877 lines integration tests ✅ Excellent
Documentation Comprehensive docstrings ✅ Excellent
Type Safety 100% (no Any types) ✅ Perfect
Error Handling Infrastructure errors preserved ✅ Excellent
ONEX Compliance Full adherence to CLAUDE.md ✅ Perfect

🎯 Acceptance Criteria Review

From PR description:

  • ✅ ModelRegistrationProjection has last_heartbeat_at field
  • ✅ SQL schema includes last_heartbeat_at column
  • ✅ Projector supports updating last_heartbeat_at via update_heartbeat()
  • ✅ ModelNodeLivenessExpired includes accurate last_heartbeat_at
  • ✅ All existing tests pass
  • ⚠️ Integration test with actual heartbeat events (deferred to follow-up)

🔧 Recommendations Summary

Must Address Before Merge:

  1. Liveness window calculation - Decide on event time vs processing time strategy (see issue Add Claude Code GitHub Workflow #2 above)

Should Address (Minor):

  1. Document union threshold increase reason
  2. Add integration test showing full event routing flow (or create follow-up ticket)

Nice to Have:

  1. Rename make_projection → create_test_projection
  2. Add comment explaining ctx variable purpose

✅ Final Verdict

APPROVE with minor recommendations

This is high-quality ONEX infrastructure code with excellent test coverage and documentation. The liveness window calculation issue (#2) should be addressed before merge, but it's a design decision rather than a bug. Once resolved, this PR is production-ready.

Suggested Next Steps:

  1. Address liveness window calculation strategy
  2. Create follow-up ticket for heartbeat event routing integration
  3. Merge after CI passes

Great work on the comprehensive test suite and adherence to ONEX patterns! 🚀


Reviewed against: CLAUDE.md ONEX Infrastructure Guidelines
Linear: OMN-1006

@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

🧹 Nitpick comments (3)
tests/integration/registration/handlers/conftest.py (1)

69-94: Prefer importing DEFAULT_LIVENESS_WINDOW_SECONDS constant.

Line 93 hardcodes 90.0 for the liveness window. For consistency and maintainability, import and use DEFAULT_LIVENESS_WINDOW_SECONDS from the handlers package instead.

🔎 Proposed refactor
 @pytest.fixture
 def heartbeat_handler(
     reader: ProjectionReaderRegistration,
     projector: ProjectorRegistration,
 ) -> HandlerNodeHeartbeat:
     """Function-scoped HandlerNodeHeartbeat instance.
 
     Creates a handler with the default liveness window (90 seconds).
     Suitable for most integration tests.
 
     Args:
         reader: ProjectionReaderRegistration fixture for state lookups.
         projector: ProjectorRegistration fixture for state updates.
 
     Returns:
         HandlerNodeHeartbeat configured with default liveness window.
     """
     from omnibase_infra.orchestrators.registration.handlers import (
+        DEFAULT_LIVENESS_WINDOW_SECONDS,
         HandlerNodeHeartbeat,
     )
 
     return HandlerNodeHeartbeat(
         projection_reader=reader,
         projector=projector,
-        liveness_window_seconds=90.0,
+        liveness_window_seconds=DEFAULT_LIVENESS_WINDOW_SECONDS,
     )
tests/integration/registration/handlers/test_handler_node_heartbeat_integration.py (2)

555-589: Consider adding a more specific assertion for rapid heartbeat test.

The test correctly handles non-deterministic ordering with the comment at lines 587-588. However, you could strengthen the test by verifying that the final last_heartbeat_at is within the expected range of event timestamps.

🔎 Optional: Add timestamp range verification
         # Final state should have the last heartbeat timestamp
         final = await reader.get_entity_state(node_id)
         assert final is not None
         assert final.last_heartbeat_at is not None
         # The last heartbeat timestamp should be one of the event timestamps
         # (exact order is non-deterministic with concurrent writes)
+        # Verify the timestamp is within the expected range
+        earliest_time = base_time
+        latest_time = base_time + timedelta(milliseconds=900)
+        assert earliest_time <= final.last_heartbeat_at <= latest_time

729-731: Consider catching a more specific exception type.

The test catches a generic Exception when verifying the model is frozen. Pydantic frozen models raise pydantic.ValidationError when attempting to modify them.

🔎 Optional: Use specific exception type
+from pydantic import ValidationError
+
         # Attempt to modify should fail
-        with pytest.raises(Exception):  # ValidationError for frozen models
+        with pytest.raises(ValidationError):
             result.success = False  # type: ignore[misc]
📜 Review details

Configuration used: defaults

Review profile: CHILL

Plan: Lite

📥 Commits

Reviewing files that changed from the base of the PR and between 9276251 and 6dfed0b.

📒 Files selected for processing (17)
  • src/omnibase_infra/nodes/node_registration_orchestrator/__init__.py
  • src/omnibase_infra/nodes/node_registration_orchestrator/contract.yaml
  • src/omnibase_infra/nodes/node_registration_orchestrator/models/model_node_liveness_expired.py
  • src/omnibase_infra/nodes/node_registration_orchestrator/node.py
  • src/omnibase_infra/orchestrators/registration/__init__.py
  • src/omnibase_infra/orchestrators/registration/handlers/__init__.py
  • src/omnibase_infra/orchestrators/registration/handlers/handler_node_heartbeat.py
  • src/omnibase_infra/projectors/projection_reader_registration.py
  • src/omnibase_infra/projectors/projector_registration.py
  • src/omnibase_infra/validation/infra_validators.py
  • tests/integration/registration/handlers/__init__.py
  • tests/integration/registration/handlers/conftest.py
  • tests/integration/registration/handlers/test_handler_node_heartbeat_integration.py
  • tests/unit/nodes/test_node_registration_orchestrator.py
  • tests/unit/nodes/test_orchestrator_decision_paths.py
  • tests/unit/nodes/test_orchestrator_no_io.py
  • tests/unit/validation/test_validator_defaults.py
✅ Files skipped from review due to trivial changes (2)
  • src/omnibase_infra/nodes/node_registration_orchestrator/models/model_node_liveness_expired.py
  • tests/integration/registration/handlers/init.py
🚧 Files skipped from review as they are similar to previous changes (3)
  • src/omnibase_infra/projectors/projection_reader_registration.py
  • src/omnibase_infra/validation/infra_validators.py
  • src/omnibase_infra/orchestrators/registration/init.py
🧰 Additional context used
📓 Path-based instructions (1)
**/*.py

📄 CodeRabbit inference engine (CLAUDE.md)

**/*.py: NEVER use Any types in Python - Always use specific types, use object for generic dispatchers instead
All data structures must be proper Pydantic models - one model per file
Use nullable type annotation X | None (PEP 604 union syntax) over Optional[X] for null types in Python
All services MUST use ModelONEXContainer for container-based dependency injection
Error classes must raise OnexError not base Exception - use raise OnexError(...) from e pattern
Never use isinstance for protocol resolution - use duck typing through protocols instead
NEVER include passwords, API keys, tokens, secrets, full connection strings with credentials, PII, or private keys in error messages or context

Files:

  • tests/integration/registration/handlers/conftest.py
  • src/omnibase_infra/projectors/projector_registration.py
  • tests/integration/registration/handlers/test_handler_node_heartbeat_integration.py
  • src/omnibase_infra/orchestrators/registration/handlers/__init__.py
  • src/omnibase_infra/orchestrators/registration/handlers/handler_node_heartbeat.py
  • tests/unit/validation/test_validator_defaults.py
  • tests/unit/nodes/test_orchestrator_no_io.py
  • tests/unit/nodes/test_orchestrator_decision_paths.py
  • tests/unit/nodes/test_node_registration_orchestrator.py
  • src/omnibase_infra/nodes/node_registration_orchestrator/node.py
  • src/omnibase_infra/nodes/node_registration_orchestrator/__init__.py
🧠 Learnings (26)
📚 Learning: 2025-12-07T17:50:13.678Z
Learnt from: CR
Repo: OmniNode-ai/omniintelligence PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-07T17:50:13.678Z
Learning: Organize tests following the structure: tests/conftest.py for shared fixtures, tests/unit/ for unit tests (no infrastructure), tests/integration/ for integration tests (requires Kafka/DBs), tests/nodes/ for node-specific tests

Applied to files:

  • tests/integration/registration/handlers/conftest.py
📚 Learning: 2025-11-30T21:55:10.298Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-11-30T21:55:10.298Z
Learning: Applies to src/omninode_bridge/nodes/**/contract.yaml : All contract YAML files for ONEX v2.0 nodes MUST define subcontract references, input/output models, and FSM configurations. Use YAML 1.2 syntax.

Applied to files:

  • src/omnibase_infra/nodes/node_registration_orchestrator/contract.yaml
📚 Learning: 2025-11-24T17:23:49.777Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: .cursor/rules/node_standards.mdc:0-0
Timestamp: 2025-11-24T17:23:49.777Z
Learning: Applies to **/node_*/v[0-9]*_[0-9]*_[0-9]*/contracts/contract_*.yaml : All ONEX node contract definitions must use the new subcontract architecture pattern, breaking down complex contracts into separate contract_actions.yaml, contract_models.yaml, contract_validation.yaml, and optional contract_cli.yaml and contract_capabilities.yaml files for separation of concerns, maintainability, reusability, modularity, and future tool-as-a-service readiness

Applied to files:

  • src/omnibase_infra/nodes/node_registration_orchestrator/contract.yaml
📚 Learning: 2025-11-24T17:23:49.777Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: .cursor/rules/node_standards.mdc:0-0
Timestamp: 2025-11-24T17:23:49.777Z
Learning: Applies to **/node_*/v[0-9]*_[0-9]*_[0-9]*/contract.yaml : All ONEX node contract definitions must follow the linked document architecture pattern with contract.yaml linking to node_config.yaml and deployment_config.yaml as associated documents

Applied to files:

  • src/omnibase_infra/nodes/node_registration_orchestrator/contract.yaml
📚 Learning: 2025-11-24T17:23:49.777Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: .cursor/rules/node_standards.mdc:0-0
Timestamp: 2025-11-24T17:23:49.777Z
Learning: Applies to **/node_*/v[0-9]*_[0-9]*_[0-9]*/contracts/contract_capabilities.yaml : All ONEX node execution capability definitions, if applicable, must be included in contract_capabilities.yaml with supported_node_types, supported_delivery_modes, and performance_constraints specifications

Applied to files:

  • src/omnibase_infra/nodes/node_registration_orchestrator/contract.yaml
📚 Learning: 2025-11-30T21:55:10.298Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-11-30T21:55:10.298Z
Learning: Applies to src/omninode_bridge/events/**/*.py : Kafka event publishing MUST use OnexEnvelopeV1 format with 13 topics for event streaming at all workflow lifecycle stages

Applied to files:

  • src/omnibase_infra/nodes/node_registration_orchestrator/contract.yaml
📚 Learning: 2025-11-24T16:32:55.606Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: .cursor/rules/node_standards.mdc:0-0
Timestamp: 2025-11-24T16:32:55.606Z
Learning: Applies to tests/unit/infrastructure/**/test_*.py : All node implementations must have comprehensive unit tests following the testing pattern in `tests/unit/infrastructure/` with tests for node initialization and node execution

Applied to files:

  • tests/integration/registration/handlers/test_handler_node_heartbeat_integration.py
📚 Learning: 2025-12-25T19:10:27.034Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-25T19:10:27.034Z
Learning: Applies to **/*adapter*.py **/*handler*.py : Use `InfraConnectionError` for connection failures, with transport-aware error codes (DATABASE_CONNECTION_ERROR for database, NETWORK_ERROR for HTTP/GRPC, SERVICE_UNAVAILABLE for Kafka/Consul/Vault/Valkey)

Applied to files:

  • src/omnibase_infra/orchestrators/registration/handlers/handler_node_heartbeat.py
📚 Learning: 2025-12-25T19:10:27.034Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-25T19:10:27.034Z
Learning: Applies to **/*adapter*.py **/*handler*.py : Use `InfraTimeoutError` for operation timeouts

Applied to files:

  • src/omnibase_infra/orchestrators/registration/handlers/handler_node_heartbeat.py
📚 Learning: 2025-12-25T19:10:27.034Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-25T19:10:27.034Z
Learning: Applies to **/*adapter*.py **/*handler*.py : Use graceful degradation for `InfraTimeoutError` - fallback to secondary data source (cache, secondary database) when primary times out

Applied to files:

  • src/omnibase_infra/orchestrators/registration/handlers/handler_node_heartbeat.py
📚 Learning: 2025-12-25T19:10:27.034Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-25T19:10:27.034Z
Learning: Applies to **/*adapter*.py **/*handler*.py : Use `InfraUnavailableError` for service unavailable conditions, including when circuit breaker is open

Applied to files:

  • src/omnibase_infra/orchestrators/registration/handlers/handler_node_heartbeat.py
📚 Learning: 2025-12-25T19:10:27.033Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-25T19:10:27.033Z
Learning: Applies to **/*adapter*.py **/*handler*.py **/*service*.py : Use `ModelInfraErrorContext` with `transport_type` when raising infrastructure errors

Applied to files:

  • src/omnibase_infra/orchestrators/registration/handlers/handler_node_heartbeat.py
📚 Learning: 2025-12-25T19:10:27.034Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-25T19:10:27.034Z
Learning: Applies to **/*adapter*.py **/*handler*.py : Use `InfraAuthenticationError` for authentication/authorization failures

Applied to files:

  • src/omnibase_infra/orchestrators/registration/handlers/handler_node_heartbeat.py
📚 Learning: 2025-12-25T19:10:27.034Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-25T19:10:27.034Z
Learning: Applies to **/*adapter*.py **/*handler*.py : Use circuit breaker pattern for `InfraUnavailableError` to prevent cascading failures - prevent requests when circuit is open, give service time to recover

Applied to files:

  • src/omnibase_infra/orchestrators/registration/handlers/handler_node_heartbeat.py
📚 Learning: 2025-11-30T21:55:10.298Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-11-30T21:55:10.298Z
Learning: Applies to src/omninode_bridge/codegen/**/*.py : Code generation service MUST auto-generate ONEX v2.0 compliant nodes with intelligent mixin injection and quality validation. Generate comprehensive test suites with 90%+ coverage.

Applied to files:

  • tests/unit/nodes/test_node_registration_orchestrator.py
📚 Learning: 2025-11-24T16:32:55.606Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: .cursor/rules/node_standards.mdc:0-0
Timestamp: 2025-11-24T16:32:55.606Z
Learning: ORCHESTRATOR Nodes must inherit from `NodeOrchestrator` or use `NodeOrchestratorService` and must coordinate workflows, manage node interactions, and handle process/event orchestration

Applied to files:

  • src/omnibase_infra/nodes/node_registration_orchestrator/node.py
  • src/omnibase_infra/nodes/node_registration_orchestrator/__init__.py
📚 Learning: 2025-12-20T04:09:41.832Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_core PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-20T04:09:41.832Z
Learning: Applies to **/*.py : Use ModelONEXContainer in node constructors for dependency injection, never use ModelContainer[T]

Applied to files:

  • src/omnibase_infra/nodes/node_registration_orchestrator/node.py
📚 Learning: 2025-11-30T21:55:10.298Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-11-30T21:55:10.298Z
Learning: Applies to src/omninode_bridge/nodes/**/*.py : All nodes in omninode_bridge MUST use omnibase_core standards (ModelServiceEffect, ModelServiceCompute for effect/compute nodes; NodeOrchestrator, NodeReducer with mixins for orchestrator/reducer nodes)

Applied to files:

  • src/omnibase_infra/nodes/node_registration_orchestrator/node.py
📚 Learning: 2025-11-30T21:55:10.298Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-11-30T21:55:10.298Z
Learning: Applies to src/omninode_bridge/nodes/**/*.py : Import mixins from omnibase_core.mixins.* and use Mixin* naming pattern (e.g., MixinHealthCheck, MixinMetrics, MixinEventBus) - never use local custom mixins unless experimental and documented

Applied to files:

  • src/omnibase_infra/nodes/node_registration_orchestrator/node.py
📚 Learning: 2025-11-24T16:32:55.606Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: .cursor/rules/node_standards.mdc:0-0
Timestamp: 2025-11-24T16:32:55.606Z
Learning: Node communication must use event-driven patterns through `ModelEventEnvelope` from `omnibase_core.models.events.model_event_envelope`

Applied to files:

  • src/omnibase_infra/nodes/node_registration_orchestrator/node.py
📚 Learning: 2025-12-25T19:10:27.034Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-25T19:10:27.034Z
Learning: Applies to nodes/**/node.py : Node archetypes and I/O models must be imported from `omnibase_core.nodes`, never defined in infra - NodeEffect, NodeCompute, NodeReducer, NodeOrchestrator

Applied to files:

  • src/omnibase_infra/nodes/node_registration_orchestrator/node.py
📚 Learning: 2025-12-08T00:48:30.737Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_spi PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-08T00:48:30.737Z
Learning: Import `omnibase_core` models and types only for type hints and runtime usage - follow the SPI → Core dependency direction

Applied to files:

  • src/omnibase_infra/nodes/node_registration_orchestrator/node.py
📚 Learning: 2025-12-25T19:10:27.033Z
Learnt from: CR
Repo: OmniNode-ai/omnibase_infra PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-25T19:10:27.033Z
Learning: Applies to **/*.py : All services MUST use `ModelONEXContainer` for container-based dependency injection

Applied to files:

  • src/omnibase_infra/nodes/node_registration_orchestrator/node.py
📚 Learning: 2025-11-24T16:32:55.606Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: .cursor/rules/node_standards.mdc:0-0
Timestamp: 2025-11-24T16:32:55.606Z
Learning: Node services must implement dependency injection using `ModelONEXContainer` from `omnibase_core.models.container.model_onex_container`

Applied to files:

  • src/omnibase_infra/nodes/node_registration_orchestrator/node.py
📚 Learning: 2025-12-06T22:21:32.649Z
Learnt from: CR
Repo: OmniNode-ai/omniagent PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-06T22:21:32.649Z
Learning: Applies to nodes/**/*.py : Use event bus mixins from `omnibase_core` for Kafka publishing instead of direct Kafka clients

Applied to files:

  • src/omnibase_infra/nodes/node_registration_orchestrator/node.py
📚 Learning: 2025-12-03T16:55:49.755Z
Learnt from: CR
Repo: OmniNode-ai/omniclaude PR: 0
File: CLAUDE.md:0-0
Timestamp: 2025-12-03T16:55:49.755Z
Learning: Applies to agents/**/*.py : Implement ONEX-compliant agent architecture with four node types: Effect (External I/O), Compute (Pure transforms), Reducer (State/persistence), and Orchestrator (Workflow coordination)

Applied to files:

  • src/omnibase_infra/nodes/node_registration_orchestrator/node.py
  • src/omnibase_infra/nodes/node_registration_orchestrator/__init__.py
🧬 Code graph analysis (4)
src/omnibase_infra/projectors/projector_registration.py (4)
src/omnibase_infra/errors/model_infra_error_context.py (1)
  • ModelInfraErrorContext (17-96)
src/omnibase_infra/enums/enum_infra_transport_type.py (1)
  • EnumInfraTransportType (28-52)
src/omnibase_infra/mixins/mixin_async_circuit_breaker.py (3)
  • _check_circuit_breaker (271-366)
  • _reset_circuit_breaker (458-533)
  • _record_circuit_failure (368-456)
src/omnibase_infra/errors/infra_errors.py (3)
  • InfraConnectionError (181-286)
  • InfraTimeoutError (289-326)
  • RuntimeHostError (35-100)
tests/integration/registration/handlers/test_handler_node_heartbeat_integration.py (10)
src/omnibase_infra/errors/infra_errors.py (2)
  • InfraConnectionError (181-286)
  • RuntimeHostError (35-100)
src/omnibase_infra/models/projection/model_registration_projection.py (1)
  • ModelRegistrationProjection (35-332)
src/omnibase_infra/models/projection/model_sequence_info.py (1)
  • ModelSequenceInfo (21-179)
src/omnibase_infra/models/registration/model_node_heartbeat_event.py (1)
  • ModelNodeHeartbeatEvent (19-95)
src/omnibase_infra/models/registration/model_node_capabilities.py (1)
  • ModelNodeCapabilities (13-167)
src/omnibase_infra/orchestrators/registration/handlers/handler_node_heartbeat.py (3)
  • HandlerNodeHeartbeat (98-305)
  • ModelHeartbeatHandlerResult (45-95)
  • handle (157-305)
src/omnibase_infra/projectors/projection_reader_registration.py (1)
  • ProjectionReaderRegistration (45-656)
tests/integration/registration/handlers/conftest.py (2)
  • heartbeat_handler (70-94)
  • heartbeat_handler_fast_window (98-122)
src/omnibase_infra/enums/enum_infra_transport_type.py (1)
  • EnumInfraTransportType (28-52)
tests/integration/projectors/conftest.py (1)
  • pg_pool (142-185)
src/omnibase_infra/orchestrators/registration/handlers/__init__.py (1)
src/omnibase_infra/orchestrators/registration/handlers/handler_node_heartbeat.py (2)
  • HandlerNodeHeartbeat (98-305)
  • ModelHeartbeatHandlerResult (45-95)
tests/unit/nodes/test_orchestrator_no_io.py (1)
tests/helpers/ast_analysis.py (1)
  • is_docstring (281-311)
🔇 Additional comments (19)
src/omnibase_infra/nodes/node_registration_orchestrator/__init__.py (1)

19-30: LGTM! Clear documentation for heartbeat handling.

The Event Handlers section effectively documents the new heartbeat handling capability with a helpful import snippet.

tests/unit/nodes/test_orchestrator_decision_paths.py (1)

529-579: LGTM! Correctly handles direct_handler events.

The logic properly skips events with direct_handler=true from workflow pattern validation since they bypass the workflow execution. The updated error message provides clear remediation options.

tests/unit/nodes/test_node_registration_orchestrator.py (1)

159-185: LGTM! Test coverage mirrors timeout handling pattern.

The test properly validates that heartbeat handling methods are present on the orchestrator class, following the same pattern as the timeout coordinator tests.

src/omnibase_infra/orchestrators/registration/handlers/__init__.py (1)

1-15: LGTM! Standard package initialization.

Clean re-export pattern for the heartbeat handler and related types.

src/omnibase_infra/nodes/node_registration_orchestrator/node.py (3)

49-74: LGTM! Clear documentation for heartbeat wiring.

The docstring provides a complete example of how to wire the heartbeat handler with the orchestrator.


112-120: LGTM! Proper TYPE_CHECKING imports and initialization.

Type hints are properly gated under TYPE_CHECKING to avoid circular dependencies, and the handler attribute is initialized correctly.

Also applies to: 203-203


260-321: LGTM! Heartbeat handling follows timeout coordinator pattern.

The implementation correctly mirrors the timeout coordinator pattern with proper delegation, error handling, and comprehensive docstrings. The RuntimeError for missing handler is consistent with the existing timeout coordinator approach.

tests/unit/nodes/test_orchestrator_no_io.py (2)

285-315: LGTM! Correctly whitelists heartbeat delegation methods.

The heartbeat handler methods are properly whitelisted as legitimate delegation points that don't violate the pure coordinator principle, following the same pattern as timeout coordination.


487-537: LGTM! Consistent whitelist across test classes.

The heartbeat methods whitelist is consistently applied, and the init statement count is appropriately updated to account for the _heartbeat_handler initialization.

tests/integration/registration/handlers/conftest.py (1)

97-122: LGTM! Fast window fixture is well-documented.

The fast window variant with 5.0 seconds is appropriate for testing deadline calculations without long waits.

tests/integration/registration/handlers/test_handler_node_heartbeat_integration.py (3)

1-66: LGTM! Well-structured test module with comprehensive documentation.

The module docstring clearly documents the test coverage areas (happy path, error scenarios, concurrency, etc.) and includes related ticket references. The graceful skip behavior for Docker availability is well-documented for CI/CD pipelines.


73-180: Well-designed test helpers with clear defaults and documentation.

The helper functions follow best practices:

  • Use keyword-only arguments for clarity (* parameter)
  • Proper type annotations with X | None syntax per coding guidelines
  • Comprehensive docstrings explaining each parameter
  • seed_projection includes an assertion to fail fast if seeding fails

599-636: Good test for connection error propagation with proper import placement.

The test correctly imports EnumInfraTransportType and ModelInfraErrorContext within the test method to construct the error context. This validates that InfraConnectionError is propagated correctly from the projector to the caller, which aligns with the handler's documented error handling behavior.

src/omnibase_infra/projectors/projector_registration.py (2)

259-332: LGTM! Upsert SQL correctly updated for last_heartbeat_at field.

The SQL changes properly:

  • Add last_heartbeat_at to the INSERT column list (line 262)
  • Include the corresponding parameter placeholder in VALUES (line 268, $9)
  • Update the column in ON CONFLICT DO UPDATE SET (line 278)
  • Pass projection.last_heartbeat_at in the params tuple (line 322)

The parameter ordering is correct and the WHERE clause logic for stale update rejection remains unchanged.


717-840: LGTM! Well-implemented update_heartbeat method following established patterns.

The new method:

  • Follows the same structure as update_ack_timeout_marker and update_liveness_timeout_marker
  • Properly implements circuit breaker pattern with lock acquisition
  • Uses appropriate error types: InfraConnectionError for connection failures, InfraTimeoutError for cancellations, RuntimeHostError for other errors (per coding guidelines)
  • Includes comprehensive logging with correlation context
  • Sets updated_at = $3 (last_heartbeat_at) which is consistent with the heartbeat update semantics
  • Docstring includes related ticket reference (OMN-1006)
src/omnibase_infra/orchestrators/registration/handlers/handler_node_heartbeat.py (4)

45-96: LGTM! Well-designed immutable result model.

The ModelHeartbeatHandlerResult model:

  • Uses frozen=True for immutability and thread safety
  • Uses extra="forbid" for strict validation
  • Uses nullable type annotation X | None per coding guidelines
  • All fields are well-documented with descriptions
  • Follows Pydantic best practices

98-156: LGTM! Handler class follows best practices.

The HandlerNodeHeartbeat class:

  • Uses constructor-based dependency injection for projection_reader and projector
  • Is stateless and thread-safe as documented
  • Provides a property for liveness_window_seconds for configuration visibility
  • Documents error handling behavior including specific exception types
  • Default liveness window (90s) matches the 3x heartbeat interval heuristic

285-305: Exception handling correctly preserves RuntimeHostError subtypes.

The exception handling:

  • Re-raises all RuntimeHostError subclasses directly (lines 285-290), preserving InfraConnectionError, InfraTimeoutError, etc.
  • Only wraps unexpected non-infrastructure errors in RuntimeHostError (lines 291-305)
  • Includes detailed logging with error type for debugging
  • This addresses the concern from the past review comment about wrapping specific error types

Based on learnings, this follows the guideline to use InfraConnectionError for connection failures and InfraTimeoutError for operation timeouts.


157-283: LGTM! Handle method implements correct heartbeat processing flow.

The method correctly:

  • Generates correlation_id if not provided for tracing
  • Returns structured result for unknown nodes (node_not_found=True)
  • Logs warnings for non-ACTIVE states but still processes the heartbeat (correct for race conditions)
  • Calculates liveness_deadline as event.timestamp + liveness_window
  • Handles the TOCTOU race condition where entity exists during lookup but is deleted before update
  • Returns comprehensive result with previous_state, timestamps, and correlation_id

Comment on lines +208 to +216
# OMN-1006: Node heartbeat events for liveness tracking
# Heartbeats update last_heartbeat_at and extend liveness_deadline in the
# registration projection. The HandlerNodeHeartbeat processes these events.
# Note: direct_handler=true indicates this event bypasses the workflow and is
# handled by a dedicated handler (HandlerNodeHeartbeat) via handle_heartbeat().
- topic: "node.heartbeat"
event_type: "ModelNodeHeartbeatEvent"
description: "Periodic heartbeat from active nodes for liveness tracking"
direct_handler: true

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 | 🔴 Critical

Critical: Topic pattern violates ONEX convention.

Line 213 uses "node.heartbeat" which doesn't follow the required ONEX topic pattern used by all other consumed events. This will break topic templating at runtime when the system attempts to replace {env} and {namespace} placeholders.

All consumed events must follow one of these patterns:

  • {env}.{namespace}.onex.evt.<slug>.v1 (external events)
  • {env}.{namespace}.onex.internal.<slug>.v1 (internal events)
🔎 Proposed fix
-  - topic: "node.heartbeat"
+  - topic: "{env}.{namespace}.onex.evt.node-heartbeat.v1"
     event_type: "ModelNodeHeartbeatEvent"
     description: "Periodic heartbeat from active nodes for liveness tracking"
     direct_handler: true
📝 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.

Suggested change
# OMN-1006: Node heartbeat events for liveness tracking
# Heartbeats update last_heartbeat_at and extend liveness_deadline in the
# registration projection. The HandlerNodeHeartbeat processes these events.
# Note: direct_handler=true indicates this event bypasses the workflow and is
# handled by a dedicated handler (HandlerNodeHeartbeat) via handle_heartbeat().
- topic: "node.heartbeat"
event_type: "ModelNodeHeartbeatEvent"
description: "Periodic heartbeat from active nodes for liveness tracking"
direct_handler: true
# OMN-1006: Node heartbeat events for liveness tracking
# Heartbeats update last_heartbeat_at and extend liveness_deadline in the
# registration projection. The HandlerNodeHeartbeat processes these events.
# Note: direct_handler=true indicates this event bypasses the workflow and is
# handled by a dedicated handler (HandlerNodeHeartbeat) via handle_heartbeat().
- topic: "{env}.{namespace}.onex.evt.node-heartbeat.v1"
event_type: "ModelNodeHeartbeatEvent"
description: "Periodic heartbeat from active nodes for liveness tracking"
direct_handler: true
🤖 Prompt for AI Agents
In src/omnibase_infra/nodes/node_registration_orchestrator/contract.yaml around
lines 208 to 216, the topic "node.heartbeat" violates the ONEX naming convention
and will break runtime templating; replace the topic value with the templated
ONEX pattern (e.g. "{env}.{namespace}.onex.evt.node.heartbeat.v1") so that {env}
and {namespace} can be substituted at runtime, keeping event_type, description
and direct_handler unchanged.

Comment thread tests/unit/validation/test_validator_defaults.py Outdated
- Fix CRITICAL: Change heartbeat topic to ONEX convention
  `node.heartbeat` → `{env}.{namespace}.onex.evt.node-heartbeat.v1`
- Fix CI failure: test_consumed_events_have_topics now passes
- Fix MINOR: Update test docstring consistency for INFRA_MAX_UNIONS
- Fix NITPICK: Add specific timestamp assertion for rapid heartbeat test

MAJOR review item verified correct: Exception handling already preserves
InfraConnectionError/InfraTimeoutError types (they inherit RuntimeHostError).
@claude

claude Bot commented Dec 25, 2025

Copy link
Copy Markdown

PR Review - feat(projection): track last_heartbeat_at in registration projection [OMN-1006]

Overall Assessment

LGTM with minor observations ✅

This is a well-architected implementation that properly extends the registration projection with heartbeat tracking. The code follows ONEX conventions closely and demonstrates excellent attention to detail in error handling, testing, and documentation.


Strengths

1. Excellent Error Handling 🎯

The heartbeat handler correctly preserves infrastructure error types:

except RuntimeHostError:
    # Re-raise all infrastructure errors directly (preserves error type)
    raise

This allows callers to catch specific subtypes (InfraConnectionError, InfraTimeoutError, etc.) for differentiated handling - exactly as specified in CLAUDE.md error patterns.

2. Comprehensive Test Coverage ✅

27 integration tests across 8 test classes is outstanding. Tests cover:

  • Happy path (ACTIVE nodes)
  • Edge cases (node not found, non-ACTIVE states)
  • Concurrency scenarios
  • Error resilience (connection failures)
  • Timestamp accuracy verification

3. Strong Type Safety 💪

  • All models use proper Pydantic with model_config = ConfigDict(frozen=True, extra="forbid")
  • No Any types present
  • Consistent use of X | None over Optional[X] (PEP 604)
  • Result model (ModelHeartbeatHandlerResult) provides complete outcome information

4. Circuit Breaker Integration 🔒

Projector's update_heartbeat() method properly uses MixinAsyncCircuitBreaker:

async with self._circuit_breaker_lock:
    await self._check_circuit_breaker("update_heartbeat", corr_id)

Follows the caller-held lock pattern from CLAUDE.md.

5. Accurate Timestamp Tracking 📅

The PR includes detailed verification comment in model_node_liveness_expired.py:

  • detected_at: Fresh UTC from RuntimeTick.now
  • liveness_deadline: Sourced from projection
  • last_heartbeat_at: From heartbeat event timestamp

This addresses the core requirement of OMN-1006.


Observations

1. SQL Schema Migration (Informational)

schema_registration_projection.sql adds last_heartbeat_at TIMESTAMPTZ column. Ensure migration scripts handle:

  • Existing rows get NULL for last_heartbeat_at (acceptable)
  • No backward compatibility issues (ONEX policy: breaking changes OK)

2. Contract YAML - direct_handler Flag (Design Note)

The new direct_handler: true flag in contract.yaml:

- topic: "{env}.{namespace}.onex.evt.node-heartbeat.v1"
  event_type: "ModelNodeHeartbeatEvent"
  description: "Periodic heartbeat from active nodes for liveness tracking"
  direct_handler: true

This is a sensible pattern for events that bypass workflow coordination. Consider documenting this pattern in docs/patterns/ if it will be used by other orchestrators.

3. Race Condition - Heartbeat Between Read and Update (Low Risk)

In handler_node_heartbeat.py:

# Look up current projection
projection = await self._projection_reader.get_entity_state(...)

# ... later ...

# Update projection via projector
updated = await self._projector.update_heartbeat(...)

There's a TOCTOU (time-of-check-time-of-use) window where:

  1. Projection exists during read
  2. Projection is deleted before update
  3. update_heartbeat() returns False

Current mitigation: Handler catches this with if not updated and returns node_not_found=True.

Assessment: Low risk - nodes don't get deleted during normal operation. Current handling is acceptable. If this becomes a concern, consider using a transaction or optimistic locking.

4. Non-ACTIVE Node Warning (Behavioral Note)

Handler logs warning but still processes heartbeats from non-ACTIVE nodes:

if not projection.current_state.is_active():
    logger.warning(...)
    # Still process the heartbeat to update tracking, but log the warning

Rationale (from comment): "This can happen during state transitions or race conditions"

Assessment: Reasonable design - heartbeat updates are idempotent and race conditions during FSM transitions are expected in distributed systems.

5. Union Validation Threshold Update (Process Note)

infra_validators.py: INFRA_MAX_UNIONS increased from 589 → 600.

This is tracked and validated correctly. The pattern validator documentation in infra_validators.py already notes that threshold increases require PR justification (satisfied by OMN-1006 adding new models).


Performance Considerations

Database Impact

Each heartbeat triggers:

  1. SELECT via get_entity_state() (indexed on entity_id, domain)
  2. UPDATE via update_heartbeat() (indexed on primary key)

Expected load: For N active nodes with 30s heartbeat interval:

  • 10 nodes: ~0.67 QPS
  • 100 nodes: ~6.7 QPS
  • 1000 nodes: ~67 QPS

Optimization opportunity (future): Consider batch heartbeat processing or single-query upsert if heartbeat volume becomes significant. Current design is fine for MVP.


Security Review

✅ No security concerns detected

  • Correlation IDs properly tracked (no PII leakage)
  • Error messages sanitized (no credentials exposed)
  • SQL uses parameterized queries (injection-safe)
  • No authentication bypass vectors
  • Circuit breaker prevents resource exhaustion

Code Quality

Naming Conventions ✅

  • Handler: HandlerNodeHeartbeat (follows pattern)
  • Result model: ModelHeartbeatHandlerResult
  • Method: update_heartbeat() (clear, idiomatic)

Documentation ✅

  • Comprehensive docstrings with examples
  • Related tickets linked
  • Error conditions documented
  • Thread safety noted

ONEX Compliance ✅

  • Follows ONEX 4-node architecture (Orchestrator → Handler pattern)
  • Circuit breaker integration via mixin
  • Container-based DI ready (handler uses projector injection)
  • Error hierarchy compliance (RuntimeHostError subtypes)

Recommendations

1. Document direct_handler Pattern (Optional)

If direct_handler flag will be used in other contracts, create:

docs/patterns/direct_handler_pattern.md

Explaining when to use direct handlers vs workflow handlers.

2. Heartbeat Metrics (Follow-up)

Consider emitting metrics for:

  • Heartbeat processing latency
  • Node liveness state distribution
  • Circuit breaker open/close events

This would aid in production observability.

3. Integration Test Parallelization (Performance)

With 27 integration tests, consider using pytest-xdist for parallel execution:

pytest -n auto tests/integration/registration/handlers/

Summary

This PR implements heartbeat tracking with production-grade quality:

  • ✅ Robust error handling
  • ✅ Comprehensive test coverage
  • ✅ ONEX convention compliance
  • ✅ Proper circuit breaker integration
  • ✅ Accurate timestamp tracking
  • ✅ Clear documentation

The implementation is ready to merge. Minor observations above are informational or for future consideration.

Great work on the detailed timestamp accuracy verification and comprehensive integration testing! 🎉


Files Reviewed: 22 files, 1716 additions, 31 deletions
Test Coverage: 27 integration tests, existing unit tests updated
ONEX Compliance: Full compliance with CLAUDE.md patterns

@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

🧹 Nitpick comments (1)
tests/integration/registration/handlers/test_handler_node_heartbeat_integration.py (1)

737-738: Consider using a more specific exception type.

The test catches a generic Exception for the frozen model mutation check. For better precision, you could use pydantic.ValidationError since Pydantic frozen models raise this on mutation attempts.

🔎 Suggested improvement
+from pydantic import ValidationError
+
 # ...
 
         # Attempt to modify should fail
-        with pytest.raises(Exception):  # ValidationError for frozen models
+        with pytest.raises(ValidationError):
             result.success = False  # type: ignore[misc]
📜 Review details

Configuration used: defaults

Review profile: CHILL

Plan: Lite

📥 Commits

Reviewing files that changed from the base of the PR and between 6dfed0b and b9f1cf4.

📒 Files selected for processing (3)
  • src/omnibase_infra/nodes/node_registration_orchestrator/contract.yaml
  • tests/integration/registration/handlers/test_handler_node_heartbeat_integration.py
  • tests/unit/validation/test_validator_defaults.py
🚧 Files skipped from review as they are similar to previous changes (2)
  • tests/unit/validation/test_validator_defaults.py
  • src/omnibase_infra/nodes/node_registration_orchestrator/contract.yaml
🧰 Additional context used
📓 Path-based instructions (1)
**/*.py

📄 CodeRabbit inference engine (CLAUDE.md)

**/*.py: NEVER use Any types in Python - Always use specific types, use object for generic dispatchers instead
All data structures must be proper Pydantic models - one model per file
Use nullable type annotation X | None (PEP 604 union syntax) over Optional[X] for null types in Python
All services MUST use ModelONEXContainer for container-based dependency injection
Error classes must raise OnexError not base Exception - use raise OnexError(...) from e pattern
Never use isinstance for protocol resolution - use duck typing through protocols instead
NEVER include passwords, API keys, tokens, secrets, full connection strings with credentials, PII, or private keys in error messages or context

Files:

  • tests/integration/registration/handlers/test_handler_node_heartbeat_integration.py
🧠 Learnings (1)
📚 Learning: 2025-11-24T16:32:55.606Z
Learnt from: CR
Repo: OmniNode-ai/omninode_bridge PR: 0
File: .cursor/rules/node_standards.mdc:0-0
Timestamp: 2025-11-24T16:32:55.606Z
Learning: Applies to tests/unit/infrastructure/**/test_*.py : All node implementations must have comprehensive unit tests following the testing pattern in `tests/unit/infrastructure/` with tests for node initialization and node execution

Applied to files:

  • tests/integration/registration/handlers/test_handler_node_heartbeat_integration.py
🔇 Additional comments (10)
tests/integration/registration/handlers/test_handler_node_heartbeat_integration.py (10)

1-66: LGTM!

Well-structured module with comprehensive docstring, proper license header, and clean imports. The use of TYPE_CHECKING for type-only imports and the pytestmark configuration for integration tests are appropriate.


68-180: LGTM!

Clean test helper factories with:

  • Proper keyword-only arguments for clarity
  • PEP 604 union syntax (X | None) per coding guidelines
  • Sensible defaults for integration testing
  • Clear docstrings documenting parameters and return values

187-217: LGTM!

Initialization tests correctly verify default and custom liveness window configuration. The redundant assertion on line 202 (assert handler.liveness_window_seconds == 90.0) after line 201 serves as documentation of the expected default value, which is acceptable.


224-347: LGTM!

Comprehensive happy path tests covering:

  • Basic heartbeat processing for ACTIVE nodes
  • Liveness deadline extension calculations with appropriate tolerance
  • Correlation ID preservation and auto-generation
  • Selective field updates (ensuring other fields remain unchanged)

Good practice using both result object assertions and direct database state verification.


354-389: LGTM!

Well-structured tests for unknown node scenarios, verifying:

  • node_not_found flag is set correctly
  • Error message contains expected context
  • Correlation ID is preserved in error responses

396-440: LGTM!

Excellent use of parametrization to test all non-ACTIVE states. The test documents the design decision that heartbeats from non-ACTIVE nodes are processed (for tracking) with a warning logged, which handles race conditions during state transitions.


447-523: LGTM!

Thorough testing of liveness window deadline calculations:

  • Default 90-second window verification
  • Event timestamp-based calculation (not current time)
  • Progressive deadline extension with consecutive heartbeats

Good test design using past timestamps to verify the calculation is based on event time.


530-596: LGTM!

Well-designed concurrency tests that:

  • Verify parallel heartbeat processing from multiple nodes
  • Acknowledge non-deterministic ordering with concurrent writes (lines 588-595)
  • Use appropriate assertion strategy for race conditions (checking final value is one of expected values)

603-713: LGTM!

Comprehensive error handling tests covering:

  • InfraConnectionError propagation (lines 606-642)
  • Unexpected error wrapping in RuntimeHostError (lines 644-674)
  • Race condition handling when entity is deleted between lookup and update (lines 676-712)

Good use of patch.object with AsyncMock for simulating infrastructure failures.


769-884: LGTM!

Excellent database state verification tests that:

  • Directly query PostgreSQL to verify column updates
  • Ensure liveness_timeout_emitted_at markers are not reset (important for timeout emission logic)
  • Verify ack_deadline remains unchanged (heartbeat should only update heartbeat-specific fields)

These tests provide strong guarantees about the atomicity and selectivity of the heartbeat update operation.

Merged OMN-949 DLQ configuration changes with OMN-1006 heartbeat handler.

Conflict resolutions:
- INFRA_MAX_UNIONS: Combined both changes (606 + 11 = 617, set to 620 with buffer)
- Updated threshold history to reflect both merges

Changes from main (OMN-949):
- DLQ message format and configuration
- Topic constants and validation
- Error sanitization utilities
- DLQ integration tests
Resolves all PR #94 review issues including critical, major, minor, and nitpicks:

Exception handling improvements:
- Add explicit catches for InfraConnectionError, InfraTimeoutError,
  InfraAuthenticationError, InfraUnavailableError before generic Exception
- Re-raise infrastructure errors directly to preserve error types
- Only wrap truly unexpected errors in RuntimeHostError

Constant export consistency:
- Export DEFAULT_LIVENESS_WINDOW_SECONDS from orchestrators/__init__.py
- Replace hardcoded 90.0 values with constant in tests/conftest.py
- Update integration tests to use constant for liveness window

New timestamp accuracy tests:
- test_heartbeat_timestamp_matches_event_timestamp_exactly
- test_liveness_deadline_calculation_precision
- test_heartbeat_preserves_utc_timezone
- test_successive_heartbeats_update_timestamps_monotonically
- test_timestamp_accuracy_for_liveness_expired_event_reporting
- test_emit_liveness_expiration_includes_last_heartbeat_at
- test_emit_liveness_expiration_handles_none_last_heartbeat_at

Documentation fixes:
- Update test docstring to reflect current union threshold (620)
@claude

claude Bot commented Dec 25, 2025

Copy link
Copy Markdown

PR Review: Add last_heartbeat_at Tracking [OMN-1006]

Summary

This PR adds heartbeat timestamp tracking to the registration projection system, enabling accurate liveness expiration reporting. The implementation is well-designed, thoroughly tested, and follows ONEX infrastructure patterns correctly.


✅ Strengths

1. Excellent Architecture & Design

  • Clean separation of concerns: Handler → Projector → Database is well-structured
  • Proper error handling: Specific infrastructure error types are preserved (InfraConnectionError, InfraTimeoutError, etc.) without wrapping
  • Strong typing throughout: No Any types, proper use of | None syntax
  • Contract-driven approach: Heartbeat event properly declared in contract.yaml with direct_handler: true

2. Outstanding Test Coverage

  • 27 integration tests covering:
    • Happy path (ACTIVE nodes)
    • Error scenarios (node not found, connection failures)
    • Concurrency and race conditions
    • Timestamp accuracy and timezone handling
    • Liveness window calculations
  • Tests use real PostgreSQL via testcontainers (proper integration testing)
  • Comprehensive edge case coverage (rapid heartbeats, non-ACTIVE nodes, etc.)

3. Proper ONEX Patterns

  • ✅ Correct enum usage: EnumRegistrationState, EnumInfraTransportType
  • ✅ Correlation ID tracking throughout the call chain
  • ✅ Circuit breaker integration in projector (update_heartbeat method)
  • ✅ Proper file/class naming: handler_node_heartbeat.py → HandlerNodeHeartbeat
  • ✅ Model naming: ModelHeartbeatHandlerResult, ModelNodeHeartbeatEvent
  • ✅ Export consistency: DEFAULT_LIVENESS_WINDOW_SECONDS re-exported at all levels

4. Database Design

  • ✅ Schema properly updated with last_heartbeat_at TIMESTAMPTZ
  • ✅ Atomic update_heartbeat() method with proper SQL (UPDATE ... RETURNING)
  • ✅ Correct projection reader mapping (_row_to_projection)
  • ✅ Upsert logic includes new field in both INSERT and UPDATE clauses

5. Documentation Quality

  • Comprehensive docstrings with usage examples
  • Thread safety notes where relevant
  • Design rationale clearly explained (e.g., timestamp accuracy verification)
  • Related tickets properly referenced

🔍 Code Quality Observations

Exception Handling (Lines 289-321)

EXCELLENT - The error handling pattern is textbook:

except (
    InfraConnectionError,
    InfraTimeoutError,
    InfraAuthenticationError,
    InfraUnavailableError,
):
    raise  # Preserve specific error types
except RuntimeHostError:
    raise  # Preserve other infrastructure errors
except Exception as e:
    # Only wrap truly unexpected errors
    raise RuntimeHostError(...) from e

This preserves error type fidelity while preventing unexpected exceptions from escaping. Callers can catch specific error types for differentiated handling.

Timestamp Handling (Lines 234-236)

heartbeat_timestamp = event.timestamp
new_liveness_deadline = heartbeat_timestamp + timedelta(
    seconds=self._liveness_window_seconds
)

CORRECT - Uses event timestamp directly (accurate source of truth). The detected_at timestamp in ModelNodeLivenessExpired comes from RuntimeTick.now, ensuring separate concerns are properly tracked.

SQL Update Method (projector_registration.py:773-813)

The update_heartbeat() method is well-designed:

  • ✅ Circuit breaker check before operation
  • ✅ Atomic UPDATE with RETURNING clause
  • ✅ Proper connection pool usage
  • ✅ Success/failure logging with correlation_id
  • ✅ Returns bool for entity existence check

Topic Naming (contract.yaml)

- topic: "{env}.{namespace}.onex.evt.node-heartbeat.v1"
  event_type: "ModelNodeHeartbeatEvent"
  description: "Periodic heartbeat from active nodes for liveness tracking"
  direct_handler: true

PERFECT - Follows ONEX topic convention: {env}.{namespace}.onex.evt.{event-name}.{version}

The direct_handler: true flag is an excellent design choice, documenting that this event bypasses workflow processing.


🎯 Minor Observations (No Action Required)

1. Validation Threshold Update

INFRA_MAX_UNIONS = 620  # Was 589, increased by 31

Rationale Documented: The increase is from merging OMN-949 (DLQ config) + OMN-1006 (heartbeat handler). The comment history shows clear tracking:

  • OMN-949: +11 unions (DLQ)
  • OMN-1006: ~11 unions (heartbeat)
  • Buffer added to prevent CI churn

This is acceptable as an infrastructure pattern with a documented migration plan to reduce unions via JsonValue (see CLAUDE.md pattern docs).

2. Constant Export Hierarchy

The DEFAULT_LIVENESS_WINDOW_SECONDS is exported through 3 levels:

  • handlers/handler_node_heartbeat.py
  • handlers/__init__.py
  • registration/__init__.py
  • orchestrators/__init__.py

Good design - This provides flexible import paths while maintaining a single source of truth. Users can import from the most convenient level.

3. Non-ACTIVE Node Handling (Lines 222-228)

if not projection.current_state.is_active():
    logger.warning(...)
    # Still process the heartbeat to update tracking, but log the warning

Pragmatic choice - Processes heartbeats even for non-ACTIVE nodes to handle race conditions gracefully. This prevents dropped heartbeats during state transitions.


🔒 Security Review

✅ Error Sanitization

All error context properly excludes sensitive data:

  • No credentials or secrets in error messages
  • Correlation IDs included (safe for tracing)
  • Operation names generic ("update_heartbeat", "handle")
  • Node IDs included (non-sensitive identifiers)

✅ SQL Injection Prevention

All queries use parameterized statements:

await conn.fetchrow(
    update_sql,
    entity_id,      # 
    domain,         # 
    last_heartbeat_at,  # 
    liveness_deadline,  # 
)

✅ UUID Validation

UUIDs are strongly typed throughout - no string manipulation vulnerabilities.


📊 Performance Considerations

Database Impact

POSITIVE:

  • last_heartbeat_at column added with minimal index overhead
  • UPDATE query is targeted (WHERE entity_id AND domain) - uses primary key
  • No new indexes required (heartbeat timestamp not queried for filtering)

Heartbeat Write Load:

  • Default 30-second heartbeat interval (from MixinNodeIntrospection)
  • 90-second liveness window (3x heartbeat interval)
  • For 1000 active nodes: ~33 writes/second to registration_projections table
  • This is well within PostgreSQL's capabilities for simple UPDATEs

Concurrency

The integration tests verify concurrent heartbeats are handled correctly (TestHandlerNodeHeartbeatConcurrency). PostgreSQL's row-level locking ensures atomic updates.


🧪 Test Quality Analysis

Integration Test Highlights

From test_handler_node_heartbeat_integration.py:

8 test classes, 27+ test methods:

  1. TestHandlerNodeHeartbeatInit - Handler initialization
  2. TestHandlerNodeHeartbeatHappyPath - Normal operation
  3. TestHandlerNodeHeartbeatNotFound - Unknown nodes
  4. TestHandlerNodeHeartbeatNonActiveNode - State validation
  5. TestHandlerNodeHeartbeatLivenessWindow - Window calculations
  6. TestHandlerNodeHeartbeatConcurrency - Race conditions
  7. TestHandlerNodeHeartbeatErrors - Error scenarios
  8. TestHandlerNodeHeartbeatDatabaseState - State verification
  9. TestHandlerNodeHeartbeatTimestampAccuracy - Timestamp precision

Particularly Strong Tests:

  • test_heartbeat_timestamp_matches_event_timestamp_exactly - Verifies timestamp fidelity
  • test_rapid_heartbeats_extend_deadline_correctly - Tests race conditions
  • test_concurrent_heartbeats_different_nodes - Multi-node concurrency
  • test_database_connection_error_propagation - Error type preservation

Test Fixture Quality

conftest.py provides comprehensive fixtures:

  • Real PostgreSQL via testcontainers
  • Proper async lifecycle management
  • Isolated test database per run
  • Helper functions for projection creation

🎨 Code Style Compliance

ONEX Naming Conventions: ✅ PERFECT

  • ✅ File: handler_node_heartbeat.py
  • ✅ Class: HandlerNodeHeartbeat
  • ✅ Model: ModelHeartbeatHandlerResult
  • ✅ Constant: DEFAULT_LIVENESS_WINDOW_SECONDS (caps with underscores)

Type Annotations: ✅ EXCELLENT

  • ✅ No Any types (ONEX zero-tolerance rule)
  • ✅ Proper use of | None (PEP 604 union syntax)
  • ✅ TYPE_CHECKING guard for imports (prevents circular imports)
  • ✅ All function signatures fully typed

Documentation: ✅ COMPREHENSIVE

  • ✅ Module docstrings with ticket references
  • ✅ Class docstrings with usage examples
  • ✅ Method docstrings with Args/Returns/Raises
  • ✅ Inline comments for complex logic

🚀 Deployment Readiness

Schema Migration

The SQL changes are idempotent and backward-compatible:

  • ALTER TABLE can be applied to existing databases safely
  • New column is nullable, so existing rows don't break
  • No data migration required

Rollout Strategy

  1. Deploy schema change (add last_heartbeat_at column)
  2. Deploy application code (heartbeat handler + projector update)
  3. Wire heartbeat handler in orchestrator initialization

Zero-downtime deployment: ✅ Possible

Monitoring Recommendations

Consider adding metrics for:

  • Heartbeat processing latency (p50, p95, p99)
  • Heartbeat update failure rate
  • Nodes with stale last_heartbeat_at (>5 minutes old)

📝 Final Verdict

Overall Assessment: EXCELLENT ✨

This is a high-quality PR that demonstrates:

  • ✅ Deep understanding of ONEX architecture patterns
  • ✅ Thorough testing methodology (27+ integration tests)
  • ✅ Proper error handling and type safety
  • ✅ Clear documentation and design rationale
  • ✅ Production-ready code with no critical issues

Recommendation: APPROVE & MERGE 🚢

This PR is ready for production deployment. The code quality, test coverage, and adherence to ONEX patterns are exemplary.


🏆 Notable Achievements

  1. Exception Handling Excellence: Preserving specific error types while preventing unexpected exceptions
  2. Test Coverage: 27+ integration tests against real PostgreSQL
  3. Documentation Quality: Timestamp accuracy verification explicitly documented
  4. Zero Technical Debt: No shortcuts, workarounds, or TODOs left behind

Great work on this implementation! 🎉


Additional Notes

  • All acceptance criteria met ✅
  • Follow-up work properly documented (heartbeat event wiring)
  • Linear issue OMN-1006 can be closed upon merge
  • No blocking issues or concerns

Reviewed against ONEX CLAUDE.md standards - Full compliance verified.

@jonahgabriel
jonahgabriel merged commit 4ba42fc into main Dec 25, 2025
6 checks passed
@jonahgabriel
jonahgabriel deleted the jonah/omn-1006-track-last_heartbeat_at-in-registration-projection branch December 25, 2025 21:27
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