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
463 changes: 463 additions & 0 deletions docs/operations/DLQ_REPLAY_RUNBOOK.md

Large diffs are not rendered by default.

2 changes: 1 addition & 1 deletion docs/operations/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ Operational runbooks and guides for deploying and managing omnibase_infra in pro

| Document | Description |
|----------|-------------|
| [DLQ Replay Guide](DLQ_REPLAY_GUIDE.md) | Dead Letter Queue replay mechanism: manual procedures, automated design, safety considerations |
| [DLQ Replay Guide](DLQ_REPLAY_RUNBOOK.md) | Dead Letter Queue replay mechanism: manual procedures, automated design, safety considerations |
| [Thread Pool Tuning](THREAD_POOL_TUNING_RUNBOOK.md) | Guide for tuning thread pool configurations in VaultAdapter and other components |

## Purpose
Expand Down
374 changes: 352 additions & 22 deletions scripts/dlq_replay.py

Large diffs are not rendered by default.

100 changes: 100 additions & 0 deletions src/omnibase_infra/dlq/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,100 @@
# SPDX-License-Identifier: MIT
# Copyright (c) 2025 OmniNode Team
"""ONEX Infrastructure DLQ (Dead Letter Queue) System.

This module provides DLQ replay tracking capabilities for message
processing in distributed systems. It enables persistent tracking
of replay attempts through PostgreSQL.

Components:
- Constants: Shared validation patterns (constants_dlq.py)
- Models: Pydantic models for DLQ tracking configuration and records
- Tracker: PostgreSQL-based tracker for replay history persistence

Constants:
- PATTERN_TABLE_NAME: Regex pattern string for PostgreSQL table name validation
- REGEX_TABLE_NAME: Pre-compiled regex for runtime validation

Models:
- ModelDlqTrackingConfig: Configuration for PostgreSQL-based tracking
- ModelDlqReplayRecord: Record of a DLQ message replay attempt
- EnumReplayStatus: Status enum for replay operations

Tracker:
- DLQReplayTracker: PostgreSQL tracker for replay history
- ServiceDlqTracking: ONEX naming convention alias
- DLQTrackingService: Backwards compatibility alias

Example - Recording Replay Attempts:
>>> from omnibase_infra.dlq import (
... DLQReplayTracker,
... ModelDlqTrackingConfig,
... ModelDlqReplayRecord,
... EnumReplayStatus,
... )
>>> from uuid import uuid4
>>> from datetime import datetime, timezone
>>>
>>> # Configure the tracker
>>> config = ModelDlqTrackingConfig(
... dsn="postgresql://user:pass@localhost:5432/mydb",
... storage_table="dlq_replay_history",
... )
>>>
>>> # Initialize and use
>>> tracker = DLQReplayTracker(config)
>>> await tracker.initialize()
>>> try:
... record = ModelDlqReplayRecord(
... id=uuid4(),
... original_message_id=uuid4(),
... replay_correlation_id=uuid4(),
... original_topic="dev.orders.command.v1",
... target_topic="dev.orders.command.v1",
... replay_status=EnumReplayStatus.COMPLETED,
... replay_timestamp=datetime.now(timezone.utc),
... success=True,
... dlq_offset=12345,
... dlq_partition=0,
... retry_count=1,
... )
... await tracker.record_replay_attempt(record)
...
... # Query replay history
... history = await tracker.get_replay_history(record.original_message_id)
... finally:
... await tracker.shutdown()

Related:
- scripts/dlq_replay.py - CLI tool for DLQ replay operations
- OMN-1032 - PostgreSQL tracking integration ticket
- OMN-949 - DLQ configuration ticket
"""

from omnibase_infra.dlq.constants_dlq import PATTERN_TABLE_NAME, REGEX_TABLE_NAME
from omnibase_infra.dlq.models import (
EnumReplayStatus,
ModelDlqReplayRecord,
ModelDlqTrackingConfig,
)
from omnibase_infra.dlq.service_dlq_tracking import (
DLQReplayTracker,
ServiceDlqTracking,
)

# Backwards compatibility alias
DLQTrackingService = DLQReplayTracker

__all__: list[str] = [
# Constants
"PATTERN_TABLE_NAME", # Regex pattern string for table name validation
"REGEX_TABLE_NAME", # Pre-compiled regex for runtime validation
# Models
"EnumReplayStatus",
"ModelDlqReplayRecord",
"ModelDlqTrackingConfig",
# Tracker (multiple names for flexibility)
"DLQReplayTracker", # Primary class name
"ServiceDlqTracking", # ONEX naming convention (service_<name>.py → Service<Name>)
"DLQTrackingService", # Backwards compatibility alias
]
57 changes: 57 additions & 0 deletions src/omnibase_infra/dlq/constants_dlq.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,57 @@
# SPDX-License-Identifier: MIT
# Copyright (c) 2025 OmniNode Team
"""DLQ Constants and Shared Validation Patterns.

This module provides shared constants for the DLQ subsystem, ensuring
consistent validation patterns across models and services.

Defense-in-Depth Note:
The PATTERN_TABLE_NAME constant is intentionally used in BOTH:
1. Pydantic model validation (ModelDlqTrackingConfig.storage_table field)
2. Runtime validation (DLQReplayTracker._validate_storage_table method)

This defense-in-depth approach ensures SQL injection prevention even if:
- Direct attribute assignment bypasses Pydantic validation
- Deserialization from untrusted sources bypasses model validation
- Future code changes inadvertently bypass config validation

DO NOT remove either validation layer. Both are intentional and required.

Related:
- model_dlq_tracking_config.py - Pydantic field pattern validation
- service_dlq_tracking.py - Runtime defense-in-depth validation
- OMN-1032 - PostgreSQL tracking integration ticket
"""

from __future__ import annotations

import re

# =============================================================================
# TABLE NAME VALIDATION
# =============================================================================

# Regex pattern string for valid PostgreSQL table names
# Must start with letter or underscore, followed by letters, digits, or underscores
# This matches PostgreSQL's identifier naming rules for unquoted identifiers
#
# Used in:
# - ModelDlqTrackingConfig: Pydantic field pattern constraint
# - DLQReplayTracker: Runtime defense-in-depth validation
#
# Pattern explanation:
# ^ - Start of string
# [a-zA-Z_] - First character must be letter (a-z, A-Z) or underscore
# [a-zA-Z0-9_]* - Subsequent characters can be letters, digits, or underscores
# $ - End of string
PATTERN_TABLE_NAME = r"^[a-zA-Z_][a-zA-Z0-9_]*$"

# Pre-compiled regex for runtime validation (avoids recompilation overhead)
# This is the compiled version of PATTERN_TABLE_NAME for use in service code
REGEX_TABLE_NAME = re.compile(PATTERN_TABLE_NAME)


__all__: list[str] = [
"PATTERN_TABLE_NAME",
"REGEX_TABLE_NAME",
]
26 changes: 26 additions & 0 deletions src/omnibase_infra/dlq/models/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
# SPDX-License-Identifier: MIT
# Copyright (c) 2025 OmniNode Team
"""ONEX Infrastructure DLQ Models.

This module provides Pydantic models for the DLQ replay tracking system,
including configuration and replay record models.

Exports:
ModelDlqTrackingConfig: Configuration for PostgreSQL-based DLQ tracking service
ModelDlqReplayRecord: Record of a DLQ message replay attempt
EnumReplayStatus: Status enum for replay operations

Related:
- scripts/dlq_replay.py - CLI tool for DLQ replay operations
- OMN-1032 - PostgreSQL tracking integration ticket
"""

from omnibase_infra.dlq.models.enum_replay_status import EnumReplayStatus
from omnibase_infra.dlq.models.model_dlq_replay_record import ModelDlqReplayRecord
from omnibase_infra.dlq.models.model_dlq_tracking_config import ModelDlqTrackingConfig

__all__: list[str] = [
"EnumReplayStatus",
"ModelDlqReplayRecord",
"ModelDlqTrackingConfig",
]
37 changes: 37 additions & 0 deletions src/omnibase_infra/dlq/models/enum_replay_status.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,37 @@
# SPDX-License-Identifier: MIT
# Copyright (c) 2025 OmniNode Team
"""DLQ Replay Status Enum.

This module provides the status enum for DLQ replay operations.

Related:
- scripts/dlq_replay.py - CLI tool that uses this enum
- OMN-1032 - PostgreSQL tracking integration ticket
"""

from __future__ import annotations

from enum import Enum


class EnumReplayStatus(str, Enum):
"""Status of a DLQ replay operation.

This enum tracks the lifecycle of a replay attempt:
- PENDING: Replay has been initiated but not yet completed
- COMPLETED: Message was successfully replayed to target topic
- FAILED: Replay attempt failed (will be recorded with error_message)
- SKIPPED: Message was intentionally not replayed (e.g., non-retryable error)

Usage:
This enum is the canonical definition for replay status tracking.
It is imported by scripts/dlq_replay.py for CLI usage.
"""

PENDING = "pending"
COMPLETED = "completed"
FAILED = "failed"
SKIPPED = "skipped"


__all__: list[str] = ["EnumReplayStatus"]
135 changes: 135 additions & 0 deletions src/omnibase_infra/dlq/models/model_dlq_replay_record.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,135 @@
# SPDX-License-Identifier: MIT
# Copyright (c) 2025 OmniNode Team
"""DLQ Replay Record Model.

This module provides the Pydantic model for DLQ replay records,
representing the state of a message replay operation stored in PostgreSQL.

Related:
- scripts/dlq_replay.py - CLI tool that uses this model
- OMN-1032 - PostgreSQL tracking integration ticket
"""

from __future__ import annotations

from datetime import datetime
from uuid import UUID

from pydantic import BaseModel, ConfigDict, Field

from omnibase_infra.dlq.models.enum_replay_status import EnumReplayStatus


class ModelDlqReplayRecord(BaseModel):
"""Record of a DLQ message replay attempt.

This model represents a single replay attempt stored in PostgreSQL,
tracking both the original message context and the replay outcome.

Table Schema:
CREATE TABLE IF NOT EXISTS dlq_replay_history (
id UUID PRIMARY KEY,
original_message_id UUID NOT NULL,
replay_correlation_id UUID NOT NULL,
original_topic VARCHAR(255) NOT NULL,
target_topic VARCHAR(255) NOT NULL,
replay_status VARCHAR(50) NOT NULL,
replay_timestamp TIMESTAMP WITH TIME ZONE NOT NULL DEFAULT now(),
success BOOLEAN NOT NULL,
error_message TEXT,
dlq_offset BIGINT NOT NULL,
dlq_partition INTEGER NOT NULL,
retry_count INTEGER NOT NULL DEFAULT 0
);

Attributes:
id: Unique identifier for this replay record (primary key).
original_message_id: Correlation ID from the original DLQ message.
This links the replay to the original failed message.
replay_correlation_id: New correlation ID assigned during replay.
Used for tracing the replayed message through the system.
original_topic: The Kafka topic where the message originally failed.
target_topic: The topic where the message was replayed to.
Usually same as original_topic, but may differ for rerouting.
replay_status: Current status of the replay operation.
replay_timestamp: When the replay was attempted (UTC).
success: Whether the replay was successful.
True for COMPLETED status, False for FAILED/SKIPPED.
error_message: Error details if replay failed. None for success.
dlq_offset: Kafka offset of the message in the DLQ topic.
dlq_partition: Kafka partition of the message in the DLQ topic.
retry_count: Number of times this message has been retried.
Includes previous attempts before this replay.

Example:
>>> from uuid import uuid4
>>> from datetime import datetime, timezone
>>> record = ModelDlqReplayRecord(
... id=uuid4(),
... original_message_id=uuid4(),
... replay_correlation_id=uuid4(),
... original_topic="dev.orders.command.v1",
... target_topic="dev.orders.command.v1",
... replay_status=EnumReplayStatus.COMPLETED,
... replay_timestamp=datetime.now(timezone.utc),
... success=True,
... dlq_offset=12345,
... dlq_partition=0,
... retry_count=1,
... )
"""

model_config = ConfigDict(
frozen=True,
extra="forbid",
from_attributes=True,
)

id: UUID = Field(
description="Unique identifier for this replay record (primary key)",
)
original_message_id: UUID = Field(
description="Correlation ID from the original DLQ message",
)
replay_correlation_id: UUID = Field(
description="New correlation ID assigned during replay",
)
original_topic: str = Field(
description="The Kafka topic where the message originally failed",
min_length=1,
max_length=255,
)
target_topic: str = Field(
description="The topic where the message was replayed to",
min_length=1,
max_length=255,
)
replay_status: EnumReplayStatus = Field(
description="Current status of the replay operation",
)
replay_timestamp: datetime = Field(
description="When the replay was attempted (UTC)",
)
success: bool = Field(
description="Whether the replay was successful",
)
error_message: str | None = Field(
default=None,
description="Error details if replay failed, None for success",
)
dlq_offset: int = Field(
description="Kafka offset of the message in the DLQ topic",
ge=0,
)
dlq_partition: int = Field(
description="Kafka partition of the message in the DLQ topic",
ge=0,
)
retry_count: int = Field(
default=0,
description="Number of times this message has been retried",
ge=0,
)


__all__: list[str] = ["ModelDlqReplayRecord"]
Loading