diff --git a/backend/src/backend/admin.py b/backend/src/backend/admin.py index f9e314c..2813faa 100644 --- a/backend/src/backend/admin.py +++ b/backend/src/backend/admin.py @@ -2,21 +2,20 @@ from __future__ import annotations from typing import Any, Literal +from common.auth.passwords import hash_password +from common.db.models import Roles, Users, WorkspaceMembers +from common.ids import new_ulid from fastapi import APIRouter, Depends, HTTPException, status from pydantic import BaseModel, ConfigDict, Field from sqlalchemy import delete, func, or_, select from sqlalchemy.ext.asyncio import AsyncSession -from common.db.models import Roles, Users, WorkspaceMembers -from common.ids import new_ulid -from common.auth.passwords import hash_password from backend.dependencies import ( RequestContext, database_session, request_context, ) - router = APIRouter(prefix="/api/v1/admin", tags=["admin"]) diff --git a/backend/src/backend/auth.py b/backend/src/backend/auth.py index 8976422..a68e11b 100644 --- a/backend/src/backend/auth.py +++ b/backend/src/backend/auth.py @@ -16,18 +16,17 @@ from __future__ import annotations from typing import Any -from fastapi import APIRouter, Depends, HTTPException, Request, Response, status -from sqlalchemy import select -from sqlalchemy.ext.asyncio import AsyncSession - -from backend.dependencies import database_session, load_user_permissions from common.auth.jwt import JwtError, issue_jwt, verify_jwt_token from common.auth.membership import resolve_is_system_admin from common.auth.passwords import verify_password from common.config import settings from common.db.models import Roles, Users, WorkspaceMembers, Workspaces from common.ids import new_ulid +from fastapi import APIRouter, Depends, HTTPException, Request, Response, status +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession +from backend.dependencies import database_session, load_user_permissions router = APIRouter(tags=["auth"]) diff --git a/backend/src/backend/dependencies.py b/backend/src/backend/dependencies.py index 1865367..5fc7103 100644 --- a/backend/src/backend/dependencies.py +++ b/backend/src/backend/dependencies.py @@ -27,12 +27,8 @@ for. from __future__ import annotations +from collections.abc import AsyncIterator from dataclasses import dataclass -from typing import AsyncIterator - -from fastapi import Depends, HTTPException, Query, Request, status -from sqlalchemy import select -from sqlalchemy.ext.asyncio import AsyncSession from common.auth.jwt import JwtError, verify_jwt_token from common.auth.membership import ( @@ -43,7 +39,9 @@ from common.auth.membership import ( from common.db import session_scope from common.db.models import Permissions, RolePermissions, Roles, Users, Workspaces from common.ids import new_ulid - +from fastapi import Depends, HTTPException, Query, Request, status +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession ACCESS_TOKEN_COOKIE = "access_token" diff --git a/backend/src/backend/jupyter.py b/backend/src/backend/jupyter.py index 230375c..97e01d7 100644 --- a/backend/src/backend/jupyter.py +++ b/backend/src/backend/jupyter.py @@ -1,5 +1,4 @@ import re -from typing import Optional from common.auth.jwt import JwtError, verify_jwt_token from common.auth.membership import MembershipError, load_active_membership @@ -19,7 +18,7 @@ security = HTTPBearer(auto_error=False) def extract_notebook_path( uri: str, workspace_id: str, -) -> Optional[str]: +) -> str | None: """Pull the relative notebook path out of the original request URI. Only ``/jupyter/{workspace_id}/notebooks/*.ipynb`` requests are @@ -83,7 +82,7 @@ async def load_active_membership_or_403( async def verify_jupyter_access( request: Request, response: Response, - auth: Optional[HTTPAuthorizationCredentials] = Depends(security), + auth: HTTPAuthorizationCredentials | None = Depends(security), session: AsyncSession = Depends(database_session), ) -> dict: """Nginx auth_request subrequest handler. diff --git a/backend/src/backend/platform.py b/backend/src/backend/platform.py index 871e05d..92f176e 100644 --- a/backend/src/backend/platform.py +++ b/backend/src/backend/platform.py @@ -65,12 +65,6 @@ import re from dataclasses import dataclass from typing import Any, Literal -from fastapi import APIRouter, Depends, HTTPException, Request, status -from pydantic import BaseModel, ConfigDict, Field -from sqlalchemy import func, insert, or_, select, update -from sqlalchemy.ext.asyncio import AsyncSession - -from backend.dependencies import current_user, database_session from common.auth.passwords import hash_password from common.db.models import ( Permissions, @@ -81,7 +75,12 @@ from common.db.models import ( Workspaces, ) from common.ids import new_ulid +from fastapi import APIRouter, Depends, HTTPException, Request, status +from pydantic import BaseModel, ConfigDict, Field +from sqlalchemy import func, insert, or_, select, update +from sqlalchemy.ext.asyncio import AsyncSession +from backend.dependencies import current_user, database_session router = APIRouter(prefix="/api/v1/platform", tags=["platform"]) @@ -1240,7 +1239,7 @@ async def patch_role_permissions( __all__ = [ - "router", "SystemAdminContext", + "router", "system_admin_context", ] \ No newline at end of file diff --git a/backend/src/backend/resources.py b/backend/src/backend/resources.py index 5863c1e..b8114f8 100644 --- a/backend/src/backend/resources.py +++ b/backend/src/backend/resources.py @@ -1,23 +1,21 @@ from __future__ import annotations -import base64 import os from datetime import UTC, datetime from pathlib import Path, PurePosixPath from typing import Any -from fastapi import APIRouter, Depends, Header, HTTPException, Query, Request, status -from sqlalchemy import or_, select -from sqlalchemy.ext.asyncio import AsyncSession - from common.db.models import DataResources, StorageObjects from common.ids import new_ulid from common.storage import workspaces_root from common.storage.schemas import ( CreateUploadRequest, DownloadUrlRequest, - ServerObjectRequest, ) +from fastapi import APIRouter, Depends, Header, HTTPException, Query, Request, status +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession + from backend.dependencies import ( RequestContext, database_session, @@ -27,11 +25,10 @@ from backend.schemas import ( CompleteResourceUploadRequest, CreateResourceUploadRequest, DownloadUrlRequest, + ResourceRelativePathRequest, ) -from backend.schemas import ResourceRelativePathRequest from backend.services.storage import ( create_download_url_payload, - create_server_object_payload, create_upload_record, soft_delete_object, upload_bytes_to_session, diff --git a/backend/src/backend/runtime_client.py b/backend/src/backend/runtime_client.py index 6137103..c2cd6f1 100644 --- a/backend/src/backend/runtime_client.py +++ b/backend/src/backend/runtime_client.py @@ -6,6 +6,7 @@ from typing import Any import httpx from loguru import logger + @dataclass(frozen=True) class RuntimeClientError(Exception): status_code: int diff --git a/backend/src/backend/schedule_runs.py b/backend/src/backend/schedule_runs.py index 6fbb131..c40e578 100644 --- a/backend/src/backend/schedule_runs.py +++ b/backend/src/backend/schedule_runs.py @@ -4,6 +4,22 @@ from datetime import UTC, datetime from typing import Any, Literal from urllib.parse import quote +from common.config import settings +from common.db.models import ( + ScheduleNodeRuns, + ScheduleRuns, + StorageObjects, +) +from common.scheduler import ( + DagTooLarge, + InvalidDag, + InvalidNodeArguments, + ScheduleNotFound, + TriggerError, + create_scheduled_run, + normalize_idempotency_key, +) +from common.schemas import StrictModel from fastapi import ( APIRouter, Depends, @@ -18,29 +34,11 @@ from pydantic import Field from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession -from common.config import settings -from common.db.models import ( - ScheduleNodeRuns, - ScheduleRuns, - StorageObjects, -) from backend.dependencies import ( RequestContext, database_session, request_context, ) -from common.schemas import StrictModel -from common.scheduler import ( - DagTooLarge, - InvalidDag, - InvalidNodeArguments, - ScheduleNotFound, - TriggerError, - create_scheduled_run, - normalize_idempotency_key, -) -from common.ids import new_ulid - router = APIRouter(tags=["schedule-runs"]) RunStatus = Literal[ diff --git a/backend/src/backend/schedule_schemas.py b/backend/src/backend/schedule_schemas.py index f614db3..1d0d705 100644 --- a/backend/src/backend/schedule_schemas.py +++ b/backend/src/backend/schedule_schemas.py @@ -3,10 +3,9 @@ from __future__ import annotations from datetime import datetime from typing import Any -from pydantic import Field, field_validator, model_validator - from common.db.models.schedules import FailurePolicy, PythonVersion, TriggerType from common.schemas import StrictModel +from pydantic import Field, field_validator, model_validator def _required_text(value: str) -> str: @@ -40,7 +39,7 @@ class CreateScheduleRequest(StrictModel): return normalized or None @model_validator(mode="after") - def validate_trigger(self) -> "CreateScheduleRequest": + def validate_trigger(self) -> CreateScheduleRequest: if self.trigger_type == "cron" and not self.cron_expression: raise ValueError("cron_expression is required for cron schedules") if self.trigger_type != "cron" and self.cron_expression: @@ -78,17 +77,17 @@ class UpdateScheduleRequest(StrictModel): return normalized or None @model_validator(mode="after") - def require_change(self) -> "UpdateScheduleRequest": + def require_change(self) -> UpdateScheduleRequest: if self.model_fields_set == {"workflow_version"}: raise ValueError("at least one schedule field must be updated") - for field in { + for field in ( "schedule_name", "trigger_type", "timezone", "enabled", "max_concurrency", "failure_policy", - }: + ): if field in self.model_fields_set and getattr(self, field) is None: raise ValueError(f"{field} cannot be null") return self @@ -158,7 +157,7 @@ class UpdateScheduleNodeRequest(StrictModel): return _required_text(value) if value is not None else None @model_validator(mode="after") - def require_change(self) -> "UpdateScheduleNodeRequest": + def require_change(self) -> UpdateScheduleNodeRequest: if self.model_fields_set == {"workflow_version"}: raise ValueError("at least one node field must be updated") for field in self.model_fields_set - {"workflow_version"}: @@ -182,7 +181,7 @@ class CreateScheduleEdgeRequest(StrictModel): return normalized or None @model_validator(mode="after") - def reject_self_edge(self) -> "CreateScheduleEdgeRequest": + def reject_self_edge(self) -> CreateScheduleEdgeRequest: if self.source_node_id == self.target_node_id: raise ValueError("an edge cannot connect a node to itself") return self diff --git a/backend/src/backend/schedules.py b/backend/src/backend/schedules.py index e163232..5976343 100644 --- a/backend/src/backend/schedules.py +++ b/backend/src/backend/schedules.py @@ -6,11 +6,6 @@ from decimal import Decimal from typing import Any from zoneinfo import ZoneInfo, ZoneInfoNotFoundError -from croniter import CroniterBadCronError, croniter -from fastapi import APIRouter, Depends, HTTPException, Query, status -from sqlalchemy import delete, func, or_, select -from sqlalchemy.ext.asyncio import AsyncSession - from common.db.models import ( ScheduleEdges, ScheduleNodeRuns, @@ -20,6 +15,11 @@ from common.db.models import ( Versions, ) from common.ids import new_ulid +from croniter import CroniterBadCronError, croniter +from fastapi import APIRouter, Depends, HTTPException, Query, status +from sqlalchemy import delete, func, or_, select +from sqlalchemy.ext.asyncio import AsyncSession + from backend.dependencies import ( RequestContext, database_session, @@ -36,7 +36,6 @@ from backend.schedule_schemas import ( WorkflowVersionRequest, ) - router = APIRouter(tags=["schedules"]) diff --git a/backend/src/backend/scripts.py b/backend/src/backend/scripts.py index a6e3903..918281d 100644 --- a/backend/src/backend/scripts.py +++ b/backend/src/backend/scripts.py @@ -1,4 +1,3 @@ -import asyncio import base64 import hashlib import json @@ -15,13 +14,12 @@ from common.db.models import ( Versions, ) from common.ids import new_ulid -from common.storage.schemas import ServerObjectRequest from common.storage import actual_bucket_name, build_storage_uri +from common.storage.schemas import ServerObjectRequest from fastapi import ( APIRouter, BackgroundTasks, Depends, - Header, HTTPException, Query, Request, @@ -46,10 +44,10 @@ from backend.schemas import ( UpdateScriptRequest, ) from backend.services.storage import ( + _resolve_unique_object_key, create_download_url_payload, create_server_object_payload, soft_delete_object, - _resolve_unique_object_key, ) router = APIRouter(tags=["scripts"]) @@ -202,9 +200,7 @@ def script_payload( workspace_prefix = f"{script.workspace_id}/" if object_key: jupyter_path = ( - object_key[len(workspace_prefix) :] - if object_key.startswith(workspace_prefix) - else object_key + object_key.removeprefix(workspace_prefix) ) else: jupyter_path = _jupyter_path(script.script_type, script.script_id) diff --git a/backend/src/backend/services/storage.py b/backend/src/backend/services/storage.py index eb513cf..7d28f89 100644 --- a/backend/src/backend/services/storage.py +++ b/backend/src/backend/services/storage.py @@ -35,21 +35,19 @@ from datetime import timedelta from pathlib import PurePosixPath from typing import Any -from fastapi import HTTPException, Request, status -from sqlalchemy import select -from sqlalchemy.ext.asyncio import AsyncSession -from sqlalchemy.exc import IntegrityError - from common.config import settings from common.db.models import StorageObjects, UploadSessions from common.ids import new_ulid +from common.storage import USAGE_TYPE_TO_PURPOSE, actual_bucket_name, build_storage_uri from common.storage.schemas import ( CreateUploadRequest, DownloadUrlRequest, ServerObjectRequest, ) -from common.storage import USAGE_TYPE_TO_PURPOSE, actual_bucket_name, build_storage_uri - +from fastapi import HTTPException, Request, status +from sqlalchemy import select +from sqlalchemy.exc import IntegrityError +from sqlalchemy.ext.asyncio import AsyncSession # ── shared low-level helpers (module-private) ──────────────────────────── @@ -91,7 +89,7 @@ def _safe_path_segment(segment: str) -> str: def _utcnow_naive() -> Any: - from datetime import datetime, UTC + from datetime import UTC, datetime return datetime.now(UTC).replace(tzinfo=None) @@ -207,9 +205,8 @@ async def create_upload_record( when the idempotency key hits an already-completed upload. """ from backend.storage_api import ( - require_workspace_member, normalized_idempotency_key, - BUCKET_FOR_USAGE, + require_workspace_member, ) workspace = await require_workspace_member( @@ -236,6 +233,7 @@ async def create_upload_record( upload = existing else: from datetime import timedelta + from backend.storage_api import utcnow bucket_name = _resolve_bucket_for_usage( diff --git a/backend/src/backend/storage_api.py b/backend/src/backend/storage_api.py index 0ff2a92..864f002 100644 --- a/backend/src/backend/storage_api.py +++ b/backend/src/backend/storage_api.py @@ -1,13 +1,10 @@ from __future__ import annotations import hashlib +from collections.abc import AsyncIterator from datetime import UTC, datetime, timedelta from pathlib import PurePosixPath -from typing import Any, AsyncIterator - -from fastapi import APIRouter, Depends, HTTPException, Request, status -from sqlalchemy import select -from sqlalchemy.ext.asyncio import AsyncSession +from typing import Any from common.config import settings from common.db import session_scope @@ -25,6 +22,10 @@ from common.storage.schemas import ( DownloadUrlRequest, ServerObjectRequest, ) +from fastapi import APIRouter, Depends, HTTPException, Request, status +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession + from backend.services.storage import ( create_download_url_payload, create_server_object_payload, @@ -44,7 +45,7 @@ def hash_bytes(value: str) -> bytes: def normalized_idempotency_key(workspace_id: str, user_id: str, value: str) -> str: digest = hashlib.sha256( - f"{workspace_id}:{user_id}:{value}".encode("utf-8") + f"{workspace_id}:{user_id}:{value}".encode() ).hexdigest() return f"v1:{digest}" diff --git a/backend/tests/test_resources.py b/backend/tests/test_resources.py index 2e534a7..750d8aa 100644 --- a/backend/tests/test_resources.py +++ b/backend/tests/test_resources.py @@ -7,12 +7,9 @@ from __future__ import annotations import datetime from types import SimpleNamespace -from unittest.mock import AsyncMock, MagicMock +from unittest.mock import MagicMock import pytest - -from types import SimpleNamespace - from backend.resources import compute_jupyter_relative_path, resource_payload from backend.services.storage import _safe_file_name, _safe_path_segment diff --git a/backend/tests/test_runtime_client_directories.py b/backend/tests/test_runtime_client_directories.py index 7601277..78627f3 100644 --- a/backend/tests/test_runtime_client_directories.py +++ b/backend/tests/test_runtime_client_directories.py @@ -15,7 +15,6 @@ from __future__ import annotations import httpx import pytest import respx - from backend.runtime_client import RuntimeClient, RuntimeClientError WORKSPACE_ID = "01HWS0000000000000000000A" diff --git a/backend/tests/test_scripts.py b/backend/tests/test_scripts.py index 32fd50d..d8caade 100644 --- a/backend/tests/test_scripts.py +++ b/backend/tests/test_scripts.py @@ -17,7 +17,6 @@ from types import SimpleNamespace from unittest.mock import AsyncMock, MagicMock import pytest - from common.config import settings from common.db.models import DataResources, Scripts, StorageObjects diff --git a/common/src/common/auth/jwt.py b/common/src/common/auth/jwt.py index 9be20ef..2759e46 100644 --- a/common/src/common/auth/jwt.py +++ b/common/src/common/auth/jwt.py @@ -25,7 +25,6 @@ from typing import Any from common.config import settings - JWT_SECRET: str = settings.jwt_secret JWT_ALGORITHM: str = "HS256" DEFAULT_TTL_SECONDS: int = 24 * 60 * 60 diff --git a/common/src/common/auth/passwords.py b/common/src/common/auth/passwords.py index f1e096f..e3a986d 100644 --- a/common/src/common/auth/passwords.py +++ b/common/src/common/auth/passwords.py @@ -17,7 +17,6 @@ import secrets from passlib.context import CryptContext - _crypt_context = CryptContext(schemes=["bcrypt"], deprecated="auto") diff --git a/common/src/common/db/base.py b/common/src/common/db/base.py index 23807ec..c861dd9 100644 --- a/common/src/common/db/base.py +++ b/common/src/common/db/base.py @@ -1,4 +1,3 @@ -# coding=utf-8 """ @Time :2026/7/29 @Author :tao.chen diff --git a/common/src/common/db/models/events.py b/common/src/common/db/models/events.py index 6d0fd94..e870eda 100644 --- a/common/src/common/db/models/events.py +++ b/common/src/common/db/models/events.py @@ -1,7 +1,6 @@ import datetime -from typing import Optional -from sqlalchemy import Index, JSON, String, text +from sqlalchemy import JSON, Index, String, text from sqlalchemy.dialects.mysql import CHAR, DATETIME, INTEGER, SMALLINT, TINYINT from sqlalchemy.orm import Mapped, mapped_column @@ -26,15 +25,15 @@ class ConsumerInbox(Base): created_at: Mapped[datetime.datetime] = mapped_column( DATETIME(fsp=3), nullable=False, server_default=text("CURRENT_TIMESTAMP(3)") ) - message_id: Mapped[Optional[str]] = mapped_column( + message_id: Mapped[str | None] = mapped_column( String(128), comment="Inbox message ID" ) - processed_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) - error_message: Mapped[Optional[str]] = mapped_column(String(2000)) + processed_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) + error_message: Mapped[str | None] = mapped_column(String(2000)) is_deleted: Mapped[int] = mapped_column( TINYINT(1), nullable=False, server_default=text("0") ) - deleted_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) + deleted_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) class OutboxEvents(Base): @@ -69,11 +68,11 @@ class OutboxEvents(Base): created_at: Mapped[datetime.datetime] = mapped_column( DATETIME(fsp=3), nullable=False, server_default=text("CURRENT_TIMESTAMP(3)") ) - trace_id: Mapped[Optional[str]] = mapped_column(String(64)) - idempotency_key: Mapped[Optional[str]] = mapped_column(String(128)) - published_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) - last_error: Mapped[Optional[str]] = mapped_column(String(2000)) + trace_id: Mapped[str | None] = mapped_column(String(64)) + idempotency_key: Mapped[str | None] = mapped_column(String(128)) + published_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) + last_error: Mapped[str | None] = mapped_column(String(2000)) is_deleted: Mapped[int] = mapped_column( TINYINT(1), nullable=False, server_default=text("0") ) - deleted_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) + deleted_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) diff --git a/common/src/common/db/models/identity.py b/common/src/common/db/models/identity.py index c22dae5..58c4364 100644 --- a/common/src/common/db/models/identity.py +++ b/common/src/common/db/models/identity.py @@ -1,5 +1,4 @@ import datetime -from typing import Optional from sqlalchemy import Index, String, text from sqlalchemy.dialects.mysql import CHAR, DATETIME, TINYINT @@ -23,11 +22,11 @@ class Permissions(Base): created_at: Mapped[datetime.datetime] = mapped_column( DATETIME(fsp=3), nullable=False, server_default=text("CURRENT_TIMESTAMP(3)") ) - description: Mapped[Optional[str]] = mapped_column(String(500)) + description: Mapped[str | None] = mapped_column(String(500)) is_deleted: Mapped[int] = mapped_column( TINYINT(1), nullable=False, server_default=text("0") ) - deleted_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) + deleted_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) class Roles(Base): @@ -54,11 +53,11 @@ class Roles(Base): nullable=False, server_default=text("CURRENT_TIMESTAMP(3) ON UPDATE CURRENT_TIMESTAMP(3)"), ) - description: Mapped[Optional[str]] = mapped_column(String(500)) + description: Mapped[str | None] = mapped_column(String(500)) is_deleted: Mapped[int] = mapped_column( TINYINT(1), nullable=False, server_default=text("0") ) - deleted_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) + deleted_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) class RolePermissions(Base): @@ -76,7 +75,7 @@ class RolePermissions(Base): is_deleted: Mapped[int] = mapped_column( TINYINT(1), nullable=False, server_default=text("0") ) - deleted_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) + deleted_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) class Users(Base): @@ -107,11 +106,11 @@ class Users(Base): nullable=False, server_default=text("CURRENT_TIMESTAMP(3) ON UPDATE CURRENT_TIMESTAMP(3)"), ) - email: Mapped[Optional[str]] = mapped_column(String(255)) - platform_role_id: Mapped[Optional[str]] = mapped_column(CHAR(26)) - avatar_uri: Mapped[Optional[str]] = mapped_column(String(1000)) - last_login_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) + email: Mapped[str | None] = mapped_column(String(255)) + platform_role_id: Mapped[str | None] = mapped_column(CHAR(26)) + avatar_uri: Mapped[str | None] = mapped_column(String(1000)) + last_login_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) is_deleted: Mapped[int] = mapped_column( TINYINT(1), nullable=False, server_default=text("0") ) - deleted_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) + deleted_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) diff --git a/common/src/common/db/models/schedules.py b/common/src/common/db/models/schedules.py index a34a54b..512ffb8 100644 --- a/common/src/common/db/models/schedules.py +++ b/common/src/common/db/models/schedules.py @@ -1,14 +1,13 @@ import datetime import decimal -from typing import Literal, Optional +from typing import Literal -from sqlalchemy import DECIMAL, Index, Integer, JSON, String, Text, text +from sqlalchemy import DECIMAL, JSON, Index, Integer, String, Text, text from sqlalchemy.dialects.mysql import BIGINT, CHAR, DATETIME, INTEGER, TINYINT from sqlalchemy.orm import Mapped, mapped_column from common.db.base import Base - TriggerType = Literal["manual", "cron", "api"] FailurePolicy = Literal["stop", "continue"] PythonVersion = Literal["3.8", "3.10", "3.12"] @@ -58,14 +57,14 @@ class Schedules(Base): nullable=False, server_default=text("CURRENT_TIMESTAMP(3) ON UPDATE CURRENT_TIMESTAMP(3)"), ) - description: Mapped[Optional[str]] = mapped_column(String(1000)) - cron_expression: Mapped[Optional[str]] = mapped_column(String(128)) - last_run_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) - next_run_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) + description: Mapped[str | None] = mapped_column(String(1000)) + cron_expression: Mapped[str | None] = mapped_column(String(128)) + last_run_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) + next_run_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) is_deleted: Mapped[int] = mapped_column( TINYINT(1), nullable=False, server_default=text("0") ) - deleted_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) + deleted_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) class ScheduleRuns(Base): @@ -112,18 +111,18 @@ class ScheduleRuns(Base): created_at: Mapped[datetime.datetime] = mapped_column( DATETIME(fsp=3), nullable=False, server_default=text("CURRENT_TIMESTAMP(3)") ) - triggered_by: Mapped[Optional[str]] = mapped_column(CHAR(26)) - started_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) - finished_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) - duration_ms: Mapped[Optional[int]] = mapped_column(BIGINT) - error_code: Mapped[Optional[str]] = mapped_column(String(64)) - error_message: Mapped[Optional[str]] = mapped_column(Text) - logs_object_id: Mapped[Optional[str]] = mapped_column(CHAR(26)) - result_object_id: Mapped[Optional[str]] = mapped_column(CHAR(26)) + triggered_by: Mapped[str | None] = mapped_column(CHAR(26)) + started_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) + finished_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) + duration_ms: Mapped[int | None] = mapped_column(BIGINT) + error_code: Mapped[str | None] = mapped_column(String(64)) + error_message: Mapped[str | None] = mapped_column(Text) + logs_object_id: Mapped[str | None] = mapped_column(CHAR(26)) + result_object_id: Mapped[str | None] = mapped_column(CHAR(26)) is_deleted: Mapped[int] = mapped_column( TINYINT(1), nullable=False, server_default=text("0") ) - deleted_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) + deleted_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) class ScheduleNodes(Base): @@ -168,14 +167,14 @@ class ScheduleNodes(Base): nullable=False, server_default=text("CURRENT_TIMESTAMP(3) ON UPDATE CURRENT_TIMESTAMP(3)"), ) - arguments_json: Mapped[Optional[dict]] = mapped_column(JSON) - env_refs_json: Mapped[Optional[dict]] = mapped_column( + arguments_json: Mapped[dict | None] = mapped_column(JSON) + env_refs_json: Mapped[dict | None] = mapped_column( JSON, comment="只存密钥引用,不存明文密钥" ) is_deleted: Mapped[int] = mapped_column( TINYINT(1), nullable=False, server_default=text("0") ) - deleted_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) + deleted_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) class ScheduleEdges(Base): @@ -200,11 +199,11 @@ class ScheduleEdges(Base): created_at: Mapped[datetime.datetime] = mapped_column( DATETIME(fsp=3), nullable=False, server_default=text("CURRENT_TIMESTAMP(3)") ) - condition_expr: Mapped[Optional[str]] = mapped_column(String(1000)) + condition_expr: Mapped[str | None] = mapped_column(String(1000)) is_deleted: Mapped[int] = mapped_column( TINYINT(1), nullable=False, server_default=text("0") ) - deleted_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) + deleted_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) class ScheduleNodeRuns(Base): @@ -240,15 +239,15 @@ class ScheduleNodeRuns(Base): created_at: Mapped[datetime.datetime] = mapped_column( DATETIME(fsp=3), nullable=False, server_default=text("CURRENT_TIMESTAMP(3)") ) - started_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) - finished_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) - duration_ms: Mapped[Optional[int]] = mapped_column(BIGINT) - exit_code: Mapped[Optional[int]] = mapped_column(Integer) - message: Mapped[Optional[str]] = mapped_column(String(2000)) - metrics_json: Mapped[Optional[dict]] = mapped_column(JSON) - logs_object_id: Mapped[Optional[str]] = mapped_column(CHAR(26)) - result_object_id: Mapped[Optional[str]] = mapped_column(CHAR(26)) + started_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) + finished_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) + duration_ms: Mapped[int | None] = mapped_column(BIGINT) + exit_code: Mapped[int | None] = mapped_column(Integer) + message: Mapped[str | None] = mapped_column(String(2000)) + metrics_json: Mapped[dict | None] = mapped_column(JSON) + logs_object_id: Mapped[str | None] = mapped_column(CHAR(26)) + result_object_id: Mapped[str | None] = mapped_column(CHAR(26)) is_deleted: Mapped[int] = mapped_column( TINYINT(1), nullable=False, server_default=text("0") ) - deleted_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) + deleted_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) diff --git a/common/src/common/db/models/scripts.py b/common/src/common/db/models/scripts.py index e956747..06a7d31 100644 --- a/common/src/common/db/models/scripts.py +++ b/common/src/common/db/models/scripts.py @@ -1,5 +1,4 @@ import datetime -from typing import Optional from sqlalchemy import Index, String, text from sqlalchemy.dialects.mysql import BIGINT, CHAR, DATETIME, INTEGER, TINYINT @@ -52,7 +51,7 @@ class Scripts(Base): is_locked: Mapped[int] = mapped_column( TINYINT(1), nullable=False, server_default=text("1") ) - deleted_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) + deleted_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) class Versions(Base): @@ -97,12 +96,12 @@ class Versions(Base): created_at: Mapped[datetime.datetime] = mapped_column( DATETIME(fsp=3), nullable=False, server_default=text("CURRENT_TIMESTAMP(3)") ) - release_note: Mapped[Optional[str]] = mapped_column(String(1000)) - schedule_hidden_at: Mapped[Optional[datetime.datetime]] = mapped_column( + release_note: Mapped[str | None] = mapped_column(String(1000)) + schedule_hidden_at: Mapped[datetime.datetime | None] = mapped_column( DATETIME(fsp=3), comment="从调度稳定版本列表移除的时间;不影响版本和运行历史", ) is_deleted: Mapped[int] = mapped_column( TINYINT(1), nullable=False, server_default=text("0") ) - deleted_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) + deleted_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) diff --git a/common/src/common/db/models/storage.py b/common/src/common/db/models/storage.py index 07ca099..b133ac0 100644 --- a/common/src/common/db/models/storage.py +++ b/common/src/common/db/models/storage.py @@ -1,7 +1,6 @@ import datetime -from typing import Optional -from sqlalchemy import BINARY, Computed, Index, JSON, String, text +from sqlalchemy import BINARY, JSON, Computed, Index, String, text from sqlalchemy.dialects.mysql import BIGINT, CHAR, DATETIME, TINYINT from sqlalchemy.orm import Mapped, mapped_column @@ -86,20 +85,20 @@ class StorageObjects(Base): nullable=False, server_default=text("CURRENT_TIMESTAMP(3) ON UPDATE CURRENT_TIMESTAMP(3)"), ) - owner_user_id: Mapped[Optional[str]] = mapped_column(CHAR(26)) - parent_object_id: Mapped[Optional[str]] = mapped_column(CHAR(26)) - relative_path: Mapped[Optional[str]] = mapped_column( + owner_user_id: Mapped[str | None] = mapped_column(CHAR(26)) + parent_object_id: Mapped[str | None] = mapped_column(CHAR(26)) + relative_path: Mapped[str | None] = mapped_column( String(1024), comment="Workspace 相对路径" ) - path_hash: Mapped[Optional[bytes]] = mapped_column( + path_hash: Mapped[bytes | None] = mapped_column( BINARY(32), comment="SHA-256(relative_path),由应用写入" ) - bucket_name: Mapped[Optional[str]] = mapped_column(String(128)) - object_key: Mapped[Optional[str]] = mapped_column(String(1024)) - object_key_hash: Mapped[Optional[bytes]] = mapped_column( + bucket_name: Mapped[str | None] = mapped_column(String(128)) + object_key: Mapped[str | None] = mapped_column(String(1024)) + object_key_hash: Mapped[bytes | None] = mapped_column( BINARY(32), comment="SHA-256(object_key),由应用写入" ) - object_key_hash_active: Mapped[Optional[bytes]] = mapped_column( + object_key_hash_active: Mapped[bytes | None] = mapped_column( BINARY(32), Computed( "CASE WHEN object_status = 'available' THEN object_key_hash ELSE NULL END", @@ -107,17 +106,17 @@ class StorageObjects(Base): ), comment="VIRTUAL generated column used by uk_storage_bucket_key_active", ) - file_extension: Mapped[Optional[str]] = mapped_column(String(32)) - mime_type: Mapped[Optional[str]] = mapped_column(String(255)) - content_hash: Mapped[Optional[str]] = mapped_column( + file_extension: Mapped[str | None] = mapped_column(String(32)) + mime_type: Mapped[str | None] = mapped_column(String(255)) + content_hash: Mapped[str | None] = mapped_column( CHAR(64), comment="SHA-256 hex" ) - object_etag: Mapped[Optional[str]] = mapped_column(String(255)) + object_etag: Mapped[str | None] = mapped_column(String(255)) is_deleted: Mapped[int] = mapped_column( TINYINT(1), nullable=False, server_default=text("0") ) - deleted_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) - trash_key: Mapped[Optional[str]] = mapped_column( + deleted_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) + trash_key: Mapped[str | None] = mapped_column( String(1100), comment=( "Path inside the trash bucket where the soft-deleted bytes " @@ -160,14 +159,14 @@ class DataResources(Base): nullable=False, server_default=text("CURRENT_TIMESTAMP(3) ON UPDATE CURRENT_TIMESTAMP(3)"), ) - description: Mapped[Optional[str]] = mapped_column(String(1000)) - schema_json: Mapped[Optional[dict]] = mapped_column( + description: Mapped[str | None] = mapped_column(String(1000)) + schema_json: Mapped[dict | None] = mapped_column( JSON, comment="字段结构、行数等可选元数据" ) is_deleted: Mapped[int] = mapped_column( TINYINT(1), nullable=False, server_default=text("0") ) - deleted_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) + deleted_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) class UploadSessions(Base): @@ -208,12 +207,12 @@ class UploadSessions(Base): nullable=False, server_default=text("CURRENT_TIMESTAMP(3) ON UPDATE CURRENT_TIMESTAMP(3)"), ) - multipart_upload_id: Mapped[Optional[str]] = mapped_column(String(255)) - expected_size_bytes: Mapped[Optional[int]] = mapped_column(BIGINT) - expected_hash: Mapped[Optional[str]] = mapped_column(CHAR(64)) - content_type: Mapped[Optional[str]] = mapped_column(String(255)) - storage_object_id: Mapped[Optional[str]] = mapped_column(CHAR(26)) - completed_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) + multipart_upload_id: Mapped[str | None] = mapped_column(String(255)) + expected_size_bytes: Mapped[int | None] = mapped_column(BIGINT) + expected_hash: Mapped[str | None] = mapped_column(CHAR(64)) + content_type: Mapped[str | None] = mapped_column(String(255)) + storage_object_id: Mapped[str | None] = mapped_column(CHAR(26)) + completed_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) # Metadata persisted at session creation so step 2 (PUT bytes) can build # the StorageObjects row without re-sending them. Replaces the # CompleteUploadRequest payload that lived between presign-PUT and head(). @@ -236,4 +235,4 @@ class UploadSessions(Base): is_deleted: Mapped[int] = mapped_column( TINYINT(1), nullable=False, server_default=text("0") ) - deleted_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) + deleted_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) diff --git a/common/src/common/db/models/workspaces.py b/common/src/common/db/models/workspaces.py index f482d3e..f029f5d 100644 --- a/common/src/common/db/models/workspaces.py +++ b/common/src/common/db/models/workspaces.py @@ -1,5 +1,4 @@ import datetime -from typing import Optional from sqlalchemy import Index, String, text from sqlalchemy.dialects.mysql import BIGINT, CHAR, DATETIME, TINYINT @@ -44,17 +43,17 @@ class Workspaces(Base): nullable=False, server_default=text("CURRENT_TIMESTAMP(3) ON UPDATE CURRENT_TIMESTAMP(3)"), ) - description: Mapped[Optional[str]] = mapped_column(String(1000)) - artifact_bucket: Mapped[Optional[str]] = mapped_column( + description: Mapped[str | None] = mapped_column(String(1000)) + artifact_bucket: Mapped[str | None] = mapped_column( String(128), comment="S3 bucket" ) - artifact_prefix: Mapped[Optional[str]] = mapped_column( + artifact_prefix: Mapped[str | None] = mapped_column( String(512), comment="S3 object key prefix" ) is_deleted: Mapped[int] = mapped_column( TINYINT(1), nullable=False, server_default=text("0") ) - deleted_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) + deleted_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) class WorkspaceMembers(Base): @@ -82,4 +81,4 @@ class WorkspaceMembers(Base): is_deleted: Mapped[int] = mapped_column( TINYINT(1), nullable=False, server_default=text("0") ) - deleted_at: Mapped[Optional[datetime.datetime]] = mapped_column(DATETIME(fsp=3)) + deleted_at: Mapped[datetime.datetime | None] = mapped_column(DATETIME(fsp=3)) diff --git a/common/src/common/scheduler/__init__.py b/common/src/common/scheduler/__init__.py index fbb997d..b6e4ad9 100644 --- a/common/src/common/scheduler/__init__.py +++ b/common/src/common/scheduler/__init__.py @@ -49,7 +49,7 @@ def build_sqlalchemy_jobstore( database_url: str, *, tablename: str = JOBSTORE_TABLE, -) -> "SQLAlchemyJobStore": +) -> SQLAlchemyJobStore: """Instantiate a :class:`SQLAlchemyJobStore` for the canonical table. The caller is responsible for ensuring APScheduler and its sync @@ -65,11 +65,11 @@ def build_sqlalchemy_jobstore( __all__ = [ + "JOBSTORE_TABLE", + "SYSTEM_CRON_USER_ID", "DagTooLarge", "InvalidDag", "InvalidNodeArguments", - "JOBSTORE_TABLE", - "SYSTEM_CRON_USER_ID", "ScheduleNotFound", "TriggerError", "build_sqlalchemy_jobstore", diff --git a/common/src/common/scheduler/trigger.py b/common/src/common/scheduler/trigger.py index 2b26179..f2eb132 100644 --- a/common/src/common/scheduler/trigger.py +++ b/common/src/common/scheduler/trigger.py @@ -35,7 +35,6 @@ from common.db.models import ( from common.eventing import add_outbox_event, schedule_event_type, utcnow from common.ids import new_ulid - # A stable user_id used for cron-triggered runs. The corresponding # ``Users`` row is seeded by the auth-bootstrap migration so any audit # query joining on ``ScheduleRuns.triggered_by`` still resolves. @@ -87,7 +86,7 @@ def normalize_idempotency_key( f"Idempotency-Key must contain at least {min_length} characters" ) digest = hashlib.sha256( - f"{workspace_id}:{schedule_id}:{normalized}".encode("utf-8") + f"{workspace_id}:{schedule_id}:{normalized}".encode() ).hexdigest() return f"run:v1:{digest}" @@ -358,10 +357,10 @@ async def create_scheduled_run( __all__ = [ + "SYSTEM_CRON_USER_ID", "DagTooLarge", "InvalidDag", "InvalidNodeArguments", - "SYSTEM_CRON_USER_ID", "ScheduleNotFound", "TriggerError", "create_scheduled_run", diff --git a/common/src/common/service_app.py b/common/src/common/service_app.py index 92f4769..2317dcb 100644 --- a/common/src/common/service_app.py +++ b/common/src/common/service_app.py @@ -1,8 +1,9 @@ from __future__ import annotations import asyncio +from collections.abc import Callable from datetime import UTC, datetime -from typing import Any, Callable +from typing import Any from fastapi import FastAPI, Response, status diff --git a/common/src/common/storage/__init__.py b/common/src/common/storage/__init__.py index e07db1c..4a4906d 100644 --- a/common/src/common/storage/__init__.py +++ b/common/src/common/storage/__init__.py @@ -22,29 +22,29 @@ from .base import AsyncStorageBackend, ObjectMeta, StorageBackend from .factory import ( PURPOSE_BUCKETS, RCLONE_REMOTE_NAME, + USAGE_TYPE_TO_PURPOSE, actual_bucket_name, - build_storage_uri, build_storage_config, + build_storage_uri, create_storage, rclone_remote_spec, - USAGE_TYPE_TO_PURPOSE, workspaces_root, ) from .registry import register_backend, registered_backends __all__ = [ - "create_storage", - "build_storage_config", - "actual_bucket_name", - "build_storage_uri", - "USAGE_TYPE_TO_PURPOSE", - "workspaces_root", - "rclone_remote_spec", - "RCLONE_REMOTE_NAME", "PURPOSE_BUCKETS", - "StorageBackend", + "RCLONE_REMOTE_NAME", + "USAGE_TYPE_TO_PURPOSE", "AsyncStorageBackend", "ObjectMeta", + "StorageBackend", + "actual_bucket_name", + "build_storage_config", + "build_storage_uri", + "create_storage", + "rclone_remote_spec", "register_backend", "registered_backends", + "workspaces_root", ] diff --git a/common/src/common/storage/backends/__init__.py b/common/src/common/storage/backends/__init__.py index 82440f2..97afac8 100644 --- a/common/src/common/storage/backends/__init__.py +++ b/common/src/common/storage/backends/__init__.py @@ -5,5 +5,7 @@ 只要在使用前 import 一次那个模块(让装饰器执行)就够了。 """ -from . import local # noqa: F401 -from . import s3 # noqa: F401 +from . import ( + local, # noqa: F401 + s3, # noqa: F401 +) diff --git a/common/src/common/storage/backends/local.py b/common/src/common/storage/backends/local.py index f3d86e3..9c2785e 100644 --- a/common/src/common/storage/backends/local.py +++ b/common/src/common/storage/backends/local.py @@ -10,9 +10,10 @@ import asyncio import os import shutil +from collections.abc import AsyncIterator, Iterable from datetime import timedelta from pathlib import Path -from typing import AsyncIterator, BinaryIO, Iterable, Optional +from typing import BinaryIO from ..base import AsyncData, AsyncStorageBackend, ObjectMeta, StorageBackend, SyncData from ..exceptions import StorageAlreadyExistsError, StorageNotFoundError @@ -52,8 +53,8 @@ class LocalStorageBackend(StorageBackend): data: SyncData, *, overwrite: bool = True, - content_type: Optional[str] = None, - metadata: Optional[dict] = None, + content_type: str | None = None, + metadata: dict | None = None, ) -> ObjectMeta: path = self._resolve(key) if path.exists() and not overwrite: @@ -107,7 +108,7 @@ class LocalStorageBackend(StorageBackend): key = str(path.relative_to(self.base_dir)).replace(os.sep, "/") yield _meta(key, path) - def get_url(self, key: str, *, expires_in: Optional[timedelta] = None) -> str: + def get_url(self, key: str, *, expires_in: timedelta | None = None) -> str: path = self._resolve(key) if not path.is_file(): raise StorageNotFoundError(f"key 不存在: {key}") @@ -146,8 +147,8 @@ class LocalAsyncStorageBackend(AsyncStorageBackend): data: AsyncData, *, overwrite: bool = True, - content_type: Optional[str] = None, - metadata: Optional[dict] = None, + content_type: str | None = None, + metadata: dict | None = None, ) -> ObjectMeta: import aiofiles @@ -229,7 +230,7 @@ class LocalAsyncStorageBackend(AsyncStorageBackend): return _iter() - async def get_url(self, key: str, *, expires_in: Optional[timedelta] = None) -> str: + async def get_url(self, key: str, *, expires_in: timedelta | None = None) -> str: path = self._resolve(key) if not await asyncio.to_thread(path.is_file): raise StorageNotFoundError(f"key 不存在: {key}") diff --git a/common/src/common/storage/backends/s3.py b/common/src/common/storage/backends/s3.py index ed82e4b..6f067a2 100644 --- a/common/src/common/storage/backends/s3.py +++ b/common/src/common/storage/backends/s3.py @@ -8,8 +8,9 @@ (aioboto3 本身依赖 botocore,异常类型从它里面拿)。 """ +from collections.abc import AsyncIterator, Iterable from datetime import timedelta -from typing import AsyncIterator, BinaryIO, Iterable, Optional +from typing import BinaryIO from ..base import AsyncData, AsyncStorageBackend, ObjectMeta, StorageBackend, SyncData from ..exceptions import ( @@ -49,10 +50,10 @@ class S3StorageBackend(StorageBackend): self, bucket: str, prefix: str = "", - region_name: Optional[str] = None, - endpoint_url: Optional[str] = None, - aws_access_key_id: Optional[str] = None, - aws_secret_access_key: Optional[str] = None, + region_name: str | None = None, + endpoint_url: str | None = None, + aws_access_key_id: str | None = None, + aws_secret_access_key: str | None = None, **_ignored, ): try: @@ -151,7 +152,7 @@ class S3StorageBackend(StorageBackend): except (self._ClientError, self._BotoCoreError) as e: raise StorageConnectionError(f"列举对象失败 prefix={prefix}: {e}") from e - def get_url(self, key: str, *, expires_in: Optional[timedelta] = None) -> str: + def get_url(self, key: str, *, expires_in: timedelta | None = None) -> str: expires_seconds = int(expires_in.total_seconds()) if expires_in else 3600 try: return self.client.generate_presigned_url( @@ -194,10 +195,10 @@ class S3AsyncStorageBackend(AsyncStorageBackend): self, bucket: str, prefix: str = "", - region_name: Optional[str] = None, - endpoint_url: Optional[str] = None, - aws_access_key_id: Optional[str] = None, - aws_secret_access_key: Optional[str] = None, + region_name: str | None = None, + endpoint_url: str | None = None, + aws_access_key_id: str | None = None, + aws_secret_access_key: str | None = None, **_ignored, ): try: @@ -254,8 +255,8 @@ class S3AsyncStorageBackend(AsyncStorageBackend): data: AsyncData, *, overwrite: bool = True, - content_type: Optional[str] = None, - metadata: Optional[dict] = None, + content_type: str | None = None, + metadata: dict | None = None, ) -> ObjectMeta: full_key = self._full_key(key) if not overwrite and await self.exists(key): @@ -388,7 +389,7 @@ class S3AsyncStorageBackend(AsyncStorageBackend): return _iter() - async def get_url(self, key: str, *, expires_in: Optional[timedelta] = None) -> str: + async def get_url(self, key: str, *, expires_in: timedelta | None = None) -> str: expires_seconds = int(expires_in.total_seconds()) if expires_in else 3600 async def _op(client): diff --git a/common/src/common/storage/base.py b/common/src/common/storage/base.py index 2dae299..30b40d3 100644 --- a/common/src/common/storage/base.py +++ b/common/src/common/storage/base.py @@ -10,9 +10,10 @@ """ from abc import ABC, abstractmethod +from collections.abc import AsyncIterator, Iterable from dataclasses import dataclass, field from datetime import timedelta -from typing import AsyncIterator, BinaryIO, Iterable, Optional, Union +from typing import BinaryIO, Union SyncData = Union[bytes, BinaryIO] AsyncData = Union[bytes, "AsyncIterator[bytes]"] @@ -24,8 +25,8 @@ class ObjectMeta: key: str size: int - last_modified: Optional[float] = None # unix timestamp - etag: Optional[str] = None + last_modified: float | None = None # unix timestamp + etag: str | None = None extra: dict = field(default_factory=dict) # 后端特有的额外信息 @@ -39,8 +40,8 @@ class StorageBackend(ABC): data: SyncData, *, overwrite: bool = True, - content_type: Optional[str] = None, - metadata: Optional[dict] = None, + content_type: str | None = None, + metadata: dict | None = None, ) -> ObjectMeta: """写入对象。overwrite=False 时 key 已存在应抛出 StorageAlreadyExistsError。 @@ -72,7 +73,7 @@ class StorageBackend(ABC): """按前缀列出对象。""" @abstractmethod - def get_url(self, key: str, *, expires_in: Optional[timedelta] = None) -> str: + def get_url(self, key: str, *, expires_in: timedelta | None = None) -> str: """获取可访问 URL;本地存储返回 file://,S3 返回预签名 URL。""" def copy(self, src_key: str, dst_key: str) -> ObjectMeta: @@ -82,7 +83,7 @@ class StorageBackend(ABC): def close(self) -> None: """释放后端持有的资源(连接池等)。不需要的后端可以不覆盖。""" - return None + return def __enter__(self) -> "StorageBackend": return self @@ -101,8 +102,8 @@ class AsyncStorageBackend(ABC): data: AsyncData, *, overwrite: bool = True, - content_type: Optional[str] = None, - metadata: Optional[dict] = None, + content_type: str | None = None, + metadata: dict | None = None, ) -> ObjectMeta: """data 可以是 bytes,也可以是异步字节流(async generator)。 @@ -138,7 +139,7 @@ class AsyncStorageBackend(ABC): """用法: `async for meta in backend.list(prefix):`。""" @abstractmethod - async def get_url(self, key: str, *, expires_in: Optional[timedelta] = None) -> str: + async def get_url(self, key: str, *, expires_in: timedelta | None = None) -> str: ... async def copy(self, src_key: str, dst_key: str) -> ObjectMeta: diff --git a/common/src/common/storage/example_usage.py b/common/src/common/storage/example_usage.py index eb2607a..4c2d06a 100644 --- a/common/src/common/storage/example_usage.py +++ b/common/src/common/storage/example_usage.py @@ -62,9 +62,8 @@ def extend_with_new_backend_demo(): """演示独立扩展一种新的存储方式(同步+异步各一个),不用改现有代码。""" import io import time - from typing import AsyncIterator, BinaryIO, Iterable, Optional - from common.storage.base import AsyncStorageBackend, ObjectMeta, StorageBackend + from common.storage.base import ObjectMeta, StorageBackend from common.storage.exceptions import StorageNotFoundError from common.storage.registry import register_backend diff --git a/common/src/common/storage/factory.py b/common/src/common/storage/factory.py index 4878a54..fb746bc 100644 --- a/common/src/common/storage/factory.py +++ b/common/src/common/storage/factory.py @@ -18,17 +18,17 @@ XxxAsyncStorageBackend 类。 """ from pathlib import Path -from typing import Any, Dict, Union +from typing import Any, Union +from .backends import local, s3 # noqa: F401 # 触发内置后端注册 from .base import AsyncStorageBackend, StorageBackend from .exceptions import StorageConfigError from .registry import get_backend_class -from .backends import local, s3 # noqa: F401 # 触发内置后端注册 AnyStorageBackend = Union[StorageBackend, AsyncStorageBackend] -def create_storage(config: Dict[str, Any]) -> AnyStorageBackend: +def create_storage(config: dict[str, Any]) -> AnyStorageBackend: """根据配置创建存储后端。 Args: @@ -111,7 +111,7 @@ def build_storage_uri(bucket_name: str, object_key: str) -> str: return f"s3://{bucket_name}/{object_key}" -def build_storage_config(bucket_name: str) -> Dict[str, Any]: +def build_storage_config(bucket_name: str) -> dict[str, Any]: """根据 ``settings.storage_backend`` 构造 ``create_storage()`` 的入参。 上层(lifespan 等)只用 ``PURPOSE_BUCKETS`` 循环调用一次, diff --git a/common/src/common/storage/registry.py b/common/src/common/storage/registry.py index 56d0b49..2d061cb 100644 --- a/common/src/common/storage/registry.py +++ b/common/src/common/storage/registry.py @@ -10,14 +10,14 @@ 只要保证模块被 import 一次即可(backends/__init__.py 里统一 import)。 """ -from typing import Dict, Tuple, Type, Union +from typing import Union from .base import AsyncStorageBackend, StorageBackend from .exceptions import StorageConfigError -BackendClass = Union[Type[StorageBackend], Type[AsyncStorageBackend]] +BackendClass = Union[type[StorageBackend], type[AsyncStorageBackend]] -_REGISTRY: Dict[Tuple[str, str], BackendClass] = {} +_REGISTRY: dict[tuple[str, str], BackendClass] = {} VALID_MODES = ("sync", "async") @@ -57,5 +57,5 @@ def get_backend_class(name: str, mode: str = "sync") -> BackendClass: ) -def registered_backends() -> Dict[Tuple[str, str], BackendClass]: +def registered_backends() -> dict[tuple[str, str], BackendClass]: return dict(_REGISTRY) diff --git a/common/src/common/storage/schemas.py b/common/src/common/storage/schemas.py index 8e6403a..b3d7b35 100644 --- a/common/src/common/storage/schemas.py +++ b/common/src/common/storage/schemas.py @@ -8,7 +8,6 @@ from pydantic import Field, field_validator from common.schemas import StrictModel - __all__ = [ "CreateUploadRequest", "DownloadUrlRequest", diff --git a/migrations/data/migrate_legacy_workspaces.py b/migrations/data/migrate_legacy_workspaces.py index da081c7..3bd59ff 100644 --- a/migrations/data/migrate_legacy_workspaces.py +++ b/migrations/data/migrate_legacy_workspaces.py @@ -10,10 +10,6 @@ from pathlib import Path from typing import Any from urllib.parse import quote -from pydantic import BaseModel, ConfigDict, TypeAdapter -from sqlalchemy import select -from sqlalchemy.ext.asyncio import AsyncSession - from common.db import create_database_engine, create_session_factory from common.db.models import ( Roles, @@ -21,6 +17,10 @@ from common.db.models import ( WorkspaceMembers, Workspaces, ) +from pydantic import BaseModel, ConfigDict, TypeAdapter +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession + from migrations.data.migrate_system_json import ( deterministic_legacy_ulid, set_changed, diff --git a/migrations/data/migrate_system_json.py b/migrations/data/migrate_system_json.py index a608ddc..1c606ae 100644 --- a/migrations/data/migrate_system_json.py +++ b/migrations/data/migrate_system_json.py @@ -9,10 +9,6 @@ from pathlib import Path from typing import Any from urllib.parse import quote -from pydantic import BaseModel, ConfigDict -from sqlalchemy import select -from sqlalchemy.ext.asyncio import AsyncSession - from common.db import create_database_engine, create_session_factory from common.db.models import ( Permissions, @@ -20,6 +16,9 @@ from common.db.models import ( Roles, Users, ) +from pydantic import BaseModel, ConfigDict +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession CROCKFORD_BASE32 = "0123456789ABCDEFGHJKMNPQRSTVWXYZ" LEGACY_PASSWORD_HASH = "!legacy-account-without-password!" @@ -76,7 +75,7 @@ class LegacySystem(BaseModel): def deterministic_legacy_ulid(entity_type: str, legacy_key: str) -> str: """Create a stable ULID-compatible ID with an epoch timestamp prefix.""" digest = hashlib.sha256( - f"model-platform-v1:{entity_type}:{legacy_key}".encode("utf-8") + f"model-platform-v1:{entity_type}:{legacy_key}".encode() ).digest() value = int.from_bytes(b"\x00" * 6 + digest[:10], byteorder="big") encoded = ["0"] * 26 diff --git a/migrations/env.py b/migrations/env.py index 4b06419..3c1de25 100644 --- a/migrations/env.py +++ b/migrations/env.py @@ -7,11 +7,10 @@ from logging.config import fileConfig from typing import Any from alembic import context +from common.db import Base from sqlalchemy import Connection, pool from sqlalchemy.ext.asyncio import async_engine_from_config -from common.db import Base - config = context.config if config.config_file_name is not None: diff --git a/migrations/versions/e1f2a3b4c5d6_rebuild_baseline.py b/migrations/versions/e1f2a3b4c5d6_rebuild_baseline.py index bc7caaa..d17f696 100644 --- a/migrations/versions/e1f2a3b4c5d6_rebuild_baseline.py +++ b/migrations/versions/e1f2a3b4c5d6_rebuild_baseline.py @@ -5,15 +5,14 @@ Revises: (none) Create Date: 2026-08-14 """ -from collections.abc import Sequence import hashlib import os +from collections.abc import Sequence -from alembic import op import sqlalchemy as sa -from sqlalchemy.dialects import mysql - +from alembic import op from common.auth.passwords import hash_password +from sqlalchemy.dialects import mysql # revision identifiers, used by Alembic. revision: str = "e1f2a3b4c5d6" @@ -68,7 +67,7 @@ CROCKFORD_BASE32 = "0123456789ABCDEFGHJKMNPQRSTVWXYZ" def _deterministic_permission_id(code: str) -> str: """Stable 26-char ULID-shaped id derived from permission_code.""" digest = hashlib.sha256( - f"model-platform-permission-v1:{code}".encode("utf-8") + f"model-platform-permission-v1:{code}".encode() ).digest() value = int.from_bytes(b"\x00" * 6 + digest[:10], byteorder="big") encoded = ["0"] * 26 diff --git a/schedule/src/schedule/context.py b/schedule/src/schedule/context.py index 5f67123..143e5c9 100644 --- a/schedule/src/schedule/context.py +++ b/schedule/src/schedule/context.py @@ -8,7 +8,6 @@ from __future__ import annotations from datetime import UTC, datetime - TERMINAL_NODE_STATES = frozenset({ "succeeded", "failed", @@ -35,8 +34,8 @@ def naive_utc(value: datetime | None) -> datetime | None: __all__ = [ - "TERMINAL_NODE_STATES", "FAILED_NODE_STATES", + "TERMINAL_NODE_STATES", "TERMINAL_RUN_STATES", "naive_utc", ] diff --git a/schedule/src/schedule/execution.py b/schedule/src/schedule/execution.py index 7d9a29f..deae579 100644 --- a/schedule/src/schedule/execution.py +++ b/schedule/src/schedule/execution.py @@ -9,7 +9,6 @@ from pathlib import Path, PurePosixPath from loguru import logger - MAX_LOG_BYTES = 4 * 1024 * 1024 diff --git a/schedule/src/schedule/executor.py b/schedule/src/schedule/executor.py index f0ff52a..3811ca1 100644 --- a/schedule/src/schedule/executor.py +++ b/schedule/src/schedule/executor.py @@ -1,4 +1,3 @@ -# coding=utf-8 """ @Time :2026/7/29 @Author :tao.chen diff --git a/schedule/src/schedule/main.py b/schedule/src/schedule/main.py index b4e174a..237816e 100644 --- a/schedule/src/schedule/main.py +++ b/schedule/src/schedule/main.py @@ -1,13 +1,14 @@ from __future__ import annotations +from collections.abc import AsyncIterator from contextlib import asynccontextmanager -from typing import Any, AsyncIterator - -from loguru import logger +from typing import Any from common.config import settings from common.db import create_database_engine, create_session_factory from common.service_app import create_service_app +from loguru import logger + from schedule.service import ( SchedulerService, build_object_store, diff --git a/schedule/src/schedule/orchestrator.py b/schedule/src/schedule/orchestrator.py index 2350f2a..61fc82b 100644 --- a/schedule/src/schedule/orchestrator.py +++ b/schedule/src/schedule/orchestrator.py @@ -16,10 +16,6 @@ import asyncio from datetime import timedelta from typing import Any -from loguru import logger -from sqlalchemy import select -from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker - from common.db import session_scope from common.db.models import ( ConsumerInbox, @@ -29,6 +25,9 @@ from common.db.models import ( ) from common.eventing import add_outbox_event, schedule_event_type, utcnow from common.ids import new_ulid +from loguru import logger +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker from schedule.context import ( FAILED_NODE_STATES, diff --git a/schedule/src/schedule/scheduler.py b/schedule/src/schedule/scheduler.py index 230ed8c..5d03941 100644 --- a/schedule/src/schedule/scheduler.py +++ b/schedule/src/schedule/scheduler.py @@ -15,20 +15,18 @@ restarts. from __future__ import annotations import asyncio +from collections.abc import Awaitable, Callable from datetime import UTC -from typing import Awaitable, Callable from zoneinfo import ZoneInfo -from loguru import logger - from apscheduler.schedulers.asyncio import AsyncIOScheduler from apscheduler.triggers.cron import CronTrigger -from sqlalchemy import select -from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker - from common.db import session_scope from common.db.models import Schedules from common.scheduler import build_sqlalchemy_jobstore +from loguru import logger +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker from schedule.context import naive_utc diff --git a/schedule/src/schedule/service.py b/schedule/src/schedule/service.py index c14a73d..6ee2a5b 100644 --- a/schedule/src/schedule/service.py +++ b/schedule/src/schedule/service.py @@ -25,17 +25,16 @@ from pathlib import Path from typing import Any import httpx -from loguru import logger -from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker - from common.config import settings +from common.ids import new_ulid from common.scheduler import ( SYSTEM_CRON_USER_ID, TriggerError, create_scheduled_run, ) -from common.ids import new_ulid from common.storage import create_storage +from loguru import logger +from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker from schedule.orchestrator import DispatchOrchestrator from schedule.scheduler import CronScheduler @@ -129,7 +128,6 @@ class SchedulerService: via the unique constraint on ``schedule_runs.idempotency_key``. """ async with self.session_factory() as session: - from sqlalchemy import select from common.db.models import Schedules diff --git a/schedule/src/schedule/storage_client.py b/schedule/src/schedule/storage_client.py index 25dd58d..72a97a9 100644 --- a/schedule/src/schedule/storage_client.py +++ b/schedule/src/schedule/storage_client.py @@ -15,9 +15,8 @@ from __future__ import annotations import base64 from typing import Any -from loguru import logger - import httpx +from loguru import logger class SchedulerStorageClient: diff --git a/schedule/src/schedule/worker.py b/schedule/src/schedule/worker.py index e727f8a..4fd9e71 100644 --- a/schedule/src/schedule/worker.py +++ b/schedule/src/schedule/worker.py @@ -13,16 +13,11 @@ worker only mutates ``ScheduleNodeRuns`` columns related to execution from __future__ import annotations -import asyncio import hashlib import json import traceback -from pathlib import Path from typing import Any -from loguru import logger -from sqlalchemy import select - from common.config import settings from common.db import session_scope from common.db.models import ( @@ -41,6 +36,8 @@ from common.eventing import ( schedule_event_type, utcnow, ) +from loguru import logger +from sqlalchemy import select from schedule.context import TERMINAL_NODE_STATES from schedule.execution import ExecutionResult, execute_artifact