diff --git a/changes/12867.feature.md b/changes/12867.feature.md new file mode 100644 index 00000000000..16b1186b589 --- /dev/null +++ b/changes/12867.feature.md @@ -0,0 +1 @@ +Add the models-layer query conditions and orders for kernel scheduling history, covering filters on id, kernel, session, phase, status transition, result and timestamps. diff --git a/changes/12868.feature.md b/changes/12868.feature.md new file mode 100644 index 00000000000..c151c6acc8f --- /dev/null +++ b/changes/12868.feature.md @@ -0,0 +1 @@ +Add the kernel scheduling-history repository search scope, which scopes a query by session, by kernel, or by both, and rejects an empty scope rather than degenerating into a system-wide search. diff --git a/changes/12869.feature.md b/changes/12869.feature.md new file mode 100644 index 00000000000..f67010c908e --- /dev/null +++ b/changes/12869.feature.md @@ -0,0 +1 @@ +Add the v2 DTOs and service actions for kernel scheduling-history search, forming the schema shared by REST v2 and GraphQL. diff --git a/changes/12887.fix.md b/changes/12887.fix.md new file mode 100644 index 00000000000..382995cb258 --- /dev/null +++ b/changes/12887.fix.md @@ -0,0 +1 @@ +Fix the legacy app_configs downgrade failing with a duplicate enum type error diff --git a/changes/12963.enhance.md b/changes/12963.enhance.md new file mode 100644 index 00000000000..abc45d7aadb --- /dev/null +++ b/changes/12963.enhance.md @@ -0,0 +1 @@ +Retire the dedicated `error_logs` cleanup timer (its rows are now purged by the DB record retention policy sweep) and redefine `clear-history` as a manual force-sweep plus `VACUUM` escape hatch over the retention policy layer. diff --git a/src/ai/backend/common/data/permission/types.py b/src/ai/backend/common/data/permission/types.py index 19c0950948d..a76d83eb0d1 100644 --- a/src/ai/backend/common/data/permission/types.py +++ b/src/ai/backend/common/data/permission/types.py @@ -174,6 +174,9 @@ class EntityType(enum.StrEnum): SESSION_DIRECT_ACCESS = "session:direct_access" SESSION_HISTORY = "session:history" SESSION_SCOPED_HISTORY = "session:scoped_history" + # Kernel sub + KERNEL_HISTORY = "kernel:history" + KERNEL_SCOPED_HISTORY = "kernel:scoped_history" # Deployment sub DEPLOYMENT_REPLICA = "deployment:replica" DEPLOYMENT_ROUTE = "deployment:route" diff --git a/src/ai/backend/common/dto/manager/v2/scheduling_history/__init__.py b/src/ai/backend/common/dto/manager/v2/scheduling_history/__init__.py index 9a923071f02..5f65f224a08 100644 --- a/src/ai/backend/common/dto/manager/v2/scheduling_history/__init__.py +++ b/src/ai/backend/common/dto/manager/v2/scheduling_history/__init__.py @@ -4,12 +4,16 @@ from ai.backend.common.dto.manager.v2.scheduling_history.request import ( AdminSearchDeploymentHistoriesInput, + AdminSearchKernelHistoriesInput, AdminSearchRouteHistoriesInput, AdminSearchSessionHistoriesInput, DeploymentHistoryFilter, DeploymentHistoryOrder, + KernelHistoryFilter, + KernelHistoryOrder, RouteHistoryFilter, RouteHistoryOrder, + ScopedSearchKernelHistoriesInput, SearchDeploymentHistoryInput, SearchRouteHistoryInput, SearchSessionHistoryInput, @@ -18,9 +22,11 @@ ) from ai.backend.common.dto.manager.v2.scheduling_history.response import ( AdminSearchDeploymentHistoriesPayload, + AdminSearchKernelHistoriesPayload, AdminSearchRouteHistoriesPayload, AdminSearchSessionHistoriesPayload, DeploymentHistoryNode, + KernelHistoryNode, ListDeploymentHistoryPayload, ListRouteHistoryPayload, ListSessionHistoryPayload, @@ -30,6 +36,8 @@ from ai.backend.common.dto.manager.v2.scheduling_history.types import ( DeploymentHistoryOrderField, DeploymentHistoryScopeDTO, + KernelHistoryOrderField, + KernelHistoryScopeDTO, OrderDirection, RouteHistoryOrderField, RouteHistoryScopeDTO, @@ -43,6 +51,8 @@ # Types "DeploymentHistoryOrderField", "DeploymentHistoryScopeDTO", + "KernelHistoryOrderField", + "KernelHistoryScopeDTO", "OrderDirection", "RouteHistoryOrderField", "RouteHistoryScopeDTO", @@ -52,12 +62,16 @@ "SubStepResultInfo", # Input models (request) "AdminSearchDeploymentHistoriesInput", + "AdminSearchKernelHistoriesInput", "AdminSearchRouteHistoriesInput", "AdminSearchSessionHistoriesInput", "DeploymentHistoryFilter", "DeploymentHistoryOrder", + "KernelHistoryFilter", + "KernelHistoryOrder", "RouteHistoryFilter", "RouteHistoryOrder", + "ScopedSearchKernelHistoriesInput", "SearchDeploymentHistoryInput", "SearchRouteHistoryInput", "SearchSessionHistoryInput", @@ -65,9 +79,11 @@ "SessionHistoryOrder", # Response models "AdminSearchDeploymentHistoriesPayload", + "AdminSearchKernelHistoriesPayload", "AdminSearchRouteHistoriesPayload", "AdminSearchSessionHistoriesPayload", "DeploymentHistoryNode", + "KernelHistoryNode", "ListDeploymentHistoryPayload", "ListRouteHistoryPayload", "ListSessionHistoryPayload", diff --git a/src/ai/backend/common/dto/manager/v2/scheduling_history/request.py b/src/ai/backend/common/dto/manager/v2/scheduling_history/request.py index ce9eac401bd..858b31f6a4d 100644 --- a/src/ai/backend/common/dto/manager/v2/scheduling_history/request.py +++ b/src/ai/backend/common/dto/manager/v2/scheduling_history/request.py @@ -11,6 +11,8 @@ from .types import ( DeploymentHistoryOrderField, + KernelHistoryOrderField, + KernelHistoryScopeDTO, OrderDirection, RouteHistoryOrderField, SchedulingResultType, @@ -19,13 +21,17 @@ __all__ = ( "AdminSearchDeploymentHistoriesInput", + "AdminSearchKernelHistoriesInput", "AdminSearchRouteHistoriesInput", "AdminSearchSessionHistoriesInput", "DeploymentHistoryFilter", "DeploymentHistoryOrder", + "KernelHistoryFilter", + "KernelHistoryOrder", "RouteHistoryFilter", "RouteHistoryOrder", "SchedulingResultFilter", + "ScopedSearchKernelHistoriesInput", "SearchDeploymentHistoryInput", "SearchRouteHistoryInput", "SearchSessionHistoryInput", @@ -86,6 +92,37 @@ class SearchSessionHistoryInput(BaseRequestModel): offset: int = Field(default=0, ge=0, description="Number of items to skip") +class KernelHistoryFilter(BaseRequestModel): + """Filter conditions for kernel scheduling history search.""" + + id: UUIDFilter | None = Field(default=None, description="Filter by history record ID") + kernel_id: UUIDFilter | None = Field(default=None, description="Filter by kernel ID") + session_id: UUIDFilter | None = Field(default=None, description="Filter by session ID") + phase: StringFilter | None = Field(default=None, description="Filter by scheduling phase") + from_status: list[str] | None = Field(default=None, description="Filter by from_status values") + to_status: list[str] | None = Field(default=None, description="Filter by to_status values") + result: SchedulingResultFilter | None = Field( + default=None, description="Filter by scheduling result" + ) + error_code: StringFilter | None = Field(default=None, description="Filter by error code") + message: StringFilter | None = Field(default=None, description="Filter by message") + created_at: DateTimeFilter | None = Field(default=None, description="Filter by created_at") + updated_at: DateTimeFilter | None = Field(default=None, description="Filter by updated_at") + AND: list[KernelHistoryFilter] | None = Field(default=None, description="AND conjunction.") + OR: list[KernelHistoryFilter] | None = Field(default=None, description="OR conjunction.") + NOT: list[KernelHistoryFilter] | None = Field(default=None, description="NOT negation.") + + +KernelHistoryFilter.model_rebuild() + + +class KernelHistoryOrder(BaseRequestModel): + """Order specification for kernel scheduling history.""" + + field: KernelHistoryOrderField = Field(description="Field to order by") + direction: OrderDirection = Field(default=OrderDirection.DESC, description="Order direction") + + class DeploymentHistoryFilter(BaseRequestModel): """Filter conditions for deployment scheduling history search.""" @@ -182,6 +219,37 @@ class AdminSearchSessionHistoriesInput(BaseRequestModel): offset: int | None = Field(default=None, description="Offset pagination: number to skip") +class AdminSearchKernelHistoriesInput(BaseRequestModel): + """Input for admin search of kernel scheduling histories.""" + + filter: KernelHistoryFilter | None = Field(default=None, description="Filter conditions") + order: list[KernelHistoryOrder] | None = Field(default=None, description="Order specifications") + first: int | None = Field(default=None, description="Cursor pagination: number of items") + after: str | None = Field(default=None, description="Cursor pagination: after cursor") + last: int | None = Field(default=None, description="Cursor pagination: last N items") + before: str | None = Field(default=None, description="Cursor pagination: before cursor") + limit: int | None = Field(default=None, description="Offset pagination: maximum items") + offset: int | None = Field(default=None, description="Offset pagination: number to skip") + + +class ScopedSearchKernelHistoriesInput(BaseRequestModel): + """Input for scoped search of kernel scheduling histories. + + Unlike the session/deployment/route scoped searches, the scope is carried in the body + because a kernel query may be scoped by session_id, kernel_id, or both. + """ + + scope: KernelHistoryScopeDTO = Field(description="Required. Restricts the rows returned.") + filter: KernelHistoryFilter | None = Field(default=None, description="Filter conditions") + order: list[KernelHistoryOrder] | None = Field(default=None, description="Order specifications") + first: int | None = Field(default=None, description="Cursor pagination: number of items") + after: str | None = Field(default=None, description="Cursor pagination: after cursor") + last: int | None = Field(default=None, description="Cursor pagination: last N items") + before: str | None = Field(default=None, description="Cursor pagination: before cursor") + limit: int | None = Field(default=None, description="Offset pagination: maximum items") + offset: int | None = Field(default=None, description="Offset pagination: number to skip") + + class AdminSearchDeploymentHistoriesInput(BaseRequestModel): """Input for admin search of deployment histories.""" diff --git a/src/ai/backend/common/dto/manager/v2/scheduling_history/response.py b/src/ai/backend/common/dto/manager/v2/scheduling_history/response.py index e667f4c1aef..f16a3556114 100644 --- a/src/ai/backend/common/dto/manager/v2/scheduling_history/response.py +++ b/src/ai/backend/common/dto/manager/v2/scheduling_history/response.py @@ -16,9 +16,11 @@ __all__ = ( "AdminSearchDeploymentHistoriesPayload", + "AdminSearchKernelHistoriesPayload", "AdminSearchRouteHistoriesPayload", "AdminSearchSessionHistoriesPayload", "DeploymentHistoryNode", + "KernelHistoryNode", "ListDeploymentHistoryPayload", "ListRouteHistoryPayload", "ListSessionHistoryPayload", @@ -46,6 +48,26 @@ class SessionHistoryNode(BaseResponseModel): updated_at: datetime = Field(description="Timestamp when the history record was last updated") +class KernelHistoryNode(BaseResponseModel): + """Node model representing a kernel scheduling history record. + + Has no ``sub_steps``: the ``kernel_scheduling_history`` table carries no such column. + """ + + id: UUID = Field(description="History record ID") + kernel_id: UUID = Field(description="Kernel ID this history belongs to") + session_id: UUID = Field(description="Session owning the kernel") + phase: str = Field(description="Scheduling phase") + from_status: str | None = Field(default=None, description="Status before transition") + to_status: str | None = Field(default=None, description="Status after transition") + result: str = Field(description="Result of the scheduling attempt") + error_code: str | None = Field(default=None, description="Error code if scheduling failed") + message: str | None = Field(default=None, description="Human-readable message or error detail") + attempts: int = Field(description="Number of scheduling attempts made") + created_at: datetime = Field(description="Timestamp when the history record was created") + updated_at: datetime = Field(description="Timestamp when the history record was last updated") + + class DeploymentHistoryNode(BaseResponseModel): """Node model representing a deployment scheduling history record. @@ -127,6 +149,15 @@ class AdminSearchSessionHistoriesPayload(BaseResponseModel): has_previous_page: bool = Field(description="Whether there is a previous page.") +class AdminSearchKernelHistoriesPayload(BaseResponseModel): + """Payload for admin and scoped search of kernel scheduling histories.""" + + items: list[KernelHistoryNode] = Field(description="List of kernel history nodes.") + total_count: int = Field(description="Total number of records matching the filter.") + has_next_page: bool = Field(description="Whether there is a next page.") + has_previous_page: bool = Field(description="Whether there is a previous page.") + + class AdminSearchDeploymentHistoriesPayload(BaseResponseModel): """Payload for admin search of deployment histories.""" diff --git a/src/ai/backend/common/dto/manager/v2/scheduling_history/types.py b/src/ai/backend/common/dto/manager/v2/scheduling_history/types.py index 2557da8eee2..43cbd5a8047 100644 --- a/src/ai/backend/common/dto/manager/v2/scheduling_history/types.py +++ b/src/ai/backend/common/dto/manager/v2/scheduling_history/types.py @@ -8,7 +8,7 @@ from enum import StrEnum from uuid import UUID -from pydantic import Field +from pydantic import Field, model_validator from ai.backend.common.api_handlers import BaseRequestModel, BaseResponseModel from ai.backend.common.dto.manager.v2.common import OrderDirection @@ -16,6 +16,8 @@ __all__ = ( "DeploymentHistoryOrderField", "DeploymentHistoryScopeDTO", + "KernelHistoryOrderField", + "KernelHistoryScopeDTO", "OrderDirection", "RouteHistoryOrderField", "RouteHistoryScopeDTO", @@ -45,6 +47,13 @@ class SessionHistoryOrderField(StrEnum): UPDATED_AT = "updated_at" +class KernelHistoryOrderField(StrEnum): + """Fields available for ordering kernel scheduling history.""" + + CREATED_AT = "created_at" + UPDATED_AT = "updated_at" + + class DeploymentHistoryOrderField(StrEnum): """Fields available for ordering deployment scheduling history.""" @@ -76,6 +85,25 @@ class SessionHistoryScopeDTO(BaseRequestModel): session_id: UUID = Field(description="Session ID to get history for.") +class KernelHistoryScopeDTO(BaseRequestModel): + """Scope for kernel scheduling history queries. + + Unlike the single-entity scopes above, a kernel query may be scoped either by the + owning session (all of its kernels) or by one kernel. At least one must be given. + """ + + session_id: UUID | None = Field( + default=None, description="Restrict to the kernels of this session." + ) + kernel_id: UUID | None = Field(default=None, description="Restrict to this kernel.") + + @model_validator(mode="after") + def _check_non_empty(self) -> KernelHistoryScopeDTO: + if self.session_id is None and self.kernel_id is None: + raise ValueError("At least one of session_id or kernel_id must be given.") + return self + + class DeploymentHistoryScopeDTO(BaseRequestModel): """Scope for deployment scheduling history queries.""" diff --git a/src/ai/backend/common/events/event_types/log/__init__.py b/src/ai/backend/common/events/event_types/log/__init__.py deleted file mode 100644 index e69de29bb2d..00000000000 diff --git a/src/ai/backend/common/events/event_types/log/anycast.py b/src/ai/backend/common/events/event_types/log/anycast.py deleted file mode 100644 index 6210c60a445..00000000000 --- a/src/ai/backend/common/events/event_types/log/anycast.py +++ /dev/null @@ -1,38 +0,0 @@ -from __future__ import annotations - -from typing import Any, Self, override - -from ai.backend.common.events.types import ( - AbstractAnycastEvent, - EventDomain, -) -from ai.backend.common.events.user_event.user_event import UserEvent - - -class DoLogCleanupEvent(AbstractAnycastEvent): - @override - def serialize(self) -> tuple[Any, ...]: - return tuple() - - @classmethod - @override - def deserialize(cls, value: tuple[Any, ...]) -> Self: - return cls() - - @classmethod - @override - def event_domain(cls) -> EventDomain: - return EventDomain.LOG - - @override - def domain_id(self) -> str | None: - return None - - @override - def user_event(self) -> UserEvent | None: - return None - - @classmethod - @override - def event_name(cls) -> str: - return "do_log_cleanup" diff --git a/src/ai/backend/common/identifier/kernel_scheduling_history.py b/src/ai/backend/common/identifier/kernel_scheduling_history.py new file mode 100644 index 00000000000..406667490cf --- /dev/null +++ b/src/ai/backend/common/identifier/kernel_scheduling_history.py @@ -0,0 +1,6 @@ +from typing import NewType +from uuid import UUID + +__all__ = ("KernelSchedulingHistoryID",) + +KernelSchedulingHistoryID = NewType("KernelSchedulingHistoryID", UUID) diff --git a/src/ai/backend/manager/api/rest/error_log/registry.py b/src/ai/backend/manager/api/rest/error_log/registry.py index 580483ebbc2..d185b466985 100644 --- a/src/ai/backend/manager/api/rest/error_log/registry.py +++ b/src/ai/backend/manager/api/rest/error_log/registry.py @@ -1,10 +1,8 @@ """Error log module registrar. -Lifecycle management (GlobalTimer for log cleanup, event dispatcher -integration) is handled by the DependencyComposer: - -* Event consumer: ``event_dispatcher.handlers.log_cleanup`` -* GlobalTimer: ``dependencies.processing.log_cleanup_timer`` +Old ``error_logs`` rows are purged by the DB record retention sweep under the +``logs`` category (BEP-1063); this module only registers the HTTP routes. A +manual immediate sweep can be triggered via the ``clear-history`` CLI. """ from __future__ import annotations diff --git a/src/ai/backend/manager/cli/__main__.py b/src/ai/backend/manager/cli/__main__.py index 680efec25c9..817449f9bc1 100644 --- a/src/ai/backend/manager/cli/__main__.py +++ b/src/ai/backend/manager/cli/__main__.py @@ -235,7 +235,11 @@ def generate_rpc_keypair(_cli_ctx: CLIContext, dst_dir: pathlib.Path, name: str) "--retention", type=str, default="1yr", - help="The retention limit. e.g., 20d, 1mo, 6mo, 1yr", + help=( + "Age boundary for the Redis kernel-statistics cleanup this command still " + "owns. e.g., 20d, 1mo, 6mo, 1yr. DB records are purged by the retention " + "policy sweep (BEP-1063), not this option." + ), ) @click.option( "-v", @@ -253,24 +257,25 @@ def generate_rpc_keypair(_cli_ctx: CLIContext, dst_dir: pathlib.Path, name: str) @click.pass_obj def clear_history(cli_ctx: CLIContext, retention: str, vacuum_full: bool) -> None: """ - Delete old records from the kernels, error_logs tables and - invoke the PostgreSQL's vaccuum operation to clear up the actual disk space. + Manual escape hatch over the automated DB record retention policy layer + (BEP-1063). It triggers an immediate retention sweep (deleting records past + each enabled policy's age boundary), cleans up terminated-kernel Redis + statistics (a non-DB store the sweep does not touch), then runs PostgreSQL + VACUUM / VACUUM FULL to reclaim disk space the sweep never reclaims. """ import asyncio - import uuid from datetime import UTC, datetime - from typing import cast import sqlalchemy as sa from more_itertools import chunked from ai.backend.common.validators import TimeDuration - from ai.backend.manager.models.error_logs import error_logs from ai.backend.manager.models.kernel import kernels - from ai.backend.manager.models.session import SessionRow from ai.backend.manager.repositories.db.engine import connect_database, vacuum_db + from ai.backend.manager.repositories.ops import DBOpsProvider + from ai.backend.manager.repositories.retention.repository import RetentionRepository - from .context import redis_ctx + from .context import config_provider_ctx, redis_ctx log = _get_logger() today = datetime.now(UTC) @@ -317,53 +322,23 @@ async def _clear_redis_history() -> None: except Exception: log.exception("Unexpected error while cleaning up redis history") - async def _clear_terminated_sessions() -> None: - async with connect_database(bootstrap_config.db, isolation_level="AUTOCOMMIT") as db: - async with db.begin() as conn: - log.info("Deleting old records...") - result = ( - await conn.scalars( - sa.select(SessionRow.id).where(SessionRow.terminated_at < expiration_date) - ) - ).all() - session_ids = cast(list[uuid.UUID], result) - if session_ids: - await conn.execute( - sa.delete(kernels).where(kernels.c.session_id.in_(session_ids)) - ) - await conn.execute(sa.delete(SessionRow).where(SessionRow.id.in_(session_ids))) - - curs = await conn.execute(sa.select(sa.func.count()).select_from(SessionRow)) - if ret := curs.fetchone(): - table_size = ret[0] - log.info( - "The number of rows of the `sessions` tables after cleanup: {}", - table_size, - ) - log.info( - "Cleaned up {:,} database records older than {}.", - len(session_ids), - expiration_date, - ) - - async def _clear_old_error_logs() -> None: - async with connect_database(bootstrap_config.db, isolation_level="AUTOCOMMIT") as db: - async with db.begin() as conn: - log.info("Deleting old error logs...") - result = await conn.execute( - sa.delete(error_logs).where(error_logs.c.created_at < expiration_date), - ) - deleted_count = result.rowcount - + async def _force_retention_sweep() -> None: + log.info("Triggering an immediate DB record retention sweep...") + async with ( + connect_database(bootstrap_config.db) as db, + config_provider_ctx(cli_ctx) as config_provider, + ): + repository = RetentionRepository(DBOpsProvider(db), config_provider) + results = await repository.sweep() + total_deleted = sum(result.deleted_count for result in results) log.info( - "Cleaned up {:,} error log records older than {}.", - deleted_count, - expiration_date, + "Retention sweep deleted {:,} record(s) across {} categor(ies).", + total_deleted, + len(results), ) asyncio.run(_clear_redis_history()) - asyncio.run(_clear_terminated_sessions()) - asyncio.run(_clear_old_error_logs()) + asyncio.run(_force_retention_sweep()) asyncio.run(vacuum_db(bootstrap_config, vacuum_full)) diff --git a/src/ai/backend/manager/cli/context.py b/src/ai/backend/manager/cli/context.py index e407341fc84..6ebeaad2f08 100644 --- a/src/ai/backend/manager/cli/context.py +++ b/src/ai/backend/manager/cli/context.py @@ -14,6 +14,7 @@ from ai.backend.logging import AbstractLogger from ai.backend.logging.types import LogLevel from ai.backend.manager.config.bootstrap import BootstrapConfig + from ai.backend.manager.config.provider import ManagerConfigProvider from ai.backend.manager.config.unified import ManagerUnifiedConfig from .types import RedisConnectionSet @@ -117,6 +118,61 @@ async def config_ctx(cli_ctx: CLIContext) -> AsyncIterator[ManagerUnifiedConfig] yield unified_config +@contextlib.asynccontextmanager +async def config_provider_ctx(cli_ctx: CLIContext) -> AsyncIterator[ManagerConfigProvider]: + """Build a full ManagerConfigProvider for one-shot CLI commands. + + Mirrors ``dependencies.config.provider.ConfigProviderDependency`` so CLI + entry points (e.g. ``clear-history``) can reuse repositories that depend on + the provider. The etcd watcher is terminated on exit. + """ + from ai.backend.common.configs.loader import ( + ConfigOverrider, + EtcdConfigLoader, + EtcdConfigWatcher, + LoaderChain, + TomlConfigLoader, + ) + from ai.backend.common.etcd import AsyncEtcd + from ai.backend.logging.types import LogLevel + from ai.backend.manager.config.loader.legacy_etcd_loader import ( + LegacyEtcdLoader, + LegacyEtcdVolumesLoader, + ) + from ai.backend.manager.config.provider import ManagerConfigProvider + + bootstrap_config = await cli_ctx.get_bootstrap_config() + etcd_config_data = bootstrap_config.etcd.to_dataclass() + async with AsyncEtcd.create_from_config(etcd_config_data) as etcd: + loaders: list[Any] = [] + if cli_ctx.config_path is not None: + loaders.append(TomlConfigLoader(cli_ctx.config_path, "manager")) + legacy_etcd_loader = LegacyEtcdLoader(etcd) + loaders.append(legacy_etcd_loader) + loaders.append(LegacyEtcdVolumesLoader(etcd)) + loaders.append(EtcdConfigLoader(etcd, prefix="ai/backend/config/common")) + loaders.append(EtcdConfigLoader(etcd, prefix="ai/backend/config/manager")) + overrides: list[tuple[tuple[str, ...], Any]] = [ + (("debug", "enabled"), cli_ctx.log_level == LogLevel.DEBUG), + ] + if cli_ctx.log_level != LogLevel.NOTSET: + overrides += [ + (("logging", "level"), cli_ctx.log_level), + (("logging", "pkg-ns", "ai.backend"), cli_ctx.log_level), + ] + loaders.append(ConfigOverrider(overrides)) + + config_provider = await ManagerConfigProvider.create( + LoaderChain(loaders), + EtcdConfigWatcher(etcd), + legacy_etcd_loader, + ) + try: + yield config_provider + finally: + await config_provider.terminate() + + @contextlib.asynccontextmanager async def redis_ctx(cli_ctx: CLIContext) -> AsyncIterator[RedisConnectionSet]: from ai.backend.common.clients.valkey_client.valkey_image.client import ValkeyImageClient diff --git a/src/ai/backend/manager/data/kernel/types.py b/src/ai/backend/manager/data/kernel/types.py index 429d4c04f51..0b4bc1d1d91 100644 --- a/src/ai/backend/manager/data/kernel/types.py +++ b/src/ai/backend/manager/data/kernel/types.py @@ -8,6 +8,7 @@ from typing import TYPE_CHECKING, Any from uuid import UUID +from ai.backend.common.identifier.kernel_scheduling_history import KernelSchedulingHistoryID from ai.backend.common.identifier.resource_group import ResourceGroupID from ai.backend.common.types import ( CIStrEnum, @@ -332,7 +333,7 @@ class KernelSchedulingPhase(StrEnum): class KernelSchedulingHistoryData: """Domain model for kernel scheduling history.""" - id: UUID + id: KernelSchedulingHistoryID kernel_id: KernelId session_id: SessionId diff --git a/src/ai/backend/manager/defs.py b/src/ai/backend/manager/defs.py index 286e83f882b..b68cacf2801 100644 --- a/src/ai/backend/manager/defs.py +++ b/src/ai/backend/manager/defs.py @@ -88,7 +88,6 @@ class LockID(enum.IntEnum): LOCKID_CHECK_PRECOND_TIMER = 192 LOCKID_START_TIMER = 198 LOCKID_SCALE_TIMER = 193 - LOCKID_LOG_CLEANUP_TIMER = 195 LOCKID_IDLE_CHECK_TIMER = 196 LOCKID_SESSION_STATUS_UPDATE_TIMER = 197 # Sokovan timers for each ScheduleType diff --git a/src/ai/backend/manager/dependencies/composer.py b/src/ai/backend/manager/dependencies/composer.py index 686a15d4e8d..9720edc9551 100644 --- a/src/ai/backend/manager/dependencies/composer.py +++ b/src/ai/backend/manager/dependencies/composer.py @@ -343,8 +343,6 @@ async def compose( registry_quota_service=domain.services_ctx.per_project_container_registries_quota, # BgtaskRegistry creation agent_client_pool=agents.agent_client_pool, - # Log cleanup timer - distributed_lock_factory=domain.distributed_lock_factory, # Lifecycle background tasks stats_monitor=monitoring.stats_monitor, pidx=setup_input.pidx, diff --git a/src/ai/backend/manager/dependencies/processing/composer.py b/src/ai/backend/manager/dependencies/processing/composer.py index 63b1cf3823b..5b3dd6c5eec 100644 --- a/src/ai/backend/manager/dependencies/processing/composer.py +++ b/src/ai/backend/manager/dependencies/processing/composer.py @@ -68,12 +68,11 @@ from ai.backend.manager.sokovan.reconciler.coordinator import ReconcilerCoordinator from ai.backend.manager.sokovan.scheduler.coordinator import ScheduleCoordinator from ai.backend.manager.sokovan.scheduling_controller import SchedulingController -from ai.backend.manager.types import DistributedLockFactory, SMTPTriggerPolicy +from ai.backend.manager.types import SMTPTriggerPolicy from .agent_lost_checker import AgentLostCheckerDependency, AgentLostCheckerInput from .bgtask_registry import BgtaskRegistryDependency, BgtaskRegistryInput from .event_dispatcher import EventDispatcherDependency, EventDispatcherInput -from .log_cleanup_timer import LogCleanupTimerDependency, LogCleanupTimerInput from .manager_status_watcher import ManagerStatusWatcherDependency, ManagerStatusWatcherInput from .processors import ProcessorsDependency, ProcessorsProviderInput from .stats_reporter import StatsReporterDependency, StatsReporterInput @@ -131,9 +130,6 @@ class ProcessingInput: # BgtaskRegistry creation (additional) agent_client_pool: AgentClientPool - # Log cleanup timer - distributed_lock_factory: DistributedLockFactory - # Lifecycle background tasks stats_monitor: StatsPluginContext pidx: int @@ -318,15 +314,6 @@ async def compose( dispatchers.dispatch(event_dispatcher) await event_dispatcher.start() - # Step 3.5: Create and start log cleanup timer - await stack.enter_dependency( - LogCleanupTimerDependency(), - LogCleanupTimerInput( - distributed_lock_factory=setup_input.distributed_lock_factory, - event_producer=setup_input.event_producer, - ), - ) - # Step 4: Create BgtaskRegistry await stack.enter_dependency( BgtaskRegistryDependency(), diff --git a/src/ai/backend/manager/dependencies/processing/log_cleanup_timer.py b/src/ai/backend/manager/dependencies/processing/log_cleanup_timer.py deleted file mode 100644 index 7ef11ce71b7..00000000000 --- a/src/ai/backend/manager/dependencies/processing/log_cleanup_timer.py +++ /dev/null @@ -1,53 +0,0 @@ -from __future__ import annotations - -from collections.abc import AsyncIterator -from contextlib import asynccontextmanager -from dataclasses import dataclass -from typing import override - -from ai.backend.common.dependencies import NonMonitorableDependencyProvider -from ai.backend.common.distributed import GlobalTimer -from ai.backend.common.events.dispatcher import EventProducer -from ai.backend.common.events.event_types.log.anycast import DoLogCleanupEvent -from ai.backend.manager.defs import LockID -from ai.backend.manager.types import DistributedLockFactory - - -@dataclass -class LogCleanupTimerInput: - """Input required for log cleanup timer setup.""" - - distributed_lock_factory: DistributedLockFactory - event_producer: EventProducer - - -class LogCleanupTimerDependency( - NonMonitorableDependencyProvider[LogCleanupTimerInput, GlobalTimer] -): - """Provides GlobalTimer lifecycle for periodic log cleanup. - - Creates a GlobalTimer that periodically fires DoLogCleanupEvent, - which is consumed by the LogCleanupEventHandler registered in Dispatchers. - """ - - @property - @override - def stage_name(self) -> str: - return "log-cleanup-timer" - - @asynccontextmanager - @override - async def provide(self, setup_input: LogCleanupTimerInput) -> AsyncIterator[GlobalTimer]: - timer = GlobalTimer( - setup_input.distributed_lock_factory(LockID.LOCKID_LOG_CLEANUP_TIMER, 20.0), - setup_input.event_producer, - lambda: DoLogCleanupEvent(), - 20.0, - initial_delay=17.0, - task_name="log_cleanup_task", - ) - await timer.join() - try: - yield timer - finally: - await timer.leave() diff --git a/src/ai/backend/manager/errors/kernel.py b/src/ai/backend/manager/errors/kernel.py index 7770116ef3c..b0db886246a 100644 --- a/src/ai/backend/manager/errors/kernel.py +++ b/src/ai/backend/manager/errors/kernel.py @@ -51,6 +51,19 @@ def error_code(self) -> ErrorCode: ) +class EmptyKernelSchedulingHistoryScope(BackendAIError, web.HTTPBadRequest): + error_type = "https://api.backend.ai/probs/empty-kernel-scheduling-history-scope" + error_title = "Kernel scheduling history scope requires session_id or kernel_id." + + @override + def error_code(self) -> ErrorCode: + return ErrorCode( + domain=ErrorDomain.KERNEL, + operation=ErrorOperation.READ, + error_detail=ErrorDetail.INVALID_PARAMETERS, + ) + + class SessionNotFound(ObjectNotFound): object_name = "session" diff --git a/src/ai/backend/manager/event_dispatcher/dispatch.py b/src/ai/backend/manager/event_dispatcher/dispatch.py index 134f8607d37..10e62119e03 100644 --- a/src/ai/backend/manager/event_dispatcher/dispatch.py +++ b/src/ai/backend/manager/event_dispatcher/dispatch.py @@ -62,7 +62,6 @@ KernelTerminatedBroadcastEvent, KernelTerminatingBroadcastEvent, ) -from ai.backend.common.events.event_types.log.anycast import DoLogCleanupEvent from ai.backend.common.events.event_types.notification.anycast import ( NotificationTriggeredEvent, ) @@ -137,7 +136,6 @@ from .handlers.idle_check import IdleCheckEventHandler from .handlers.image import ImageEventHandler from .handlers.kernel import KernelEventHandler -from .handlers.log_cleanup import LogCleanupEventHandler from .handlers.notification import NotificationEventHandler from .handlers.service_catalog import ServiceCatalogEventHandler from .handlers.session import SessionEventHandler @@ -185,7 +183,6 @@ class Dispatchers: _artifact_event_handler: ArtifactEventHandler _artifact_registry_event_handler: ArtifactRegistryEventHandler _service_catalog_event_handler: ServiceCatalogEventHandler - _log_cleanup_event_handler: LogCleanupEventHandler stream_cleanup_handler: StreamCleanupEventHandler def __init__(self, args: DispatcherArgs) -> None: @@ -247,7 +244,6 @@ def __init__(self, args: DispatcherArgs) -> None: args.config_provider, ) self._service_catalog_event_handler = ServiceCatalogEventHandler(args.db) - self._log_cleanup_event_handler = LogCleanupEventHandler(args.etcd, args.db) self.stream_cleanup_handler = StreamCleanupEventHandler(args.db) def dispatch(self, event_dispatcher: EventDispatcher) -> None: @@ -266,7 +262,6 @@ def dispatch(self, event_dispatcher: EventDispatcher) -> None: self._dispatch_artifact_events(event_dispatcher) self._dispatch_artifact_registry_events(event_dispatcher) self._dispatch_service_catalog_events(event_dispatcher) - self._dispatch_log_cleanup_events(event_dispatcher) self._dispatch_session_broadcast_propagation(event_dispatcher) self._dispatch_stream_cleanup_events(event_dispatcher) @@ -648,17 +643,6 @@ def _dispatch_service_catalog_events( name="service-catalog.sweep", ) - def _dispatch_log_cleanup_events( - self, - event_dispatcher: EventDispatcher, - ) -> None: - event_dispatcher.consume( - DoLogCleanupEvent, - None, - self._log_cleanup_event_handler.handle_log_cleanup, - name="log_cleanup", - ) - def _dispatch_session_broadcast_propagation( self, event_dispatcher: EventDispatcher, diff --git a/src/ai/backend/manager/event_dispatcher/handlers/log_cleanup.py b/src/ai/backend/manager/event_dispatcher/handlers/log_cleanup.py deleted file mode 100644 index 12490fc47c4..00000000000 --- a/src/ai/backend/manager/event_dispatcher/handlers/log_cleanup.py +++ /dev/null @@ -1,55 +0,0 @@ -import datetime as dt -import logging -from datetime import UTC, datetime - -import sqlalchemy as sa -from dateutil.relativedelta import relativedelta - -from ai.backend.common import validators as tx -from ai.backend.common.etcd import AsyncEtcd -from ai.backend.common.events.event_types.log.anycast import DoLogCleanupEvent -from ai.backend.common.types import AgentId -from ai.backend.logging import BraceStyleAdapter -from ai.backend.manager.models.error_logs import error_logs -from ai.backend.manager.models.utils import ExtendedAsyncSAEngine - -log = BraceStyleAdapter(logging.getLogger(__spec__.name)) - - -class LogCleanupEventHandler: - _etcd: AsyncEtcd - _db: ExtendedAsyncSAEngine - - def __init__( - self, - etcd: AsyncEtcd, - db: ExtendedAsyncSAEngine, - ) -> None: - self._etcd = etcd - self._db = db - - async def handle_log_cleanup( - self, - _context: None, - _source: AgentId, - _event: DoLogCleanupEvent, - ) -> None: - raw_lifetime = await self._etcd.get("config/logs/error/retention") - if raw_lifetime is None: - raw_lifetime = "90d" - lifetime: dt.timedelta | relativedelta - try: - lifetime = tx.TimeDuration().check(raw_lifetime) - except ValueError: - lifetime = dt.timedelta(days=90) - log.warning( - "Failed to parse the error log retention period ({}) read from etcd; " - "falling back to 90 days", - raw_lifetime, - ) - boundary = datetime.now(UTC) - lifetime - async with self._db.begin() as conn: - query = sa.delete(error_logs).where(error_logs.c.created_at < boundary) - result = await conn.execute(query) - if result.rowcount > 0: - log.info("Cleaned up {} log(s) filed before {}", result.rowcount, boundary) diff --git a/src/ai/backend/manager/models/alembic/versions/84d5c6daf8cc_drop_legacy_app_configs_table.py b/src/ai/backend/manager/models/alembic/versions/84d5c6daf8cc_drop_legacy_app_configs_table.py index eb8a26fa72e..13e4e665c3a 100644 --- a/src/ai/backend/manager/models/alembic/versions/84d5c6daf8cc_drop_legacy_app_configs_table.py +++ b/src/ai/backend/manager/models/alembic/versions/84d5c6daf8cc_drop_legacy_app_configs_table.py @@ -36,7 +36,6 @@ def downgrade() -> None: # Recreate the predecessor table to allow `alembic downgrade` to # complete cleanly. Existing-row restoration is not attempted. app_config_scope_type = sa.Enum("DOMAIN", "PROJECT", "USER", name="app_config_scope_type") - app_config_scope_type.create(op.get_bind(), checkfirst=True) op.create_table( "app_configs", IDColumn(), diff --git a/src/ai/backend/manager/models/scheduling_history/conditions.py b/src/ai/backend/manager/models/scheduling_history/conditions.py index 6ea430ac8ce..89a71af8e62 100644 --- a/src/ai/backend/manager/models/scheduling_history/conditions.py +++ b/src/ai/backend/manager/models/scheduling_history/conditions.py @@ -10,9 +10,9 @@ import sqlalchemy as sa from ai.backend.common.data.filter_specs import StringMatchSpec, UUIDEqualMatchSpec, UUIDInMatchSpec +from ai.backend.common.identifier.kernel_scheduling_history import KernelSchedulingHistoryID from ai.backend.common.types import KernelId, SessionId from ai.backend.manager.data.deployment.types import RouteStatus -from ai.backend.manager.data.kernel.types import KernelSchedulingPhase from ai.backend.manager.data.session.types import SchedulingResult, SessionStatus from ai.backend.manager.models.clauses import QueryCondition from ai.backend.manager.models.condition_utils import make_string_in_factory @@ -400,6 +400,32 @@ def inner() -> sa.sql.expression.ColumnElement[bool]: class KernelSchedulingHistoryConditions: """Query conditions for kernel scheduling history.""" + # UUID filter conditions for history id + @staticmethod + def by_id_filter(spec: UUIDEqualMatchSpec) -> QueryCondition: + def inner() -> sa.sql.expression.ColumnElement[bool]: + if spec.negated: + return KernelSchedulingHistoryRow.id != spec.value + return KernelSchedulingHistoryRow.id == spec.value + + return inner + + @staticmethod + def by_id_in(spec: UUIDInMatchSpec) -> QueryCondition: + def inner() -> sa.sql.expression.ColumnElement[bool]: + if spec.negated: + return KernelSchedulingHistoryRow.id.notin_(spec.values) + return KernelSchedulingHistoryRow.id.in_(spec.values) + + return inner + + @staticmethod + def by_ids(ids: Collection[KernelSchedulingHistoryID]) -> QueryCondition: + def inner() -> sa.sql.expression.ColumnElement[bool]: + return KernelSchedulingHistoryRow.id.in_(ids) + + return inner + @staticmethod def by_kernel_id(kernel_id: KernelId) -> QueryCondition: def inner() -> sa.sql.expression.ColumnElement[bool]: @@ -421,86 +447,243 @@ def inner() -> sa.sql.expression.ColumnElement[bool]: return inner + @staticmethod + def by_results(results: list[SchedulingResult]) -> QueryCondition: + def inner() -> sa.sql.expression.ColumnElement[bool]: + return KernelSchedulingHistoryRow.result.in_([str(r) for r in results]) + + return inner + + @staticmethod + def by_result_not_equals(result: SchedulingResult) -> QueryCondition: + def inner() -> sa.sql.expression.ColumnElement[bool]: + return KernelSchedulingHistoryRow.result != str(result) + + return inner + + @staticmethod + def by_result_not_in(results: list[SchedulingResult]) -> QueryCondition: + def inner() -> sa.sql.expression.ColumnElement[bool]: + return KernelSchedulingHistoryRow.result.not_in([str(r) for r in results]) + + return inner + + @staticmethod + def by_from_statuses(statuses: list[str]) -> QueryCondition: + def inner() -> sa.sql.expression.ColumnElement[bool]: + return KernelSchedulingHistoryRow.from_status.in_(statuses) + + return inner + + @staticmethod + def by_to_statuses(statuses: list[str]) -> QueryCondition: + def inner() -> sa.sql.expression.ColumnElement[bool]: + return KernelSchedulingHistoryRow.to_status.in_(statuses) + + return inner + + # UUID filter conditions for kernel_id + @staticmethod + def by_kernel_id_filter(spec: UUIDEqualMatchSpec) -> QueryCondition: + def inner() -> sa.sql.expression.ColumnElement[bool]: + if spec.negated: + return KernelSchedulingHistoryRow.kernel_id != spec.value + return KernelSchedulingHistoryRow.kernel_id == spec.value + + return inner + + @staticmethod + def by_kernel_id_in(spec: UUIDInMatchSpec) -> QueryCondition: + def inner() -> sa.sql.expression.ColumnElement[bool]: + if spec.negated: + return KernelSchedulingHistoryRow.kernel_id.notin_(spec.values) + return KernelSchedulingHistoryRow.kernel_id.in_(spec.values) + + return inner + + # UUID filter conditions for session_id + @staticmethod + def by_session_id_filter(spec: UUIDEqualMatchSpec) -> QueryCondition: + def inner() -> sa.sql.expression.ColumnElement[bool]: + if spec.negated: + return KernelSchedulingHistoryRow.session_id != spec.value + return KernelSchedulingHistoryRow.session_id == spec.value + + return inner + + @staticmethod + def by_session_id_in(spec: UUIDInMatchSpec) -> QueryCondition: + def inner() -> sa.sql.expression.ColumnElement[bool]: + if spec.negated: + return KernelSchedulingHistoryRow.session_id.notin_(spec.values) + return KernelSchedulingHistoryRow.session_id.in_(spec.values) + + return inner + + # String filter conditions for error_code + @staticmethod + def by_error_code_contains(spec: StringMatchSpec) -> QueryCondition: + def inner() -> sa.sql.expression.ColumnElement[bool]: + if spec.case_insensitive: + condition = KernelSchedulingHistoryRow.error_code.ilike(f"%{spec.value}%") + else: + condition = KernelSchedulingHistoryRow.error_code.like(f"%{spec.value}%") + if spec.negated: + condition = sa.not_(condition) + return condition + + return inner + + @staticmethod + def by_error_code_equals(spec: StringMatchSpec) -> QueryCondition: + def inner() -> sa.sql.expression.ColumnElement[bool]: + if spec.case_insensitive: + condition = ( + sa.func.lower(KernelSchedulingHistoryRow.error_code) == spec.value.lower() + ) + else: + condition = KernelSchedulingHistoryRow.error_code == spec.value + if spec.negated: + condition = sa.not_(condition) + return condition + + return inner + + @staticmethod + def by_error_code_starts_with(spec: StringMatchSpec) -> QueryCondition: + def inner() -> sa.sql.expression.ColumnElement[bool]: + if spec.case_insensitive: + condition = KernelSchedulingHistoryRow.error_code.ilike(f"{spec.value}%") + else: + condition = KernelSchedulingHistoryRow.error_code.like(f"{spec.value}%") + if spec.negated: + condition = sa.not_(condition) + return condition + + return inner + + @staticmethod + def by_error_code_ends_with(spec: StringMatchSpec) -> QueryCondition: + def inner() -> sa.sql.expression.ColumnElement[bool]: + if spec.case_insensitive: + condition = KernelSchedulingHistoryRow.error_code.ilike(f"%{spec.value}") + else: + condition = KernelSchedulingHistoryRow.error_code.like(f"%{spec.value}") + if spec.negated: + condition = sa.not_(condition) + return condition + + return inner + # String filter conditions for phase @staticmethod def by_phase_contains(spec: StringMatchSpec) -> QueryCondition: def inner() -> sa.sql.expression.ColumnElement[bool]: - col = cast(sa.ColumnElement[str], KernelSchedulingHistoryRow.phase) if spec.case_insensitive: - col = sa.func.lower(col) - pattern = f"%{spec.value.lower()}%" + condition = KernelSchedulingHistoryRow.phase.ilike(f"%{spec.value}%") else: - pattern = f"%{spec.value}%" - expr = col.like(pattern) - return ~expr if spec.negated else expr + condition = KernelSchedulingHistoryRow.phase.like(f"%{spec.value}%") + if spec.negated: + condition = sa.not_(condition) + return condition return inner @staticmethod def by_phase_equals(spec: StringMatchSpec) -> QueryCondition: def inner() -> sa.sql.expression.ColumnElement[bool]: - col = cast(sa.ColumnElement[str], KernelSchedulingHistoryRow.phase) if spec.case_insensitive: - col = sa.func.lower(col) - val = spec.value.lower() + condition = sa.func.lower(KernelSchedulingHistoryRow.phase) == spec.value.lower() else: - val = spec.value + condition = KernelSchedulingHistoryRow.phase == spec.value if spec.negated: - return col != val - return col == val + condition = sa.not_(condition) + return condition return inner @staticmethod def by_phase_starts_with(spec: StringMatchSpec) -> QueryCondition: def inner() -> sa.sql.expression.ColumnElement[bool]: - col = cast(sa.ColumnElement[str], KernelSchedulingHistoryRow.phase) if spec.case_insensitive: - col = sa.func.lower(col) - pattern = f"{spec.value.lower()}%" + condition = KernelSchedulingHistoryRow.phase.ilike(f"{spec.value}%") else: - pattern = f"{spec.value}%" - expr = col.like(pattern) - return ~expr if spec.negated else expr + condition = KernelSchedulingHistoryRow.phase.like(f"{spec.value}%") + if spec.negated: + condition = sa.not_(condition) + return condition return inner @staticmethod def by_phase_ends_with(spec: StringMatchSpec) -> QueryCondition: def inner() -> sa.sql.expression.ColumnElement[bool]: - col = cast(sa.ColumnElement[str], KernelSchedulingHistoryRow.phase) if spec.case_insensitive: - col = sa.func.lower(col) - pattern = f"%{spec.value.lower()}" + condition = KernelSchedulingHistoryRow.phase.ilike(f"%{spec.value}") else: - pattern = f"%{spec.value}" - expr = col.like(pattern) - return ~expr if spec.negated else expr + condition = KernelSchedulingHistoryRow.phase.like(f"%{spec.value}") + if spec.negated: + condition = sa.not_(condition) + return condition return inner + # String filter conditions for message @staticmethod - def by_from_status(phase: KernelSchedulingPhase) -> QueryCondition: + def by_message_contains(spec: StringMatchSpec) -> QueryCondition: def inner() -> sa.sql.expression.ColumnElement[bool]: - return KernelSchedulingHistoryRow.from_status == str(phase) + if spec.case_insensitive: + condition = KernelSchedulingHistoryRow.message.ilike(f"%{spec.value}%") + else: + condition = KernelSchedulingHistoryRow.message.like(f"%{spec.value}%") + if spec.negated: + condition = sa.not_(condition) + return condition return inner @staticmethod - def by_to_status(phase: KernelSchedulingPhase) -> QueryCondition: + def by_message_equals(spec: StringMatchSpec) -> QueryCondition: def inner() -> sa.sql.expression.ColumnElement[bool]: - return KernelSchedulingHistoryRow.to_status == str(phase) + if spec.case_insensitive: + condition = sa.func.lower(KernelSchedulingHistoryRow.message) == spec.value.lower() + else: + condition = KernelSchedulingHistoryRow.message == spec.value + if spec.negated: + condition = sa.not_(condition) + return condition return inner @staticmethod - def by_error_code(error_code: str) -> QueryCondition: + def by_message_starts_with(spec: StringMatchSpec) -> QueryCondition: def inner() -> sa.sql.expression.ColumnElement[bool]: - return KernelSchedulingHistoryRow.error_code == error_code + if spec.case_insensitive: + condition = KernelSchedulingHistoryRow.message.ilike(f"{spec.value}%") + else: + condition = KernelSchedulingHistoryRow.message.like(f"{spec.value}%") + if spec.negated: + condition = sa.not_(condition) + return condition + + return inner + + @staticmethod + def by_message_ends_with(spec: StringMatchSpec) -> QueryCondition: + def inner() -> sa.sql.expression.ColumnElement[bool]: + if spec.case_insensitive: + condition = KernelSchedulingHistoryRow.message.ilike(f"%{spec.value}") + else: + condition = KernelSchedulingHistoryRow.message.like(f"%{spec.value}") + if spec.negated: + condition = sa.not_(condition) + return condition return inner by_phase_in = staticmethod(make_string_in_factory(KernelSchedulingHistoryRow.phase)) + by_error_code_in = staticmethod(make_string_in_factory(KernelSchedulingHistoryRow.error_code)) + by_message_in = staticmethod(make_string_in_factory(KernelSchedulingHistoryRow.message)) @staticmethod def by_cursor_forward(cursor_id: str) -> QueryCondition: @@ -532,6 +715,49 @@ def inner() -> sa.sql.expression.ColumnElement[bool]: return inner + # DateTime filter conditions + @staticmethod + def by_created_at_before(dt: datetime) -> QueryCondition: + def inner() -> sa.sql.expression.ColumnElement[bool]: + return KernelSchedulingHistoryRow.created_at < dt + + return inner + + @staticmethod + def by_created_at_after(dt: datetime) -> QueryCondition: + def inner() -> sa.sql.expression.ColumnElement[bool]: + return KernelSchedulingHistoryRow.created_at > dt + + return inner + + @staticmethod + def by_created_at_equals(dt: datetime) -> QueryCondition: + def inner() -> sa.sql.expression.ColumnElement[bool]: + return KernelSchedulingHistoryRow.created_at == dt + + return inner + + @staticmethod + def by_updated_at_before(dt: datetime) -> QueryCondition: + def inner() -> sa.sql.expression.ColumnElement[bool]: + return KernelSchedulingHistoryRow.updated_at < dt + + return inner + + @staticmethod + def by_updated_at_after(dt: datetime) -> QueryCondition: + def inner() -> sa.sql.expression.ColumnElement[bool]: + return KernelSchedulingHistoryRow.updated_at > dt + + return inner + + @staticmethod + def by_updated_at_equals(dt: datetime) -> QueryCondition: + def inner() -> sa.sql.expression.ColumnElement[bool]: + return KernelSchedulingHistoryRow.updated_at == dt + + return inner + class DeploymentHistoryConditions: """Query conditions for deployment history.""" diff --git a/src/ai/backend/manager/models/scheduling_history/orders.py b/src/ai/backend/manager/models/scheduling_history/orders.py index 7af40f8a6cc..beba5b5ad96 100644 --- a/src/ai/backend/manager/models/scheduling_history/orders.py +++ b/src/ai/backend/manager/models/scheduling_history/orders.py @@ -9,6 +9,7 @@ from ai.backend.common.dto.manager.v2.scheduling_history.types import ( DeploymentHistoryOrderField, + KernelHistoryOrderField, OrderDirection, RouteHistoryOrderField, SessionHistoryOrderField, @@ -63,6 +64,15 @@ def resolve_session_order(field: SessionHistoryOrderField, direction: OrderDirec # ========== Kernel Scheduling History ========== +KERNEL_ORDER_FIELD_MAP: dict[KernelHistoryOrderField, _OrderColumn] = { + KernelHistoryOrderField.CREATED_AT: KernelSchedulingHistoryRow.created_at, + KernelHistoryOrderField.UPDATED_AT: KernelSchedulingHistoryRow.updated_at, +} + +KERNEL_DEFAULT_FORWARD_ORDER: QueryOrder = KernelSchedulingHistoryRow.created_at.desc() +KERNEL_DEFAULT_BACKWARD_ORDER: QueryOrder = KernelSchedulingHistoryRow.created_at.asc() +KERNEL_TIEBREAKER_ORDER: QueryOrder = KernelSchedulingHistoryRow.id.asc() + class KernelSchedulingHistoryOrders: """Order factories used by GQL KernelSchedulingHistoryOrderBy.to_query_order().""" @@ -80,6 +90,14 @@ def updated_at(ascending: bool = True) -> QueryOrder: return KernelSchedulingHistoryRow.updated_at.desc() +def resolve_kernel_order(field: KernelHistoryOrderField, direction: OrderDirection) -> QueryOrder: + """Resolve a DTO order field + direction to a SQLAlchemy order expression.""" + col = KERNEL_ORDER_FIELD_MAP[field] + if direction == OrderDirection.DESC: + return col.desc() + return col.asc() + + # ========== Deployment History ========== DEPLOYMENT_ORDER_FIELD_MAP: dict[DeploymentHistoryOrderField, _OrderColumn] = { diff --git a/src/ai/backend/manager/models/scheduling_history/row.py b/src/ai/backend/manager/models/scheduling_history/row.py index 09bd2aa902d..8a3b2ffb424 100644 --- a/src/ai/backend/manager/models/scheduling_history/row.py +++ b/src/ai/backend/manager/models/scheduling_history/row.py @@ -7,6 +7,7 @@ from sqlalchemy.orm import Mapped, mapped_column from ai.backend.common.data.model_deployment.types import ModelDeploymentStatus +from ai.backend.common.identifier.kernel_scheduling_history import KernelSchedulingHistoryID from ai.backend.common.identifier.replica import ReplicaID from ai.backend.common.types import KernelId, SessionId from ai.backend.manager.data.deployment.types import ( @@ -110,8 +111,11 @@ def to_data(self) -> SessionSchedulingHistoryData: class KernelSchedulingHistoryRow(Base): # type: ignore[misc] __tablename__ = "kernel_scheduling_history" - id: Mapped[uuid.UUID] = mapped_column( - "id", GUID, primary_key=True, server_default=sa.text("uuid_generate_v4()") + id: Mapped[KernelSchedulingHistoryID] = mapped_column( + "id", + GUID(KernelSchedulingHistoryID), + primary_key=True, + server_default=sa.text("uuid_generate_v4()"), ) kernel_id: Mapped[uuid.UUID] = mapped_column("kernel_id", GUID, nullable=False, index=True) session_id: Mapped[uuid.UUID] = mapped_column("session_id", GUID, nullable=False, index=True) diff --git a/src/ai/backend/manager/repositories/scheduling_history/__init__.py b/src/ai/backend/manager/repositories/scheduling_history/__init__.py index fe14354a5c3..a6063fea478 100644 --- a/src/ai/backend/manager/repositories/scheduling_history/__init__.py +++ b/src/ai/backend/manager/repositories/scheduling_history/__init__.py @@ -8,6 +8,7 @@ from .repository import SchedulingHistoryRepository from .types import ( DeploymentHistorySearchScope, + KernelSchedulingHistorySearchScope, RouteHistorySearchScope, SessionSchedulingHistorySearchScope, ) @@ -16,6 +17,7 @@ "DeploymentHistoryCreatorSpec", "DeploymentHistorySearchScope", "KernelSchedulingHistoryCreatorSpec", + "KernelSchedulingHistorySearchScope", "RouteHistoryCreatorSpec", "RouteHistorySearchScope", "SchedulingHistoryRepositories", diff --git a/src/ai/backend/manager/repositories/scheduling_history/db_source/db_source.py b/src/ai/backend/manager/repositories/scheduling_history/db_source/db_source.py index 8e2ed5aa254..a4224356250 100644 --- a/src/ai/backend/manager/repositories/scheduling_history/db_source/db_source.py +++ b/src/ai/backend/manager/repositories/scheduling_history/db_source/db_source.py @@ -28,6 +28,7 @@ ) from ai.backend.manager.repositories.scheduling_history.types import ( DeploymentHistorySearchScope, + KernelSchedulingHistorySearchScope, RouteHistorySearchScope, SessionSchedulingHistorySearchScope, ) @@ -94,7 +95,7 @@ async def search_session_scoped_history( has_previous_page=result.has_previous_page, ) - # ========== Kernel History ========== + # ========== Kernel History (Admin) ========== async def search_kernel_history( self, @@ -119,6 +120,28 @@ async def search_kernel_history( has_previous_page=result.has_previous_page, ) + # ========== Kernel History (Scoped) ========== + + async def search_kernel_scoped_history( + self, + querier: BatchQuerier, + scope: KernelSchedulingHistorySearchScope, + ) -> KernelSchedulingHistoryListResult: + """Search kernel scheduling history within scope.""" + async with self._db.begin_readonly_session() as db_sess: + query = sa.select(KernelSchedulingHistoryRow) + + result = await execute_batch_querier(db_sess, query, querier, scopes=[scope]) + + items = [row.KernelSchedulingHistoryRow.to_data() for row in result.rows] + + return KernelSchedulingHistoryListResult( + items=items, + total_count=result.total_count, + has_next_page=result.has_next_page, + has_previous_page=result.has_previous_page, + ) + # ========== Deployment History (Admin) ========== async def search_deployment_history( diff --git a/src/ai/backend/manager/repositories/scheduling_history/repository.py b/src/ai/backend/manager/repositories/scheduling_history/repository.py index ed6308adea4..a97440c1482 100644 --- a/src/ai/backend/manager/repositories/scheduling_history/repository.py +++ b/src/ai/backend/manager/repositories/scheduling_history/repository.py @@ -26,6 +26,7 @@ from .db_source import SchedulingHistoryDBSource from .types import ( DeploymentHistorySearchScope, + KernelSchedulingHistorySearchScope, RouteHistorySearchScope, SessionSchedulingHistorySearchScope, ) @@ -81,7 +82,7 @@ async def search_session_scoped_history( """Search session scheduling history within scope.""" return await self._db_source.search_session_scoped_history(querier, scope) - # ========== Kernel History ========== + # ========== Kernel History (Admin) ========== @scheduling_history_repository_resilience.apply() async def search_kernel_history( @@ -91,6 +92,17 @@ async def search_kernel_history( """Search kernel scheduling history with pagination.""" return await self._db_source.search_kernel_history(querier) + # ========== Kernel History (Scoped) ========== + + @scheduling_history_repository_resilience.apply() + async def search_kernel_scoped_history( + self, + querier: BatchQuerier, + scope: KernelSchedulingHistorySearchScope, + ) -> KernelSchedulingHistoryListResult: + """Search kernel scheduling history within scope.""" + return await self._db_source.search_kernel_scoped_history(querier, scope) + # ========== Deployment History (Admin) ========== @scheduling_history_repository_resilience.apply() diff --git a/src/ai/backend/manager/repositories/scheduling_history/types.py b/src/ai/backend/manager/repositories/scheduling_history/types.py index 150834e4043..b7169bd69bf 100644 --- a/src/ai/backend/manager/repositories/scheduling_history/types.py +++ b/src/ai/backend/manager/repositories/scheduling_history/types.py @@ -6,16 +6,25 @@ from typing import Any, override from uuid import UUID +import sqlalchemy as sa + from ai.backend.common.data.filter_specs import UUIDEqualMatchSpec from ai.backend.common.identifier.replica import ReplicaID +from ai.backend.common.types import KernelId, SessionId from ai.backend.manager.errors.deployment import EndpointNotFound -from ai.backend.manager.errors.kernel import SessionNotFound +from ai.backend.manager.errors.kernel import ( + EmptyKernelSchedulingHistoryScope, + KernelNotFound, + SessionNotFound, +) from ai.backend.manager.errors.service import RouteNotFound from ai.backend.manager.models.clauses import QueryCondition from ai.backend.manager.models.endpoint import EndpointRow +from ai.backend.manager.models.kernel.row import KernelRow from ai.backend.manager.models.routing import RoutingRow from ai.backend.manager.models.scheduling_history.conditions import ( DeploymentHistoryConditions, + KernelSchedulingHistoryConditions, RouteHistoryConditions, SessionSchedulingHistoryConditions, ) @@ -24,6 +33,7 @@ __all__ = ( "SessionSchedulingHistorySearchScope", + "KernelSchedulingHistorySearchScope", "DeploymentHistorySearchScope", "RouteHistorySearchScope", ) @@ -62,6 +72,73 @@ def existence_checks(self) -> list[ExistenceCheck[Any]]: ] +# Kernel Scheduling History Scope + + +@dataclass(frozen=True) +class KernelSchedulingHistorySearchScope(SearchScope): + """Scope for kernel scheduling history search. + + Either axis may be given; when both are, they intersect. At least one is required — + an empty scope would degenerate into an unscoped (admin) search. + """ + + session_id: SessionId | None = None + """Restrict to the kernels of this session.""" + + kernel_id: KernelId | None = None + """Restrict to this kernel.""" + + def __post_init__(self) -> None: + if self.session_id is None and self.kernel_id is None: + raise EmptyKernelSchedulingHistoryScope() + + @override + def to_condition(self) -> QueryCondition: + """Convert scope to a query condition for KernelSchedulingHistoryRow.""" + conditions: list[QueryCondition] = [] + if self.session_id is not None: + conditions.append( + KernelSchedulingHistoryConditions.by_session_id_filter( + UUIDEqualMatchSpec(value=self.session_id, negated=False) + ) + ) + if self.kernel_id is not None: + conditions.append( + KernelSchedulingHistoryConditions.by_kernel_id_filter( + UUIDEqualMatchSpec(value=self.kernel_id, negated=False) + ) + ) + + def inner() -> sa.sql.expression.ColumnElement[bool]: + return sa.and_(*(cond() for cond in conditions)) + + return inner + + @property + @override + def existence_checks(self) -> list[ExistenceCheck[Any]]: + """Check that each scoped entity exists.""" + checks: list[ExistenceCheck[Any]] = [] + if self.session_id is not None: + checks.append( + ExistenceCheck( + column=SessionRow.id, + value=self.session_id, + error=SessionNotFound(str(self.session_id)), + ) + ) + if self.kernel_id is not None: + checks.append( + ExistenceCheck( + column=KernelRow.id, + value=self.kernel_id, + error=KernelNotFound(str(self.kernel_id)), + ) + ) + return checks + + # Deployment History Scope diff --git a/src/ai/backend/manager/services/scheduling_history/actions/__init__.py b/src/ai/backend/manager/services/scheduling_history/actions/__init__.py index 82b26def324..5c350b977fd 100644 --- a/src/ai/backend/manager/services/scheduling_history/actions/__init__.py +++ b/src/ai/backend/manager/services/scheduling_history/actions/__init__.py @@ -9,6 +9,14 @@ SearchDeploymentScopedHistoryAction, SearchDeploymentScopedHistoryActionResult, ) +from .search_kernel_history import ( + SearchKernelHistoryAction, + SearchKernelHistoryActionResult, +) +from .search_kernel_scoped_history import ( + SearchKernelScopedHistoryAction, + SearchKernelScopedHistoryActionResult, +) from .search_route_history import ( SearchRouteHistoryAction, SearchRouteHistoryActionResult, @@ -31,6 +39,8 @@ # Admin actions "SearchSessionHistoryAction", "SearchSessionHistoryActionResult", + "SearchKernelHistoryAction", + "SearchKernelHistoryActionResult", "SearchDeploymentHistoryAction", "SearchDeploymentHistoryActionResult", "SearchRouteHistoryAction", @@ -38,6 +48,8 @@ # Scoped actions (added in 26.2.0) "SearchSessionScopedHistoryAction", "SearchSessionScopedHistoryActionResult", + "SearchKernelScopedHistoryAction", + "SearchKernelScopedHistoryActionResult", "SearchDeploymentScopedHistoryAction", "SearchDeploymentScopedHistoryActionResult", "SearchRouteScopedHistoryAction", diff --git a/src/ai/backend/manager/services/scheduling_history/actions/search_kernel_history.py b/src/ai/backend/manager/services/scheduling_history/actions/search_kernel_history.py new file mode 100644 index 00000000000..a76d8a359c6 --- /dev/null +++ b/src/ai/backend/manager/services/scheduling_history/actions/search_kernel_history.py @@ -0,0 +1,51 @@ +from __future__ import annotations + +from dataclasses import dataclass +from typing import override + +from ai.backend.common.data.permission.types import EntityType +from ai.backend.manager.actions.action import BaseActionResult +from ai.backend.manager.actions.action.global_action import BaseGlobalAction +from ai.backend.manager.actions.types import ActionOperationType +from ai.backend.manager.data.kernel.types import KernelSchedulingHistoryData +from ai.backend.manager.repositories.base import BatchQuerier + + +@dataclass +class SearchKernelHistoryAction(BaseGlobalAction): + """Action to search kernel scheduling history (admin API). + + System-wide and unscoped: authorization is the SUPERADMIN role gate rather + than RBAC scope resolution, so this runs through ``GlobalActionProcessor``. + The scoped counterpart stays on the RBAC path. + """ + + querier: BatchQuerier + + @override + @classmethod + def entity_type(cls) -> EntityType: + return EntityType.KERNEL_HISTORY + + @override + @classmethod + def operation_type(cls) -> ActionOperationType: + return ActionOperationType.SEARCH + + @override + def entity_id(self) -> str | None: + return None + + +@dataclass +class SearchKernelHistoryActionResult(BaseActionResult): + """Result of searching kernel scheduling history.""" + + histories: list[KernelSchedulingHistoryData] + total_count: int + has_next_page: bool + has_previous_page: bool + + @override + def entity_id(self) -> str | None: + return None diff --git a/src/ai/backend/manager/services/scheduling_history/actions/search_kernel_scoped_history.py b/src/ai/backend/manager/services/scheduling_history/actions/search_kernel_scoped_history.py new file mode 100644 index 00000000000..d2f181234f6 --- /dev/null +++ b/src/ai/backend/manager/services/scheduling_history/actions/search_kernel_scoped_history.py @@ -0,0 +1,58 @@ +from __future__ import annotations + +from dataclasses import dataclass +from typing import override + +from ai.backend.common.data.permission.types import EntityType +from ai.backend.manager.actions.action import BaseActionResult +from ai.backend.manager.actions.types import ActionOperationType +from ai.backend.manager.data.kernel.types import KernelSchedulingHistoryData +from ai.backend.manager.repositories.base import BatchQuerier +from ai.backend.manager.repositories.scheduling_history.types import ( + KernelSchedulingHistorySearchScope, +) + +from .base import SchedulingHistoryAction + + +@dataclass +class SearchKernelScopedHistoryAction(SchedulingHistoryAction): + """Action to search kernel scheduling history within a scope. + + This is the scoped version used by entity-scoped APIs. The scope is required and + narrows the query to one session's kernels, one kernel, or their intersection. + """ + + scope: KernelSchedulingHistorySearchScope + querier: BatchQuerier + + @override + @classmethod + def entity_type(cls) -> EntityType: + return EntityType.KERNEL_SCOPED_HISTORY + + @override + @classmethod + def operation_type(cls) -> ActionOperationType: + return ActionOperationType.SEARCH + + @override + def entity_id(self) -> str | None: + # The narrower axis identifies the query best; the scope guarantees one is set. + if self.scope.kernel_id is not None: + return str(self.scope.kernel_id) + return str(self.scope.session_id) + + +@dataclass +class SearchKernelScopedHistoryActionResult(BaseActionResult): + """Result of searching kernel scheduling history within scope.""" + + histories: list[KernelSchedulingHistoryData] + total_count: int + has_next_page: bool + has_previous_page: bool + + @override + def entity_id(self) -> str | None: + return None diff --git a/src/ai/backend/manager/services/scheduling_history/processors.py b/src/ai/backend/manager/services/scheduling_history/processors.py index 84cb133a677..210fa7b50e4 100644 --- a/src/ai/backend/manager/services/scheduling_history/processors.py +++ b/src/ai/backend/manager/services/scheduling_history/processors.py @@ -4,6 +4,7 @@ from ai.backend.manager.actions.monitors.monitor import ActionMonitor from ai.backend.manager.actions.processor import ActionProcessor +from ai.backend.manager.actions.processor.global_action import GlobalActionProcessor from ai.backend.manager.actions.types import AbstractProcessorPackage, ActionSpec from ai.backend.manager.actions.validators import ActionValidators @@ -12,6 +13,10 @@ SearchDeploymentHistoryActionResult, SearchDeploymentScopedHistoryAction, SearchDeploymentScopedHistoryActionResult, + SearchKernelHistoryAction, + SearchKernelHistoryActionResult, + SearchKernelScopedHistoryAction, + SearchKernelScopedHistoryActionResult, SearchRouteHistoryAction, SearchRouteHistoryActionResult, SearchRouteScopedHistoryAction, @@ -31,6 +36,9 @@ class SchedulingHistoryProcessors(AbstractProcessorPackage): search_session_history: ActionProcessor[ SearchSessionHistoryAction, SearchSessionHistoryActionResult ] + search_kernel_history: GlobalActionProcessor[ + SearchKernelHistoryAction, SearchKernelHistoryActionResult + ] search_deployment_history: ActionProcessor[ SearchDeploymentHistoryAction, SearchDeploymentHistoryActionResult ] @@ -40,6 +48,9 @@ class SchedulingHistoryProcessors(AbstractProcessorPackage): search_session_scoped_history: ActionProcessor[ SearchSessionScopedHistoryAction, SearchSessionScopedHistoryActionResult ] + search_kernel_scoped_history: ActionProcessor[ + SearchKernelScopedHistoryAction, SearchKernelScopedHistoryActionResult + ] search_deployment_scoped_history: ActionProcessor[ SearchDeploymentScopedHistoryAction, SearchDeploymentScopedHistoryActionResult ] @@ -57,6 +68,9 @@ def __init__( self.search_session_history = ActionProcessor( service.search_session_history, action_monitors ) + self.search_kernel_history = GlobalActionProcessor( + service.search_kernel_history, action_monitors + ) self.search_deployment_history = ActionProcessor( service.search_deployment_history, action_monitors ) @@ -66,6 +80,9 @@ def __init__( self.search_session_scoped_history = ActionProcessor( service.search_session_scoped_history, action_monitors ) + self.search_kernel_scoped_history = ActionProcessor( + service.search_kernel_scoped_history, action_monitors + ) self.search_deployment_scoped_history = ActionProcessor( service.search_deployment_scoped_history, action_monitors ) @@ -78,10 +95,12 @@ def supported_actions(self) -> list[ActionSpec]: return [ # Admin actions SearchSessionHistoryAction.spec(), + SearchKernelHistoryAction.spec(), SearchDeploymentHistoryAction.spec(), SearchRouteHistoryAction.spec(), # Scoped actions (added in 26.2.0) SearchSessionScopedHistoryAction.spec(), + SearchKernelScopedHistoryAction.spec(), SearchDeploymentScopedHistoryAction.spec(), SearchRouteScopedHistoryAction.spec(), ] diff --git a/src/ai/backend/manager/services/scheduling_history/service.py b/src/ai/backend/manager/services/scheduling_history/service.py index 19fc6781277..9dae4c42c7a 100644 --- a/src/ai/backend/manager/services/scheduling_history/service.py +++ b/src/ai/backend/manager/services/scheduling_history/service.py @@ -10,6 +10,14 @@ SearchDeploymentScopedHistoryAction, SearchDeploymentScopedHistoryActionResult, ) +from .actions.search_kernel_history import ( + SearchKernelHistoryAction, + SearchKernelHistoryActionResult, +) +from .actions.search_kernel_scoped_history import ( + SearchKernelScopedHistoryAction, + SearchKernelScopedHistoryActionResult, +) from .actions.search_route_history import ( SearchRouteHistoryAction, SearchRouteHistoryActionResult, @@ -54,6 +62,22 @@ async def search_session_history( has_previous_page=result.has_previous_page, ) + async def search_kernel_history( + self, + action: SearchKernelHistoryAction, + ) -> SearchKernelHistoryActionResult: + """Searches kernel scheduling history (admin API).""" + result = await self._repository.search_kernel_history( + querier=action.querier, + ) + + return SearchKernelHistoryActionResult( + histories=result.items, + total_count=result.total_count, + has_next_page=result.has_next_page, + has_previous_page=result.has_previous_page, + ) + async def search_deployment_history( self, action: SearchDeploymentHistoryAction, @@ -105,6 +129,23 @@ async def search_session_scoped_history( has_previous_page=result.has_previous_page, ) + async def search_kernel_scoped_history( + self, + action: SearchKernelScopedHistoryAction, + ) -> SearchKernelScopedHistoryActionResult: + """Searches kernel scheduling history within scope.""" + result = await self._repository.search_kernel_scoped_history( + querier=action.querier, + scope=action.scope, + ) + + return SearchKernelScopedHistoryActionResult( + histories=result.items, + total_count=result.total_count, + has_next_page=result.has_next_page, + has_previous_page=result.has_previous_page, + ) + async def search_deployment_scoped_history( self, action: SearchDeploymentScopedHistoryAction, diff --git a/tests/unit/common/dto/manager/v2/scheduling_history/test_types.py b/tests/unit/common/dto/manager/v2/scheduling_history/test_types.py index 769918cb9a5..54567f01639 100644 --- a/tests/unit/common/dto/manager/v2/scheduling_history/test_types.py +++ b/tests/unit/common/dto/manager/v2/scheduling_history/test_types.py @@ -3,10 +3,17 @@ from __future__ import annotations import json +import uuid +from dataclasses import dataclass from datetime import UTC, datetime +import pytest +from pydantic import ValidationError + from ai.backend.common.dto.manager.v2.scheduling_history.types import ( DeploymentHistoryOrderField, + KernelHistoryOrderField, + KernelHistoryScopeDTO, OrderDirection, RouteHistoryOrderField, SchedulingResultType, @@ -14,6 +21,9 @@ SubStepResultInfo, ) +_SESSION_ID = uuid.UUID("11111111-1111-1111-1111-111111111111") +_KERNEL_ID = uuid.UUID("22222222-2222-2222-2222-222222222222") + class TestOrderDirection: """Tests for OrderDirection enum.""" @@ -111,6 +121,56 @@ def test_enum_members_count(self) -> None: assert len(list(RouteHistoryOrderField)) == 2 +class TestKernelHistoryOrderField: + """Tests for KernelHistoryOrderField enum.""" + + def test_created_at_value(self) -> None: + assert KernelHistoryOrderField.CREATED_AT.value == "created_at" + + def test_updated_at_value(self) -> None: + assert KernelHistoryOrderField.UPDATED_AT.value == "updated_at" + + def test_enum_members_count(self) -> None: + assert len(list(KernelHistoryOrderField)) == 2 + + +@dataclass(frozen=True) +class _KernelScopeCase: + session_id: uuid.UUID | None + kernel_id: uuid.UUID | None + + +class TestKernelHistoryScopeDTO: + """Tests for KernelHistoryScopeDTO, whose axes are optional but not both-empty.""" + + def test_empty_scope_is_rejected(self) -> None: + with pytest.raises(ValidationError): + KernelHistoryScopeDTO() + + @pytest.mark.parametrize( + "case", + [ + _KernelScopeCase(session_id=_SESSION_ID, kernel_id=None), + _KernelScopeCase(session_id=None, kernel_id=_KERNEL_ID), + _KernelScopeCase(session_id=_SESSION_ID, kernel_id=_KERNEL_ID), + ], + ids=lambda case: f"session={case.session_id is not None}-kernel={case.kernel_id is not None}", + ) + def test_any_non_empty_axis_combination_is_accepted(self, case: _KernelScopeCase) -> None: + scope = KernelHistoryScopeDTO(session_id=case.session_id, kernel_id=case.kernel_id) + + assert scope.session_id == case.session_id + assert scope.kernel_id == case.kernel_id + + def test_serialization_round_trip(self) -> None: + scope = KernelHistoryScopeDTO(session_id=_SESSION_ID, kernel_id=_KERNEL_ID) + + restored = KernelHistoryScopeDTO.model_validate_json(scope.model_dump_json()) + + assert restored.session_id == _SESSION_ID + assert restored.kernel_id == _KERNEL_ID + + class TestSubStepResultInfo: """Tests for SubStepResultInfo Pydantic model.""" diff --git a/tests/unit/manager/dependencies/processing/test_composer.py b/tests/unit/manager/dependencies/processing/test_composer.py index 40873440c71..bbacc98a8bc 100644 --- a/tests/unit/manager/dependencies/processing/test_composer.py +++ b/tests/unit/manager/dependencies/processing/test_composer.py @@ -60,7 +60,6 @@ def _make_processing_input() -> ProcessingInput: appproxy_client_pool=MagicMock(), prometheus_client=MagicMock(), agent_client_pool=MagicMock(), - distributed_lock_factory=MagicMock(), stats_monitor=mock_stats_monitor, pidx=0, ) diff --git a/tests/unit/manager/models/scheduling_history/BUILD b/tests/unit/manager/models/scheduling_history/BUILD deleted file mode 100644 index dabf212d7e7..00000000000 --- a/tests/unit/manager/models/scheduling_history/BUILD +++ /dev/null @@ -1 +0,0 @@ -python_tests() diff --git a/tests/unit/manager/models/scheduling_history/test_conditions.py b/tests/unit/manager/models/scheduling_history/test_conditions.py deleted file mode 100644 index 8677133bd2e..00000000000 --- a/tests/unit/manager/models/scheduling_history/test_conditions.py +++ /dev/null @@ -1,196 +0,0 @@ -from __future__ import annotations - -from dataclasses import dataclass - -import pytest - -from ai.backend.common.data.filter_specs import StringInMatchSpec, StringMatchSpec -from ai.backend.manager.data.kernel.types import KernelSchedulingPhase -from ai.backend.manager.models.clauses import QueryCondition -from ai.backend.manager.models.scheduling_history.conditions import ( - KernelSchedulingHistoryConditions, -) - - -@dataclass(frozen=True) -class _ConditionCase: - label: str - condition: QueryCondition - expected_sql: str - - -class TestKernelSchedulingHistoryConditions: - @pytest.mark.parametrize( - "case", - [ - _ConditionCase( - label="by_from_status", - condition=KernelSchedulingHistoryConditions.by_from_status( - KernelSchedulingPhase.PULLING - ), - expected_sql="kernel_scheduling_history.from_status = 'PULLING'", - ), - _ConditionCase( - label="by_to_status", - condition=KernelSchedulingHistoryConditions.by_to_status( - KernelSchedulingPhase.RUNNING - ), - expected_sql="kernel_scheduling_history.to_status = 'RUNNING'", - ), - _ConditionCase( - label="by_phase_contains", - condition=KernelSchedulingHistoryConditions.by_phase_contains( - StringMatchSpec(value="PULLING", case_insensitive=False, negated=False) - ), - expected_sql="kernel_scheduling_history.phase LIKE '%PULLING%'", - ), - _ConditionCase( - label="by_phase_contains-negated", - condition=KernelSchedulingHistoryConditions.by_phase_contains( - StringMatchSpec(value="PULLING", case_insensitive=False, negated=True) - ), - expected_sql="kernel_scheduling_history.phase NOT LIKE '%PULLING%'", - ), - _ConditionCase( - label="by_phase_contains-case_insensitive", - condition=KernelSchedulingHistoryConditions.by_phase_contains( - StringMatchSpec(value="PULLING", case_insensitive=True, negated=False) - ), - expected_sql="lower(kernel_scheduling_history.phase) LIKE '%pulling%'", - ), - _ConditionCase( - label="by_phase_contains-case_insensitive-negated", - condition=KernelSchedulingHistoryConditions.by_phase_contains( - StringMatchSpec(value="PULLING", case_insensitive=True, negated=True) - ), - expected_sql="lower(kernel_scheduling_history.phase) NOT LIKE '%pulling%'", - ), - _ConditionCase( - label="by_phase_equals", - condition=KernelSchedulingHistoryConditions.by_phase_equals( - StringMatchSpec(value="PULLING", case_insensitive=False, negated=False) - ), - expected_sql="kernel_scheduling_history.phase = 'PULLING'", - ), - _ConditionCase( - label="by_phase_equals-negated", - condition=KernelSchedulingHistoryConditions.by_phase_equals( - StringMatchSpec(value="PULLING", case_insensitive=False, negated=True) - ), - expected_sql="kernel_scheduling_history.phase != 'PULLING'", - ), - _ConditionCase( - label="by_phase_equals-case_insensitive", - condition=KernelSchedulingHistoryConditions.by_phase_equals( - StringMatchSpec(value="PULLING", case_insensitive=True, negated=False) - ), - expected_sql="lower(kernel_scheduling_history.phase) = 'pulling'", - ), - _ConditionCase( - label="by_phase_equals-case_insensitive-negated", - condition=KernelSchedulingHistoryConditions.by_phase_equals( - StringMatchSpec(value="PULLING", case_insensitive=True, negated=True) - ), - expected_sql="lower(kernel_scheduling_history.phase) != 'pulling'", - ), - _ConditionCase( - label="by_phase_starts_with", - condition=KernelSchedulingHistoryConditions.by_phase_starts_with( - StringMatchSpec(value="PULLING", case_insensitive=False, negated=False) - ), - expected_sql="kernel_scheduling_history.phase LIKE 'PULLING%'", - ), - _ConditionCase( - label="by_phase_starts_with-negated", - condition=KernelSchedulingHistoryConditions.by_phase_starts_with( - StringMatchSpec(value="PULLING", case_insensitive=False, negated=True) - ), - expected_sql="kernel_scheduling_history.phase NOT LIKE 'PULLING%'", - ), - _ConditionCase( - label="by_phase_starts_with-case_insensitive", - condition=KernelSchedulingHistoryConditions.by_phase_starts_with( - StringMatchSpec(value="PULLING", case_insensitive=True, negated=False) - ), - expected_sql="lower(kernel_scheduling_history.phase) LIKE 'pulling%'", - ), - _ConditionCase( - label="by_phase_starts_with-case_insensitive-negated", - condition=KernelSchedulingHistoryConditions.by_phase_starts_with( - StringMatchSpec(value="PULLING", case_insensitive=True, negated=True) - ), - expected_sql="lower(kernel_scheduling_history.phase) NOT LIKE 'pulling%'", - ), - _ConditionCase( - label="by_phase_ends_with", - condition=KernelSchedulingHistoryConditions.by_phase_ends_with( - StringMatchSpec(value="PULLING", case_insensitive=False, negated=False) - ), - expected_sql="kernel_scheduling_history.phase LIKE '%PULLING'", - ), - _ConditionCase( - label="by_phase_ends_with-negated", - condition=KernelSchedulingHistoryConditions.by_phase_ends_with( - StringMatchSpec(value="PULLING", case_insensitive=False, negated=True) - ), - expected_sql="kernel_scheduling_history.phase NOT LIKE '%PULLING'", - ), - _ConditionCase( - label="by_phase_ends_with-case_insensitive", - condition=KernelSchedulingHistoryConditions.by_phase_ends_with( - StringMatchSpec(value="PULLING", case_insensitive=True, negated=False) - ), - expected_sql="lower(kernel_scheduling_history.phase) LIKE '%pulling'", - ), - _ConditionCase( - label="by_phase_ends_with-case_insensitive-negated", - condition=KernelSchedulingHistoryConditions.by_phase_ends_with( - StringMatchSpec(value="PULLING", case_insensitive=True, negated=True) - ), - expected_sql="lower(kernel_scheduling_history.phase) NOT LIKE '%pulling'", - ), - _ConditionCase( - label="by_phase_in", - condition=KernelSchedulingHistoryConditions.by_phase_in( - StringInMatchSpec( - values=["PULLING", "CREATING"], case_insensitive=False, negated=False - ) - ), - expected_sql="kernel_scheduling_history.phase IN ('PULLING', 'CREATING')", - ), - _ConditionCase( - label="by_phase_in-negated", - condition=KernelSchedulingHistoryConditions.by_phase_in( - StringInMatchSpec( - values=["PULLING", "CREATING"], case_insensitive=False, negated=True - ) - ), - expected_sql="(kernel_scheduling_history.phase NOT IN ('PULLING', 'CREATING'))", - ), - _ConditionCase( - label="by_phase_in-case_insensitive", - condition=KernelSchedulingHistoryConditions.by_phase_in( - StringInMatchSpec( - values=["PULLING", "CREATING"], case_insensitive=True, negated=False - ) - ), - expected_sql="lower(kernel_scheduling_history.phase) IN ('pulling', 'creating')", - ), - _ConditionCase( - label="by_phase_in-case_insensitive-negated", - condition=KernelSchedulingHistoryConditions.by_phase_in( - StringInMatchSpec( - values=["PULLING", "CREATING"], case_insensitive=True, negated=True - ) - ), - expected_sql=( - "(lower(kernel_scheduling_history.phase) NOT IN ('pulling', 'creating'))" - ), - ), - ], - ids=lambda case: case.label, - ) - def test_condition_compiles_to_expected_sql(self, case: _ConditionCase) -> None: - sql = str(case.condition().compile(compile_kwargs={"literal_binds": True})) - - assert sql == case.expected_sql diff --git a/tests/unit/manager/repositories/scheduling_history/test_kernel_scheduling_history_scope.py b/tests/unit/manager/repositories/scheduling_history/test_kernel_scheduling_history_scope.py new file mode 100644 index 00000000000..18e0ac5fe34 --- /dev/null +++ b/tests/unit/manager/repositories/scheduling_history/test_kernel_scheduling_history_scope.py @@ -0,0 +1,137 @@ +"""Tests for KernelSchedulingHistorySearchScope. + +The scope is the one piece of the kernel history stack with real branching: it accepts +two optional axes, must reject an empty scope, and must intersect when both are given. +""" + +from __future__ import annotations + +import uuid +from dataclasses import dataclass, field + +import pytest + +from ai.backend.common.types import KernelId, SessionId +from ai.backend.manager.errors.kernel import EmptyKernelSchedulingHistoryScope +from ai.backend.manager.models.kernel import KernelRow +from ai.backend.manager.models.session import SessionRow +from ai.backend.manager.repositories.scheduling_history.types import ( + KernelSchedulingHistorySearchScope, +) + +_SESSION_ID = SessionId(uuid.UUID("11111111-1111-1111-1111-111111111111")) +_KERNEL_ID = KernelId(uuid.UUID("22222222-2222-2222-2222-222222222222")) + + +@dataclass(frozen=True) +class _ExistenceCheckExpectation: + column_name: str + value: uuid.UUID + + +@dataclass(frozen=True) +class _ScopeCase: + session_id: SessionId | None + kernel_id: KernelId | None + expected_sql_fragments: list[str] = field(default_factory=list) + expected_checks: list[_ExistenceCheckExpectation] = field(default_factory=list) + + +def _scope_case_id(case: _ScopeCase) -> str: + return f"session={case.session_id is not None}-kernel={case.kernel_id is not None}" + + +class TestKernelSchedulingHistorySearchScope: + """Test cases for the kernel scheduling history search scope.""" + + def test_empty_scope_is_rejected(self) -> None: + with pytest.raises(EmptyKernelSchedulingHistoryScope): + KernelSchedulingHistorySearchScope() + + @pytest.mark.parametrize( + "case", + [ + _ScopeCase( + session_id=_SESSION_ID, + kernel_id=None, + expected_sql_fragments=["kernel_scheduling_history.session_id ="], + ), + _ScopeCase( + session_id=None, + kernel_id=_KERNEL_ID, + expected_sql_fragments=["kernel_scheduling_history.kernel_id ="], + ), + _ScopeCase( + session_id=_SESSION_ID, + kernel_id=_KERNEL_ID, + expected_sql_fragments=[ + "kernel_scheduling_history.session_id =", + "kernel_scheduling_history.kernel_id =", + ], + ), + ], + ids=_scope_case_id, + ) + def test_to_condition_filters_on_each_given_axis(self, case: _ScopeCase) -> None: + scope = KernelSchedulingHistorySearchScope( + session_id=case.session_id, kernel_id=case.kernel_id + ) + + rendered = str(scope.to_condition()()) + + for fragment in case.expected_sql_fragments: + assert fragment in rendered + + def test_to_condition_intersects_both_axes(self) -> None: + scope = KernelSchedulingHistorySearchScope(session_id=_SESSION_ID, kernel_id=_KERNEL_ID) + + rendered = str(scope.to_condition()()) + + assert " AND " in rendered + + @pytest.mark.parametrize( + "case", + [ + _ScopeCase( + session_id=_SESSION_ID, + kernel_id=None, + expected_checks=[ + _ExistenceCheckExpectation(column_name="id", value=_SESSION_ID), + ], + ), + _ScopeCase( + session_id=None, + kernel_id=_KERNEL_ID, + expected_checks=[ + _ExistenceCheckExpectation(column_name="id", value=_KERNEL_ID), + ], + ), + _ScopeCase( + session_id=_SESSION_ID, + kernel_id=_KERNEL_ID, + expected_checks=[ + _ExistenceCheckExpectation(column_name="id", value=_SESSION_ID), + _ExistenceCheckExpectation(column_name="id", value=_KERNEL_ID), + ], + ), + ], + ids=_scope_case_id, + ) + def test_existence_checks_cover_each_given_axis(self, case: _ScopeCase) -> None: + scope = KernelSchedulingHistorySearchScope( + session_id=case.session_id, kernel_id=case.kernel_id + ) + + checks = scope.existence_checks + + assert [ + _ExistenceCheckExpectation(column_name=c.column.key, value=c.value) for c in checks + ] == case.expected_checks + + def test_existence_checks_point_at_the_owning_tables(self) -> None: + scope = KernelSchedulingHistorySearchScope(session_id=_SESSION_ID, kernel_id=_KERNEL_ID) + + checks = scope.existence_checks + + assert checks[0].column is SessionRow.id + assert checks[1].column is KernelRow.id