Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions src/omnibase_infra/idempotency/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -73,10 +73,15 @@
ModelIdempotencyStoreMetrics,
ModelPostgresIdempotencyStoreConfig,
)
from omnibase_infra.idempotency.protocol_idempotency_store import (
ProtocolIdempotencyStore,
)
from omnibase_infra.idempotency.store_inmemory import InMemoryIdempotencyStore
from omnibase_infra.idempotency.store_postgres import PostgresIdempotencyStore

__all__ = [
# Protocol
"ProtocolIdempotencyStore",
# Models
"ModelIdempotencyCheckResult",
"ModelIdempotencyGuardConfig",
Expand Down
181 changes: 181 additions & 0 deletions src/omnibase_infra/idempotency/protocol_idempotency_store.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,181 @@
# SPDX-License-Identifier: MIT
# Copyright (c) 2025 OmniNode Team
"""Protocol definition for Idempotency Store.

This module defines the ProtocolIdempotencyStore protocol that all idempotency
store implementations must follow. The protocol defines the contract for
message deduplication in distributed systems.

Migration Note:
This protocol is defined locally in omnibase_infra because it is not
available in omnibase_spi versions 0.4.0/0.4.1. This is a TEMPORARY
definition that should be migrated to omnibase_spi in a future release.

Migration Path (OMN-1000):
1. When omnibase_spi 0.5.0+ is released with ProtocolIdempotencyStore,
update pyproject.toml to require the new version
2. Update imports in InMemoryIdempotencyStore and PostgresIdempotencyStore
to use: `from omnibase_spi.protocols import ProtocolIdempotencyStore`
3. Remove this local protocol definition
4. Run tests to verify compatibility

The protocol contract is intentionally designed to match the expected
omnibase_spi interface to ensure a smooth migration.

Protocol Methods:
- check_and_record: Atomically check if message was processed and record if not
- is_processed: Check if a message was already processed (read-only)
- mark_processed: Mark a message as processed (upsert)
- cleanup_expired: Remove entries older than TTL

Implementations:
- InMemoryIdempotencyStore: In-memory store for testing (OMN-945)
- PostgresIdempotencyStore: Production PostgreSQL store (OMN-945)

Security Considerations:
- Thread Safety: All implementations MUST be safe for concurrent access.
Multiple coroutines may call check_and_record simultaneously with the
same message_id. Implementations must use appropriate synchronization
(e.g., asyncio.Lock for in-memory, database transactions for PostgreSQL).

- Atomicity: The check_and_record method MUST provide atomic check-and-set
semantics. When multiple callers race with the same (domain, message_id),
exactly ONE caller must receive True. This prevents duplicate processing
in concurrent scenarios.

- Domain Isolation: Messages are namespaced by domain to prevent cross-tenant
conflicts. The (domain, message_id) tuple forms the unique key. Different
domains can safely use the same message_id without collision. This is
critical for multi-tenant deployments where tenant data must be isolated.

- Correlation ID Usage: The correlation_id parameter is used ONLY for
distributed tracing and observability. It is NOT used for authentication
or authorization. Security-sensitive operations must implement their own
authentication layer; do not rely on correlation_id for access control.

- Input Validation: Implementations should validate that message_id is a
valid UUID. Domain names should be sanitized if used in storage backends
(e.g., as part of database keys or Redis key prefixes).
"""

from __future__ import annotations

from datetime import datetime
from typing import Protocol, runtime_checkable
from uuid import UUID


@runtime_checkable
class ProtocolIdempotencyStore(Protocol):
"""Protocol for idempotency store implementations.

Defines the contract for message deduplication stores that track processed
messages and prevent duplicate processing in distributed systems.

All implementations must provide atomic check-and-record semantics to
ensure exactly-once processing guarantees.

Key Properties:
- Thread-safe: All operations must be safe for concurrent access
- Atomic: check_and_record must provide atomic check-and-set semantics
- Domain-isolated: Messages can be namespaced by domain for isolated deduplication

Example:
>>> store: ProtocolIdempotencyStore = InMemoryIdempotencyStore()
>>> message_id = uuid4()
>>> is_new = await store.check_and_record(message_id, domain="orders")
>>> if is_new:
... # Process the message
... pass
"""

async def check_and_record(
self,
message_id: UUID,
domain: str | None = None,
correlation_id: UUID | None = None,
) -> bool:
"""Atomically check if message was processed and record if not.

This is the primary idempotency operation. It must be atomic to ensure
that when multiple coroutines call this method simultaneously with the
same (domain, message_id), exactly ONE caller receives True.

Args:
message_id: Unique identifier for the message.
domain: Optional domain namespace for isolated deduplication.
Messages with the same message_id but different domains are
treated as distinct messages.
correlation_id: Optional correlation ID for distributed tracing.
Stored with the record for observability purposes.

Returns:
True if message is new (should be processed).
False if message is duplicate (should be skipped).
"""
...

async def is_processed(
self,
message_id: UUID,
domain: str | None = None,
) -> bool:
"""Check if a message was already processed.

Read-only check that does not modify the store. Useful for querying
message status without affecting the idempotency state.

Args:
message_id: Unique identifier for the message.
domain: Optional domain namespace.

Returns:
True if the message has been processed.
False if the message has not been processed.
"""
...

async def mark_processed(
self,
message_id: UUID,
domain: str | None = None,
correlation_id: UUID | None = None,
processed_at: datetime | None = None,
) -> None:
"""Mark a message as processed.

Records a message as processed without checking if it already exists.
If the record already exists, updates it with the new values.

This is an upsert operation - it will create a new record if one
doesn't exist, or update the existing record if it does.

Args:
message_id: Unique identifier for the message.
domain: Optional domain namespace for isolated deduplication.
correlation_id: Optional correlation ID for tracing.
processed_at: Optional timestamp of when processing occurred.
If None, implementations should use the current UTC time.
"""
...

async def cleanup_expired(
self,
ttl_seconds: int,
) -> int:
"""Remove entries older than TTL.

Cleans up old idempotency records based on their processed_at timestamp.
This prevents unbounded storage growth.

Args:
ttl_seconds: Time-to-live in seconds. Records older than this
value (based on processed_at timestamp) are removed.

Returns:
Number of entries removed.
"""
...


__all__ = ["ProtocolIdempotencyStore"]
5 changes: 2 additions & 3 deletions src/omnibase_infra/idempotency/store_inmemory.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,12 +20,11 @@
from datetime import UTC, datetime
from uuid import UUID

from omnibase_spi.protocols.storage.protocol_idempotency_store import (
from omnibase_infra.idempotency.models import ModelIdempotencyRecord
from omnibase_infra.idempotency.protocol_idempotency_store import (
ProtocolIdempotencyStore,
)

from omnibase_infra.idempotency.models import ModelIdempotencyRecord


class InMemoryIdempotencyStore(ProtocolIdempotencyStore):
"""In-memory idempotency store for testing.
Expand Down
6 changes: 3 additions & 3 deletions src/omnibase_infra/idempotency/store_postgres.py
Original file line number Diff line number Diff line change
Expand Up @@ -67,9 +67,6 @@
from uuid import UUID, uuid4

import asyncpg
from omnibase_spi.protocols.storage.protocol_idempotency_store import (
ProtocolIdempotencyStore,
)

from omnibase_infra.enums import EnumInfraTransportType
from omnibase_infra.errors import (
Expand All @@ -84,6 +81,9 @@
ModelIdempotencyStoreMetrics,
ModelPostgresIdempotencyStoreConfig,
)
from omnibase_infra.idempotency.protocol_idempotency_store import (
ProtocolIdempotencyStore,
)

logger = logging.getLogger(__name__)

Expand Down
3 changes: 2 additions & 1 deletion src/omnibase_infra/mixins/mixin_node_introspection.py
Original file line number Diff line number Diff line change
Expand Up @@ -1415,7 +1415,8 @@ async def _publish_heartbeat(self) -> bool:
node_id=node_id,
node_type=node_type,
uptime_seconds=uptime_seconds,
# TODO(OMN-XXX): Implement active operation tracking
# TODO(ACTIVE-OP-TRACKING): Implement active operation tracking
# Ticket: Create Linear ticket for active operation tracking implementation
# Currently hardcoded to 0. Full implementation requires:
# - Operation counter increment/decrement around async operations
# - Thread-safe counter for concurrent operations
Expand Down
53 changes: 43 additions & 10 deletions src/omnibase_infra/runtime/runtime_host_process.py
Original file line number Diff line number Diff line change
Expand Up @@ -64,12 +64,12 @@ async def main() -> None:
if TYPE_CHECKING:
from omnibase_core.types import JsonValue
from omnibase_spi.protocols.handlers.protocol_handler import ProtocolHandler
from omnibase_spi.protocols.storage.protocol_idempotency_store import (
ProtocolIdempotencyStore,
)

from omnibase_infra.event_bus.models import ModelEventMessage
from omnibase_infra.idempotency import ModelIdempotencyGuardConfig
from omnibase_infra.idempotency.protocol_idempotency_store import (
ProtocolIdempotencyStore,
)

# Expose wire_default_handlers as wire_handlers for test patching compatibility
# Tests patch "omnibase_infra.runtime.runtime_host_process.wire_handlers"
Expand Down Expand Up @@ -1364,6 +1364,18 @@ async def _initialize_idempotency_store(self) -> None:
self._idempotency_store = None
self._idempotency_config = None

# =========================================================================
# WARNING: FAIL-OPEN BEHAVIOR
# =========================================================================
# This method implements FAIL-OPEN semantics: if the idempotency store
# is unavailable or errors, messages are ALLOWED THROUGH for processing.
#
# This is an intentional design decision prioritizing availability over
# exactly-once guarantees. See docstring below for full trade-off analysis.
#
# IMPORTANT: Downstream handlers MUST be designed for at-least-once delivery
# and implement their own idempotency for critical operations.
# =========================================================================
async def _check_idempotency(
self,
envelope: dict[str, object],
Expand All @@ -1375,10 +1387,27 @@ async def _check_idempotency(
idempotency store. If duplicate detected, publishes a duplicate
response and returns False.

Fail-Open Behavior:
If the idempotency store is unavailable or throws an error,
the message is allowed through (logged with warning). This
prioritizes availability over exactly-once semantics.
Fail-Open Semantics:
This method implements **fail-open** error handling: if the
idempotency store is unavailable or throws an error, the message
is allowed through for processing (with a warning log).

**Design Rationale**: In distributed event-driven systems, the
idempotency store (e.g., Redis/Valkey) is a supporting service,
not a critical path dependency. A temporary store outage should
not halt message processing entirely, as this would cascade into
broader system unavailability.

**Trade-offs**:
- Pro: High availability - processing continues during store outages
- Pro: Graceful degradation - system remains functional
- Con: May result in duplicate message processing during outages
- Con: Downstream handlers must be designed for at-least-once delivery

**Mitigation**: Handlers consuming messages should implement their
own idempotency logic for critical operations (e.g., using database
constraints or transaction guards) to ensure correctness even when
duplicates slip through.

Args:
envelope: Validated envelope dict.
Expand Down Expand Up @@ -1437,6 +1466,7 @@ async def _check_idempotency(
message_id=message_id,
correlation_id=correlation_id,
)
# duplicate_response is already a dict from _create_duplicate_response
await self._publish_envelope_safe(
duplicate_response, self._output_topic
)
Expand All @@ -1445,12 +1475,15 @@ async def _check_idempotency(
return True

except Exception as e:
# Idempotency check failure - log and allow processing
# (fail-open for availability, may result in duplicate processing)
# FAIL-OPEN: Allow message through on idempotency store errors.
# Rationale: Availability over exactly-once. Store outages should not
# halt processing. Downstream handlers must tolerate duplicates.
# See docstring for full trade-off analysis.
logger.warning(
"Idempotency check failed, allowing message through",
"Idempotency check failed, allowing message through (fail-open)",
extra={
"error": str(e),
"error_type": type(e).__name__,
"message_id": str(message_id),
"domain": domain,
"correlation_id": str(correlation_id),
Expand Down
8 changes: 6 additions & 2 deletions src/omnibase_infra/validation/infra_validators.py
Original file line number Diff line number Diff line change
Expand Up @@ -327,7 +327,7 @@ def get_union_exemptions() -> list[ExemptionPattern]:
# This is a COUNT threshold, not a violation threshold. The validator counts all
# unions including the ONEX-preferred `X | None` patterns, which are valid.
#
# Current baseline (515 unions as of 2025-12-22):
# Current baseline (544 unions as of 2025-12-23):
# - Most unions are legitimate `X | None` nullable patterns
# - These are NOT flagged as violations, just counted
# - Actual violations (primitive soup, Union[X,None] syntax) are reported separately
Expand All @@ -336,9 +336,13 @@ def get_union_exemptions() -> list[ExemptionPattern]:
# - 491 (2025-12-21): Initial baseline with DispatcherFunc | ContextAwareDispatcherFunc
# - 515 (2025-12-22): OMN-990 MessageDispatchEngine + OMN-947 snapshots (~24 unions added)
# - 540 (2025-12-23): OMN-950 comprehensive reducer tests (~25 unions from type annotations)
# - 544 (2025-12-23): OMN-954 effect idempotency and retry tests (PR #78) (~14 unions added)
# - nodes/effects/ module with protocol and model definitions
# - Legitimate X | None nullable patterns for optional fields
#
# Threshold: 555 (11 buffer above 544 baseline for codebase growth)
# Target: Reduce to <200 through dict[str, object] -> JsonValue migration.
INFRA_MAX_UNIONS = 540
INFRA_MAX_UNIONS = 555

# Maximum allowed architecture violations in infrastructure code.
# Set to 0 (strict enforcement) to ensure one-model-per-file principle is always followed.
Expand Down
Loading