update: ruff check --fix
This commit is contained in:
@@ -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"])
|
||||
|
||||
|
||||
|
||||
@@ -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"])
|
||||
|
||||
|
||||
@@ -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"
|
||||
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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",
|
||||
]
|
||||
@@ -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,
|
||||
|
||||
@@ -6,6 +6,7 @@ from typing import Any
|
||||
import httpx
|
||||
from loguru import logger
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class RuntimeClientError(Exception):
|
||||
status_code: int
|
||||
|
||||
@@ -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[
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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"])
|
||||
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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}"
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -15,7 +15,6 @@ from __future__ import annotations
|
||||
import httpx
|
||||
import pytest
|
||||
import respx
|
||||
|
||||
from backend.runtime_client import RuntimeClient, RuntimeClientError
|
||||
|
||||
WORKSPACE_ID = "01HWS0000000000000000000A"
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -17,7 +17,6 @@ import secrets
|
||||
|
||||
from passlib.context import CryptContext
|
||||
|
||||
|
||||
_crypt_context = CryptContext(schemes=["bcrypt"], deprecated="auto")
|
||||
|
||||
|
||||
|
||||
@@ -1,4 +1,3 @@
|
||||
# coding=utf-8
|
||||
"""
|
||||
@Time :2026/7/29
|
||||
@Author :tao.chen
|
||||
|
||||
@@ -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))
|
||||
|
||||
@@ -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))
|
||||
|
||||
@@ -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))
|
||||
|
||||
@@ -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))
|
||||
|
||||
@@ -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))
|
||||
|
||||
@@ -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))
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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",
|
||||
]
|
||||
|
||||
@@ -5,5 +5,7 @@
|
||||
只要在使用前 import 一次那个模块(让装饰器执行)就够了。
|
||||
"""
|
||||
|
||||
from . import local # noqa: F401
|
||||
from . import s3 # noqa: F401
|
||||
from . import (
|
||||
local, # noqa: F401
|
||||
s3, # noqa: F401
|
||||
)
|
||||
|
||||
@@ -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}")
|
||||
|
||||
@@ -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):
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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`` 循环调用一次,
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -8,7 +8,6 @@ from pydantic import Field, field_validator
|
||||
|
||||
from common.schemas import StrictModel
|
||||
|
||||
|
||||
__all__ = [
|
||||
"CreateUploadRequest",
|
||||
"DownloadUrlRequest",
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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
|
||||
|
||||
+1
-2
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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",
|
||||
]
|
||||
|
||||
@@ -9,7 +9,6 @@ from pathlib import Path, PurePosixPath
|
||||
|
||||
from loguru import logger
|
||||
|
||||
|
||||
MAX_LOG_BYTES = 4 * 1024 * 1024
|
||||
|
||||
|
||||
|
||||
@@ -1,4 +1,3 @@
|
||||
# coding=utf-8
|
||||
"""
|
||||
@Time :2026/7/29
|
||||
@Author :tao.chen
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user