diff --git a/src/omnibase_infra/idempotency/__init__.py b/src/omnibase_infra/idempotency/__init__.py index 10e3741b0c..c83af49a6b 100644 --- a/src/omnibase_infra/idempotency/__init__.py +++ b/src/omnibase_infra/idempotency/__init__.py @@ -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", diff --git a/src/omnibase_infra/idempotency/protocol_idempotency_store.py b/src/omnibase_infra/idempotency/protocol_idempotency_store.py new file mode 100644 index 0000000000..acf5c0bbdf --- /dev/null +++ b/src/omnibase_infra/idempotency/protocol_idempotency_store.py @@ -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"] diff --git a/src/omnibase_infra/idempotency/store_inmemory.py b/src/omnibase_infra/idempotency/store_inmemory.py index 403cd5fe1b..f5490b2547 100644 --- a/src/omnibase_infra/idempotency/store_inmemory.py +++ b/src/omnibase_infra/idempotency/store_inmemory.py @@ -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. diff --git a/src/omnibase_infra/idempotency/store_postgres.py b/src/omnibase_infra/idempotency/store_postgres.py index a569d0ef8b..7586667640 100644 --- a/src/omnibase_infra/idempotency/store_postgres.py +++ b/src/omnibase_infra/idempotency/store_postgres.py @@ -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 ( @@ -84,6 +81,9 @@ ModelIdempotencyStoreMetrics, ModelPostgresIdempotencyStoreConfig, ) +from omnibase_infra.idempotency.protocol_idempotency_store import ( + ProtocolIdempotencyStore, +) logger = logging.getLogger(__name__) diff --git a/src/omnibase_infra/mixins/mixin_node_introspection.py b/src/omnibase_infra/mixins/mixin_node_introspection.py index c240f2d6db..2093868f1e 100644 --- a/src/omnibase_infra/mixins/mixin_node_introspection.py +++ b/src/omnibase_infra/mixins/mixin_node_introspection.py @@ -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 diff --git a/src/omnibase_infra/runtime/runtime_host_process.py b/src/omnibase_infra/runtime/runtime_host_process.py index 70b12f76a2..6344fa0491 100644 --- a/src/omnibase_infra/runtime/runtime_host_process.py +++ b/src/omnibase_infra/runtime/runtime_host_process.py @@ -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" @@ -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], @@ -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. @@ -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 ) @@ -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), diff --git a/src/omnibase_infra/validation/infra_validators.py b/src/omnibase_infra/validation/infra_validators.py index 8b369a36eb..d68bafd123 100644 --- a/src/omnibase_infra/validation/infra_validators.py +++ b/src/omnibase_infra/validation/infra_validators.py @@ -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 @@ -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. diff --git a/tests/unit/idempotency/test_protocol_idempotency_store.py b/tests/unit/idempotency/test_protocol_idempotency_store.py new file mode 100644 index 0000000000..807b6f3aea --- /dev/null +++ b/tests/unit/idempotency/test_protocol_idempotency_store.py @@ -0,0 +1,453 @@ +# SPDX-License-Identifier: MIT +# Copyright (c) 2025 OmniNode Team +"""Unit tests for ProtocolIdempotencyStore protocol definition. + +Verifies that the protocol is correctly defined and that all implementations +conform to the protocol contract. + +This file tests: +1. Protocol definition correctness (runtime_checkable, required methods) +2. Protocol method signatures (parameter types, return types) +3. Implementation conformance for InMemoryIdempotencyStore and PostgresIdempotencyStore + +Ticket: OMN-945 +""" + +from __future__ import annotations + +import inspect +from datetime import datetime +from typing import get_type_hints +from uuid import UUID + +import pytest + +from omnibase_infra.idempotency import ( + InMemoryIdempotencyStore, + PostgresIdempotencyStore, + ProtocolIdempotencyStore, +) + + +class TestProtocolDefinition: + """Tests for ProtocolIdempotencyStore definition.""" + + def test_protocol_is_runtime_checkable(self) -> None: + """Protocol should be decorated with @runtime_checkable. + + The protocol uses typing.runtime_checkable decorator, which adds + __protocol_attrs__ to the class for isinstance() checking support. + """ + # runtime_checkable protocols have _is_runtime_protocol attribute + assert hasattr(ProtocolIdempotencyStore, "_is_runtime_protocol") + assert ProtocolIdempotencyStore._is_runtime_protocol is True + + def test_protocol_has_required_methods(self) -> None: + """Protocol should define all required methods. + + The protocol must define: + - check_and_record: Atomic check-and-set for idempotency + - is_processed: Read-only check + - mark_processed: Upsert operation + - cleanup_expired: TTL-based cleanup + """ + required_methods = [ + "check_and_record", + "is_processed", + "mark_processed", + "cleanup_expired", + ] + for method in required_methods: + assert hasattr(ProtocolIdempotencyStore, method), ( + f"Protocol missing method: {method}" + ) + + def test_protocol_methods_are_async(self) -> None: + """All protocol methods should be coroutine functions. + + The protocol is designed for async I/O operations, so all methods + must be async (coroutine functions). + """ + async_methods = [ + "check_and_record", + "is_processed", + "mark_processed", + "cleanup_expired", + ] + for method_name in async_methods: + method = getattr(ProtocolIdempotencyStore, method_name) + assert inspect.iscoroutinefunction(method), f"{method_name} should be async" + + +class TestProtocolMethodSignatures: + """Tests for protocol method signatures.""" + + def test_check_and_record_signature(self) -> None: + """check_and_record should have correct parameter and return types. + + Expected signature: + async def check_and_record( + self, + message_id: UUID, + domain: str | None = None, + correlation_id: UUID | None = None, + ) -> bool + """ + method = ProtocolIdempotencyStore.check_and_record + sig = inspect.signature(method) + params = sig.parameters + + # Check parameter names exist + assert "message_id" in params + assert "domain" in params + assert "correlation_id" in params + + # Check default values + assert params["domain"].default is None + assert params["correlation_id"].default is None + + # Check return annotation + hints = get_type_hints(method) + assert hints.get("return") is bool + + def test_is_processed_signature(self) -> None: + """is_processed should have correct parameter and return types. + + Expected signature: + async def is_processed( + self, + message_id: UUID, + domain: str | None = None, + ) -> bool + """ + method = ProtocolIdempotencyStore.is_processed + sig = inspect.signature(method) + params = sig.parameters + + # Check parameter names exist + assert "message_id" in params + assert "domain" in params + + # Check default values + assert params["domain"].default is None + + # Check return annotation + hints = get_type_hints(method) + assert hints.get("return") is bool + + def test_mark_processed_signature(self) -> None: + """mark_processed should have correct parameter and return types. + + Expected signature: + async def mark_processed( + self, + message_id: UUID, + domain: str | None = None, + correlation_id: UUID | None = None, + processed_at: datetime | None = None, + ) -> None + """ + method = ProtocolIdempotencyStore.mark_processed + sig = inspect.signature(method) + params = sig.parameters + + # Check parameter names exist + assert "message_id" in params + assert "domain" in params + assert "correlation_id" in params + assert "processed_at" in params + + # Check default values + assert params["domain"].default is None + assert params["correlation_id"].default is None + assert params["processed_at"].default is None + + # Check return annotation + hints = get_type_hints(method) + assert hints.get("return") is type(None) + + def test_cleanup_expired_signature(self) -> None: + """cleanup_expired should have correct parameter and return types. + + Expected signature: + async def cleanup_expired( + self, + ttl_seconds: int, + ) -> int + """ + method = ProtocolIdempotencyStore.cleanup_expired + sig = inspect.signature(method) + params = sig.parameters + + # Check parameter names exist + assert "ttl_seconds" in params + + # ttl_seconds should be required (no default) + assert params["ttl_seconds"].default is inspect.Parameter.empty + + # Check return annotation + hints = get_type_hints(method) + assert hints.get("return") is int + + +class TestProtocolConformance: + """Tests for implementation conformance to ProtocolIdempotencyStore.""" + + def test_inmemory_conforms_to_protocol(self) -> None: + """InMemoryIdempotencyStore should conform to ProtocolIdempotencyStore. + + Verifies that isinstance() check passes for InMemoryIdempotencyStore, + which requires the class to implement all protocol methods with + compatible signatures. + """ + store = InMemoryIdempotencyStore() + assert isinstance(store, ProtocolIdempotencyStore) + + def test_postgres_conforms_to_protocol(self) -> None: + """PostgresIdempotencyStore should conform to ProtocolIdempotencyStore. + + Verifies that isinstance() check passes for PostgresIdempotencyStore. + Note: This test uses a mock configuration since we're only testing + protocol conformance, not actual database connectivity. + """ + from omnibase_infra.idempotency import ModelPostgresIdempotencyStoreConfig + + config = ModelPostgresIdempotencyStoreConfig( + dsn="postgresql://user:pass@localhost:5432/testdb", + ) + store = PostgresIdempotencyStore(config) + assert isinstance(store, ProtocolIdempotencyStore) + + def test_inmemory_has_all_protocol_methods(self) -> None: + """InMemoryIdempotencyStore should implement all protocol methods. + + Verifies that all required protocol methods are present and callable. + """ + store = InMemoryIdempotencyStore() + required_methods = [ + "check_and_record", + "is_processed", + "mark_processed", + "cleanup_expired", + ] + for method_name in required_methods: + assert hasattr(store, method_name), f"Missing method: {method_name}" + method = getattr(store, method_name) + assert callable(method), f"Method {method_name} is not callable" + assert inspect.iscoroutinefunction(method), ( + f"Method {method_name} should be async" + ) + + def test_postgres_has_all_protocol_methods(self) -> None: + """PostgresIdempotencyStore should implement all protocol methods. + + Verifies that all required protocol methods are present and callable. + """ + from omnibase_infra.idempotency import ModelPostgresIdempotencyStoreConfig + + config = ModelPostgresIdempotencyStoreConfig( + dsn="postgresql://user:pass@localhost:5432/testdb", + ) + store = PostgresIdempotencyStore(config) + required_methods = [ + "check_and_record", + "is_processed", + "mark_processed", + "cleanup_expired", + ] + for method_name in required_methods: + assert hasattr(store, method_name), f"Missing method: {method_name}" + method = getattr(store, method_name) + assert callable(method), f"Method {method_name} is not callable" + assert inspect.iscoroutinefunction(method), ( + f"Method {method_name} should be async" + ) + + +class TestProtocolTypeAnnotations: + """Tests for protocol type annotations correctness.""" + + @staticmethod + def _is_optional_type(hint: type, expected_type: type) -> bool: + """Check if a type hint is expected_type | None. + + In Python 3.10+, `X | None` creates a types.UnionType, not typing.Union. + This helper handles both cases for compatibility. + + Args: + hint: The type hint to check. + expected_type: The expected non-None type (e.g., str, UUID, datetime). + + Returns: + True if hint is expected_type | None, False otherwise. + """ + import types + + # Check for Python 3.10+ union type (X | None) + if isinstance(hint, types.UnionType): + args = hint.__args__ + return len(args) == 2 and expected_type in args and type(None) in args + + # Check for typing.Union (Optional[X]) + if hasattr(hint, "__origin__") and hasattr(hint, "__args__"): + from typing import Union + + if hint.__origin__ is Union: + args = hint.__args__ + return len(args) == 2 and expected_type in args and type(None) in args + + return False + + def test_check_and_record_type_hints(self) -> None: + """check_and_record type hints should be correctly defined.""" + hints = get_type_hints(ProtocolIdempotencyStore.check_and_record) + + # message_id should be UUID + assert hints["message_id"] is UUID + + # domain should be str | None + assert self._is_optional_type(hints["domain"], str), ( + f"Expected str | None, got {hints['domain']}" + ) + + # correlation_id should be UUID | None + assert self._is_optional_type(hints["correlation_id"], UUID), ( + f"Expected UUID | None, got {hints['correlation_id']}" + ) + + # Return type should be bool + assert hints["return"] is bool + + def test_is_processed_type_hints(self) -> None: + """is_processed type hints should be correctly defined.""" + hints = get_type_hints(ProtocolIdempotencyStore.is_processed) + + # message_id should be UUID + assert hints["message_id"] is UUID + + # domain should be str | None + assert self._is_optional_type(hints["domain"], str), ( + f"Expected str | None, got {hints['domain']}" + ) + + # Return type should be bool + assert hints["return"] is bool + + def test_mark_processed_type_hints(self) -> None: + """mark_processed type hints should be correctly defined.""" + hints = get_type_hints(ProtocolIdempotencyStore.mark_processed) + + # message_id should be UUID + assert hints["message_id"] is UUID + + # domain should be str | None + assert self._is_optional_type(hints["domain"], str), ( + f"Expected str | None, got {hints['domain']}" + ) + + # correlation_id should be UUID | None + assert self._is_optional_type(hints["correlation_id"], UUID), ( + f"Expected UUID | None, got {hints['correlation_id']}" + ) + + # processed_at should be datetime | None + assert self._is_optional_type(hints["processed_at"], datetime), ( + f"Expected datetime | None, got {hints['processed_at']}" + ) + + # Return type should be None + assert hints["return"] is type(None) + + def test_cleanup_expired_type_hints(self) -> None: + """cleanup_expired type hints should be correctly defined.""" + hints = get_type_hints(ProtocolIdempotencyStore.cleanup_expired) + + # ttl_seconds should be int + assert hints["ttl_seconds"] is int + + # Return type should be int + assert hints["return"] is int + + +class TestNonConformingImplementation: + """Tests for classes that should NOT conform to the protocol.""" + + def test_empty_class_does_not_conform(self) -> None: + """An empty class should not pass isinstance check. + + Verifies that the protocol properly rejects classes that don't + implement the required methods. + """ + + class EmptyStore: + pass + + store = EmptyStore() + assert not isinstance(store, ProtocolIdempotencyStore) + + def test_partial_implementation_does_not_conform(self) -> None: + """A class with only some methods should not conform. + + Verifies that partial implementations are rejected by the + protocol's isinstance check. + """ + + class PartialStore: + async def check_and_record( + self, + message_id: UUID, + domain: str | None = None, + correlation_id: UUID | None = None, + ) -> bool: + return True + + # Missing: is_processed, mark_processed, cleanup_expired + + store = PartialStore() + assert not isinstance(store, ProtocolIdempotencyStore) + + def test_sync_methods_do_not_conform(self) -> None: + """A class with sync (non-async) methods should not conform. + + The protocol requires async methods, so sync implementations + should not pass the isinstance check. + """ + + class SyncStore: + def check_and_record( + self, + message_id: UUID, + domain: str | None = None, + correlation_id: UUID | None = None, + ) -> bool: + return True + + def is_processed( + self, + message_id: UUID, + domain: str | None = None, + ) -> bool: + return False + + def mark_processed( + self, + message_id: UUID, + domain: str | None = None, + correlation_id: UUID | None = None, + processed_at: datetime | None = None, + ) -> None: + pass + + def cleanup_expired( + self, + ttl_seconds: int, + ) -> int: + return 0 + + store = SyncStore() + # Note: runtime_checkable only checks method existence, not signatures. + # This is a known limitation of typing.Protocol. + # The sync implementation will pass isinstance() but fail at runtime. + # This test documents the expected behavior rather than a strict check. + # In practice, type checkers like mypy will catch this. + assert isinstance(store, ProtocolIdempotencyStore) diff --git a/tests/unit/idempotency/test_store_inmemory.py b/tests/unit/idempotency/test_store_inmemory.py index 2233a78326..5a793dfdbd 100644 --- a/tests/unit/idempotency/test_store_inmemory.py +++ b/tests/unit/idempotency/test_store_inmemory.py @@ -16,12 +16,13 @@ from uuid import uuid4 import pytest -from omnibase_spi.protocols.storage.protocol_idempotency_store import ( + +from omnibase_infra.idempotency import ( + InMemoryIdempotencyStore, + ModelIdempotencyRecord, ProtocolIdempotencyStore, ) -from omnibase_infra.idempotency import InMemoryIdempotencyStore, ModelIdempotencyRecord - class TestInMemoryIdempotencyStoreProtocol: """Test protocol conformance.""" diff --git a/tests/unit/validation/test_validator_defaults.py b/tests/unit/validation/test_validator_defaults.py index 3198d79841..49e0c05206 100644 --- a/tests/unit/validation/test_validator_defaults.py +++ b/tests/unit/validation/test_validator_defaults.py @@ -40,7 +40,7 @@ def test_infra_max_unions_constant(self) -> None: OMN-983: Strict validation mode enabled. - Current baseline (~540 unions as of 2025-12-23): + Current baseline (544 unions as of 2025-12-23): - Most unions are legitimate `X | None` nullable patterns (ONEX-preferred) - These are counted but NOT flagged as violations - Actual violations (primitive soup, Union[X,None] syntax) are reported separately @@ -49,11 +49,13 @@ def test_infra_max_unions_constant(self) -> None: - 491 (2025-12-21): Initial baseline with DispatcherFunc | ContextAwareDispatcherFunc - 515 (2025-12-22): OMN-990 MessageDispatchEngine + OMN-947 snapshots - 540 (2025-12-23): OMN-950 comprehensive reducer tests + - 544 (2025-12-23): OMN-954 effect idempotency and retry tests (PR #78) + Threshold: 555 (11 buffer above 544 baseline for codebase growth) Target: Reduce to <200 through ongoing dict[str, object] -> JsonValue migration. """ - assert INFRA_MAX_UNIONS == 540, ( - "INFRA_MAX_UNIONS should be 540 (OMN-950 reducer tests)" + assert INFRA_MAX_UNIONS == 555, ( + "INFRA_MAX_UNIONS should be 555 (11 buffer above 544 baseline)" ) def test_infra_max_violations_constant(self) -> None: @@ -492,12 +494,12 @@ def test_union_count_within_threshold(self) -> None: the threshold, it indicates new code added unions without using proper typed patterns from omnibase_core. - Current baseline (~513 unions as of 2025-12-22): + Current baseline (~544 unions as of 2025-12-23): - Most unions are legitimate `X | None` nullable patterns (ONEX-preferred) - These are counted but NOT flagged as violations - Actual violations (primitive soup, Union[X,None] syntax) are reported separately - Threshold: INFRA_MAX_UNIONS (515) - buffer above baseline. + Threshold: INFRA_MAX_UNIONS (555) - buffer above baseline. Target: Reduce to <200 through ongoing dict[str, object] -> JsonValue migration. """ result = validate_infra_union_usage()