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
8 changes: 7 additions & 1 deletion .claude/settings.local.json
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,13 @@
"mcp__serena__replace_regex",
"mcp__archon__assess_code_quality",
"Bash(agent-onex-coordinator)",
"mcp__archon__health_check"
"mcp__archon__health_check",
"mcp__serena__list_dir",
"mcp__serena__create_text_file",
"Read(//Volumes/PRO-G40/Code/omnibase_core/src/omnibase_core/enums/intelligence/**)",
"Read(//Volumes/PRO-G40/Code/omnibase_core/src/omnibase_core/models/metrics/**)",
"Read(//Volumes/PRO-G40/Code/omnibase_core/src/omnibase_core/models/core/**)",
"Read(//Volumes/PRO-G40/Code/omnibase_core/src/**)"
],
"deny": [],
"ask": []
Expand Down
3 changes: 3 additions & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,9 @@ rich = "^13.7.0"
cryptography = "^41.0.0"
jinja2 = "^3.1.0"

# Tracing and observability dependencies
sqlparse = "^0.4.4" # For secure SQL query sanitization in tracing

[tool.poetry.group.dev.dependencies]
pytest = "^8.4.0"
pytest-asyncio = "^0.25.0"
Expand Down
14 changes: 7 additions & 7 deletions src/omnibase_infra/infrastructure/container.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@
import time
from typing import Callable, Optional, Type, TypeVar, Union, Dict, List

from omnibase_core.core.onex_container import ModelONEXContainer as ONEXContainer
from omnibase_core.core.onex_container import ModelONEXContainer
from omnibase_core.protocol.protocol_event_bus import ProtocolEventBus
from omnibase_core.model.core.model_onex_event import ModelOnexEvent
from omnibase_core.utils.generation.utility_schema_loader import UtilitySchemaLoader
Expand Down Expand Up @@ -618,7 +618,7 @@ def _event_to_topic(self, event: ModelOnexEvent) -> str:
return topic


def create_infrastructure_container() -> ONEXContainer:
def create_infrastructure_container() -> ModelONEXContainer:
"""
Create infrastructure container with all shared dependencies.

Expand All @@ -628,10 +628,10 @@ def create_infrastructure_container() -> ONEXContainer:
- "Everything needs to be resolved by duck typing"

Returns:
Configured ONEXContainer with infrastructure dependencies
Configured ModelONEXContainer with infrastructure dependencies
"""
# Create base ONEX container
container = ONEXContainer()
container = ModelONEXContainer()

# Set up all shared dependencies for infrastructure services
_setup_infrastructure_dependencies(container)
Expand All @@ -642,7 +642,7 @@ def create_infrastructure_container() -> ONEXContainer:
return container


def _setup_infrastructure_dependencies(container: ONEXContainer):
def _setup_infrastructure_dependencies(container: ModelONEXContainer):
"""Set up all dependencies needed by infrastructure services."""

# Get logger for container setup
Expand Down Expand Up @@ -687,13 +687,13 @@ def _setup_infrastructure_dependencies(container: ONEXContainer):
logger.info(" PostgreSQL connection manager skipped (environment not configured)")


def _register_service(container: ONEXContainer, service_name: str, service_instance):
def _register_service(container: ModelONEXContainer, service_name: str, service_instance):
"""Register a service in the container for later retrieval."""
# Use the ONEX container's native service registration
container.register_service(service_name, service_instance)


def _bind_infrastructure_get_service_method(container: ONEXContainer):
def _bind_infrastructure_get_service_method(container: ModelONEXContainer):
"""Configure infrastructure container with proper dependency injection."""
# The ModelONEXContainer should handle get_service natively
# We just need to register our services properly in the container
Expand Down
33 changes: 20 additions & 13 deletions src/omnibase_infra/infrastructure/distributed_tracing.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@
import os
import time
from contextlib import asynccontextmanager
from typing import Dict, Any, Optional, Union, AsyncIterator
from typing import Optional, Union, AsyncIterator
from uuid import UUID, uuid4
from datetime import datetime

Expand All @@ -26,7 +26,7 @@
from opentelemetry.instrumentation.asyncpg import AsyncPGInstrumentor
from opentelemetry.instrumentation.kafka import KafkaInstrumentor
from opentelemetry.trace.status import Status, StatusCode
from opentelemetry.trace import Span, SpanKind
from opentelemetry.trace import Span, SpanKind, Tracer
from opentelemetry.context import Context

OPENTELEMETRY_AVAILABLE = True
Expand All @@ -39,6 +39,7 @@
from omnibase_core.core.errors.onex_error import OnexError
from omnibase_core.core.errors.onex_error import CoreErrorCode
from omnibase_core.model.core.model_onex_event import ModelOnexEvent
from omnibase_infra.models.tracing.model_span_attributes import ModelSpanAttributes

from ..security.audit_logger import AuditLogger, AuditEvent, AuditEventType, AuditSeverity

Expand Down Expand Up @@ -122,7 +123,7 @@ def __init__(self, config: Optional[TracingConfiguration] = None):

# Tracing components
self.tracer_provider: Optional[TracerProvider] = None
self.tracer: Optional[Any] = None # OpenTelemetry tracer
self.tracer: Optional[Tracer] = None # OpenTelemetry tracer
self.is_initialized = False

# Integration with audit logging
Expand Down Expand Up @@ -208,7 +209,7 @@ async def trace_operation(
correlation_id: Optional[Union[str, UUID]] = None,
parent_context: Optional[Context] = None,
span_kind: Optional[SpanKind] = SpanKind.INTERNAL,
attributes: Optional[Dict[str, Any]] = None
attributes: Optional[ModelSpanAttributes] = None
) -> AsyncIterator[Span]:
"""
Create a trace span for an operation with automatic error handling.
Expand Down Expand Up @@ -239,15 +240,21 @@ async def trace_operation(

try:
# Create span
base_attributes = {
"correlation_id": correlation_str,
"environment": self.config.environment,
"service.name": self.config.service_name,
}

# Add provided attributes if available
if attributes:
attribute_dict = attributes.dict(exclude_none=True)
base_attributes.update(attribute_dict)

span = self.tracer.start_span(
name=operation_name,
kind=span_kind,
attributes={
"correlation_id": correlation_str,
"environment": self.config.environment,
"service.name": self.config.service_name,
**(attributes or {})
}
attributes=base_attributes
)

# Set correlation ID in baggage for propagation
Expand Down Expand Up @@ -285,12 +292,12 @@ async def trace_operation(
if token:
context.detach(token)

def _create_noop_span(self) -> Any:
def _create_noop_span(self) -> object:
"""Create a no-op span when tracing is disabled."""
class NoOpSpan:
def set_attribute(self, key: str, value: Any) -> None:
def set_attribute(self, key: str, value: Union[str, int, float, bool]) -> None:
pass
def set_status(self, status: Any) -> None:
def set_status(self, status: object) -> None:
pass
def record_exception(self, exception: Exception) -> None:
pass
Expand Down
5 changes: 5 additions & 0 deletions src/omnibase_infra/models/circuit_breaker/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
"""Circuit Breaker Models Package.

Shared models for circuit breaker operations and configurations.
Used by circuit breaker nodes and related infrastructure components.
"""
Original file line number Diff line number Diff line change
@@ -0,0 +1,89 @@
"""Circuit Breaker Metrics Model.

Shared model for circuit breaker metrics and performance data.
Used across circuit breaker nodes and observability systems.
"""

from pydantic import BaseModel, Field
from typing import Optional
from datetime import datetime


class ModelCircuitBreakerMetrics(BaseModel):
"""Model for circuit breaker metrics tracking."""

total_events: int = Field(
default=0,
ge=0,
description="Total number of events processed"
)

successful_events: int = Field(
default=0,
ge=0,
description="Number of successfully processed events"
)

failed_events: int = Field(
default=0,
ge=0,
description="Number of failed events"
)

queued_events: int = Field(
default=0,
ge=0,
description="Number of events currently queued"
)

dropped_events: int = Field(
default=0,
ge=0,
description="Number of events dropped due to capacity limits"
)

dead_letter_events: int = Field(
default=0,
ge=0,
description="Number of events in dead letter queue"
)

circuit_opens: int = Field(
default=0,
ge=0,
description="Number of times circuit has opened"
)

circuit_closes: int = Field(
default=0,
ge=0,
description="Number of times circuit has closed"
)

last_failure: Optional[datetime] = Field(
default=None,
description="Timestamp of last failure"
)

last_success: Optional[datetime] = Field(
default=None,
description="Timestamp of last success"
)

success_rate_percent: float = Field(
default=100.0,
ge=0.0,
le=100.0,
description="Success rate percentage"
)

average_response_time_ms: float = Field(
default=0.0,
ge=0.0,
description="Average response time in milliseconds"
)

class Config:
json_encoders = {
datetime: lambda v: v.isoformat()
}
Loading