From 283011ff631731c4d6ca948a2e7b652da56468c0 Mon Sep 17 00:00:00 2001 From: HardMax71 Date: Sun, 8 Feb 2026 20:23:05 +0100 Subject: [PATCH 1/3] fixed enum imports (now using import from enums/init, not from separate enums/xx files) --- backend/app/api/routes/admin/events.py | 2 +- backend/app/api/routes/admin/users.py | 2 +- backend/app/api/routes/dlq.py | 2 +- backend/app/api/routes/events.py | 4 +- backend/app/api/routes/execution.py | 4 +- backend/app/api/routes/notifications.py | 2 +- backend/app/api/routes/replay.py | 2 +- backend/app/api/routes/saga.py | 2 +- backend/app/core/metrics/execution.py | 2 +- backend/app/core/providers.py | 2 +- backend/app/db/docs/event.py | 2 +- backend/app/db/docs/execution.py | 3 +- backend/app/db/docs/notification.py | 6 +- backend/app/db/docs/replay.py | 3 +- backend/app/db/docs/saga.py | 2 +- backend/app/db/docs/user.py | 2 +- backend/app/db/docs/user_settings.py | 3 +- .../admin/admin_events_repository.py | 3 +- backend/app/db/repositories/dlq_repository.py | 2 +- .../app/db/repositories/event_repository.py | 2 +- .../repositories/notification_repository.py | 3 +- .../app/db/repositories/replay_repository.py | 2 +- .../app/db/repositories/saga_repository.py | 2 +- .../app/db/repositories/user_repository.py | 2 +- .../repositories/user_settings_repository.py | 2 +- backend/app/dlq/manager.py | 2 +- backend/app/dlq/models.py | 2 +- backend/app/domain/admin/replay_models.py | 3 +- backend/app/domain/admin/replay_updates.py | 2 +- backend/app/domain/enums/__init__.py | 23 ++++- backend/app/domain/enums/kafka.py | 85 ------------------- backend/app/domain/events/event_models.py | 2 +- backend/app/domain/events/typed.py | 16 ++-- backend/app/domain/execution/models.py | 3 +- backend/app/domain/notification/models.py | 6 +- backend/app/domain/replay/models.py | 4 +- backend/app/domain/saga/models.py | 2 +- backend/app/domain/sse/models.py | 2 +- backend/app/domain/user/__init__.py | 2 +- backend/app/domain/user/settings_models.py | 4 +- backend/app/domain/user/user_models.py | 2 +- backend/app/events/core/producer.py | 2 +- backend/app/events/handlers.py | 4 +- backend/app/infrastructure/kafka/mappings.py | 38 ++++++++- backend/app/infrastructure/kafka/topics.py | 2 +- backend/app/schemas_pydantic/admin_events.py | 2 +- backend/app/schemas_pydantic/events.py | 3 +- backend/app/schemas_pydantic/execution.py | 3 +- .../app/schemas_pydantic/health_dashboard.py | 2 +- backend/app/schemas_pydantic/notification.py | 6 +- backend/app/schemas_pydantic/replay.py | 4 +- backend/app/schemas_pydantic/replay_models.py | 4 +- backend/app/schemas_pydantic/saga.py | 2 +- backend/app/schemas_pydantic/sse.py | 5 +- backend/app/schemas_pydantic/user.py | 2 +- backend/app/schemas_pydantic/user_settings.py | 4 +- .../services/admin/admin_events_service.py | 2 +- .../app/services/admin/admin_user_service.py | 4 +- backend/app/services/auth_service.py | 2 +- backend/app/services/coordinator/__init__.py | 2 +- .../app/services/coordinator/coordinator.py | 3 +- backend/app/services/event_replay/__init__.py | 2 +- .../services/event_replay/replay_service.py | 2 +- backend/app/services/event_service.py | 3 +- backend/app/services/execution_service.py | 3 +- .../app/services/grafana_alert_processor.py | 3 +- backend/app/services/k8s_worker/worker.py | 2 +- backend/app/services/kafka_event_service.py | 2 +- backend/app/services/notification_service.py | 7 +- backend/app/services/pod_monitor/config.py | 2 +- .../app/services/pod_monitor/event_mapper.py | 3 +- .../services/result_processor/processor.py | 4 +- backend/app/services/saga/__init__.py | 2 +- .../app/services/saga/saga_orchestrator.py | 2 +- backend/app/services/sse/redis_bus.py | 2 +- backend/app/services/sse/sse_service.py | 4 +- backend/app/services/user_settings_service.py | 3 +- backend/tests/conftest.py | 2 +- backend/tests/e2e/conftest.py | 4 +- .../db/repositories/test_dlq_repository.py | 2 +- .../repositories/test_execution_repository.py | 2 +- backend/tests/e2e/dlq/test_dlq_discard.py | 2 +- backend/tests/e2e/dlq/test_dlq_manager.py | 3 +- backend/tests/e2e/dlq/test_dlq_retry.py | 2 +- .../notifications/test_notification_sse.py | 2 +- .../result_processor/test_result_processor.py | 2 +- .../services/admin/test_admin_user_service.py | 2 +- .../coordinator/test_execution_coordinator.py | 2 +- .../events/test_kafka_event_service.py | 3 +- .../execution/test_execution_service.py | 3 +- .../test_notification_service.py | 5 +- .../services/replay/test_replay_service.py | 2 +- .../e2e/services/saga/test_saga_service.py | 3 +- .../tests/e2e/services/sse/test_redis_bus.py | 3 +- backend/tests/e2e/test_admin_events_routes.py | 3 +- backend/tests/e2e/test_admin_users_routes.py | 2 +- backend/tests/e2e/test_auth_routes.py | 2 +- backend/tests/e2e/test_dlq_routes.py | 2 +- backend/tests/e2e/test_events_routes.py | 2 +- backend/tests/e2e/test_execution_routes.py | 3 +- .../tests/e2e/test_k8s_worker_create_pod.py | 2 +- .../tests/e2e/test_notifications_routes.py | 2 +- backend/tests/e2e/test_replay_routes.py | 3 +- backend/tests/e2e/test_saga_routes.py | 2 +- .../tests/e2e/test_user_settings_routes.py | 2 +- .../test_execution_and_events_metrics.py | 2 +- .../unit/core/metrics/test_metrics_classes.py | 2 +- backend/tests/unit/core/test_security.py | 2 +- .../events/test_event_schema_coverage.py | 2 +- .../unit/events/test_mappings_and_types.py | 3 +- .../schemas_pydantic/test_events_schemas.py | 2 +- .../test_notification_schemas.py | 2 +- .../coordinator/test_coordinator_queue.py | 2 +- .../services/pod_monitor/test_event_mapper.py | 3 +- .../result_processor/test_processor.py | 3 +- .../services/saga/test_saga_comprehensive.py | 2 +- .../saga/test_saga_orchestrator_unit.py | 2 +- .../unit/services/sse/test_sse_service.py | 3 +- .../tests/unit/services/test_pod_builder.py | 2 +- backend/workers/run_coordinator.py | 2 +- backend/workers/run_dlq_processor.py | 2 +- backend/workers/run_k8s_worker.py | 2 +- backend/workers/run_pod_monitor.py | 2 +- backend/workers/run_result_processor.py | 2 +- backend/workers/run_saga_orchestrator.py | 2 +- 125 files changed, 190 insertions(+), 288 deletions(-) diff --git a/backend/app/api/routes/admin/events.py b/backend/app/api/routes/admin/events.py index 2aeeb014..88207d33 100644 --- a/backend/app/api/routes/admin/events.py +++ b/backend/app/api/routes/admin/events.py @@ -8,7 +8,7 @@ from app.api.dependencies import admin_user from app.core.correlation import CorrelationContext -from app.domain.enums.events import EventType +from app.domain.enums import EventType from app.domain.events.event_models import EventFilter from app.domain.replay import ReplayFilter from app.domain.user import User diff --git a/backend/app/api/routes/admin/users.py b/backend/app/api/routes/admin/users.py index c76dab70..4d9df577 100644 --- a/backend/app/api/routes/admin/users.py +++ b/backend/app/api/routes/admin/users.py @@ -6,7 +6,7 @@ from app.api.dependencies import admin_user from app.db.repositories.admin.admin_user_repository import AdminUserRepository -from app.domain.enums.user import UserRole +from app.domain.enums import UserRole from app.domain.rate_limit import RateLimitRule, UserRateLimit from app.domain.user import User from app.domain.user import UserUpdate as DomainUserUpdate diff --git a/backend/app/api/routes/dlq.py b/backend/app/api/routes/dlq.py index 9132631e..ecd5a335 100644 --- a/backend/app/api/routes/dlq.py +++ b/backend/app/api/routes/dlq.py @@ -9,7 +9,7 @@ from app.dlq import RetryPolicy from app.dlq.manager import DLQManager from app.dlq.models import DLQMessageStatus -from app.domain.enums.events import EventType +from app.domain.enums import EventType from app.schemas_pydantic.common import ErrorResponse from app.schemas_pydantic.dlq import ( DLQBatchRetryResponse, diff --git a/backend/app/api/routes/events.py b/backend/app/api/routes/events.py index f1f81900..6382f794 100644 --- a/backend/app/api/routes/events.py +++ b/backend/app/api/routes/events.py @@ -9,9 +9,7 @@ from app.api.dependencies import admin_user, current_user from app.core.correlation import CorrelationContext from app.core.utils import get_client_ip -from app.domain.enums.common import SortOrder -from app.domain.enums.events import EventType -from app.domain.enums.user import UserRole +from app.domain.enums import EventType, SortOrder, UserRole from app.domain.events.event_models import EventFilter from app.domain.events.typed import BaseEvent, DomainEvent, EventMetadata from app.domain.user import User diff --git a/backend/app/api/routes/execution.py b/backend/app/api/routes/execution.py index e0f18822..200a4c4e 100644 --- a/backend/app/api/routes/execution.py +++ b/backend/app/api/routes/execution.py @@ -9,9 +9,7 @@ from app.api.dependencies import admin_user, current_user from app.core.tracing import EventAttributes, add_span_attributes from app.core.utils import get_client_ip -from app.domain.enums.events import EventType -from app.domain.enums.execution import ExecutionStatus -from app.domain.enums.user import UserRole +from app.domain.enums import EventType, ExecutionStatus, UserRole from app.domain.events.typed import BaseEvent, DomainEvent, EventMetadata from app.domain.exceptions import DomainError from app.domain.idempotency import KeyStrategy diff --git a/backend/app/api/routes/notifications.py b/backend/app/api/routes/notifications.py index 1b29e746..1549a35e 100644 --- a/backend/app/api/routes/notifications.py +++ b/backend/app/api/routes/notifications.py @@ -4,7 +4,7 @@ from dishka.integrations.fastapi import DishkaRoute from fastapi import APIRouter, Query, Request, Response -from app.domain.enums.notification import NotificationChannel, NotificationStatus +from app.domain.enums import NotificationChannel, NotificationStatus from app.schemas_pydantic.notification import ( DeleteNotificationResponse, NotificationListResponse, diff --git a/backend/app/api/routes/replay.py b/backend/app/api/routes/replay.py index 44c37661..9e519c55 100644 --- a/backend/app/api/routes/replay.py +++ b/backend/app/api/routes/replay.py @@ -5,7 +5,7 @@ from fastapi import APIRouter, Depends, Query from app.api.dependencies import admin_user -from app.domain.enums.replay import ReplayStatus +from app.domain.enums import ReplayStatus from app.domain.replay import ReplayConfig from app.schemas_pydantic.replay import ( CleanupResponse, diff --git a/backend/app/api/routes/saga.py b/backend/app/api/routes/saga.py index 243f87eb..b75ced32 100644 --- a/backend/app/api/routes/saga.py +++ b/backend/app/api/routes/saga.py @@ -4,7 +4,7 @@ from dishka.integrations.fastapi import DishkaRoute from fastapi import APIRouter, Query, Request -from app.domain.enums.saga import SagaState +from app.domain.enums import SagaState from app.schemas_pydantic.common import ErrorResponse from app.schemas_pydantic.saga import ( SagaCancellationResponse, diff --git a/backend/app/core/metrics/execution.py b/backend/app/core/metrics/execution.py index a2f3b74a..4d7f59a9 100644 --- a/backend/app/core/metrics/execution.py +++ b/backend/app/core/metrics/execution.py @@ -1,5 +1,5 @@ from app.core.metrics.base import BaseMetrics -from app.domain.enums.execution import ExecutionStatus +from app.domain.enums import ExecutionStatus class ExecutionMetrics(BaseMetrics): diff --git a/backend/app/core/providers.py b/backend/app/core/providers.py index 0aaccf4c..b9cd6db3 100644 --- a/backend/app/core/providers.py +++ b/backend/app/core/providers.py @@ -46,7 +46,7 @@ from app.db.repositories.user_settings_repository import UserSettingsRepository from app.dlq.manager import DLQManager from app.dlq.models import RetryPolicy, RetryStrategy -from app.domain.enums.kafka import KafkaTopic +from app.domain.enums import KafkaTopic from app.domain.rate_limit import RateLimitConfig from app.domain.saga.models import SagaConfig from app.events.core import UnifiedProducer diff --git a/backend/app/db/docs/event.py b/backend/app/db/docs/event.py index f5a65609..d8c34b02 100644 --- a/backend/app/db/docs/event.py +++ b/backend/app/db/docs/event.py @@ -6,7 +6,7 @@ from pydantic import ConfigDict, Field from pymongo import ASCENDING, DESCENDING, IndexModel -from app.domain.enums.events import EventType +from app.domain.enums import EventType from app.domain.events.typed import EventMetadata diff --git a/backend/app/db/docs/execution.py b/backend/app/db/docs/execution.py index 80724e35..b893da0c 100644 --- a/backend/app/db/docs/execution.py +++ b/backend/app/db/docs/execution.py @@ -5,8 +5,7 @@ from pydantic import BaseModel, ConfigDict, Field from pymongo import IndexModel -from app.domain.enums.execution import ExecutionStatus -from app.domain.enums.storage import ExecutionErrorType +from app.domain.enums import ExecutionErrorType, ExecutionStatus # Pydantic model required here because Beanie embedded documents must be Pydantic BaseModel subclasses. diff --git a/backend/app/db/docs/notification.py b/backend/app/db/docs/notification.py index ee70ff6d..fbc200b3 100644 --- a/backend/app/db/docs/notification.py +++ b/backend/app/db/docs/notification.py @@ -6,11 +6,7 @@ from pydantic import ConfigDict, Field, field_validator from pymongo import ASCENDING, DESCENDING, IndexModel -from app.domain.enums.notification import ( - NotificationChannel, - NotificationSeverity, - NotificationStatus, -) +from app.domain.enums import NotificationChannel, NotificationSeverity, NotificationStatus class NotificationDocument(Document): diff --git a/backend/app/db/docs/replay.py b/backend/app/db/docs/replay.py index 0e2debac..4a503ace 100644 --- a/backend/app/db/docs/replay.py +++ b/backend/app/db/docs/replay.py @@ -6,8 +6,7 @@ from pydantic import BaseModel, ConfigDict, Field from pymongo import IndexModel -from app.domain.enums.events import EventType -from app.domain.enums.replay import ReplayStatus, ReplayTarget, ReplayType +from app.domain.enums import EventType, ReplayStatus, ReplayTarget, ReplayType class ReplayFilter(BaseModel): diff --git a/backend/app/db/docs/saga.py b/backend/app/db/docs/saga.py index c762d6c8..da220ef5 100644 --- a/backend/app/db/docs/saga.py +++ b/backend/app/db/docs/saga.py @@ -6,7 +6,7 @@ from pydantic import ConfigDict, Field from pymongo import ASCENDING, IndexModel -from app.domain.enums.saga import SagaState +from app.domain.enums import SagaState class SagaDocument(Document): diff --git a/backend/app/db/docs/user.py b/backend/app/db/docs/user.py index 27b8cdc5..3eb01059 100644 --- a/backend/app/db/docs/user.py +++ b/backend/app/db/docs/user.py @@ -4,7 +4,7 @@ from beanie import Document, Indexed from pydantic import ConfigDict, EmailStr, Field -from app.domain.enums.user import UserRole +from app.domain.enums import UserRole class UserDocument(Document): diff --git a/backend/app/db/docs/user_settings.py b/backend/app/db/docs/user_settings.py index 499c9f8c..96fe6c91 100644 --- a/backend/app/db/docs/user_settings.py +++ b/backend/app/db/docs/user_settings.py @@ -4,8 +4,7 @@ from beanie import Document, Indexed from pydantic import BaseModel, ConfigDict, Field, field_validator -from app.domain.enums.common import Theme -from app.domain.enums.notification import NotificationChannel +from app.domain.enums import NotificationChannel, Theme class NotificationSettings(BaseModel): diff --git a/backend/app/db/repositories/admin/admin_events_repository.py b/backend/app/db/repositories/admin/admin_events_repository.py index 380ebcea..c2bdfa0a 100644 --- a/backend/app/db/repositories/admin/admin_events_repository.py +++ b/backend/app/db/repositories/admin/admin_events_repository.py @@ -13,8 +13,7 @@ ) from app.domain.admin import ExecutionResultSummary, ReplaySessionData, ReplaySessionStatusDetail from app.domain.admin.replay_updates import ReplaySessionUpdate -from app.domain.enums.events import EventType -from app.domain.enums.replay import ReplayStatus +from app.domain.enums import EventType, ReplayStatus from app.domain.events import ( DomainEvent, DomainEventAdapter, diff --git a/backend/app/db/repositories/dlq_repository.py b/backend/app/db/repositories/dlq_repository.py index fee84466..1a94a617 100644 --- a/backend/app/db/repositories/dlq_repository.py +++ b/backend/app/db/repositories/dlq_repository.py @@ -18,7 +18,7 @@ EventTypeStatistic, TopicStatistic, ) -from app.domain.enums.events import EventType +from app.domain.enums import EventType class DLQRepository: diff --git a/backend/app/db/repositories/event_repository.py b/backend/app/db/repositories/event_repository.py index 17598020..d7ffe2b7 100644 --- a/backend/app/db/repositories/event_repository.py +++ b/backend/app/db/repositories/event_repository.py @@ -11,7 +11,7 @@ from app.core.tracing import EventAttributes from app.core.tracing.utils import add_span_attributes from app.db.docs import EventArchiveDocument, EventDocument -from app.domain.enums.events import EventType +from app.domain.enums import EventType from app.domain.events import ( ArchivedEvent, DomainEvent, diff --git a/backend/app/db/repositories/notification_repository.py b/backend/app/db/repositories/notification_repository.py index 4202d7c7..e38c2a6f 100644 --- a/backend/app/db/repositories/notification_repository.py +++ b/backend/app/db/repositories/notification_repository.py @@ -5,8 +5,7 @@ from beanie.operators import GTE, LTE, ElemMatch, In, NotIn, Or from app.db.docs import NotificationDocument, NotificationSubscriptionDocument, UserDocument -from app.domain.enums.notification import NotificationChannel, NotificationStatus -from app.domain.enums.user import UserRole +from app.domain.enums import NotificationChannel, NotificationStatus, UserRole from app.domain.notification import ( DomainNotification, DomainNotificationCreate, diff --git a/backend/app/db/repositories/replay_repository.py b/backend/app/db/repositories/replay_repository.py index c69593b9..998c4cb0 100644 --- a/backend/app/db/repositories/replay_repository.py +++ b/backend/app/db/repositories/replay_repository.py @@ -7,7 +7,7 @@ from app.db.docs import EventDocument, ReplaySessionDocument from app.domain.admin.replay_updates import ReplaySessionUpdate -from app.domain.enums.replay import ReplayStatus +from app.domain.enums import ReplayStatus from app.domain.replay.models import ReplayFilter, ReplaySessionState diff --git a/backend/app/db/repositories/saga_repository.py b/backend/app/db/repositories/saga_repository.py index a1d3f3e8..f2c1e435 100644 --- a/backend/app/db/repositories/saga_repository.py +++ b/backend/app/db/repositories/saga_repository.py @@ -8,7 +8,7 @@ from monggregate import Pipeline, S from app.db.docs import ExecutionDocument, SagaDocument -from app.domain.enums.saga import SagaState +from app.domain.enums import SagaState from app.domain.saga import Saga, SagaFilter, SagaListResult diff --git a/backend/app/db/repositories/user_repository.py b/backend/app/db/repositories/user_repository.py index ce0dddd8..7e72377d 100644 --- a/backend/app/db/repositories/user_repository.py +++ b/backend/app/db/repositories/user_repository.py @@ -5,7 +5,7 @@ from beanie.operators import Eq, Or, RegEx from app.db.docs import UserDocument -from app.domain.enums.user import UserRole +from app.domain.enums import UserRole from app.domain.user import DomainUserCreate, DomainUserUpdate, User, UserListResult diff --git a/backend/app/db/repositories/user_settings_repository.py b/backend/app/db/repositories/user_settings_repository.py index 69718e25..019ee9f4 100644 --- a/backend/app/db/repositories/user_settings_repository.py +++ b/backend/app/db/repositories/user_settings_repository.py @@ -6,7 +6,7 @@ from beanie.operators import GT, LTE, Eq, In from app.db.docs import EventDocument, UserSettingsDocument, UserSettingsSnapshotDocument -from app.domain.enums.events import EventType +from app.domain.enums import EventType from app.domain.user.settings_models import DomainUserSettings, DomainUserSettingsChangedEvent diff --git a/backend/app/dlq/manager.py b/backend/app/dlq/manager.py index 41b2a69b..a3724b88 100644 --- a/backend/app/dlq/manager.py +++ b/backend/app/dlq/manager.py @@ -16,7 +16,7 @@ RetryPolicy, RetryStrategy, ) -from app.domain.enums.kafka import KafkaTopic +from app.domain.enums import KafkaTopic from app.domain.events.typed import ( DLQMessageDiscardedEvent, DLQMessageReceivedEvent, diff --git a/backend/app/dlq/models.py b/backend/app/dlq/models.py index 66961243..7df828b4 100644 --- a/backend/app/dlq/models.py +++ b/backend/app/dlq/models.py @@ -5,7 +5,7 @@ from pydantic import BaseModel, ConfigDict, Field from app.core.utils import StringEnum -from app.domain.enums.events import EventType +from app.domain.enums import EventType from app.domain.events.typed import DomainEvent diff --git a/backend/app/domain/admin/replay_models.py b/backend/app/domain/admin/replay_models.py index e88e49d5..9b5ef7e5 100644 --- a/backend/app/domain/admin/replay_models.py +++ b/backend/app/domain/admin/replay_models.py @@ -4,8 +4,7 @@ from pydantic import ConfigDict from pydantic.dataclasses import dataclass -from app.domain.enums.execution import ExecutionStatus -from app.domain.enums.replay import ReplayStatus +from app.domain.enums import ExecutionStatus, ReplayStatus from app.domain.events.event_models import EventSummary from app.domain.replay.models import ReplayFilter, ReplaySessionState diff --git a/backend/app/domain/admin/replay_updates.py b/backend/app/domain/admin/replay_updates.py index 24c034d3..5ada0b9a 100644 --- a/backend/app/domain/admin/replay_updates.py +++ b/backend/app/domain/admin/replay_updates.py @@ -2,7 +2,7 @@ from pydantic import BaseModel, ConfigDict -from app.domain.enums.replay import ReplayStatus +from app.domain.enums import ReplayStatus class ReplaySessionUpdate(BaseModel): diff --git a/backend/app/domain/enums/__init__.py b/backend/app/domain/enums/__init__.py index 145ec7db..cfdbe8b9 100644 --- a/backend/app/domain/enums/__init__.py +++ b/backend/app/domain/enums/__init__.py @@ -1,20 +1,31 @@ -from app.domain.enums.common import ErrorType, SortOrder, Theme +from app.domain.enums.auth import LoginMethod, SettingsType +from app.domain.enums.common import Environment, ErrorType, SortOrder, Theme +from app.domain.enums.events import EventType from app.domain.enums.execution import ExecutionStatus, QueuePriority from app.domain.enums.health import AlertSeverity, AlertStatus, ComponentStatus +from app.domain.enums.kafka import GroupId, KafkaTopic from app.domain.enums.notification import ( NotificationChannel, NotificationSeverity, NotificationStatus, ) +from app.domain.enums.replay import ReplayStatus, ReplayTarget, ReplayType from app.domain.enums.saga import SagaState from app.domain.enums.sse import SSEControlEvent +from app.domain.enums.storage import ExecutionErrorType, StorageType from app.domain.enums.user import UserRole __all__ = [ + # Auth + "LoginMethod", + "SettingsType", # Common + "Environment", "ErrorType", "SortOrder", "Theme", + # Events + "EventType", # Execution "ExecutionStatus", "QueuePriority", @@ -22,14 +33,24 @@ "AlertSeverity", "AlertStatus", "ComponentStatus", + # Kafka + "GroupId", + "KafkaTopic", # Notification "NotificationChannel", "NotificationSeverity", "NotificationStatus", + # Replay + "ReplayStatus", + "ReplayTarget", + "ReplayType", # Saga "SagaState", # SSE "SSEControlEvent", + # Storage + "ExecutionErrorType", + "StorageType", # User "UserRole", ] diff --git a/backend/app/domain/enums/kafka.py b/backend/app/domain/enums/kafka.py index e1eceeb7..d248bd8b 100644 --- a/backend/app/domain/enums/kafka.py +++ b/backend/app/domain/enums/kafka.py @@ -1,5 +1,4 @@ from app.core.utils import StringEnum -from app.domain.enums.events import EventType class KafkaTopic(StringEnum): @@ -64,87 +63,3 @@ class GroupId(StringEnum): NOTIFICATION_SERVICE = "notification-service" DLQ_PROCESSOR = "dlq-processor" DLQ_MANAGER = "dlq-manager" - - -# Consumer group topic subscriptions -CONSUMER_GROUP_SUBSCRIPTIONS: dict[GroupId, set[KafkaTopic]] = { - GroupId.EXECUTION_COORDINATOR: { - KafkaTopic.EXECUTION_EVENTS, - KafkaTopic.EXECUTION_RESULTS, - }, - GroupId.K8S_WORKER: { - KafkaTopic.SAGA_COMMANDS, # Receives CreatePodCommand/DeletePodCommand from coordinator - }, - GroupId.POD_MONITOR: { - KafkaTopic.POD_EVENTS, - KafkaTopic.POD_STATUS_UPDATES, - }, - GroupId.RESULT_PROCESSOR: { - KafkaTopic.EXECUTION_EVENTS, # Listens for COMPLETED/FAILED/TIMEOUT, publishes to EXECUTION_RESULTS - }, - GroupId.SAGA_ORCHESTRATOR: { - # Orchestrator is triggered by domain events, specifically EXECUTION_REQUESTED, - # and emits commands on SAGA_COMMANDS. - KafkaTopic.EXECUTION_EVENTS, - KafkaTopic.SAGA_COMMANDS, - }, - GroupId.WEBSOCKET_GATEWAY: { - KafkaTopic.EXECUTION_EVENTS, - KafkaTopic.EXECUTION_RESULTS, - KafkaTopic.POD_EVENTS, - KafkaTopic.POD_STATUS_UPDATES, - }, - GroupId.NOTIFICATION_SERVICE: { - KafkaTopic.NOTIFICATION_EVENTS, - KafkaTopic.EXECUTION_EVENTS, - }, - GroupId.DLQ_PROCESSOR: { - KafkaTopic.DEAD_LETTER_QUEUE, - }, -} - -# Consumer group event filters -CONSUMER_GROUP_EVENTS: dict[GroupId, set[EventType]] = { - GroupId.EXECUTION_COORDINATOR: { - EventType.EXECUTION_REQUESTED, - EventType.EXECUTION_COMPLETED, - EventType.EXECUTION_FAILED, - EventType.EXECUTION_CANCELLED, - }, - GroupId.K8S_WORKER: { - EventType.EXECUTION_STARTED, - }, - GroupId.POD_MONITOR: { - EventType.POD_CREATED, - EventType.POD_RUNNING, - EventType.POD_SUCCEEDED, - EventType.POD_FAILED, - }, - GroupId.RESULT_PROCESSOR: { - EventType.EXECUTION_COMPLETED, - EventType.EXECUTION_FAILED, - EventType.EXECUTION_TIMEOUT, - }, - GroupId.SAGA_ORCHESTRATOR: { - EventType.EXECUTION_REQUESTED, - EventType.EXECUTION_COMPLETED, - EventType.EXECUTION_FAILED, - EventType.EXECUTION_TIMEOUT, - }, - GroupId.WEBSOCKET_GATEWAY: { - EventType.EXECUTION_REQUESTED, - EventType.EXECUTION_STARTED, - EventType.EXECUTION_COMPLETED, - EventType.EXECUTION_FAILED, - EventType.POD_CREATED, - EventType.POD_RUNNING, - EventType.RESULT_STORED, - }, - GroupId.NOTIFICATION_SERVICE: { - EventType.NOTIFICATION_CREATED, - EventType.EXECUTION_COMPLETED, - EventType.EXECUTION_FAILED, - EventType.EXECUTION_TIMEOUT, - }, - GroupId.DLQ_PROCESSOR: set(), -} diff --git a/backend/app/domain/events/event_models.py b/backend/app/domain/events/event_models.py index 6e152d41..3bdf58e7 100644 --- a/backend/app/domain/events/event_models.py +++ b/backend/app/domain/events/event_models.py @@ -6,7 +6,7 @@ from pydantic.dataclasses import dataclass from app.core.utils import StringEnum -from app.domain.enums.events import EventType +from app.domain.enums import EventType from app.domain.events.typed import DomainEvent MongoQueryValue = str | dict[str, str | list[str] | float | datetime] diff --git a/backend/app/domain/events/typed.py b/backend/app/domain/events/typed.py index fc9aa14e..f0a1e7ee 100644 --- a/backend/app/domain/events/typed.py +++ b/backend/app/domain/events/typed.py @@ -4,12 +4,16 @@ from pydantic import BaseModel, ConfigDict, Discriminator, Field, TypeAdapter -from app.domain.enums.auth import LoginMethod -from app.domain.enums.common import Environment -from app.domain.enums.events import EventType -from app.domain.enums.execution import QueuePriority -from app.domain.enums.notification import NotificationChannel, NotificationSeverity -from app.domain.enums.storage import ExecutionErrorType, StorageType +from app.domain.enums import ( + Environment, + EventType, + ExecutionErrorType, + LoginMethod, + NotificationChannel, + NotificationSeverity, + QueuePriority, + StorageType, +) class ResourceUsageDomain(BaseModel): diff --git a/backend/app/domain/execution/models.py b/backend/app/domain/execution/models.py index df8e38d7..45f04943 100644 --- a/backend/app/domain/execution/models.py +++ b/backend/app/domain/execution/models.py @@ -5,8 +5,7 @@ from pydantic import BaseModel, ConfigDict, Field -from app.domain.enums.execution import ExecutionStatus -from app.domain.enums.storage import ExecutionErrorType +from app.domain.enums import ExecutionErrorType, ExecutionStatus from app.domain.events.typed import EventMetadata, ResourceUsageDomain diff --git a/backend/app/domain/notification/models.py b/backend/app/domain/notification/models.py index 9b50326d..9a7ff4e3 100644 --- a/backend/app/domain/notification/models.py +++ b/backend/app/domain/notification/models.py @@ -6,11 +6,7 @@ from pydantic import BaseModel, ConfigDict, Field -from app.domain.enums.notification import ( - NotificationChannel, - NotificationSeverity, - NotificationStatus, -) +from app.domain.enums import NotificationChannel, NotificationSeverity, NotificationStatus class DomainNotification(BaseModel): diff --git a/backend/app/domain/replay/models.py b/backend/app/domain/replay/models.py index 2d1e3c9f..2aedbb7c 100644 --- a/backend/app/domain/replay/models.py +++ b/backend/app/domain/replay/models.py @@ -4,9 +4,7 @@ from pydantic import BaseModel, ConfigDict, Field -from app.domain.enums.events import EventType -from app.domain.enums.kafka import KafkaTopic -from app.domain.enums.replay import ReplayStatus, ReplayTarget, ReplayType +from app.domain.enums import EventType, KafkaTopic, ReplayStatus, ReplayTarget, ReplayType class ReplayError(BaseModel): diff --git a/backend/app/domain/saga/models.py b/backend/app/domain/saga/models.py index f95434be..f245a889 100644 --- a/backend/app/domain/saga/models.py +++ b/backend/app/domain/saga/models.py @@ -4,7 +4,7 @@ from pydantic import BaseModel, ConfigDict, Field -from app.domain.enums.saga import SagaState +from app.domain.enums import SagaState class Saga(BaseModel): diff --git a/backend/app/domain/sse/models.py b/backend/app/domain/sse/models.py index 8d954f17..2a783dc4 100644 --- a/backend/app/domain/sse/models.py +++ b/backend/app/domain/sse/models.py @@ -4,7 +4,7 @@ from pydantic import BaseModel, ConfigDict -from app.domain.enums.execution import ExecutionStatus +from app.domain.enums import ExecutionStatus class SSEExecutionStatusDomain(BaseModel): diff --git a/backend/app/domain/user/__init__.py b/backend/app/domain/user/__init__.py index 0eda945a..b110d774 100644 --- a/backend/app/domain/user/__init__.py +++ b/backend/app/domain/user/__init__.py @@ -1,4 +1,4 @@ -from app.domain.enums.user import UserRole +from app.domain.enums import UserRole from .exceptions import ( AdminAccessRequiredError, diff --git a/backend/app/domain/user/settings_models.py b/backend/app/domain/user/settings_models.py index 30af0354..6bf5f391 100644 --- a/backend/app/domain/user/settings_models.py +++ b/backend/app/domain/user/settings_models.py @@ -5,9 +5,7 @@ from pydantic import BaseModel, ConfigDict, Field -from app.domain.enums.common import Theme -from app.domain.enums.events import EventType -from app.domain.enums.notification import NotificationChannel +from app.domain.enums import EventType, NotificationChannel, Theme class DomainNotificationSettings(BaseModel): diff --git a/backend/app/domain/user/user_models.py b/backend/app/domain/user/user_models.py index 80d2e28e..1f9b16ce 100644 --- a/backend/app/domain/user/user_models.py +++ b/backend/app/domain/user/user_models.py @@ -4,7 +4,7 @@ from pydantic import BaseModel, ConfigDict from app.core.utils import StringEnum -from app.domain.enums.user import UserRole +from app.domain.enums import UserRole EMAIL_PATTERN = re.compile(r"^[a-zA-Z0-9._%+-]+@[a-zA-Z0-9.-]+\.[a-zA-Z]{2,}$") diff --git a/backend/app/events/core/producer.py b/backend/app/events/core/producer.py index 2e727aaa..17940d33 100644 --- a/backend/app/events/core/producer.py +++ b/backend/app/events/core/producer.py @@ -9,7 +9,7 @@ from app.core.tracing.utils import inject_trace_context from app.db.repositories.event_repository import EventRepository from app.dlq.models import DLQMessageStatus -from app.domain.enums.kafka import KafkaTopic +from app.domain.enums import KafkaTopic from app.domain.events.typed import DomainEvent from app.infrastructure.kafka.mappings import EVENT_TYPE_TO_TOPIC from app.settings import Settings diff --git a/backend/app/events/handlers.py b/backend/app/events/handlers.py index 19040925..684a71f8 100644 --- a/backend/app/events/handlers.py +++ b/backend/app/events/handlers.py @@ -13,8 +13,7 @@ from app.core.tracing.utils import extract_trace_context, get_tracer from app.dlq.manager import DLQManager from app.dlq.models import DLQMessage, DLQMessageStatus -from app.domain.enums.events import EventType -from app.domain.enums.kafka import CONSUMER_GROUP_SUBSCRIPTIONS, GroupId, KafkaTopic +from app.domain.enums import EventType, GroupId, KafkaTopic from app.domain.events.typed import ( CreatePodCommandEvent, DeletePodCommandEvent, @@ -26,6 +25,7 @@ ExecutionTimeoutEvent, ) from app.domain.idempotency import KeyStrategy +from app.infrastructure.kafka.mappings import CONSUMER_GROUP_SUBSCRIPTIONS from app.services.coordinator.coordinator import ExecutionCoordinator from app.services.idempotency import IdempotencyManager from app.services.k8s_worker import KubernetesWorker diff --git a/backend/app/infrastructure/kafka/mappings.py b/backend/app/infrastructure/kafka/mappings.py index e0a41100..2a9795de 100644 --- a/backend/app/infrastructure/kafka/mappings.py +++ b/backend/app/infrastructure/kafka/mappings.py @@ -1,8 +1,7 @@ from functools import lru_cache from typing import get_args, get_origin -from app.domain.enums.events import EventType -from app.domain.enums.kafka import KafkaTopic +from app.domain.enums import EventType, GroupId, KafkaTopic # EventType -> KafkaTopic routing EVENT_TYPE_TO_TOPIC: dict[EventType, KafkaTopic] = { @@ -102,3 +101,38 @@ def get_topic_for_event(event_type: EventType) -> KafkaTopic: def get_event_types_for_topic(topic: KafkaTopic) -> list[EventType]: """Get all event types that publish to a given topic.""" return [et for et, t in EVENT_TYPE_TO_TOPIC.items() if t == topic] + + +CONSUMER_GROUP_SUBSCRIPTIONS: dict[GroupId, set[KafkaTopic]] = { + GroupId.EXECUTION_COORDINATOR: { + KafkaTopic.EXECUTION_EVENTS, + KafkaTopic.EXECUTION_RESULTS, + }, + GroupId.K8S_WORKER: { + KafkaTopic.SAGA_COMMANDS, + }, + GroupId.POD_MONITOR: { + KafkaTopic.POD_EVENTS, + KafkaTopic.POD_STATUS_UPDATES, + }, + GroupId.RESULT_PROCESSOR: { + KafkaTopic.EXECUTION_EVENTS, + }, + GroupId.SAGA_ORCHESTRATOR: { + KafkaTopic.EXECUTION_EVENTS, + KafkaTopic.SAGA_COMMANDS, + }, + GroupId.WEBSOCKET_GATEWAY: { + KafkaTopic.EXECUTION_EVENTS, + KafkaTopic.EXECUTION_RESULTS, + KafkaTopic.POD_EVENTS, + KafkaTopic.POD_STATUS_UPDATES, + }, + GroupId.NOTIFICATION_SERVICE: { + KafkaTopic.NOTIFICATION_EVENTS, + KafkaTopic.EXECUTION_EVENTS, + }, + GroupId.DLQ_PROCESSOR: { + KafkaTopic.DEAD_LETTER_QUEUE, + }, +} diff --git a/backend/app/infrastructure/kafka/topics.py b/backend/app/infrastructure/kafka/topics.py index be5ae6d8..a664bb19 100644 --- a/backend/app/infrastructure/kafka/topics.py +++ b/backend/app/infrastructure/kafka/topics.py @@ -1,6 +1,6 @@ from typing import Any -from app.domain.enums.kafka import KafkaTopic +from app.domain.enums import KafkaTopic def get_all_topics() -> set[KafkaTopic]: diff --git a/backend/app/schemas_pydantic/admin_events.py b/backend/app/schemas_pydantic/admin_events.py index 4c8c7e61..9cceb9b0 100644 --- a/backend/app/schemas_pydantic/admin_events.py +++ b/backend/app/schemas_pydantic/admin_events.py @@ -2,7 +2,7 @@ from pydantic import BaseModel, ConfigDict, Field, computed_field -from app.domain.enums.events import EventType +from app.domain.enums import EventType from app.domain.events.event_models import EventSummary from app.domain.events.typed import DomainEvent from app.domain.replay import ReplayError diff --git a/backend/app/schemas_pydantic/events.py b/backend/app/schemas_pydantic/events.py index 93d560b4..2d040d25 100644 --- a/backend/app/schemas_pydantic/events.py +++ b/backend/app/schemas_pydantic/events.py @@ -4,8 +4,7 @@ from pydantic import BaseModel, ConfigDict, Field, field_validator -from app.domain.enums.common import Environment, SortOrder -from app.domain.enums.events import EventType +from app.domain.enums import Environment, EventType, SortOrder from app.domain.events.typed import ContainerStatusInfo, DomainEvent diff --git a/backend/app/schemas_pydantic/execution.py b/backend/app/schemas_pydantic/execution.py index 9331baf8..1854b87e 100644 --- a/backend/app/schemas_pydantic/execution.py +++ b/backend/app/schemas_pydantic/execution.py @@ -5,8 +5,7 @@ from pydantic import BaseModel, ConfigDict, Field, model_validator -from app.domain.enums.execution import ExecutionStatus -from app.domain.enums.storage import ExecutionErrorType +from app.domain.enums import ExecutionErrorType, ExecutionStatus from app.runtime_registry import SUPPORTED_RUNTIMES diff --git a/backend/app/schemas_pydantic/health_dashboard.py b/backend/app/schemas_pydantic/health_dashboard.py index 634bde38..6767aafd 100644 --- a/backend/app/schemas_pydantic/health_dashboard.py +++ b/backend/app/schemas_pydantic/health_dashboard.py @@ -2,7 +2,7 @@ from pydantic import BaseModel, ConfigDict, Field -from app.domain.enums.health import AlertSeverity +from app.domain.enums import AlertSeverity class HealthAlert(BaseModel): diff --git a/backend/app/schemas_pydantic/notification.py b/backend/app/schemas_pydantic/notification.py index 949e87b0..99888046 100644 --- a/backend/app/schemas_pydantic/notification.py +++ b/backend/app/schemas_pydantic/notification.py @@ -4,11 +4,7 @@ from pydantic import BaseModel, ConfigDict, Field, field_validator -from app.domain.enums.notification import ( - NotificationChannel, - NotificationSeverity, - NotificationStatus, -) +from app.domain.enums import NotificationChannel, NotificationSeverity, NotificationStatus # Templates are removed in the unified model diff --git a/backend/app/schemas_pydantic/replay.py b/backend/app/schemas_pydantic/replay.py index e498bbda..3728d834 100644 --- a/backend/app/schemas_pydantic/replay.py +++ b/backend/app/schemas_pydantic/replay.py @@ -2,9 +2,7 @@ from pydantic import BaseModel, ConfigDict, Field, computed_field -from app.domain.enums.events import EventType -from app.domain.enums.kafka import KafkaTopic -from app.domain.enums.replay import ReplayStatus, ReplayTarget, ReplayType +from app.domain.enums import EventType, KafkaTopic, ReplayStatus, ReplayTarget, ReplayType from app.domain.replay import ReplayFilter diff --git a/backend/app/schemas_pydantic/replay_models.py b/backend/app/schemas_pydantic/replay_models.py index 3d3a9b46..0663fe4d 100644 --- a/backend/app/schemas_pydantic/replay_models.py +++ b/backend/app/schemas_pydantic/replay_models.py @@ -4,9 +4,7 @@ from pydantic import BaseModel, ConfigDict, Field -from app.domain.enums.events import EventType -from app.domain.enums.kafka import KafkaTopic -from app.domain.enums.replay import ReplayStatus, ReplayTarget, ReplayType +from app.domain.enums import EventType, KafkaTopic, ReplayStatus, ReplayTarget, ReplayType from app.domain.replay import ReplayError diff --git a/backend/app/schemas_pydantic/saga.py b/backend/app/schemas_pydantic/saga.py index 130c8296..dc6beb19 100644 --- a/backend/app/schemas_pydantic/saga.py +++ b/backend/app/schemas_pydantic/saga.py @@ -2,7 +2,7 @@ from pydantic import BaseModel, ConfigDict -from app.domain.enums.saga import SagaState +from app.domain.enums import SagaState class SagaStatusResponse(BaseModel): diff --git a/backend/app/schemas_pydantic/sse.py b/backend/app/schemas_pydantic/sse.py index b6d9874d..e79711c0 100644 --- a/backend/app/schemas_pydantic/sse.py +++ b/backend/app/schemas_pydantic/sse.py @@ -3,10 +3,7 @@ from pydantic import BaseModel, Field -from app.domain.enums.events import EventType -from app.domain.enums.execution import ExecutionStatus -from app.domain.enums.notification import NotificationSeverity, NotificationStatus -from app.domain.enums.sse import SSEControlEvent +from app.domain.enums import EventType, ExecutionStatus, NotificationSeverity, NotificationStatus, SSEControlEvent from app.schemas_pydantic.execution import ExecutionResult, ResourceUsage # Type variable for generic Redis message parsing diff --git a/backend/app/schemas_pydantic/user.py b/backend/app/schemas_pydantic/user.py index 8ed324bf..d89d706c 100644 --- a/backend/app/schemas_pydantic/user.py +++ b/backend/app/schemas_pydantic/user.py @@ -3,7 +3,7 @@ from pydantic import BaseModel, ConfigDict, EmailStr, Field -from app.domain.enums.user import UserRole +from app.domain.enums import UserRole from app.domain.rate_limit import EndpointGroup, EndpointUsageStats, RateLimitAlgorithm diff --git a/backend/app/schemas_pydantic/user_settings.py b/backend/app/schemas_pydantic/user_settings.py index e16f1849..83a96eb8 100644 --- a/backend/app/schemas_pydantic/user_settings.py +++ b/backend/app/schemas_pydantic/user_settings.py @@ -3,9 +3,7 @@ from pydantic import BaseModel, ConfigDict, Field, field_validator -from app.domain.enums.common import Theme -from app.domain.enums.events import EventType -from app.domain.enums.notification import NotificationChannel +from app.domain.enums import EventType, NotificationChannel, Theme class NotificationSettings(BaseModel): diff --git a/backend/app/services/admin/admin_events_service.py b/backend/app/services/admin/admin_events_service.py index f3380578..c0de6533 100644 --- a/backend/app/services/admin/admin_events_service.py +++ b/backend/app/services/admin/admin_events_service.py @@ -11,7 +11,7 @@ from app.db.repositories.admin import AdminEventsRepository from app.domain.admin import ReplaySessionStatusDetail from app.domain.admin.replay_updates import ReplaySessionUpdate -from app.domain.enums.replay import ReplayStatus, ReplayTarget, ReplayType +from app.domain.enums import ReplayStatus, ReplayTarget, ReplayType from app.domain.events.event_models import ( EventBrowseResult, EventDetail, diff --git a/backend/app/services/admin/admin_user_service.py b/backend/app/services/admin/admin_user_service.py index 415fc684..95aa232e 100644 --- a/backend/app/services/admin/admin_user_service.py +++ b/backend/app/services/admin/admin_user_service.py @@ -5,9 +5,7 @@ from app.core.security import SecurityService from app.db.repositories.admin.admin_user_repository import AdminUserRepository from app.domain.admin import AdminUserOverviewDomain, DerivedCountsDomain, RateLimitSummaryDomain -from app.domain.enums.events import EventType -from app.domain.enums.execution import ExecutionStatus -from app.domain.enums.user import UserRole +from app.domain.enums import EventType, ExecutionStatus, UserRole from app.domain.rate_limit import RateLimitUpdateResult, UserRateLimit, UserRateLimitsResult from app.domain.user import DomainUserCreate, PasswordReset, User, UserDeleteResult, UserListResult, UserUpdate from app.schemas_pydantic.user import UserCreate diff --git a/backend/app/services/auth_service.py b/backend/app/services/auth_service.py index 8d851861..c4c88c6f 100644 --- a/backend/app/services/auth_service.py +++ b/backend/app/services/auth_service.py @@ -4,7 +4,7 @@ from app.core.security import SecurityService from app.db.repositories.user_repository import UserRepository -from app.domain.enums.user import UserRole +from app.domain.enums import UserRole from app.domain.user import AdminAccessRequiredError, AuthenticationRequiredError, InvalidCredentialsError, User diff --git a/backend/app/services/coordinator/__init__.py b/backend/app/services/coordinator/__init__.py index 03ccbd8d..7f925273 100644 --- a/backend/app/services/coordinator/__init__.py +++ b/backend/app/services/coordinator/__init__.py @@ -1,4 +1,4 @@ -from app.domain.enums.execution import QueuePriority +from app.domain.enums import QueuePriority from app.services.coordinator.coordinator import ExecutionCoordinator, QueueRejectError __all__ = [ diff --git a/backend/app/services/coordinator/coordinator.py b/backend/app/services/coordinator/coordinator.py index d0bcffbc..a2c6bd19 100644 --- a/backend/app/services/coordinator/coordinator.py +++ b/backend/app/services/coordinator/coordinator.py @@ -7,8 +7,7 @@ from app.core.metrics import CoordinatorMetrics from app.db.repositories.execution_repository import ExecutionRepository -from app.domain.enums.execution import QueuePriority -from app.domain.enums.storage import ExecutionErrorType +from app.domain.enums import ExecutionErrorType, QueuePriority from app.domain.events.typed import ( CreatePodCommandEvent, EventMetadata, diff --git a/backend/app/services/event_replay/__init__.py b/backend/app/services/event_replay/__init__.py index aab4eb0d..e0d45753 100644 --- a/backend/app/services/event_replay/__init__.py +++ b/backend/app/services/event_replay/__init__.py @@ -1,4 +1,4 @@ -from app.domain.enums.replay import ReplayStatus, ReplayTarget, ReplayType +from app.domain.enums import ReplayStatus, ReplayTarget, ReplayType from app.domain.replay import ReplayConfig, ReplayFilter from app.services.event_replay.replay_service import EventReplayService diff --git a/backend/app/services/event_replay/replay_service.py b/backend/app/services/event_replay/replay_service.py index 84468380..a42e0f71 100644 --- a/backend/app/services/event_replay/replay_service.py +++ b/backend/app/services/event_replay/replay_service.py @@ -13,7 +13,7 @@ from app.core.metrics import ReplayMetrics from app.db.repositories.replay_repository import ReplayRepository from app.domain.admin.replay_updates import ReplaySessionUpdate -from app.domain.enums.replay import ReplayStatus, ReplayTarget +from app.domain.enums import ReplayStatus, ReplayTarget from app.domain.events.typed import DomainEvent, DomainEventAdapter from app.domain.replay import ( CleanupResult, diff --git a/backend/app/services/event_service.py b/backend/app/services/event_service.py index 1418df26..c274de9e 100644 --- a/backend/app/services/event_service.py +++ b/backend/app/services/event_service.py @@ -2,8 +2,7 @@ from typing import Any from app.db.repositories.event_repository import EventRepository -from app.domain.enums.events import EventType -from app.domain.enums.user import UserRole +from app.domain.enums import EventType, UserRole from app.domain.events import ( ArchivedEvent, DomainEvent, diff --git a/backend/app/services/execution_service.py b/backend/app/services/execution_service.py index 59a7b556..96734880 100644 --- a/backend/app/services/execution_service.py +++ b/backend/app/services/execution_service.py @@ -8,8 +8,7 @@ from app.core.metrics import ExecutionMetrics from app.db.repositories.event_repository import EventRepository from app.db.repositories.execution_repository import ExecutionRepository -from app.domain.enums.events import EventType -from app.domain.enums.execution import ExecutionStatus, QueuePriority +from app.domain.enums import EventType, ExecutionStatus, QueuePriority from app.domain.events.typed import ( DomainEvent, EventMetadata, diff --git a/backend/app/services/grafana_alert_processor.py b/backend/app/services/grafana_alert_processor.py index 2513cfbb..1c6cd59a 100644 --- a/backend/app/services/grafana_alert_processor.py +++ b/backend/app/services/grafana_alert_processor.py @@ -3,8 +3,7 @@ import logging from typing import Any -from app.domain.enums.notification import NotificationSeverity -from app.domain.enums.user import UserRole +from app.domain.enums import NotificationSeverity, UserRole from app.schemas_pydantic.grafana import GrafanaAlertItem, GrafanaWebhook from app.services.notification_service import NotificationService diff --git a/backend/app/services/k8s_worker/worker.py b/backend/app/services/k8s_worker/worker.py index 3d6e4bda..3fdf71fc 100644 --- a/backend/app/services/k8s_worker/worker.py +++ b/backend/app/services/k8s_worker/worker.py @@ -8,7 +8,7 @@ from kubernetes_asyncio.client.rest import ApiException from app.core.metrics import EventMetrics, ExecutionMetrics, KubernetesMetrics -from app.domain.enums.storage import ExecutionErrorType +from app.domain.enums import ExecutionErrorType from app.domain.events.typed import ( CreatePodCommandEvent, DeletePodCommandEvent, diff --git a/backend/app/services/kafka_event_service.py b/backend/app/services/kafka_event_service.py index deca49a3..d344f7cd 100644 --- a/backend/app/services/kafka_event_service.py +++ b/backend/app/services/kafka_event_service.py @@ -8,7 +8,7 @@ from app.core.correlation import CorrelationContext from app.core.metrics import EventMetrics -from app.domain.enums.events import EventType +from app.domain.enums import EventType from app.domain.events import DomainEventAdapter from app.domain.events.typed import DomainEvent, EventMetadata from app.events.core import UnifiedProducer diff --git a/backend/app/services/notification_service.py b/backend/app/services/notification_service.py index 06fbeccf..2185fe54 100644 --- a/backend/app/services/notification_service.py +++ b/backend/app/services/notification_service.py @@ -10,12 +10,7 @@ from app.core.metrics import NotificationMetrics from app.core.tracing.utils import add_span_attributes from app.db.repositories.notification_repository import NotificationRepository -from app.domain.enums.notification import ( - NotificationChannel, - NotificationSeverity, - NotificationStatus, -) -from app.domain.enums.user import UserRole +from app.domain.enums import NotificationChannel, NotificationSeverity, NotificationStatus, UserRole from app.domain.events.typed import ( EventMetadata, ExecutionCompletedEvent, diff --git a/backend/app/services/pod_monitor/config.py b/backend/app/services/pod_monitor/config.py index f862f016..5e8c358d 100644 --- a/backend/app/services/pod_monitor/config.py +++ b/backend/app/services/pod_monitor/config.py @@ -1,7 +1,7 @@ import os from dataclasses import dataclass, field -from app.domain.enums.events import EventType +from app.domain.enums import EventType from app.infrastructure.kafka import get_topic_for_event from app.services.pod_monitor.event_mapper import PodPhase diff --git a/backend/app/services/pod_monitor/event_mapper.py b/backend/app/services/pod_monitor/event_mapper.py index 0d4a0bac..8b5b37a8 100644 --- a/backend/app/services/pod_monitor/event_mapper.py +++ b/backend/app/services/pod_monitor/event_mapper.py @@ -5,8 +5,7 @@ from kubernetes_asyncio import client as k8s_client -from app.domain.enums.kafka import GroupId -from app.domain.enums.storage import ExecutionErrorType +from app.domain.enums import ExecutionErrorType, GroupId from app.domain.events.typed import ( ContainerStatusInfo, DomainEvent, diff --git a/backend/app/services/result_processor/processor.py b/backend/app/services/result_processor/processor.py index 55909b62..a9f4a1bf 100644 --- a/backend/app/services/result_processor/processor.py +++ b/backend/app/services/result_processor/processor.py @@ -2,9 +2,7 @@ from app.core.metrics import ExecutionMetrics from app.db.repositories.execution_repository import ExecutionRepository -from app.domain.enums.execution import ExecutionStatus -from app.domain.enums.kafka import GroupId -from app.domain.enums.storage import ExecutionErrorType, StorageType +from app.domain.enums import ExecutionErrorType, ExecutionStatus, GroupId, StorageType from app.domain.events.typed import ( DomainEvent, EventMetadata, diff --git a/backend/app/services/saga/__init__.py b/backend/app/services/saga/__init__.py index de5ad07c..3a5c87a0 100644 --- a/backend/app/services/saga/__init__.py +++ b/backend/app/services/saga/__init__.py @@ -1,4 +1,4 @@ -from app.domain.enums.saga import SagaState +from app.domain.enums import SagaState from app.domain.saga.models import SagaConfig, SagaInstance from app.services.saga.execution_saga import ( AllocateResourcesStep, diff --git a/backend/app/services/saga/saga_orchestrator.py b/backend/app/services/saga/saga_orchestrator.py index 5fead192..c3f1180d 100644 --- a/backend/app/services/saga/saga_orchestrator.py +++ b/backend/app/services/saga/saga_orchestrator.py @@ -8,7 +8,7 @@ from app.core.tracing.utils import get_tracer from app.db.repositories.resource_allocation_repository import ResourceAllocationRepository from app.db.repositories.saga_repository import SagaRepository -from app.domain.enums.saga import SagaState +from app.domain.enums import SagaState from app.domain.events.typed import ( DomainEvent, EventMetadata, diff --git a/backend/app/services/sse/redis_bus.py b/backend/app/services/sse/redis_bus.py index 1426056e..d7375a85 100644 --- a/backend/app/services/sse/redis_bus.py +++ b/backend/app/services/sse/redis_bus.py @@ -6,7 +6,7 @@ import redis.asyncio as redis from pydantic import BaseModel -from app.domain.enums.events import EventType +from app.domain.enums import EventType from app.domain.events.typed import DomainEvent from app.schemas_pydantic.sse import RedisNotificationMessage, RedisSSEMessage diff --git a/backend/app/services/sse/sse_service.py b/backend/app/services/sse/sse_service.py index 2edf8bef..62071b79 100644 --- a/backend/app/services/sse/sse_service.py +++ b/backend/app/services/sse/sse_service.py @@ -6,9 +6,7 @@ from app.core.metrics import ConnectionMetrics from app.db.repositories.sse_repository import SSERepository -from app.domain.enums.events import EventType -from app.domain.enums.notification import NotificationChannel -from app.domain.enums.sse import SSEControlEvent +from app.domain.enums import EventType, NotificationChannel, SSEControlEvent from app.schemas_pydantic.execution import ExecutionResult from app.schemas_pydantic.notification import NotificationResponse from app.schemas_pydantic.sse import ( diff --git a/backend/app/services/user_settings_service.py b/backend/app/services/user_settings_service.py index a2a607be..49dce6da 100644 --- a/backend/app/services/user_settings_service.py +++ b/backend/app/services/user_settings_service.py @@ -5,8 +5,7 @@ from cachetools import TTLCache from app.db.repositories.user_settings_repository import UserSettingsRepository -from app.domain.enums import Theme -from app.domain.enums.events import EventType +from app.domain.enums import EventType, Theme from app.domain.user import ( DomainEditorSettings, DomainNotificationSettings, diff --git a/backend/tests/conftest.py b/backend/tests/conftest.py index 843f0359..37252e3d 100644 --- a/backend/tests/conftest.py +++ b/backend/tests/conftest.py @@ -7,7 +7,7 @@ import pytest import pytest_asyncio import redis.asyncio as redis -from app.domain.enums.execution import QueuePriority +from app.domain.enums import QueuePriority from app.domain.events.typed import EventMetadata, ExecutionRequestedEvent from app.main import create_app from app.settings import Settings diff --git a/backend/tests/e2e/conftest.py b/backend/tests/e2e/conftest.py index d0d767fb..deaf14bc 100644 --- a/backend/tests/e2e/conftest.py +++ b/backend/tests/e2e/conftest.py @@ -9,9 +9,7 @@ import pytest_asyncio from aiokafka import AIOKafkaConsumer from app.db.docs.saga import SagaDocument -from app.domain.enums.events import EventType -from app.domain.enums.kafka import KafkaTopic -from app.domain.enums.user import UserRole +from app.domain.enums import EventType, KafkaTopic, UserRole from app.domain.events.typed import DomainEvent, DomainEventAdapter from app.schemas_pydantic.execution import ExecutionRequest, ExecutionResponse from app.schemas_pydantic.notification import NotificationListResponse, NotificationResponse diff --git a/backend/tests/e2e/db/repositories/test_dlq_repository.py b/backend/tests/e2e/db/repositories/test_dlq_repository.py index 9464d087..7197f08b 100644 --- a/backend/tests/e2e/db/repositories/test_dlq_repository.py +++ b/backend/tests/e2e/db/repositories/test_dlq_repository.py @@ -5,7 +5,7 @@ from app.db.docs import DLQMessageDocument from app.db.repositories.dlq_repository import DLQRepository from app.dlq import DLQMessageStatus -from app.domain.enums.events import EventType +from app.domain.enums import EventType pytestmark = pytest.mark.e2e diff --git a/backend/tests/e2e/db/repositories/test_execution_repository.py b/backend/tests/e2e/db/repositories/test_execution_repository.py index ce701bd8..a6db17e3 100644 --- a/backend/tests/e2e/db/repositories/test_execution_repository.py +++ b/backend/tests/e2e/db/repositories/test_execution_repository.py @@ -3,7 +3,7 @@ import pytest from app.db.repositories.execution_repository import ExecutionRepository -from app.domain.enums.execution import ExecutionStatus +from app.domain.enums import ExecutionStatus from app.domain.execution import DomainExecutionCreate, DomainExecutionUpdate _test_logger = logging.getLogger("test.db.repositories.execution_repository") diff --git a/backend/tests/e2e/dlq/test_dlq_discard.py b/backend/tests/e2e/dlq/test_dlq_discard.py index 2c4650f4..39c023fb 100644 --- a/backend/tests/e2e/dlq/test_dlq_discard.py +++ b/backend/tests/e2e/dlq/test_dlq_discard.py @@ -6,7 +6,7 @@ from app.db.docs import DLQMessageDocument from app.db.repositories.dlq_repository import DLQRepository from app.dlq.models import DLQMessageStatus -from app.domain.enums.kafka import KafkaTopic +from app.domain.enums import KafkaTopic from dishka import AsyncContainer from tests.conftest import make_execution_requested_event diff --git a/backend/tests/e2e/dlq/test_dlq_manager.py b/backend/tests/e2e/dlq/test_dlq_manager.py index 9bc5ab35..35a90adf 100644 --- a/backend/tests/e2e/dlq/test_dlq_manager.py +++ b/backend/tests/e2e/dlq/test_dlq_manager.py @@ -11,8 +11,7 @@ from app.db.repositories.dlq_repository import DLQRepository from app.dlq.manager import DLQManager from app.dlq.models import DLQMessage -from app.domain.enums.events import EventType -from app.domain.enums.kafka import KafkaTopic +from app.domain.enums import EventType, KafkaTopic from app.domain.events.typed import DLQMessageReceivedEvent, DomainEventAdapter from app.settings import Settings from dishka import AsyncContainer diff --git a/backend/tests/e2e/dlq/test_dlq_retry.py b/backend/tests/e2e/dlq/test_dlq_retry.py index d01fefe7..06425c58 100644 --- a/backend/tests/e2e/dlq/test_dlq_retry.py +++ b/backend/tests/e2e/dlq/test_dlq_retry.py @@ -6,7 +6,7 @@ from app.db.docs import DLQMessageDocument from app.db.repositories.dlq_repository import DLQRepository from app.dlq.models import DLQMessageStatus -from app.domain.enums.kafka import KafkaTopic +from app.domain.enums import KafkaTopic from dishka import AsyncContainer from tests.conftest import make_execution_requested_event diff --git a/backend/tests/e2e/notifications/test_notification_sse.py b/backend/tests/e2e/notifications/test_notification_sse.py index e3007a6d..1a34f840 100644 --- a/backend/tests/e2e/notifications/test_notification_sse.py +++ b/backend/tests/e2e/notifications/test_notification_sse.py @@ -2,7 +2,7 @@ from uuid import uuid4 import pytest -from app.domain.enums.notification import NotificationChannel, NotificationSeverity +from app.domain.enums import NotificationChannel, NotificationSeverity from app.schemas_pydantic.sse import RedisNotificationMessage from app.services.notification_service import NotificationService from app.services.sse.redis_bus import SSERedisBus diff --git a/backend/tests/e2e/result_processor/test_result_processor.py b/backend/tests/e2e/result_processor/test_result_processor.py index 1d57f851..e9a2a43c 100644 --- a/backend/tests/e2e/result_processor/test_result_processor.py +++ b/backend/tests/e2e/result_processor/test_result_processor.py @@ -4,7 +4,7 @@ from app.core.metrics import ExecutionMetrics from app.db.docs import ExecutionDocument from app.db.repositories.execution_repository import ExecutionRepository -from app.domain.enums.execution import ExecutionStatus +from app.domain.enums import ExecutionStatus from app.domain.events.typed import ( EventMetadata, ExecutionCompletedEvent, diff --git a/backend/tests/e2e/services/admin/test_admin_user_service.py b/backend/tests/e2e/services/admin/test_admin_user_service.py index 1fbf3223..eb399c57 100644 --- a/backend/tests/e2e/services/admin/test_admin_user_service.py +++ b/backend/tests/e2e/services/admin/test_admin_user_service.py @@ -1,6 +1,6 @@ import pytest from app.db.docs import UserDocument -from app.domain.enums.user import UserRole +from app.domain.enums import UserRole from app.services.admin import AdminUserService from dishka import AsyncContainer diff --git a/backend/tests/e2e/services/coordinator/test_execution_coordinator.py b/backend/tests/e2e/services/coordinator/test_execution_coordinator.py index 8335c4b0..26a6cdba 100644 --- a/backend/tests/e2e/services/coordinator/test_execution_coordinator.py +++ b/backend/tests/e2e/services/coordinator/test_execution_coordinator.py @@ -1,5 +1,5 @@ import pytest -from app.domain.enums.execution import QueuePriority +from app.domain.enums import QueuePriority from app.services.coordinator.coordinator import ExecutionCoordinator from dishka import AsyncContainer from tests.conftest import make_execution_requested_event diff --git a/backend/tests/e2e/services/events/test_kafka_event_service.py b/backend/tests/e2e/services/events/test_kafka_event_service.py index 1a02e800..d2781fd9 100644 --- a/backend/tests/e2e/services/events/test_kafka_event_service.py +++ b/backend/tests/e2e/services/events/test_kafka_event_service.py @@ -1,7 +1,6 @@ import pytest from app.db.repositories import EventRepository -from app.domain.enums.events import EventType -from app.domain.enums.execution import ExecutionStatus +from app.domain.enums import EventType, ExecutionStatus from app.services.kafka_event_service import KafkaEventService from dishka import AsyncContainer diff --git a/backend/tests/e2e/services/execution/test_execution_service.py b/backend/tests/e2e/services/execution/test_execution_service.py index 8ae06e85..d17ae90d 100644 --- a/backend/tests/e2e/services/execution/test_execution_service.py +++ b/backend/tests/e2e/services/execution/test_execution_service.py @@ -1,8 +1,7 @@ import uuid import pytest -from app.domain.enums.events import EventType -from app.domain.enums.execution import ExecutionStatus +from app.domain.enums import EventType, ExecutionStatus from app.domain.execution import ResourceLimitsDomain from app.domain.execution.exceptions import ExecutionNotFoundError from app.services.execution_service import ExecutionService diff --git a/backend/tests/e2e/services/notifications/test_notification_service.py b/backend/tests/e2e/services/notifications/test_notification_service.py index 9e7a7c8c..ffa63932 100644 --- a/backend/tests/e2e/services/notifications/test_notification_service.py +++ b/backend/tests/e2e/services/notifications/test_notification_service.py @@ -2,10 +2,7 @@ import pytest from app.db.repositories import NotificationRepository -from app.domain.enums.notification import ( - NotificationChannel, - NotificationSeverity, -) +from app.domain.enums import NotificationChannel, NotificationSeverity from app.domain.notification import ( DomainNotificationListResult, NotificationNotFoundError, diff --git a/backend/tests/e2e/services/replay/test_replay_service.py b/backend/tests/e2e/services/replay/test_replay_service.py index 81c5c9bd..e7cbb584 100644 --- a/backend/tests/e2e/services/replay/test_replay_service.py +++ b/backend/tests/e2e/services/replay/test_replay_service.py @@ -1,5 +1,5 @@ import pytest -from app.domain.enums.replay import ReplayStatus, ReplayTarget, ReplayType +from app.domain.enums import ReplayStatus, ReplayTarget, ReplayType from app.domain.replay.exceptions import ReplaySessionNotFoundError from app.services.event_replay import EventReplayService, ReplayConfig, ReplayFilter from dishka import AsyncContainer diff --git a/backend/tests/e2e/services/saga/test_saga_service.py b/backend/tests/e2e/services/saga/test_saga_service.py index 1d936625..f41888ac 100644 --- a/backend/tests/e2e/services/saga/test_saga_service.py +++ b/backend/tests/e2e/services/saga/test_saga_service.py @@ -3,8 +3,7 @@ import pytest from app.db.repositories import ExecutionRepository, SagaRepository -from app.domain.enums import SagaState -from app.domain.enums.user import UserRole +from app.domain.enums import SagaState, UserRole from app.domain.execution import DomainExecutionCreate from app.domain.saga.exceptions import SagaAccessDeniedError, SagaNotFoundError from app.domain.saga.models import Saga, SagaListResult diff --git a/backend/tests/e2e/services/sse/test_redis_bus.py b/backend/tests/e2e/services/sse/test_redis_bus.py index 8d0ac726..0c712a78 100644 --- a/backend/tests/e2e/services/sse/test_redis_bus.py +++ b/backend/tests/e2e/services/sse/test_redis_bus.py @@ -5,8 +5,7 @@ import pytest import redis.asyncio as redis_async -from app.domain.enums.events import EventType -from app.domain.enums.notification import NotificationSeverity, NotificationStatus +from app.domain.enums import EventType, NotificationSeverity, NotificationStatus from app.domain.events.typed import EventMetadata, ExecutionCompletedEvent from app.schemas_pydantic.sse import RedisNotificationMessage, RedisSSEMessage from app.services.sse.redis_bus import SSERedisBus diff --git a/backend/tests/e2e/test_admin_events_routes.py b/backend/tests/e2e/test_admin_events_routes.py index 3413d032..5af0087b 100644 --- a/backend/tests/e2e/test_admin_events_routes.py +++ b/backend/tests/e2e/test_admin_events_routes.py @@ -3,8 +3,7 @@ import pytest import pytest_asyncio from app.db.repositories.event_repository import EventRepository -from app.domain.enums.events import EventType -from app.domain.enums.replay import ReplayStatus +from app.domain.enums import EventType, ReplayStatus from app.domain.events.typed import DomainEvent from app.schemas_pydantic.admin_events import ( EventBrowseRequest, diff --git a/backend/tests/e2e/test_admin_users_routes.py b/backend/tests/e2e/test_admin_users_routes.py index db4edead..7e6ece05 100644 --- a/backend/tests/e2e/test_admin_users_routes.py +++ b/backend/tests/e2e/test_admin_users_routes.py @@ -1,7 +1,7 @@ import uuid import pytest -from app.domain.enums.user import UserRole +from app.domain.enums import UserRole from app.schemas_pydantic.admin_user_overview import AdminUserOverview from app.schemas_pydantic.user import ( DeleteUserResponse, diff --git a/backend/tests/e2e/test_auth_routes.py b/backend/tests/e2e/test_auth_routes.py index aa6f825e..0f2b3a78 100644 --- a/backend/tests/e2e/test_auth_routes.py +++ b/backend/tests/e2e/test_auth_routes.py @@ -1,7 +1,7 @@ import uuid import pytest -from app.domain.enums.user import UserRole +from app.domain.enums import UserRole from app.schemas_pydantic.user import ( LoginResponse, MessageResponse, diff --git a/backend/tests/e2e/test_dlq_routes.py b/backend/tests/e2e/test_dlq_routes.py index bd396c59..ad7130eb 100644 --- a/backend/tests/e2e/test_dlq_routes.py +++ b/backend/tests/e2e/test_dlq_routes.py @@ -2,7 +2,7 @@ import pytest_asyncio from app.db.docs.dlq import DLQMessageDocument from app.dlq.models import DLQMessageStatus, RetryStrategy -from app.domain.enums.events import EventType +from app.domain.enums import EventType from app.schemas_pydantic.dlq import ( DLQBatchRetryResponse, DLQMessageDetail, diff --git a/backend/tests/e2e/test_events_routes.py b/backend/tests/e2e/test_events_routes.py index 5a6050ce..96e60f03 100644 --- a/backend/tests/e2e/test_events_routes.py +++ b/backend/tests/e2e/test_events_routes.py @@ -1,5 +1,5 @@ import pytest -from app.domain.enums.events import EventType +from app.domain.enums import EventType from app.domain.events.typed import DomainEvent, ExecutionRequestedEvent from app.schemas_pydantic.events import ( DeleteEventResponse, diff --git a/backend/tests/e2e/test_execution_routes.py b/backend/tests/e2e/test_execution_routes.py index 1e5f1373..3facdeda 100644 --- a/backend/tests/e2e/test_execution_routes.py +++ b/backend/tests/e2e/test_execution_routes.py @@ -8,8 +8,7 @@ import asyncio import pytest -from app.domain.enums.events import EventType -from app.domain.enums.execution import ExecutionStatus +from app.domain.enums import EventType, ExecutionStatus from app.domain.events.typed import ExecutionDomainEvent from app.schemas_pydantic.execution import ( CancelExecutionRequest, diff --git a/backend/tests/e2e/test_k8s_worker_create_pod.py b/backend/tests/e2e/test_k8s_worker_create_pod.py index 91d9c0dd..ecf48a07 100644 --- a/backend/tests/e2e/test_k8s_worker_create_pod.py +++ b/backend/tests/e2e/test_k8s_worker_create_pod.py @@ -3,7 +3,7 @@ import pytest from app.core.metrics import EventMetrics -from app.domain.enums.execution import QueuePriority +from app.domain.enums import QueuePriority from app.domain.events.typed import CreatePodCommandEvent, EventMetadata from app.events.core import UnifiedProducer from app.services.k8s_worker import KubernetesWorker diff --git a/backend/tests/e2e/test_notifications_routes.py b/backend/tests/e2e/test_notifications_routes.py index 43565dda..7812a0d8 100644 --- a/backend/tests/e2e/test_notifications_routes.py +++ b/backend/tests/e2e/test_notifications_routes.py @@ -1,5 +1,5 @@ import pytest -from app.domain.enums.notification import NotificationChannel, NotificationSeverity, NotificationStatus +from app.domain.enums import NotificationChannel, NotificationSeverity, NotificationStatus from app.schemas_pydantic.execution import ExecutionResponse from app.schemas_pydantic.notification import ( DeleteNotificationResponse, diff --git a/backend/tests/e2e/test_replay_routes.py b/backend/tests/e2e/test_replay_routes.py index e19be9df..61456b2d 100644 --- a/backend/tests/e2e/test_replay_routes.py +++ b/backend/tests/e2e/test_replay_routes.py @@ -1,6 +1,5 @@ import pytest -from app.domain.enums.events import EventType -from app.domain.enums.replay import ReplayStatus, ReplayTarget, ReplayType +from app.domain.enums import EventType, ReplayStatus, ReplayTarget, ReplayType from app.domain.replay import ReplayFilter from app.schemas_pydantic.replay import ( CleanupResponse, diff --git a/backend/tests/e2e/test_saga_routes.py b/backend/tests/e2e/test_saga_routes.py index 8b36a6aa..df8d2efc 100644 --- a/backend/tests/e2e/test_saga_routes.py +++ b/backend/tests/e2e/test_saga_routes.py @@ -1,5 +1,5 @@ import pytest -from app.domain.enums.saga import SagaState +from app.domain.enums import SagaState from app.schemas_pydantic.execution import ExecutionRequest, ExecutionResponse from app.schemas_pydantic.saga import ( SagaCancellationResponse, diff --git a/backend/tests/e2e/test_user_settings_routes.py b/backend/tests/e2e/test_user_settings_routes.py index 46225bd9..b9eb4bdf 100644 --- a/backend/tests/e2e/test_user_settings_routes.py +++ b/backend/tests/e2e/test_user_settings_routes.py @@ -1,5 +1,5 @@ import pytest -from app.domain.enums.common import Theme +from app.domain.enums import Theme from app.schemas_pydantic.user_settings import ( EditorSettings, NotificationSettings, diff --git a/backend/tests/unit/core/metrics/test_execution_and_events_metrics.py b/backend/tests/unit/core/metrics/test_execution_and_events_metrics.py index fdd09bdc..49e4f7d3 100644 --- a/backend/tests/unit/core/metrics/test_execution_and_events_metrics.py +++ b/backend/tests/unit/core/metrics/test_execution_and_events_metrics.py @@ -1,6 +1,6 @@ import pytest from app.core.metrics import EventMetrics, ExecutionMetrics -from app.domain.enums.execution import ExecutionStatus +from app.domain.enums import ExecutionStatus from app.settings import Settings pytestmark = pytest.mark.unit diff --git a/backend/tests/unit/core/metrics/test_metrics_classes.py b/backend/tests/unit/core/metrics/test_metrics_classes.py index ce9e7f8f..12fba98f 100644 --- a/backend/tests/unit/core/metrics/test_metrics_classes.py +++ b/backend/tests/unit/core/metrics/test_metrics_classes.py @@ -13,7 +13,7 @@ ReplayMetrics, SecurityMetrics, ) -from app.domain.enums.execution import ExecutionStatus +from app.domain.enums import ExecutionStatus from app.settings import Settings pytestmark = pytest.mark.unit diff --git a/backend/tests/unit/core/test_security.py b/backend/tests/unit/core/test_security.py index bef94a0b..27bd4386 100644 --- a/backend/tests/unit/core/test_security.py +++ b/backend/tests/unit/core/test_security.py @@ -5,7 +5,7 @@ import jwt import pytest from app.core.security import SecurityService -from app.domain.enums.user import UserRole +from app.domain.enums import UserRole from app.domain.user import InvalidCredentialsError from app.settings import Settings from jwt.exceptions import InvalidTokenError diff --git a/backend/tests/unit/domain/events/test_event_schema_coverage.py b/backend/tests/unit/domain/events/test_event_schema_coverage.py index 5888ed92..9055a8a1 100644 --- a/backend/tests/unit/domain/events/test_event_schema_coverage.py +++ b/backend/tests/unit/domain/events/test_event_schema_coverage.py @@ -11,7 +11,7 @@ from typing import get_args -from app.domain.enums.events import EventType +from app.domain.enums import EventType from app.domain.events.typed import BaseEvent, DomainEvent, DomainEventAdapter diff --git a/backend/tests/unit/events/test_mappings_and_types.py b/backend/tests/unit/events/test_mappings_and_types.py index cdbefb24..bc7cfb10 100644 --- a/backend/tests/unit/events/test_mappings_and_types.py +++ b/backend/tests/unit/events/test_mappings_and_types.py @@ -1,5 +1,4 @@ -from app.domain.enums.events import EventType -from app.domain.enums.kafka import KafkaTopic +from app.domain.enums import EventType, KafkaTopic from app.infrastructure.kafka.mappings import ( get_event_class_for_type, get_event_types_for_topic, diff --git a/backend/tests/unit/schemas_pydantic/test_events_schemas.py b/backend/tests/unit/schemas_pydantic/test_events_schemas.py index 38d17179..b2fb4b9e 100644 --- a/backend/tests/unit/schemas_pydantic/test_events_schemas.py +++ b/backend/tests/unit/schemas_pydantic/test_events_schemas.py @@ -1,5 +1,5 @@ import pytest -from app.domain.enums.common import SortOrder +from app.domain.enums import SortOrder from app.schemas_pydantic.events import EventFilterRequest diff --git a/backend/tests/unit/schemas_pydantic/test_notification_schemas.py b/backend/tests/unit/schemas_pydantic/test_notification_schemas.py index b50603f1..16b58e3a 100644 --- a/backend/tests/unit/schemas_pydantic/test_notification_schemas.py +++ b/backend/tests/unit/schemas_pydantic/test_notification_schemas.py @@ -1,7 +1,7 @@ from datetime import UTC, datetime, timedelta import pytest -from app.domain.enums.notification import NotificationChannel, NotificationSeverity, NotificationStatus +from app.domain.enums import NotificationChannel, NotificationSeverity, NotificationStatus from app.schemas_pydantic.notification import Notification, NotificationBatch diff --git a/backend/tests/unit/services/coordinator/test_coordinator_queue.py b/backend/tests/unit/services/coordinator/test_coordinator_queue.py index 82e47be4..cc6467c8 100644 --- a/backend/tests/unit/services/coordinator/test_coordinator_queue.py +++ b/backend/tests/unit/services/coordinator/test_coordinator_queue.py @@ -3,7 +3,7 @@ import pytest from app.core.metrics import CoordinatorMetrics -from app.domain.enums.execution import QueuePriority +from app.domain.enums import QueuePriority from app.domain.events.typed import ExecutionRequestedEvent from app.services.coordinator.coordinator import ExecutionCoordinator, QueueRejectError diff --git a/backend/tests/unit/services/pod_monitor/test_event_mapper.py b/backend/tests/unit/services/pod_monitor/test_event_mapper.py index 13e84484..bd4b6281 100644 --- a/backend/tests/unit/services/pod_monitor/test_event_mapper.py +++ b/backend/tests/unit/services/pod_monitor/test_event_mapper.py @@ -5,8 +5,7 @@ import pytest from kubernetes_asyncio.client import V1Pod, V1PodCondition -from app.domain.enums.events import EventType -from app.domain.enums.storage import ExecutionErrorType +from app.domain.enums import EventType, ExecutionErrorType from app.domain.events.typed import ( EventMetadata, ExecutionCompletedEvent, diff --git a/backend/tests/unit/services/result_processor/test_processor.py b/backend/tests/unit/services/result_processor/test_processor.py index 7199bd30..8ec741d3 100644 --- a/backend/tests/unit/services/result_processor/test_processor.py +++ b/backend/tests/unit/services/result_processor/test_processor.py @@ -3,8 +3,7 @@ import pytest from app.core.metrics import ExecutionMetrics -from app.domain.enums.execution import ExecutionStatus -from app.domain.enums.storage import ExecutionErrorType +from app.domain.enums import ExecutionErrorType, ExecutionStatus from app.domain.events.typed import ( EventMetadata, ExecutionCompletedEvent, diff --git a/backend/tests/unit/services/saga/test_saga_comprehensive.py b/backend/tests/unit/services/saga/test_saga_comprehensive.py index d5eea475..2c88b581 100644 --- a/backend/tests/unit/services/saga/test_saga_comprehensive.py +++ b/backend/tests/unit/services/saga/test_saga_comprehensive.py @@ -6,7 +6,7 @@ """ import pytest -from app.domain.enums.saga import SagaState +from app.domain.enums import SagaState from app.domain.events.typed import DomainEvent, ExecutionRequestedEvent from app.domain.saga.models import Saga from app.services.saga.saga_step import CompensationStep, SagaContext, SagaStep diff --git a/backend/tests/unit/services/saga/test_saga_orchestrator_unit.py b/backend/tests/unit/services/saga/test_saga_orchestrator_unit.py index eb0228b4..33ce5d92 100644 --- a/backend/tests/unit/services/saga/test_saga_orchestrator_unit.py +++ b/backend/tests/unit/services/saga/test_saga_orchestrator_unit.py @@ -3,7 +3,7 @@ import pytest from app.db.repositories.resource_allocation_repository import ResourceAllocationRepository from app.db.repositories.saga_repository import SagaRepository -from app.domain.enums.saga import SagaState +from app.domain.enums import SagaState from app.domain.events.typed import DomainEvent from app.domain.saga import DomainResourceAllocation, DomainResourceAllocationCreate from app.domain.saga.models import Saga, SagaConfig diff --git a/backend/tests/unit/services/sse/test_sse_service.py b/backend/tests/unit/services/sse/test_sse_service.py index 17dd77c9..855c35ad 100644 --- a/backend/tests/unit/services/sse/test_sse_service.py +++ b/backend/tests/unit/services/sse/test_sse_service.py @@ -8,8 +8,7 @@ import pytest from app.core.metrics import ConnectionMetrics from app.db.repositories.sse_repository import SSERepository -from app.domain.enums.events import EventType -from app.domain.enums.execution import ExecutionStatus +from app.domain.enums import EventType, ExecutionStatus from app.domain.events import ResourceUsageDomain from app.domain.execution import DomainExecution from app.domain.sse import SSEExecutionStatusDomain diff --git a/backend/tests/unit/services/test_pod_builder.py b/backend/tests/unit/services/test_pod_builder.py index b7b0ee9b..e01951a5 100644 --- a/backend/tests/unit/services/test_pod_builder.py +++ b/backend/tests/unit/services/test_pod_builder.py @@ -2,7 +2,7 @@ from uuid import uuid4 import pytest -from app.domain.enums.execution import QueuePriority +from app.domain.enums import QueuePriority from app.domain.events.typed import CreatePodCommandEvent, EventMetadata from app.services.k8s_worker import PodBuilder from kubernetes_asyncio import client as k8s_client diff --git a/backend/workers/run_coordinator.py b/backend/workers/run_coordinator.py index 2528d6a3..e1dc2d1e 100644 --- a/backend/workers/run_coordinator.py +++ b/backend/workers/run_coordinator.py @@ -5,7 +5,7 @@ from app.core.logging import setup_logger from app.core.tracing import init_tracing from app.db.docs import ALL_DOCUMENTS -from app.domain.enums.kafka import GroupId +from app.domain.enums import GroupId from app.events.handlers import register_coordinator_subscriber from app.settings import Settings from beanie import init_beanie diff --git a/backend/workers/run_dlq_processor.py b/backend/workers/run_dlq_processor.py index e8dd5862..8ae604ed 100644 --- a/backend/workers/run_dlq_processor.py +++ b/backend/workers/run_dlq_processor.py @@ -6,7 +6,7 @@ from app.core.tracing import init_tracing from app.db.docs import ALL_DOCUMENTS from app.dlq.manager import DLQManager -from app.domain.enums.kafka import GroupId +from app.domain.enums import GroupId from app.events.handlers import register_dlq_subscriber from app.settings import Settings from beanie import init_beanie diff --git a/backend/workers/run_k8s_worker.py b/backend/workers/run_k8s_worker.py index 9e150a04..5ddec7f8 100644 --- a/backend/workers/run_k8s_worker.py +++ b/backend/workers/run_k8s_worker.py @@ -5,7 +5,7 @@ from app.core.logging import setup_logger from app.core.tracing import init_tracing from app.db.docs import ALL_DOCUMENTS -from app.domain.enums.kafka import GroupId +from app.domain.enums import GroupId from app.events.handlers import register_k8s_worker_subscriber from app.services.k8s_worker import KubernetesWorker from app.settings import Settings diff --git a/backend/workers/run_pod_monitor.py b/backend/workers/run_pod_monitor.py index 854d0e19..95baf0e3 100644 --- a/backend/workers/run_pod_monitor.py +++ b/backend/workers/run_pod_monitor.py @@ -5,7 +5,7 @@ from app.core.logging import setup_logger from app.core.tracing import init_tracing from app.db.docs import ALL_DOCUMENTS -from app.domain.enums.kafka import GroupId +from app.domain.enums import GroupId from app.services.pod_monitor.monitor import PodMonitor from app.settings import Settings from beanie import init_beanie diff --git a/backend/workers/run_result_processor.py b/backend/workers/run_result_processor.py index a7f789a3..e96b04e6 100644 --- a/backend/workers/run_result_processor.py +++ b/backend/workers/run_result_processor.py @@ -5,7 +5,7 @@ from app.core.logging import setup_logger from app.core.tracing import init_tracing from app.db.docs import ALL_DOCUMENTS -from app.domain.enums.kafka import GroupId +from app.domain.enums import GroupId from app.events.handlers import register_result_processor_subscriber from app.settings import Settings from beanie import init_beanie diff --git a/backend/workers/run_saga_orchestrator.py b/backend/workers/run_saga_orchestrator.py index 16ad502f..e07b0334 100644 --- a/backend/workers/run_saga_orchestrator.py +++ b/backend/workers/run_saga_orchestrator.py @@ -5,7 +5,7 @@ from app.core.logging import setup_logger from app.core.tracing import init_tracing from app.db.docs import ALL_DOCUMENTS -from app.domain.enums.kafka import GroupId +from app.domain.enums import GroupId from app.events.handlers import register_saga_subscriber from app.services.saga import SagaOrchestrator from app.settings import Settings From 583c661378851654a40a2a7b0d4e069361312bf1 Mon Sep 17 00:00:00 2001 From: HardMax71 Date: Sun, 8 Feb 2026 20:46:49 +0100 Subject: [PATCH 2/3] fixed enum like imports --- backend/app/api/routes/admin/events.py | 2 +- backend/app/api/routes/admin/users.py | 2 +- backend/app/api/routes/dlq.py | 2 +- backend/app/api/routes/events.py | 3 +- backend/app/api/routes/execution.py | 2 +- backend/app/api/routes/sse.py | 2 +- backend/app/api/routes/user_settings.py | 4 +- backend/app/core/middlewares/rate_limit.py | 2 +- backend/app/core/providers.py | 19 +++++---- backend/app/db/docs/dlq.py | 2 +- backend/app/db/docs/event.py | 2 +- backend/app/db/repositories/__init__.py | 6 +++ .../admin/admin_events_repository.py | 5 +-- .../app/db/repositories/replay_repository.py | 4 +- .../repositories/user_settings_repository.py | 2 +- backend/app/dlq/manager.py | 4 +- backend/app/dlq/models.py | 2 +- backend/app/domain/admin/__init__.py | 3 ++ backend/app/domain/admin/replay_models.py | 4 +- backend/app/domain/events/__init__.py | 12 ++++++ backend/app/domain/execution/models.py | 2 +- backend/app/events/core/producer.py | 4 +- backend/app/events/handlers.py | 4 +- backend/app/infrastructure/kafka/__init__.py | 2 +- backend/app/schemas_pydantic/admin_events.py | 3 +- .../schemas_pydantic/admin_user_overview.py | 2 +- backend/app/schemas_pydantic/dlq.py | 2 +- backend/app/schemas_pydantic/events.py | 2 +- .../services/admin/admin_events_service.py | 7 ++-- .../services/admin/admin_settings_service.py | 2 +- .../app/services/admin/admin_user_service.py | 2 +- backend/app/services/auth_service.py | 2 +- .../app/services/coordinator/coordinator.py | 4 +- .../services/event_replay/replay_service.py | 6 +-- backend/app/services/event_service.py | 2 +- backend/app/services/execution_service.py | 40 ++++++++----------- .../idempotency/idempotency_manager.py | 2 +- .../app/services/k8s_worker/pod_builder.py | 2 +- backend/app/services/k8s_worker/worker.py | 2 +- backend/app/services/kafka_event_service.py | 3 +- .../app/services/notification_scheduler.py | 2 +- backend/app/services/notification_service.py | 6 +-- .../app/services/pod_monitor/event_mapper.py | 2 +- backend/app/services/pod_monitor/monitor.py | 2 +- .../services/result_processor/processor.py | 4 +- backend/app/services/saga/__init__.py | 2 +- backend/app/services/saga/execution_saga.py | 4 +- .../app/services/saga/saga_orchestrator.py | 7 ++-- backend/app/services/saga/saga_service.py | 6 ++- backend/app/services/saga/saga_step.py | 2 +- backend/app/services/sse/__init__.py | 8 ++++ backend/app/services/sse/redis_bus.py | 2 +- backend/app/services/sse/sse_service.py | 2 +- backend/app/services/user_settings_service.py | 2 +- backend/tests/conftest.py | 2 +- backend/tests/e2e/conftest.py | 2 +- .../tests/e2e/core/test_dishka_lifespan.py | 2 +- .../test_admin_settings_repository.py | 2 +- .../db/repositories/test_dlq_repository.py | 2 +- .../repositories/test_execution_repository.py | 2 +- .../test_saved_script_repository.py | 2 +- backend/tests/e2e/dlq/test_dlq_discard.py | 2 +- backend/tests/e2e/dlq/test_dlq_manager.py | 4 +- backend/tests/e2e/dlq/test_dlq_retry.py | 2 +- .../notifications/test_notification_sse.py | 2 +- .../result_processor/test_result_processor.py | 4 +- .../services/replay/test_replay_service.py | 2 +- .../e2e/services/saga/test_saga_service.py | 3 +- .../sse/test_partitioned_event_router.py | 2 +- .../tests/e2e/services/sse/test_redis_bus.py | 4 +- .../test_user_settings_service.py | 2 +- backend/tests/e2e/test_admin_events_routes.py | 4 +- backend/tests/e2e/test_events_routes.py | 2 +- backend/tests/e2e/test_execution_routes.py | 2 +- .../tests/e2e/test_k8s_worker_create_pod.py | 2 +- backend/tests/unit/core/test_csrf.py | 2 +- .../events/test_event_schema_coverage.py | 2 +- .../tests/unit/events/test_metadata_model.py | 2 +- .../coordinator/test_coordinator_queue.py | 2 +- .../idempotency/test_idempotency_manager.py | 2 +- .../services/pod_monitor/test_event_mapper.py | 2 +- .../unit/services/pod_monitor/test_monitor.py | 2 +- .../result_processor/test_processor.py | 2 +- .../saga/test_execution_saga_steps.py | 4 +- .../services/saga/test_saga_comprehensive.py | 4 +- .../saga/test_saga_orchestrator_unit.py | 8 ++-- .../services/saga/test_saga_step_and_base.py | 2 +- .../services/sse/test_kafka_redis_bridge.py | 4 +- .../unit/services/sse/test_sse_service.py | 5 +-- .../tests/unit/services/test_pod_builder.py | 2 +- 90 files changed, 170 insertions(+), 156 deletions(-) diff --git a/backend/app/api/routes/admin/events.py b/backend/app/api/routes/admin/events.py index 88207d33..de2dd28d 100644 --- a/backend/app/api/routes/admin/events.py +++ b/backend/app/api/routes/admin/events.py @@ -9,7 +9,7 @@ from app.api.dependencies import admin_user from app.core.correlation import CorrelationContext from app.domain.enums import EventType -from app.domain.events.event_models import EventFilter +from app.domain.events import EventFilter from app.domain.replay import ReplayFilter from app.domain.user import User from app.schemas_pydantic.admin_events import ( diff --git a/backend/app/api/routes/admin/users.py b/backend/app/api/routes/admin/users.py index 4d9df577..5eb253fa 100644 --- a/backend/app/api/routes/admin/users.py +++ b/backend/app/api/routes/admin/users.py @@ -5,7 +5,7 @@ from fastapi import APIRouter, Depends, HTTPException, Query from app.api.dependencies import admin_user -from app.db.repositories.admin.admin_user_repository import AdminUserRepository +from app.db.repositories import AdminUserRepository from app.domain.enums import UserRole from app.domain.rate_limit import RateLimitRule, UserRateLimit from app.domain.user import User diff --git a/backend/app/api/routes/dlq.py b/backend/app/api/routes/dlq.py index ecd5a335..ebb07e2b 100644 --- a/backend/app/api/routes/dlq.py +++ b/backend/app/api/routes/dlq.py @@ -5,7 +5,7 @@ from fastapi import APIRouter, Depends, HTTPException, Query from app.api.dependencies import current_user -from app.db.repositories.dlq_repository import DLQRepository +from app.db.repositories import DLQRepository from app.dlq import RetryPolicy from app.dlq.manager import DLQManager from app.dlq.models import DLQMessageStatus diff --git a/backend/app/api/routes/events.py b/backend/app/api/routes/events.py index 6382f794..bd6ef80a 100644 --- a/backend/app/api/routes/events.py +++ b/backend/app/api/routes/events.py @@ -10,8 +10,7 @@ from app.core.correlation import CorrelationContext from app.core.utils import get_client_ip from app.domain.enums import EventType, SortOrder, UserRole -from app.domain.events.event_models import EventFilter -from app.domain.events.typed import BaseEvent, DomainEvent, EventMetadata +from app.domain.events import BaseEvent, DomainEvent, EventFilter, EventMetadata from app.domain.user import User from app.schemas_pydantic.common import ErrorResponse from app.schemas_pydantic.events import ( diff --git a/backend/app/api/routes/execution.py b/backend/app/api/routes/execution.py index 200a4c4e..19f99120 100644 --- a/backend/app/api/routes/execution.py +++ b/backend/app/api/routes/execution.py @@ -10,7 +10,7 @@ from app.core.tracing import EventAttributes, add_span_attributes from app.core.utils import get_client_ip from app.domain.enums import EventType, ExecutionStatus, UserRole -from app.domain.events.typed import BaseEvent, DomainEvent, EventMetadata +from app.domain.events import BaseEvent, DomainEvent, EventMetadata from app.domain.exceptions import DomainError from app.domain.idempotency import KeyStrategy from app.domain.user import User diff --git a/backend/app/api/routes/sse.py b/backend/app/api/routes/sse.py index 6d22d0ea..cc05e10c 100644 --- a/backend/app/api/routes/sse.py +++ b/backend/app/api/routes/sse.py @@ -6,7 +6,7 @@ from app.schemas_pydantic.notification import NotificationResponse from app.schemas_pydantic.sse import SSEExecutionEventData from app.services.auth_service import AuthService -from app.services.sse.sse_service import SSEService +from app.services.sse import SSEService router = APIRouter(prefix="/events", tags=["sse"], route_class=DishkaRoute) diff --git a/backend/app/api/routes/user_settings.py b/backend/app/api/routes/user_settings.py index 0677a7a3..a00d7279 100644 --- a/backend/app/api/routes/user_settings.py +++ b/backend/app/api/routes/user_settings.py @@ -5,11 +5,11 @@ from fastapi import APIRouter, Depends, Query from app.api.dependencies import current_user -from app.domain.user import User -from app.domain.user.settings_models import ( +from app.domain.user import ( DomainEditorSettings, DomainNotificationSettings, DomainUserSettingsUpdate, + User, ) from app.schemas_pydantic.user_settings import ( EditorSettings, diff --git a/backend/app/core/middlewares/rate_limit.py b/backend/app/core/middlewares/rate_limit.py index 56b2da62..0a958c3a 100644 --- a/backend/app/core/middlewares/rate_limit.py +++ b/backend/app/core/middlewares/rate_limit.py @@ -7,7 +7,7 @@ from app.core.utils import get_client_ip from app.domain.rate_limit import RateLimitStatus -from app.domain.user.user_models import User +from app.domain.user import User from app.services.rate_limit_service import RateLimitService from app.settings import Settings diff --git a/backend/app/core/providers.py b/backend/app/core/providers.py index b9cd6db3..4f303a31 100644 --- a/backend/app/core/providers.py +++ b/backend/app/core/providers.py @@ -29,26 +29,26 @@ from app.core.security import SecurityService from app.core.tracing import TracerManager from app.db.repositories import ( + AdminEventsRepository, + AdminSettingsRepository, + AdminUserRepository, + DLQRepository, EventRepository, ExecutionRepository, NotificationRepository, + ReplayRepository, + ResourceAllocationRepository, SagaRepository, SavedScriptRepository, SSERepository, UserRepository, + UserSettingsRepository, ) -from app.db.repositories.admin.admin_events_repository import AdminEventsRepository -from app.db.repositories.admin.admin_settings_repository import AdminSettingsRepository -from app.db.repositories.admin.admin_user_repository import AdminUserRepository -from app.db.repositories.dlq_repository import DLQRepository -from app.db.repositories.replay_repository import ReplayRepository -from app.db.repositories.resource_allocation_repository import ResourceAllocationRepository -from app.db.repositories.user_settings_repository import UserSettingsRepository from app.dlq.manager import DLQManager from app.dlq.models import RetryPolicy, RetryStrategy from app.domain.enums import KafkaTopic from app.domain.rate_limit import RateLimitConfig -from app.domain.saga.models import SagaConfig +from app.domain.saga import SagaConfig from app.events.core import UnifiedProducer from app.services.admin import AdminEventsService, AdminSettingsService, AdminUserService from app.services.auth_service import AuthService @@ -72,8 +72,7 @@ from app.services.saga import SagaOrchestrator from app.services.saga.saga_service import SagaService from app.services.saved_script_service import SavedScriptService -from app.services.sse.redis_bus import SSERedisBus -from app.services.sse.sse_service import SSEService +from app.services.sse import SSERedisBus, SSEService from app.services.user_settings_service import UserSettingsService from app.settings import Settings diff --git a/backend/app/db/docs/dlq.py b/backend/app/db/docs/dlq.py index 71e2f7a3..0b63235f 100644 --- a/backend/app/db/docs/dlq.py +++ b/backend/app/db/docs/dlq.py @@ -5,7 +5,7 @@ from pymongo import ASCENDING, DESCENDING, IndexModel from app.dlq.models import DLQMessageStatus -from app.domain.events.typed import DomainEvent +from app.domain.events import DomainEvent class DLQMessageDocument(Document): diff --git a/backend/app/db/docs/event.py b/backend/app/db/docs/event.py index d8c34b02..2ce1a126 100644 --- a/backend/app/db/docs/event.py +++ b/backend/app/db/docs/event.py @@ -7,7 +7,7 @@ from pymongo import ASCENDING, DESCENDING, IndexModel from app.domain.enums import EventType -from app.domain.events.typed import EventMetadata +from app.domain.events import EventMetadata class EventDocument(Document): diff --git a/backend/app/db/repositories/__init__.py b/backend/app/db/repositories/__init__.py index 1e985797..d5406541 100644 --- a/backend/app/db/repositories/__init__.py +++ b/backend/app/db/repositories/__init__.py @@ -1,9 +1,12 @@ +from app.db.repositories.admin.admin_events_repository import AdminEventsRepository from app.db.repositories.admin.admin_settings_repository import AdminSettingsRepository from app.db.repositories.admin.admin_user_repository import AdminUserRepository +from app.db.repositories.dlq_repository import DLQRepository from app.db.repositories.event_repository import EventRepository from app.db.repositories.execution_repository import ExecutionRepository from app.db.repositories.notification_repository import NotificationRepository from app.db.repositories.replay_repository import ReplayRepository +from app.db.repositories.resource_allocation_repository import ResourceAllocationRepository from app.db.repositories.saga_repository import SagaRepository from app.db.repositories.saved_script_repository import SavedScriptRepository from app.db.repositories.sse_repository import SSERepository @@ -11,12 +14,15 @@ from app.db.repositories.user_settings_repository import UserSettingsRepository __all__ = [ + "AdminEventsRepository", "AdminSettingsRepository", "AdminUserRepository", + "DLQRepository", "EventRepository", "ExecutionRepository", "NotificationRepository", "ReplayRepository", + "ResourceAllocationRepository", "SagaRepository", "SavedScriptRepository", "SSERepository", diff --git a/backend/app/db/repositories/admin/admin_events_repository.py b/backend/app/db/repositories/admin/admin_events_repository.py index c2bdfa0a..b5d43a98 100644 --- a/backend/app/db/repositories/admin/admin_events_repository.py +++ b/backend/app/db/repositories/admin/admin_events_repository.py @@ -11,8 +11,7 @@ ExecutionDocument, ReplaySessionDocument, ) -from app.domain.admin import ExecutionResultSummary, ReplaySessionData, ReplaySessionStatusDetail -from app.domain.admin.replay_updates import ReplaySessionUpdate +from app.domain.admin import ExecutionResultSummary, ReplaySessionData, ReplaySessionStatusDetail, ReplaySessionUpdate from app.domain.enums import EventType, ReplayStatus from app.domain.events import ( DomainEvent, @@ -27,7 +26,7 @@ HourlyEventCount, UserEventCount, ) -from app.domain.replay.models import ReplayFilter, ReplaySessionState +from app.domain.replay import ReplayFilter, ReplaySessionState class AdminEventsRepository: diff --git a/backend/app/db/repositories/replay_repository.py b/backend/app/db/repositories/replay_repository.py index 998c4cb0..ee449bb3 100644 --- a/backend/app/db/repositories/replay_repository.py +++ b/backend/app/db/repositories/replay_repository.py @@ -6,9 +6,9 @@ from beanie.operators import LT, In from app.db.docs import EventDocument, ReplaySessionDocument -from app.domain.admin.replay_updates import ReplaySessionUpdate +from app.domain.admin import ReplaySessionUpdate from app.domain.enums import ReplayStatus -from app.domain.replay.models import ReplayFilter, ReplaySessionState +from app.domain.replay import ReplayFilter, ReplaySessionState class ReplayRepository: diff --git a/backend/app/db/repositories/user_settings_repository.py b/backend/app/db/repositories/user_settings_repository.py index 019ee9f4..8cb5e1fd 100644 --- a/backend/app/db/repositories/user_settings_repository.py +++ b/backend/app/db/repositories/user_settings_repository.py @@ -7,7 +7,7 @@ from app.db.docs import EventDocument, UserSettingsDocument, UserSettingsSnapshotDocument from app.domain.enums import EventType -from app.domain.user.settings_models import DomainUserSettings, DomainUserSettingsChangedEvent +from app.domain.user import DomainUserSettings, DomainUserSettingsChangedEvent class UserSettingsRepository: diff --git a/backend/app/dlq/manager.py b/backend/app/dlq/manager.py index a3724b88..cead164b 100644 --- a/backend/app/dlq/manager.py +++ b/backend/app/dlq/manager.py @@ -6,7 +6,7 @@ from app.core.metrics import DLQMetrics from app.core.tracing.utils import inject_trace_context -from app.db.repositories.dlq_repository import DLQRepository +from app.db.repositories import DLQRepository from app.dlq.models import ( DLQBatchRetryResult, DLQMessage, @@ -17,7 +17,7 @@ RetryStrategy, ) from app.domain.enums import KafkaTopic -from app.domain.events.typed import ( +from app.domain.events import ( DLQMessageDiscardedEvent, DLQMessageReceivedEvent, DLQMessageRetriedEvent, diff --git a/backend/app/dlq/models.py b/backend/app/dlq/models.py index 7df828b4..85e91a09 100644 --- a/backend/app/dlq/models.py +++ b/backend/app/dlq/models.py @@ -6,7 +6,7 @@ from app.core.utils import StringEnum from app.domain.enums import EventType -from app.domain.events.typed import DomainEvent +from app.domain.events import DomainEvent class DLQMessageStatus(StringEnum): diff --git a/backend/app/domain/admin/__init__.py b/backend/app/domain/admin/__init__.py index 8733ce14..43004a3b 100644 --- a/backend/app/domain/admin/__init__.py +++ b/backend/app/domain/admin/__init__.py @@ -9,6 +9,7 @@ ReplaySessionStatusDetail, ReplaySessionStatusInfo, ) +from .replay_updates import ReplaySessionUpdate from .settings_models import ( AuditAction, AuditLogEntry, @@ -37,4 +38,6 @@ "ReplaySessionData", "ReplaySessionStatusDetail", "ReplaySessionStatusInfo", + # Replay Updates + "ReplaySessionUpdate", ] diff --git a/backend/app/domain/admin/replay_models.py b/backend/app/domain/admin/replay_models.py index 9b5ef7e5..68a80019 100644 --- a/backend/app/domain/admin/replay_models.py +++ b/backend/app/domain/admin/replay_models.py @@ -5,8 +5,8 @@ from pydantic.dataclasses import dataclass from app.domain.enums import ExecutionStatus, ReplayStatus -from app.domain.events.event_models import EventSummary -from app.domain.replay.models import ReplayFilter, ReplaySessionState +from app.domain.events import EventSummary +from app.domain.replay import ReplayFilter, ReplaySessionState @dataclass(config=ConfigDict(from_attributes=True)) diff --git a/backend/app/domain/events/__init__.py b/backend/app/domain/events/__init__.py index 2a9bc41c..b6dc3b56 100644 --- a/backend/app/domain/events/__init__.py +++ b/backend/app/domain/events/__init__.py @@ -28,6 +28,10 @@ ContainerStatusInfo, CreatePodCommandEvent, DeletePodCommandEvent, + # DLQ Events + DLQMessageDiscardedEvent, + DLQMessageReceivedEvent, + DLQMessageRetriedEvent, DomainEvent, DomainEventAdapter, EventMetadata, @@ -35,6 +39,8 @@ ExecutionAcceptedEvent, ExecutionCancelledEvent, ExecutionCompletedEvent, + # Type aliases + ExecutionDomainEvent, ExecutionFailedEvent, ExecutionQueuedEvent, ExecutionRequestedEvent, @@ -181,8 +187,14 @@ # Resource Events "ResourceLimitExceededEvent", "QuotaExceededEvent", + # DLQ Events + "DLQMessageReceivedEvent", + "DLQMessageRetriedEvent", + "DLQMessageDiscardedEvent", # System Events "SystemErrorEvent", "ServiceUnhealthyEvent", "ServiceRecoveredEvent", + # Type aliases + "ExecutionDomainEvent", ] diff --git a/backend/app/domain/execution/models.py b/backend/app/domain/execution/models.py index 45f04943..335c0ec0 100644 --- a/backend/app/domain/execution/models.py +++ b/backend/app/domain/execution/models.py @@ -6,7 +6,7 @@ from pydantic import BaseModel, ConfigDict, Field from app.domain.enums import ExecutionErrorType, ExecutionStatus -from app.domain.events.typed import EventMetadata, ResourceUsageDomain +from app.domain.events import EventMetadata, ResourceUsageDomain class DomainExecution(BaseModel): diff --git a/backend/app/events/core/producer.py b/backend/app/events/core/producer.py index 17940d33..771bea94 100644 --- a/backend/app/events/core/producer.py +++ b/backend/app/events/core/producer.py @@ -7,10 +7,10 @@ from app.core.metrics import EventMetrics from app.core.tracing.utils import inject_trace_context -from app.db.repositories.event_repository import EventRepository +from app.db.repositories import EventRepository from app.dlq.models import DLQMessageStatus from app.domain.enums import KafkaTopic -from app.domain.events.typed import DomainEvent +from app.domain.events import DomainEvent from app.infrastructure.kafka.mappings import EVENT_TYPE_TO_TOPIC from app.settings import Settings diff --git a/backend/app/events/handlers.py b/backend/app/events/handlers.py index 684a71f8..e22aee8d 100644 --- a/backend/app/events/handlers.py +++ b/backend/app/events/handlers.py @@ -14,7 +14,7 @@ from app.dlq.manager import DLQManager from app.dlq.models import DLQMessage, DLQMessageStatus from app.domain.enums import EventType, GroupId, KafkaTopic -from app.domain.events.typed import ( +from app.domain.events import ( CreatePodCommandEvent, DeletePodCommandEvent, DomainEvent, @@ -32,7 +32,7 @@ from app.services.notification_service import NotificationService from app.services.result_processor.processor import ResultProcessor from app.services.saga import SagaOrchestrator -from app.services.sse.redis_bus import SSERedisBus +from app.services.sse import SSERedisBus from app.settings import Settings diff --git a/backend/app/infrastructure/kafka/__init__.py b/backend/app/infrastructure/kafka/__init__.py index 97295a56..fae49311 100644 --- a/backend/app/infrastructure/kafka/__init__.py +++ b/backend/app/infrastructure/kafka/__init__.py @@ -1,4 +1,4 @@ -from app.domain.events.typed import DomainEvent, EventMetadata +from app.domain.events import DomainEvent, EventMetadata from app.infrastructure.kafka.mappings import get_event_class_for_type, get_topic_for_event from app.infrastructure.kafka.topics import get_all_topics, get_topic_configs diff --git a/backend/app/schemas_pydantic/admin_events.py b/backend/app/schemas_pydantic/admin_events.py index 9cceb9b0..63fc3d3f 100644 --- a/backend/app/schemas_pydantic/admin_events.py +++ b/backend/app/schemas_pydantic/admin_events.py @@ -3,8 +3,7 @@ from pydantic import BaseModel, ConfigDict, Field, computed_field from app.domain.enums import EventType -from app.domain.events.event_models import EventSummary -from app.domain.events.typed import DomainEvent +from app.domain.events import DomainEvent, EventSummary from app.domain.replay import ReplayError from app.schemas_pydantic.events import EventTypeCountSchema, HourlyEventCountSchema from app.schemas_pydantic.execution import ExecutionResult diff --git a/backend/app/schemas_pydantic/admin_user_overview.py b/backend/app/schemas_pydantic/admin_user_overview.py index 52301677..08c34921 100644 --- a/backend/app/schemas_pydantic/admin_user_overview.py +++ b/backend/app/schemas_pydantic/admin_user_overview.py @@ -2,7 +2,7 @@ from pydantic import BaseModel, ConfigDict -from app.domain.events.typed import DomainEvent +from app.domain.events import DomainEvent from app.schemas_pydantic.events import EventStatistics from app.schemas_pydantic.user import UserResponse diff --git a/backend/app/schemas_pydantic/dlq.py b/backend/app/schemas_pydantic/dlq.py index 4093d03f..d1a8137a 100644 --- a/backend/app/schemas_pydantic/dlq.py +++ b/backend/app/schemas_pydantic/dlq.py @@ -10,7 +10,7 @@ RetryStrategy, TopicStatistic, ) -from app.domain.events.typed import DomainEvent +from app.domain.events import DomainEvent class DLQStats(BaseModel): diff --git a/backend/app/schemas_pydantic/events.py b/backend/app/schemas_pydantic/events.py index 2d040d25..b02533ed 100644 --- a/backend/app/schemas_pydantic/events.py +++ b/backend/app/schemas_pydantic/events.py @@ -5,7 +5,7 @@ from pydantic import BaseModel, ConfigDict, Field, field_validator from app.domain.enums import Environment, EventType, SortOrder -from app.domain.events.typed import ContainerStatusInfo, DomainEvent +from app.domain.events import ContainerStatusInfo, DomainEvent class EventTypeCountSchema(BaseModel): diff --git a/backend/app/services/admin/admin_events_service.py b/backend/app/services/admin/admin_events_service.py index c0de6533..2e9cee40 100644 --- a/backend/app/services/admin/admin_events_service.py +++ b/backend/app/services/admin/admin_events_service.py @@ -8,11 +8,10 @@ from beanie.odm.enums import SortDirection -from app.db.repositories.admin import AdminEventsRepository -from app.domain.admin import ReplaySessionStatusDetail -from app.domain.admin.replay_updates import ReplaySessionUpdate +from app.db.repositories import AdminEventsRepository +from app.domain.admin import ReplaySessionStatusDetail, ReplaySessionUpdate from app.domain.enums import ReplayStatus, ReplayTarget, ReplayType -from app.domain.events.event_models import ( +from app.domain.events import ( EventBrowseResult, EventDetail, EventExportRow, diff --git a/backend/app/services/admin/admin_settings_service.py b/backend/app/services/admin/admin_settings_service.py index 674ad8de..f16dedfc 100644 --- a/backend/app/services/admin/admin_settings_service.py +++ b/backend/app/services/admin/admin_settings_service.py @@ -1,6 +1,6 @@ import logging -from app.db.repositories.admin.admin_settings_repository import AdminSettingsRepository +from app.db.repositories import AdminSettingsRepository from app.domain.admin import SystemSettings diff --git a/backend/app/services/admin/admin_user_service.py b/backend/app/services/admin/admin_user_service.py index 95aa232e..16abd974 100644 --- a/backend/app/services/admin/admin_user_service.py +++ b/backend/app/services/admin/admin_user_service.py @@ -3,7 +3,7 @@ from datetime import datetime, timedelta, timezone from app.core.security import SecurityService -from app.db.repositories.admin.admin_user_repository import AdminUserRepository +from app.db.repositories import AdminUserRepository from app.domain.admin import AdminUserOverviewDomain, DerivedCountsDomain, RateLimitSummaryDomain from app.domain.enums import EventType, ExecutionStatus, UserRole from app.domain.rate_limit import RateLimitUpdateResult, UserRateLimit, UserRateLimitsResult diff --git a/backend/app/services/auth_service.py b/backend/app/services/auth_service.py index c4c88c6f..b1cdf066 100644 --- a/backend/app/services/auth_service.py +++ b/backend/app/services/auth_service.py @@ -3,7 +3,7 @@ from fastapi import Request from app.core.security import SecurityService -from app.db.repositories.user_repository import UserRepository +from app.db.repositories import UserRepository from app.domain.enums import UserRole from app.domain.user import AdminAccessRequiredError, AuthenticationRequiredError, InvalidCredentialsError, User diff --git a/backend/app/services/coordinator/coordinator.py b/backend/app/services/coordinator/coordinator.py index a2c6bd19..7e0b4de5 100644 --- a/backend/app/services/coordinator/coordinator.py +++ b/backend/app/services/coordinator/coordinator.py @@ -6,9 +6,9 @@ from uuid import uuid4 from app.core.metrics import CoordinatorMetrics -from app.db.repositories.execution_repository import ExecutionRepository +from app.db.repositories import ExecutionRepository from app.domain.enums import ExecutionErrorType, QueuePriority -from app.domain.events.typed import ( +from app.domain.events import ( CreatePodCommandEvent, EventMetadata, ExecutionAcceptedEvent, diff --git a/backend/app/services/event_replay/replay_service.py b/backend/app/services/event_replay/replay_service.py index a42e0f71..6c17b455 100644 --- a/backend/app/services/event_replay/replay_service.py +++ b/backend/app/services/event_replay/replay_service.py @@ -11,10 +11,10 @@ from pydantic import ValidationError from app.core.metrics import ReplayMetrics -from app.db.repositories.replay_repository import ReplayRepository -from app.domain.admin.replay_updates import ReplaySessionUpdate +from app.db.repositories import ReplayRepository +from app.domain.admin import ReplaySessionUpdate from app.domain.enums import ReplayStatus, ReplayTarget -from app.domain.events.typed import DomainEvent, DomainEventAdapter +from app.domain.events import DomainEvent, DomainEventAdapter from app.domain.replay import ( CleanupResult, ReplayConfig, diff --git a/backend/app/services/event_service.py b/backend/app/services/event_service.py index c274de9e..62d60f0c 100644 --- a/backend/app/services/event_service.py +++ b/backend/app/services/event_service.py @@ -1,7 +1,7 @@ from datetime import datetime from typing import Any -from app.db.repositories.event_repository import EventRepository +from app.db.repositories import EventRepository from app.domain.enums import EventType, UserRole from app.domain.events import ( ArchivedEvent, diff --git a/backend/app/services/execution_service.py b/backend/app/services/execution_service.py index 96734880..c3bf44ce 100644 --- a/backend/app/services/execution_service.py +++ b/backend/app/services/execution_service.py @@ -2,14 +2,13 @@ from contextlib import contextmanager from datetime import datetime from time import time -from typing import Any, Generator, TypeAlias +from typing import Any, Generator from app.core.correlation import CorrelationContext from app.core.metrics import ExecutionMetrics -from app.db.repositories.event_repository import EventRepository -from app.db.repositories.execution_repository import ExecutionRepository +from app.db.repositories import EventRepository, ExecutionRepository from app.domain.enums import EventType, ExecutionStatus, QueuePriority -from app.domain.events.typed import ( +from app.domain.events import ( DomainEvent, EventMetadata, ExecutionCancelledEvent, @@ -27,13 +26,6 @@ from app.runtime_registry import RUNTIME_REGISTRY from app.settings import Settings -# Type aliases for better readability -UserId: TypeAlias = str -EventFilter: TypeAlias = list[EventType] | None -TimeRange: TypeAlias = tuple[datetime | None, datetime | None] -ExecutionQuery: TypeAlias = dict[str, Any] -ExecutionStats: TypeAlias = dict[str, Any] - class ExecutionService: """ @@ -288,7 +280,7 @@ async def get_execution_result(self, execution_id: str) -> DomainExecution: async def get_execution_events( self, execution_id: str, - event_types: EventFilter = None, + event_types: list[EventType] | None = None, limit: int = 100, ) -> list[DomainEvent]: """ @@ -320,7 +312,7 @@ async def get_execution_events( async def get_user_executions( self, - user_id: UserId, + user_id: str, status: ExecutionStatus | None = None, lang: str | None = None, start_time: datetime | None = None, @@ -363,7 +355,7 @@ async def get_user_executions( async def count_user_executions( self, - user_id: UserId, + user_id: str, status: ExecutionStatus | None = None, lang: str | None = None, start_time: datetime | None = None, @@ -387,12 +379,12 @@ async def count_user_executions( def _build_user_query( self, - user_id: UserId, + user_id: str, status: ExecutionStatus | None = None, lang: str | None = None, start_time: datetime | None = None, end_time: datetime | None = None, - ) -> ExecutionQuery: + ) -> dict[str, Any]: """ Build MongoDB query for user executions. @@ -406,7 +398,7 @@ def _build_user_query( Returns: MongoDB query dictionary. """ - query: ExecutionQuery = {"user_id": str(user_id)} + query: dict[str, Any] = {"user_id": str(user_id)} if status: query["status"] = status @@ -471,8 +463,8 @@ async def _publish_deletion_event(self, execution_id: str) -> None: ) async def get_execution_stats( - self, user_id: UserId | None = None, time_range: TimeRange = (None, None) - ) -> ExecutionStats: + self, user_id: str | None = None, time_range: tuple[datetime | None, datetime | None] = (None, None) + ) -> dict[str, Any]: """ Get execution statistics. @@ -493,7 +485,9 @@ async def get_execution_stats( return self._calculate_stats(executions) - def _build_stats_query(self, user_id: UserId | None, time_range: TimeRange) -> ExecutionQuery: + def _build_stats_query( + self, user_id: str | None, time_range: tuple[datetime | None, datetime | None] + ) -> dict[str, Any]: """ Build query for statistics. @@ -504,7 +498,7 @@ def _build_stats_query(self, user_id: UserId | None, time_range: TimeRange) -> E Returns: MongoDB query dictionary. """ - query: ExecutionQuery = {} + query: dict[str, Any] = {} if user_id: query["user_id"] = str(user_id) @@ -520,7 +514,7 @@ def _build_stats_query(self, user_id: UserId | None, time_range: TimeRange) -> E return query - def _calculate_stats(self, executions: list[DomainExecution]) -> ExecutionStats: + def _calculate_stats(self, executions: list[DomainExecution]) -> dict[str, Any]: """ Calculate statistics from executions. @@ -530,7 +524,7 @@ def _calculate_stats(self, executions: list[DomainExecution]) -> ExecutionStats: Returns: Statistics dictionary. """ - stats: ExecutionStats = { + stats: dict[str, Any] = { "total": len(executions), "by_status": {}, "by_language": {}, diff --git a/backend/app/services/idempotency/idempotency_manager.py b/backend/app/services/idempotency/idempotency_manager.py index 41dc64ac..aa56f220 100644 --- a/backend/app/services/idempotency/idempotency_manager.py +++ b/backend/app/services/idempotency/idempotency_manager.py @@ -7,7 +7,7 @@ from pymongo.errors import DuplicateKeyError from app.core.metrics import DatabaseMetrics -from app.domain.events.typed import BaseEvent +from app.domain.events import BaseEvent from app.domain.idempotency import IdempotencyRecord, IdempotencyStatus, KeyStrategy from app.services.idempotency.redis_repository import RedisIdempotencyRepository diff --git a/backend/app/services/k8s_worker/pod_builder.py b/backend/app/services/k8s_worker/pod_builder.py index 4aaa0cf0..8b9dcb3e 100644 --- a/backend/app/services/k8s_worker/pod_builder.py +++ b/backend/app/services/k8s_worker/pod_builder.py @@ -1,6 +1,6 @@ from kubernetes_asyncio import client as k8s_client -from app.domain.events.typed import CreatePodCommandEvent +from app.domain.events import CreatePodCommandEvent from app.settings import Settings diff --git a/backend/app/services/k8s_worker/worker.py b/backend/app/services/k8s_worker/worker.py index 3fdf71fc..a2ff22d3 100644 --- a/backend/app/services/k8s_worker/worker.py +++ b/backend/app/services/k8s_worker/worker.py @@ -9,7 +9,7 @@ from app.core.metrics import EventMetrics, ExecutionMetrics, KubernetesMetrics from app.domain.enums import ExecutionErrorType -from app.domain.events.typed import ( +from app.domain.events import ( CreatePodCommandEvent, DeletePodCommandEvent, ExecutionFailedEvent, diff --git a/backend/app/services/kafka_event_service.py b/backend/app/services/kafka_event_service.py index d344f7cd..018227f8 100644 --- a/backend/app/services/kafka_event_service.py +++ b/backend/app/services/kafka_event_service.py @@ -9,8 +9,7 @@ from app.core.correlation import CorrelationContext from app.core.metrics import EventMetrics from app.domain.enums import EventType -from app.domain.events import DomainEventAdapter -from app.domain.events.typed import DomainEvent, EventMetadata +from app.domain.events import DomainEvent, DomainEventAdapter, EventMetadata from app.events.core import UnifiedProducer from app.settings import Settings diff --git a/backend/app/services/notification_scheduler.py b/backend/app/services/notification_scheduler.py index 4e16e9ce..91bbcd81 100644 --- a/backend/app/services/notification_scheduler.py +++ b/backend/app/services/notification_scheduler.py @@ -1,6 +1,6 @@ import logging -from app.db.repositories.notification_repository import NotificationRepository +from app.db.repositories import NotificationRepository from app.services.notification_service import NotificationService diff --git a/backend/app/services/notification_service.py b/backend/app/services/notification_service.py index 2185fe54..a39f0c36 100644 --- a/backend/app/services/notification_service.py +++ b/backend/app/services/notification_service.py @@ -9,9 +9,9 @@ from app.core.metrics import NotificationMetrics from app.core.tracing.utils import add_span_attributes -from app.db.repositories.notification_repository import NotificationRepository +from app.db.repositories import NotificationRepository from app.domain.enums import NotificationChannel, NotificationSeverity, NotificationStatus, UserRole -from app.domain.events.typed import ( +from app.domain.events import ( EventMetadata, ExecutionCompletedEvent, ExecutionFailedEvent, @@ -31,7 +31,7 @@ ) from app.schemas_pydantic.sse import RedisNotificationMessage from app.services.kafka_event_service import KafkaEventService -from app.services.sse.redis_bus import SSERedisBus +from app.services.sse import SSERedisBus from app.settings import Settings # Constants diff --git a/backend/app/services/pod_monitor/event_mapper.py b/backend/app/services/pod_monitor/event_mapper.py index 8b5b37a8..2abd5f80 100644 --- a/backend/app/services/pod_monitor/event_mapper.py +++ b/backend/app/services/pod_monitor/event_mapper.py @@ -6,7 +6,7 @@ from kubernetes_asyncio import client as k8s_client from app.domain.enums import ExecutionErrorType, GroupId -from app.domain.events.typed import ( +from app.domain.events import ( ContainerStatusInfo, DomainEvent, EventMetadata, diff --git a/backend/app/services/pod_monitor/monitor.py b/backend/app/services/pod_monitor/monitor.py index b639ead0..88a56814 100644 --- a/backend/app/services/pod_monitor/monitor.py +++ b/backend/app/services/pod_monitor/monitor.py @@ -9,7 +9,7 @@ from app.core.metrics import KubernetesMetrics from app.core.utils import StringEnum -from app.domain.events.typed import DomainEvent +from app.domain.events import DomainEvent from app.services.kafka_event_service import KafkaEventService from app.services.pod_monitor.config import PodMonitorConfig from app.services.pod_monitor.event_mapper import PodEventMapper diff --git a/backend/app/services/result_processor/processor.py b/backend/app/services/result_processor/processor.py index a9f4a1bf..2abea59d 100644 --- a/backend/app/services/result_processor/processor.py +++ b/backend/app/services/result_processor/processor.py @@ -1,9 +1,9 @@ import logging from app.core.metrics import ExecutionMetrics -from app.db.repositories.execution_repository import ExecutionRepository +from app.db.repositories import ExecutionRepository from app.domain.enums import ExecutionErrorType, ExecutionStatus, GroupId, StorageType -from app.domain.events.typed import ( +from app.domain.events import ( DomainEvent, EventMetadata, ExecutionCompletedEvent, diff --git a/backend/app/services/saga/__init__.py b/backend/app/services/saga/__init__.py index 3a5c87a0..d5b27436 100644 --- a/backend/app/services/saga/__init__.py +++ b/backend/app/services/saga/__init__.py @@ -1,5 +1,5 @@ from app.domain.enums import SagaState -from app.domain.saga.models import SagaConfig, SagaInstance +from app.domain.saga import SagaConfig, SagaInstance from app.services.saga.execution_saga import ( AllocateResourcesStep, CreatePodStep, diff --git a/backend/app/services/saga/execution_saga.py b/backend/app/services/saga/execution_saga.py index 063a0ac1..57e76fd1 100644 --- a/backend/app/services/saga/execution_saga.py +++ b/backend/app/services/saga/execution_saga.py @@ -1,8 +1,8 @@ import logging from typing import Any -from app.db.repositories.resource_allocation_repository import ResourceAllocationRepository -from app.domain.events.typed import CreatePodCommandEvent, DeletePodCommandEvent, EventMetadata, ExecutionRequestedEvent +from app.db.repositories import ResourceAllocationRepository +from app.domain.events import CreatePodCommandEvent, DeletePodCommandEvent, EventMetadata, ExecutionRequestedEvent from app.domain.saga import DomainResourceAllocationCreate from app.events.core import UnifiedProducer diff --git a/backend/app/services/saga/saga_orchestrator.py b/backend/app/services/saga/saga_orchestrator.py index c3f1180d..65afadfa 100644 --- a/backend/app/services/saga/saga_orchestrator.py +++ b/backend/app/services/saga/saga_orchestrator.py @@ -6,10 +6,9 @@ from app.core.tracing import EventAttributes from app.core.tracing.utils import get_tracer -from app.db.repositories.resource_allocation_repository import ResourceAllocationRepository -from app.db.repositories.saga_repository import SagaRepository +from app.db.repositories import ResourceAllocationRepository, SagaRepository from app.domain.enums import SagaState -from app.domain.events.typed import ( +from app.domain.events import ( DomainEvent, EventMetadata, ExecutionCompletedEvent, @@ -19,7 +18,7 @@ SagaCancelledEvent, SagaStartedEvent, ) -from app.domain.saga.models import Saga, SagaConfig +from app.domain.saga import Saga, SagaConfig from app.events.core import UnifiedProducer from .execution_saga import ExecutionSaga diff --git a/backend/app/services/saga/saga_service.py b/backend/app/services/saga/saga_service.py index 12cba068..dc45f381 100644 --- a/backend/app/services/saga/saga_service.py +++ b/backend/app/services/saga/saga_service.py @@ -2,12 +2,14 @@ from app.db.repositories import ExecutionRepository, SagaRepository from app.domain.enums import SagaState, UserRole -from app.domain.saga.exceptions import ( +from app.domain.saga import ( + Saga, SagaAccessDeniedError, + SagaFilter, SagaInvalidStateError, + SagaListResult, SagaNotFoundError, ) -from app.domain.saga.models import Saga, SagaFilter, SagaListResult from app.schemas_pydantic.user import User from app.services.saga import SagaOrchestrator diff --git a/backend/app/services/saga/saga_step.py b/backend/app/services/saga/saga_step.py index 81d65e1e..fe3ab191 100644 --- a/backend/app/services/saga/saga_step.py +++ b/backend/app/services/saga/saga_step.py @@ -3,7 +3,7 @@ from fastapi.encoders import jsonable_encoder -from app.domain.events.typed import DomainEvent +from app.domain.events import DomainEvent T = TypeVar("T", bound=DomainEvent) diff --git a/backend/app/services/sse/__init__.py b/backend/app/services/sse/__init__.py index e69de29b..a8f9db84 100644 --- a/backend/app/services/sse/__init__.py +++ b/backend/app/services/sse/__init__.py @@ -0,0 +1,8 @@ +from app.services.sse.redis_bus import SSERedisBus, SSERedisSubscription +from app.services.sse.sse_service import SSEService + +__all__ = [ + "SSERedisBus", + "SSERedisSubscription", + "SSEService", +] diff --git a/backend/app/services/sse/redis_bus.py b/backend/app/services/sse/redis_bus.py index d7375a85..23408fdb 100644 --- a/backend/app/services/sse/redis_bus.py +++ b/backend/app/services/sse/redis_bus.py @@ -7,7 +7,7 @@ from pydantic import BaseModel from app.domain.enums import EventType -from app.domain.events.typed import DomainEvent +from app.domain.events import DomainEvent from app.schemas_pydantic.sse import RedisNotificationMessage, RedisSSEMessage T = TypeVar("T", bound=BaseModel) diff --git a/backend/app/services/sse/sse_service.py b/backend/app/services/sse/sse_service.py index 62071b79..a8f7c622 100644 --- a/backend/app/services/sse/sse_service.py +++ b/backend/app/services/sse/sse_service.py @@ -5,7 +5,7 @@ from typing import Any from app.core.metrics import ConnectionMetrics -from app.db.repositories.sse_repository import SSERepository +from app.db.repositories import SSERepository from app.domain.enums import EventType, NotificationChannel, SSEControlEvent from app.schemas_pydantic.execution import ExecutionResult from app.schemas_pydantic.notification import NotificationResponse diff --git a/backend/app/services/user_settings_service.py b/backend/app/services/user_settings_service.py index 49dce6da..28ede701 100644 --- a/backend/app/services/user_settings_service.py +++ b/backend/app/services/user_settings_service.py @@ -4,7 +4,7 @@ from cachetools import TTLCache -from app.db.repositories.user_settings_repository import UserSettingsRepository +from app.db.repositories import UserSettingsRepository from app.domain.enums import EventType, Theme from app.domain.user import ( DomainEditorSettings, diff --git a/backend/tests/conftest.py b/backend/tests/conftest.py index 37252e3d..6245af02 100644 --- a/backend/tests/conftest.py +++ b/backend/tests/conftest.py @@ -8,7 +8,7 @@ import pytest_asyncio import redis.asyncio as redis from app.domain.enums import QueuePriority -from app.domain.events.typed import EventMetadata, ExecutionRequestedEvent +from app.domain.events import EventMetadata, ExecutionRequestedEvent from app.main import create_app from app.settings import Settings from dishka import AsyncContainer diff --git a/backend/tests/e2e/conftest.py b/backend/tests/e2e/conftest.py index deaf14bc..57fbd95e 100644 --- a/backend/tests/e2e/conftest.py +++ b/backend/tests/e2e/conftest.py @@ -10,7 +10,7 @@ from aiokafka import AIOKafkaConsumer from app.db.docs.saga import SagaDocument from app.domain.enums import EventType, KafkaTopic, UserRole -from app.domain.events.typed import DomainEvent, DomainEventAdapter +from app.domain.events import DomainEvent, DomainEventAdapter from app.schemas_pydantic.execution import ExecutionRequest, ExecutionResponse from app.schemas_pydantic.notification import NotificationListResponse, NotificationResponse from app.schemas_pydantic.saga import SagaStatusResponse diff --git a/backend/tests/e2e/core/test_dishka_lifespan.py b/backend/tests/e2e/core/test_dishka_lifespan.py index 1781a099..6059017c 100644 --- a/backend/tests/e2e/core/test_dishka_lifespan.py +++ b/backend/tests/e2e/core/test_dishka_lifespan.py @@ -3,7 +3,7 @@ import pytest import redis.asyncio as aioredis from app.db.docs import UserDocument -from app.services.sse.redis_bus import SSERedisBus +from app.services.sse import SSERedisBus from app.settings import Settings from dishka import AsyncContainer from fastapi import FastAPI diff --git a/backend/tests/e2e/db/repositories/test_admin_settings_repository.py b/backend/tests/e2e/db/repositories/test_admin_settings_repository.py index cce154c9..49a9e4c8 100644 --- a/backend/tests/e2e/db/repositories/test_admin_settings_repository.py +++ b/backend/tests/e2e/db/repositories/test_admin_settings_repository.py @@ -1,6 +1,6 @@ import pytest from app.db.docs import AuditLogDocument -from app.db.repositories.admin.admin_settings_repository import AdminSettingsRepository +from app.db.repositories import AdminSettingsRepository from app.domain.admin import SystemSettings from dishka import AsyncContainer diff --git a/backend/tests/e2e/db/repositories/test_dlq_repository.py b/backend/tests/e2e/db/repositories/test_dlq_repository.py index 7197f08b..2965cd8b 100644 --- a/backend/tests/e2e/db/repositories/test_dlq_repository.py +++ b/backend/tests/e2e/db/repositories/test_dlq_repository.py @@ -3,7 +3,7 @@ import pytest from app.db.docs import DLQMessageDocument -from app.db.repositories.dlq_repository import DLQRepository +from app.db.repositories import DLQRepository from app.dlq import DLQMessageStatus from app.domain.enums import EventType diff --git a/backend/tests/e2e/db/repositories/test_execution_repository.py b/backend/tests/e2e/db/repositories/test_execution_repository.py index a6db17e3..60694fd1 100644 --- a/backend/tests/e2e/db/repositories/test_execution_repository.py +++ b/backend/tests/e2e/db/repositories/test_execution_repository.py @@ -2,7 +2,7 @@ from uuid import uuid4 import pytest -from app.db.repositories.execution_repository import ExecutionRepository +from app.db.repositories import ExecutionRepository from app.domain.enums import ExecutionStatus from app.domain.execution import DomainExecutionCreate, DomainExecutionUpdate diff --git a/backend/tests/e2e/db/repositories/test_saved_script_repository.py b/backend/tests/e2e/db/repositories/test_saved_script_repository.py index 5be6a5f6..4ce7c040 100644 --- a/backend/tests/e2e/db/repositories/test_saved_script_repository.py +++ b/backend/tests/e2e/db/repositories/test_saved_script_repository.py @@ -1,5 +1,5 @@ import pytest -from app.db.repositories.saved_script_repository import SavedScriptRepository +from app.db.repositories import SavedScriptRepository from app.domain.saved_script import DomainSavedScriptCreate, DomainSavedScriptUpdate from dishka import AsyncContainer diff --git a/backend/tests/e2e/dlq/test_dlq_discard.py b/backend/tests/e2e/dlq/test_dlq_discard.py index 39c023fb..5c6388a1 100644 --- a/backend/tests/e2e/dlq/test_dlq_discard.py +++ b/backend/tests/e2e/dlq/test_dlq_discard.py @@ -4,7 +4,7 @@ import pytest from app.db.docs import DLQMessageDocument -from app.db.repositories.dlq_repository import DLQRepository +from app.db.repositories import DLQRepository from app.dlq.models import DLQMessageStatus from app.domain.enums import KafkaTopic from dishka import AsyncContainer diff --git a/backend/tests/e2e/dlq/test_dlq_manager.py b/backend/tests/e2e/dlq/test_dlq_manager.py index 35a90adf..efa7d537 100644 --- a/backend/tests/e2e/dlq/test_dlq_manager.py +++ b/backend/tests/e2e/dlq/test_dlq_manager.py @@ -8,11 +8,11 @@ from aiokafka import AIOKafkaConsumer from app.core.metrics import DLQMetrics from app.core.providers import _default_retry_policies, _default_retry_policy -from app.db.repositories.dlq_repository import DLQRepository +from app.db.repositories import DLQRepository from app.dlq.manager import DLQManager from app.dlq.models import DLQMessage from app.domain.enums import EventType, KafkaTopic -from app.domain.events.typed import DLQMessageReceivedEvent, DomainEventAdapter +from app.domain.events import DLQMessageReceivedEvent, DomainEventAdapter from app.settings import Settings from dishka import AsyncContainer from faststream.kafka import KafkaBroker diff --git a/backend/tests/e2e/dlq/test_dlq_retry.py b/backend/tests/e2e/dlq/test_dlq_retry.py index 06425c58..6429ac1b 100644 --- a/backend/tests/e2e/dlq/test_dlq_retry.py +++ b/backend/tests/e2e/dlq/test_dlq_retry.py @@ -4,7 +4,7 @@ import pytest from app.db.docs import DLQMessageDocument -from app.db.repositories.dlq_repository import DLQRepository +from app.db.repositories import DLQRepository from app.dlq.models import DLQMessageStatus from app.domain.enums import KafkaTopic from dishka import AsyncContainer diff --git a/backend/tests/e2e/notifications/test_notification_sse.py b/backend/tests/e2e/notifications/test_notification_sse.py index 1a34f840..e9e1f4b7 100644 --- a/backend/tests/e2e/notifications/test_notification_sse.py +++ b/backend/tests/e2e/notifications/test_notification_sse.py @@ -5,7 +5,7 @@ from app.domain.enums import NotificationChannel, NotificationSeverity from app.schemas_pydantic.sse import RedisNotificationMessage from app.services.notification_service import NotificationService -from app.services.sse.redis_bus import SSERedisBus +from app.services.sse import SSERedisBus from dishka import AsyncContainer pytestmark = [pytest.mark.e2e, pytest.mark.redis] diff --git a/backend/tests/e2e/result_processor/test_result_processor.py b/backend/tests/e2e/result_processor/test_result_processor.py index e9a2a43c..3e36b1f8 100644 --- a/backend/tests/e2e/result_processor/test_result_processor.py +++ b/backend/tests/e2e/result_processor/test_result_processor.py @@ -3,9 +3,9 @@ import pytest from app.core.metrics import ExecutionMetrics from app.db.docs import ExecutionDocument -from app.db.repositories.execution_repository import ExecutionRepository +from app.db.repositories import ExecutionRepository from app.domain.enums import ExecutionStatus -from app.domain.events.typed import ( +from app.domain.events import ( EventMetadata, ExecutionCompletedEvent, ResourceUsageDomain, diff --git a/backend/tests/e2e/services/replay/test_replay_service.py b/backend/tests/e2e/services/replay/test_replay_service.py index e7cbb584..27046e37 100644 --- a/backend/tests/e2e/services/replay/test_replay_service.py +++ b/backend/tests/e2e/services/replay/test_replay_service.py @@ -1,6 +1,6 @@ import pytest from app.domain.enums import ReplayStatus, ReplayTarget, ReplayType -from app.domain.replay.exceptions import ReplaySessionNotFoundError +from app.domain.replay import ReplaySessionNotFoundError from app.services.event_replay import EventReplayService, ReplayConfig, ReplayFilter from dishka import AsyncContainer diff --git a/backend/tests/e2e/services/saga/test_saga_service.py b/backend/tests/e2e/services/saga/test_saga_service.py index f41888ac..e91bd844 100644 --- a/backend/tests/e2e/services/saga/test_saga_service.py +++ b/backend/tests/e2e/services/saga/test_saga_service.py @@ -5,8 +5,7 @@ from app.db.repositories import ExecutionRepository, SagaRepository from app.domain.enums import SagaState, UserRole from app.domain.execution import DomainExecutionCreate -from app.domain.saga.exceptions import SagaAccessDeniedError, SagaNotFoundError -from app.domain.saga.models import Saga, SagaListResult +from app.domain.saga import Saga, SagaAccessDeniedError, SagaListResult, SagaNotFoundError from app.schemas_pydantic.user import User from app.services.execution_service import ExecutionService from app.services.saga.saga_service import SagaService diff --git a/backend/tests/e2e/services/sse/test_partitioned_event_router.py b/backend/tests/e2e/services/sse/test_partitioned_event_router.py index 76548d80..11a18091 100644 --- a/backend/tests/e2e/services/sse/test_partitioned_event_router.py +++ b/backend/tests/e2e/services/sse/test_partitioned_event_router.py @@ -5,7 +5,7 @@ import pytest import redis.asyncio as redis from app.schemas_pydantic.sse import RedisSSEMessage -from app.services.sse.redis_bus import SSERedisBus +from app.services.sse import SSERedisBus from app.settings import Settings from tests.conftest import make_execution_requested_event diff --git a/backend/tests/e2e/services/sse/test_redis_bus.py b/backend/tests/e2e/services/sse/test_redis_bus.py index 0c712a78..fffea4a0 100644 --- a/backend/tests/e2e/services/sse/test_redis_bus.py +++ b/backend/tests/e2e/services/sse/test_redis_bus.py @@ -6,9 +6,9 @@ import pytest import redis.asyncio as redis_async from app.domain.enums import EventType, NotificationSeverity, NotificationStatus -from app.domain.events.typed import EventMetadata, ExecutionCompletedEvent +from app.domain.events import EventMetadata, ExecutionCompletedEvent from app.schemas_pydantic.sse import RedisNotificationMessage, RedisSSEMessage -from app.services.sse.redis_bus import SSERedisBus +from app.services.sse import SSERedisBus pytestmark = pytest.mark.e2e diff --git a/backend/tests/e2e/services/user_settings/test_user_settings_service.py b/backend/tests/e2e/services/user_settings/test_user_settings_service.py index 5cf14a9f..11a2dda9 100644 --- a/backend/tests/e2e/services/user_settings/test_user_settings_service.py +++ b/backend/tests/e2e/services/user_settings/test_user_settings_service.py @@ -3,7 +3,7 @@ import pytest from app.domain.enums import Theme -from app.domain.user.settings_models import ( +from app.domain.user import ( DomainEditorSettings, DomainNotificationSettings, DomainSettingsHistoryEntry, diff --git a/backend/tests/e2e/test_admin_events_routes.py b/backend/tests/e2e/test_admin_events_routes.py index 5af0087b..4f0dd72e 100644 --- a/backend/tests/e2e/test_admin_events_routes.py +++ b/backend/tests/e2e/test_admin_events_routes.py @@ -2,9 +2,9 @@ import pytest import pytest_asyncio -from app.db.repositories.event_repository import EventRepository +from app.db.repositories import EventRepository from app.domain.enums import EventType, ReplayStatus -from app.domain.events.typed import DomainEvent +from app.domain.events import DomainEvent from app.schemas_pydantic.admin_events import ( EventBrowseRequest, EventBrowseResponse, diff --git a/backend/tests/e2e/test_events_routes.py b/backend/tests/e2e/test_events_routes.py index 96e60f03..7c774a0b 100644 --- a/backend/tests/e2e/test_events_routes.py +++ b/backend/tests/e2e/test_events_routes.py @@ -1,6 +1,6 @@ import pytest from app.domain.enums import EventType -from app.domain.events.typed import DomainEvent, ExecutionRequestedEvent +from app.domain.events import DomainEvent, ExecutionRequestedEvent from app.schemas_pydantic.events import ( DeleteEventResponse, EventListResponse, diff --git a/backend/tests/e2e/test_execution_routes.py b/backend/tests/e2e/test_execution_routes.py index 3facdeda..f73979d9 100644 --- a/backend/tests/e2e/test_execution_routes.py +++ b/backend/tests/e2e/test_execution_routes.py @@ -9,7 +9,7 @@ import pytest from app.domain.enums import EventType, ExecutionStatus -from app.domain.events.typed import ExecutionDomainEvent +from app.domain.events import ExecutionDomainEvent from app.schemas_pydantic.execution import ( CancelExecutionRequest, CancelResponse, diff --git a/backend/tests/e2e/test_k8s_worker_create_pod.py b/backend/tests/e2e/test_k8s_worker_create_pod.py index ecf48a07..6ec9984d 100644 --- a/backend/tests/e2e/test_k8s_worker_create_pod.py +++ b/backend/tests/e2e/test_k8s_worker_create_pod.py @@ -4,7 +4,7 @@ import pytest from app.core.metrics import EventMetrics from app.domain.enums import QueuePriority -from app.domain.events.typed import CreatePodCommandEvent, EventMetadata +from app.domain.events import CreatePodCommandEvent, EventMetadata from app.events.core import UnifiedProducer from app.services.k8s_worker import KubernetesWorker from app.settings import Settings diff --git a/backend/tests/unit/core/test_csrf.py b/backend/tests/unit/core/test_csrf.py index d733f5cf..bc20f7d2 100644 --- a/backend/tests/unit/core/test_csrf.py +++ b/backend/tests/unit/core/test_csrf.py @@ -1,6 +1,6 @@ import pytest from app.core.security import SecurityService -from app.domain.user.exceptions import CSRFValidationError +from app.domain.user import CSRFValidationError from app.settings import Settings from starlette.requests import Request diff --git a/backend/tests/unit/domain/events/test_event_schema_coverage.py b/backend/tests/unit/domain/events/test_event_schema_coverage.py index 9055a8a1..7d386f1a 100644 --- a/backend/tests/unit/domain/events/test_event_schema_coverage.py +++ b/backend/tests/unit/domain/events/test_event_schema_coverage.py @@ -12,7 +12,7 @@ from typing import get_args from app.domain.enums import EventType -from app.domain.events.typed import BaseEvent, DomainEvent, DomainEventAdapter +from app.domain.events import BaseEvent, DomainEvent, DomainEventAdapter def get_domain_event_classes() -> dict[EventType, type]: diff --git a/backend/tests/unit/events/test_metadata_model.py b/backend/tests/unit/events/test_metadata_model.py index f237a263..a99ac70e 100644 --- a/backend/tests/unit/events/test_metadata_model.py +++ b/backend/tests/unit/events/test_metadata_model.py @@ -1,4 +1,4 @@ -from app.domain.events.typed import EventMetadata +from app.domain.events import EventMetadata def test_metadata_creation() -> None: diff --git a/backend/tests/unit/services/coordinator/test_coordinator_queue.py b/backend/tests/unit/services/coordinator/test_coordinator_queue.py index cc6467c8..f31cfbe8 100644 --- a/backend/tests/unit/services/coordinator/test_coordinator_queue.py +++ b/backend/tests/unit/services/coordinator/test_coordinator_queue.py @@ -4,7 +4,7 @@ import pytest from app.core.metrics import CoordinatorMetrics from app.domain.enums import QueuePriority -from app.domain.events.typed import ExecutionRequestedEvent +from app.domain.events import ExecutionRequestedEvent from app.services.coordinator.coordinator import ExecutionCoordinator, QueueRejectError from tests.conftest import make_execution_requested_event diff --git a/backend/tests/unit/services/idempotency/test_idempotency_manager.py b/backend/tests/unit/services/idempotency/test_idempotency_manager.py index 98d2c433..96d5c44f 100644 --- a/backend/tests/unit/services/idempotency/test_idempotency_manager.py +++ b/backend/tests/unit/services/idempotency/test_idempotency_manager.py @@ -3,7 +3,7 @@ import pytest from app.core.metrics import DatabaseMetrics -from app.domain.events.typed import BaseEvent +from app.domain.events import BaseEvent from app.domain.idempotency import KeyStrategy from app.services.idempotency.idempotency_manager import ( IdempotencyConfig, diff --git a/backend/tests/unit/services/pod_monitor/test_event_mapper.py b/backend/tests/unit/services/pod_monitor/test_event_mapper.py index bd4b6281..6094a5b3 100644 --- a/backend/tests/unit/services/pod_monitor/test_event_mapper.py +++ b/backend/tests/unit/services/pod_monitor/test_event_mapper.py @@ -6,7 +6,7 @@ from kubernetes_asyncio.client import V1Pod, V1PodCondition from app.domain.enums import EventType, ExecutionErrorType -from app.domain.events.typed import ( +from app.domain.events import ( EventMetadata, ExecutionCompletedEvent, ExecutionFailedEvent, diff --git a/backend/tests/unit/services/pod_monitor/test_monitor.py b/backend/tests/unit/services/pod_monitor/test_monitor.py index 14f0a61d..36fbd9ed 100644 --- a/backend/tests/unit/services/pod_monitor/test_monitor.py +++ b/backend/tests/unit/services/pod_monitor/test_monitor.py @@ -5,7 +5,7 @@ import pytest from app.core.metrics import EventMetrics, KubernetesMetrics -from app.domain.events.typed import ( +from app.domain.events import ( DomainEvent, EventMetadata, ExecutionCompletedEvent, diff --git a/backend/tests/unit/services/result_processor/test_processor.py b/backend/tests/unit/services/result_processor/test_processor.py index 8ec741d3..bf1b8eb3 100644 --- a/backend/tests/unit/services/result_processor/test_processor.py +++ b/backend/tests/unit/services/result_processor/test_processor.py @@ -4,7 +4,7 @@ import pytest from app.core.metrics import ExecutionMetrics from app.domain.enums import ExecutionErrorType, ExecutionStatus -from app.domain.events.typed import ( +from app.domain.events import ( EventMetadata, ExecutionCompletedEvent, ExecutionFailedEvent, diff --git a/backend/tests/unit/services/saga/test_execution_saga_steps.py b/backend/tests/unit/services/saga/test_execution_saga_steps.py index f0349f12..84d06b5d 100644 --- a/backend/tests/unit/services/saga/test_execution_saga_steps.py +++ b/backend/tests/unit/services/saga/test_execution_saga_steps.py @@ -1,6 +1,6 @@ import pytest -from app.db.repositories.resource_allocation_repository import ResourceAllocationRepository -from app.domain.events.typed import DomainEvent, ExecutionRequestedEvent +from app.db.repositories import ResourceAllocationRepository +from app.domain.events import DomainEvent, ExecutionRequestedEvent from app.domain.saga import DomainResourceAllocation, DomainResourceAllocationCreate from app.events.core import UnifiedProducer from app.services.saga.execution_saga import ( diff --git a/backend/tests/unit/services/saga/test_saga_comprehensive.py b/backend/tests/unit/services/saga/test_saga_comprehensive.py index 2c88b581..7a94d67c 100644 --- a/backend/tests/unit/services/saga/test_saga_comprehensive.py +++ b/backend/tests/unit/services/saga/test_saga_comprehensive.py @@ -7,8 +7,8 @@ import pytest from app.domain.enums import SagaState -from app.domain.events.typed import DomainEvent, ExecutionRequestedEvent -from app.domain.saga.models import Saga +from app.domain.events import DomainEvent, ExecutionRequestedEvent +from app.domain.saga import Saga from app.services.saga.saga_step import CompensationStep, SagaContext, SagaStep from tests.conftest import make_execution_requested_event diff --git a/backend/tests/unit/services/saga/test_saga_orchestrator_unit.py b/backend/tests/unit/services/saga/test_saga_orchestrator_unit.py index 33ce5d92..86fa4fdd 100644 --- a/backend/tests/unit/services/saga/test_saga_orchestrator_unit.py +++ b/backend/tests/unit/services/saga/test_saga_orchestrator_unit.py @@ -1,12 +1,10 @@ import logging import pytest -from app.db.repositories.resource_allocation_repository import ResourceAllocationRepository -from app.db.repositories.saga_repository import SagaRepository +from app.db.repositories import ResourceAllocationRepository, SagaRepository from app.domain.enums import SagaState -from app.domain.events.typed import DomainEvent -from app.domain.saga import DomainResourceAllocation, DomainResourceAllocationCreate -from app.domain.saga.models import Saga, SagaConfig +from app.domain.events import DomainEvent +from app.domain.saga import DomainResourceAllocation, DomainResourceAllocationCreate, Saga, SagaConfig from app.events.core import UnifiedProducer from app.services.saga.execution_saga import ExecutionSaga from app.services.saga.saga_orchestrator import SagaOrchestrator diff --git a/backend/tests/unit/services/saga/test_saga_step_and_base.py b/backend/tests/unit/services/saga/test_saga_step_and_base.py index 850ab187..33eea711 100644 --- a/backend/tests/unit/services/saga/test_saga_step_and_base.py +++ b/backend/tests/unit/services/saga/test_saga_step_and_base.py @@ -1,5 +1,5 @@ import pytest -from app.domain.events.typed import SystemErrorEvent +from app.domain.events import SystemErrorEvent from app.services.saga.saga_step import CompensationStep, SagaContext, SagaStep pytestmark = pytest.mark.unit diff --git a/backend/tests/unit/services/sse/test_kafka_redis_bridge.py b/backend/tests/unit/services/sse/test_kafka_redis_bridge.py index f3122dc4..ce0b281e 100644 --- a/backend/tests/unit/services/sse/test_kafka_redis_bridge.py +++ b/backend/tests/unit/services/sse/test_kafka_redis_bridge.py @@ -1,8 +1,8 @@ import logging import pytest -from app.domain.events.typed import DomainEvent, EventMetadata, ExecutionStartedEvent -from app.services.sse.redis_bus import SSERedisBus +from app.domain.events import DomainEvent, EventMetadata, ExecutionStartedEvent +from app.services.sse import SSERedisBus pytestmark = pytest.mark.unit diff --git a/backend/tests/unit/services/sse/test_sse_service.py b/backend/tests/unit/services/sse/test_sse_service.py index 855c35ad..a43a82ea 100644 --- a/backend/tests/unit/services/sse/test_sse_service.py +++ b/backend/tests/unit/services/sse/test_sse_service.py @@ -7,13 +7,12 @@ import pytest from app.core.metrics import ConnectionMetrics -from app.db.repositories.sse_repository import SSERepository +from app.db.repositories import SSERepository from app.domain.enums import EventType, ExecutionStatus from app.domain.events import ResourceUsageDomain from app.domain.execution import DomainExecution from app.domain.sse import SSEExecutionStatusDomain -from app.services.sse.redis_bus import SSERedisBus, SSERedisSubscription -from app.services.sse.sse_service import SSEService +from app.services.sse import SSERedisBus, SSERedisSubscription, SSEService from app.settings import Settings from pydantic import BaseModel diff --git a/backend/tests/unit/services/test_pod_builder.py b/backend/tests/unit/services/test_pod_builder.py index e01951a5..528f1b17 100644 --- a/backend/tests/unit/services/test_pod_builder.py +++ b/backend/tests/unit/services/test_pod_builder.py @@ -3,7 +3,7 @@ import pytest from app.domain.enums import QueuePriority -from app.domain.events.typed import CreatePodCommandEvent, EventMetadata +from app.domain.events import CreatePodCommandEvent, EventMetadata from app.services.k8s_worker import PodBuilder from kubernetes_asyncio import client as k8s_client From 45aaf6a8a3c140c13551313a5fd29c4060187150 Mon Sep 17 00:00:00 2001 From: HardMax71 Date: Sun, 8 Feb 2026 20:53:07 +0100 Subject: [PATCH 3/3] dlq processor + manager -> dlq manager only --- backend/app/domain/enums/kafka.py | 1 - backend/app/infrastructure/kafka/mappings.py | 2 +- 2 files changed, 1 insertion(+), 2 deletions(-) diff --git a/backend/app/domain/enums/kafka.py b/backend/app/domain/enums/kafka.py index d248bd8b..fbb5d692 100644 --- a/backend/app/domain/enums/kafka.py +++ b/backend/app/domain/enums/kafka.py @@ -61,5 +61,4 @@ class GroupId(StringEnum): EVENT_STORE_CONSUMER = "event-store-consumer" WEBSOCKET_GATEWAY = "websocket-gateway" NOTIFICATION_SERVICE = "notification-service" - DLQ_PROCESSOR = "dlq-processor" DLQ_MANAGER = "dlq-manager" diff --git a/backend/app/infrastructure/kafka/mappings.py b/backend/app/infrastructure/kafka/mappings.py index 2a9795de..c7597db2 100644 --- a/backend/app/infrastructure/kafka/mappings.py +++ b/backend/app/infrastructure/kafka/mappings.py @@ -132,7 +132,7 @@ def get_event_types_for_topic(topic: KafkaTopic) -> list[EventType]: KafkaTopic.NOTIFICATION_EVENTS, KafkaTopic.EXECUTION_EVENTS, }, - GroupId.DLQ_PROCESSOR: { + GroupId.DLQ_MANAGER: { KafkaTopic.DEAD_LETTER_QUEUE, }, }