From 99e72cdcba88ed775376c146e6cbb71d500c78cb Mon Sep 17 00:00:00 2001 From: Xiaoxia AI Date: Mon, 15 Jun 2026 19:34:47 +0800 Subject: [PATCH] feat: persist ingest flow in postgres --- apps/api/app/api/routes/asset_libraries.py | 6 +-- apps/api/app/api/routes/assets.py | 6 +-- apps/api/app/api/routes/ingest_jobs.py | 4 +- apps/api/app/api/routes/projects.py | 6 +-- apps/api/app/api/routes/upload.py | 6 +-- apps/api/app/db.py | 13 +++--- apps/api/app/dependencies.py | 34 +++++++------- apps/worker/worker_app/core/config.py | 5 ++ apps/worker/worker_app/db.py | 12 +++++ apps/worker/worker_app/tasks/ingest.py | 29 ++++++------ packages/adapters/sqlalchemy_impl/__init__.py | 17 +++++++ .../asset_library_repository.py | 16 +++---- .../sqlalchemy_impl/asset_repository.py | 23 +++++----- .../sqlalchemy_impl/ingest_job_repository.py | 2 +- .../sqlalchemy_impl/project_repository.py | 14 +++--- packages/adapters/sqlalchemy_impl/session.py | 46 +++++++++++++++++++ 16 files changed, 161 insertions(+), 78 deletions(-) create mode 100644 apps/worker/worker_app/db.py create mode 100644 packages/adapters/sqlalchemy_impl/session.py diff --git a/apps/api/app/api/routes/asset_libraries.py b/apps/api/app/api/routes/asset_libraries.py index f26e6394a..ef1cfbcb4 100644 --- a/apps/api/app/api/routes/asset_libraries.py +++ b/apps/api/app/api/routes/asset_libraries.py @@ -2,7 +2,7 @@ from fastapi import APIRouter, Depends from app.dependencies import get_asset_library_repository from app.schemas.asset_library import AssetLibraryResponse, CreateAssetLibraryRequest, ListAssetLibrariesResponse -from packages.adapters.in_memory import InMemoryAssetLibraryRepository +from packages.adapters.sqlalchemy_impl import SQLAlchemyAssetLibraryRepository from packages.application import CreateAssetLibraryCommand, CreateAssetLibraryUseCase, ListAssetLibrariesUseCase from packages.domain import AssetLibraryKind @@ -13,7 +13,7 @@ router = APIRouter() def list_asset_libraries( project_id: str, kind: str | None = None, - asset_library_repository: InMemoryAssetLibraryRepository = Depends(get_asset_library_repository), + asset_library_repository: SQLAlchemyAssetLibraryRepository = Depends(get_asset_library_repository), ) -> ListAssetLibrariesResponse: use_case = ListAssetLibrariesUseCase(asset_library_repository) parsed_kind = AssetLibraryKind(kind) if kind else None @@ -35,7 +35,7 @@ def list_asset_libraries( @router.post("", response_model=AssetLibraryResponse) def create_asset_library( request: CreateAssetLibraryRequest, - asset_library_repository: InMemoryAssetLibraryRepository = Depends(get_asset_library_repository), + asset_library_repository: SQLAlchemyAssetLibraryRepository = Depends(get_asset_library_repository), ) -> AssetLibraryResponse: use_case = CreateAssetLibraryUseCase(asset_library_repository) item = use_case.execute( diff --git a/apps/api/app/api/routes/assets.py b/apps/api/app/api/routes/assets.py index e222729aa..95a22cd74 100644 --- a/apps/api/app/api/routes/assets.py +++ b/apps/api/app/api/routes/assets.py @@ -2,7 +2,7 @@ from fastapi import APIRouter, Depends from app.dependencies import get_asset_repository from app.schemas.asset import AssetResponse, CreateAssetRequest, ListAssetsResponse -from packages.adapters.in_memory import InMemoryAssetRepository +from packages.adapters.sqlalchemy_impl import SQLAlchemyAssetRepository from packages.application import CreateAssetCommand, CreateAssetUseCase, ListAssetsUseCase router = APIRouter() @@ -11,7 +11,7 @@ router = APIRouter() @router.get("", response_model=ListAssetsResponse) def list_assets( library_id: str, - asset_repository: InMemoryAssetRepository = Depends(get_asset_repository), + asset_repository: SQLAlchemyAssetRepository = Depends(get_asset_repository), ) -> ListAssetsResponse: use_case = ListAssetsUseCase(asset_repository) items = use_case.execute(library_id) @@ -35,7 +35,7 @@ def list_assets( @router.post("", response_model=AssetResponse) def create_asset( request: CreateAssetRequest, - asset_repository: InMemoryAssetRepository = Depends(get_asset_repository), + asset_repository: SQLAlchemyAssetRepository = Depends(get_asset_repository), ) -> AssetResponse: use_case = CreateAssetUseCase(asset_repository) item = use_case.execute( diff --git a/apps/api/app/api/routes/ingest_jobs.py b/apps/api/app/api/routes/ingest_jobs.py index 0e34ac45d..04ad38749 100644 --- a/apps/api/app/api/routes/ingest_jobs.py +++ b/apps/api/app/api/routes/ingest_jobs.py @@ -3,7 +3,7 @@ from fastapi import APIRouter, Depends from app.core.celery_app import celery_app from app.dependencies import get_ingest_job_repository from app.schemas.ingest_job import IngestJobResponse, SubmitIngestJobRequest -from packages.adapters.in_memory import InMemoryIngestJobRepository +from packages.adapters.sqlalchemy_impl import SQLAlchemyIngestJobRepository from packages.application import SubmitIngestJobCommand, SubmitIngestJobUseCase router = APIRouter() @@ -12,7 +12,7 @@ router = APIRouter() @router.post("", response_model=IngestJobResponse) def submit_ingest_job( request: SubmitIngestJobRequest, - ingest_job_repository: InMemoryIngestJobRepository = Depends(get_ingest_job_repository), + ingest_job_repository: SQLAlchemyIngestJobRepository = Depends(get_ingest_job_repository), ) -> IngestJobResponse: use_case = SubmitIngestJobUseCase(ingest_job_repository) job = use_case.execute( diff --git a/apps/api/app/api/routes/projects.py b/apps/api/app/api/routes/projects.py index 8af550475..245f6254d 100644 --- a/apps/api/app/api/routes/projects.py +++ b/apps/api/app/api/routes/projects.py @@ -2,8 +2,8 @@ from fastapi import APIRouter, Depends from app.dependencies import get_project_repository from app.schemas.project import CreateProjectRequest, ListProjectsResponse, ProjectResponse +from packages.adapters.sqlalchemy_impl import SQLAlchemyProjectRepository from packages.application import CreateProjectCommand, CreateProjectUseCase, ListProjectsUseCase -from packages.adapters.in_memory import InMemoryProjectRepository router = APIRouter() @@ -11,7 +11,7 @@ router = APIRouter() @router.get("", response_model=ListProjectsResponse) def list_projects( workspace_id: str, - project_repository: InMemoryProjectRepository = Depends(get_project_repository), + project_repository: SQLAlchemyProjectRepository = Depends(get_project_repository), ) -> ListProjectsResponse: use_case = ListProjectsUseCase(project_repository) projects = use_case.execute(workspace_id) @@ -31,7 +31,7 @@ def list_projects( @router.post("", response_model=ProjectResponse) def create_project( request: CreateProjectRequest, - project_repository: InMemoryProjectRepository = Depends(get_project_repository), + project_repository: SQLAlchemyProjectRepository = Depends(get_project_repository), ) -> ProjectResponse: use_case = CreateProjectUseCase(project_repository) project = use_case.execute( diff --git a/apps/api/app/api/routes/upload.py b/apps/api/app/api/routes/upload.py index e6ba2d5b2..c0a3d2820 100644 --- a/apps/api/app/api/routes/upload.py +++ b/apps/api/app/api/routes/upload.py @@ -3,9 +3,9 @@ from uuid import uuid4 from app.core.celery_app import celery_app from app.dependencies import get_ingest_job_repository -from app.core.storage import get_minio_service, MinIOService +from app.core.storage import MinIOService, get_minio_service from app.schemas.upload import UploadAssetResponse -from packages.adapters.in_memory import InMemoryIngestJobRepository +from packages.adapters.sqlalchemy_impl import SQLAlchemyIngestJobRepository from packages.application import SubmitIngestJobCommand, SubmitIngestJobUseCase router = APIRouter() @@ -17,7 +17,7 @@ async def upload_asset( workspace_id: str = Form(..., description="工作空间 ID"), project_id: str = Form(..., description="项目 ID"), library_id: str = Form(..., description="资产库 ID"), - ingest_job_repository: InMemoryIngestJobRepository = Depends(get_ingest_job_repository), + ingest_job_repository: SQLAlchemyIngestJobRepository = Depends(get_ingest_job_repository), storage_service: MinIOService = Depends(get_minio_service), ) -> UploadAssetResponse: """ diff --git a/apps/api/app/db.py b/apps/api/app/db.py index 32f6cb3d7..93980d99f 100644 --- a/apps/api/app/db.py +++ b/apps/api/app/db.py @@ -1,21 +1,22 @@ -from sqlalchemy import create_engine -from sqlalchemy.orm import sessionmaker, Session +from collections.abc import Generator + +from sqlalchemy.orm import Session from app.core.database import get_database_settings +from packages.adapters.sqlalchemy_impl import build_session_factory, initialize_database settings = get_database_settings() -engine = create_engine( +engine, SessionLocal = build_session_factory( settings.database_url, pool_size=settings.pool_size, max_overflow=settings.max_overflow, pool_timeout=settings.pool_timeout, pool_recycle=settings.pool_recycle, ) -SessionLocal = sessionmaker(autocommit=False, autoflush=False, bind=engine) +initialize_database(engine) -def get_db() -> Session: - """Dependency for database session.""" +def get_db() -> Generator[Session, None, None]: db = SessionLocal() try: yield db diff --git a/apps/api/app/dependencies.py b/apps/api/app/dependencies.py index 63f583736..dcb94dba2 100644 --- a/apps/api/app/dependencies.py +++ b/apps/api/app/dependencies.py @@ -1,28 +1,26 @@ -from functools import lru_cache +from fastapi import Depends +from sqlalchemy.orm import Session -from packages.adapters.in_memory import ( - InMemoryAssetLibraryRepository, - InMemoryAssetRepository, - InMemoryIngestJobRepository, - InMemoryProjectRepository, +from app.db import get_db +from packages.adapters.sqlalchemy_impl import ( + SQLAlchemyAssetLibraryRepository, + SQLAlchemyAssetRepository, + SQLAlchemyIngestJobRepository, + SQLAlchemyProjectRepository, ) -@lru_cache(maxsize=1) -def get_project_repository() -> InMemoryProjectRepository: - return InMemoryProjectRepository() +def get_project_repository(db: Session = Depends(get_db)) -> SQLAlchemyProjectRepository: + return SQLAlchemyProjectRepository(db) -@lru_cache(maxsize=1) -def get_asset_library_repository() -> InMemoryAssetLibraryRepository: - return InMemoryAssetLibraryRepository() +def get_asset_library_repository(db: Session = Depends(get_db)) -> SQLAlchemyAssetLibraryRepository: + return SQLAlchemyAssetLibraryRepository(db) -@lru_cache(maxsize=1) -def get_asset_repository() -> InMemoryAssetRepository: - return InMemoryAssetRepository() +def get_asset_repository(db: Session = Depends(get_db)) -> SQLAlchemyAssetRepository: + return SQLAlchemyAssetRepository(db) -@lru_cache(maxsize=1) -def get_ingest_job_repository() -> InMemoryIngestJobRepository: - return InMemoryIngestJobRepository() +def get_ingest_job_repository(db: Session = Depends(get_db)) -> SQLAlchemyIngestJobRepository: + return SQLAlchemyIngestJobRepository(db) diff --git a/apps/worker/worker_app/core/config.py b/apps/worker/worker_app/core/config.py index 9b395ae1d..a79244bdd 100644 --- a/apps/worker/worker_app/core/config.py +++ b/apps/worker/worker_app/core/config.py @@ -9,6 +9,11 @@ class WorkerSettings(BaseSettings): result_backend: str = "redis://redis:6379/1" worker_concurrency: int = 4 worker_max_tasks_per_child: int = 1000 + database_url: str = "postgresql+psycopg://postgres:postgres@postgres:5432/xiaoxia_saas" + database_pool_size: int = 20 + database_max_overflow: int = 40 + database_pool_timeout: int = 30 + database_pool_recycle: int = 3600 model_config = SettingsConfigDict( env_file=".env", diff --git a/apps/worker/worker_app/db.py b/apps/worker/worker_app/db.py new file mode 100644 index 000000000..ca58fe311 --- /dev/null +++ b/apps/worker/worker_app/db.py @@ -0,0 +1,12 @@ +from worker_app.core.config import get_settings +from packages.adapters.sqlalchemy_impl import build_session_factory, initialize_database + +settings = get_settings() +engine, SessionLocal = build_session_factory( + settings.database_url, + pool_size=settings.database_pool_size, + max_overflow=settings.database_max_overflow, + pool_timeout=settings.database_pool_timeout, + pool_recycle=settings.database_pool_recycle, +) +initialize_database(engine) diff --git a/apps/worker/worker_app/tasks/ingest.py b/apps/worker/worker_app/tasks/ingest.py index 007f94e5f..a49bf3e22 100644 --- a/apps/worker/worker_app/tasks/ingest.py +++ b/apps/worker/worker_app/tasks/ingest.py @@ -1,11 +1,9 @@ +from datetime import datetime, timezone + +from packages.adapters.sqlalchemy_impl import SQLAlchemyAssetRepository, SQLAlchemyIngestJobRepository +from packages.domain import Asset, IngestJobStatus from worker_app.celery_app import celery_app -from packages.domain import Asset, IngestJob, IngestJobStatus -from packages.adapters.in_memory import InMemoryAssetRepository, InMemoryIngestJobRepository - - -@celery_app.task(name="worker.healthcheck") -def healthcheck() -> dict: - return {"ok": True, "service": "worker"} +from worker_app.db import SessionLocal @celery_app.task(name="worker.ingest_asset") @@ -20,10 +18,10 @@ def ingest_asset(job_id: str) -> dict: 4. Update IngestJob status to COMPLETED 5. Return result """ - # TODO: Replace with real repository injection - job_repo = InMemoryIngestJobRepository() - asset_repo = InMemoryAssetRepository() - + db = SessionLocal() + job_repo = SQLAlchemyIngestJobRepository(db) + asset_repo = SQLAlchemyAssetRepository(db) + job = job_repo.get(job_id) if job is None: return {"status": "failed", "error": "job not found"} @@ -31,6 +29,7 @@ def ingest_asset(job_id: str) -> dict: try: # Update job status to PROCESSING job.status = IngestJobStatus.PROCESSING + job.updated_at = datetime.now(timezone.utc) job_repo.update(job) # Mock metadata extraction (in real implementation: use ffprobe, Pillow, etc.) @@ -60,8 +59,9 @@ def ingest_asset(job_id: str) -> dict: # Update job status to COMPLETED job.status = IngestJobStatus.COMPLETED job.result_asset_id = asset.id + job.updated_at = datetime.now(timezone.utc) job_repo.update(job) - + return { "status": "completed", "job_id": job.id, @@ -71,10 +71,13 @@ def ingest_asset(job_id: str) -> dict: # Update job status to FAILED job.status = IngestJobStatus.FAILED job.error_message = str(e) + job.updated_at = datetime.now(timezone.utc) job_repo.update(job) - + return { "status": "failed", "job_id": job.id, "error": str(e), } + finally: + db.close() diff --git a/packages/adapters/sqlalchemy_impl/__init__.py b/packages/adapters/sqlalchemy_impl/__init__.py index 60e96fbbb..53a79cc54 100644 --- a/packages/adapters/sqlalchemy_impl/__init__.py +++ b/packages/adapters/sqlalchemy_impl/__init__.py @@ -1 +1,18 @@ """SQLAlchemy-based repository implementations.""" + +from .asset_library_repository import SQLAlchemyAssetLibraryRepository +from .asset_repository import SQLAlchemyAssetRepository +from .ingest_job_repository import SQLAlchemyIngestJobRepository +from .project_repository import SQLAlchemyProjectRepository +from .session import Base, build_engine, build_session_factory, initialize_database + +__all__ = [ + "Base", + "SQLAlchemyAssetLibraryRepository", + "SQLAlchemyAssetRepository", + "SQLAlchemyIngestJobRepository", + "SQLAlchemyProjectRepository", + "build_engine", + "build_session_factory", + "initialize_database", +] diff --git a/packages/adapters/sqlalchemy_impl/asset_library_repository.py b/packages/adapters/sqlalchemy_impl/asset_library_repository.py index a4bc0d7a0..c4e25afb0 100644 --- a/packages/adapters/sqlalchemy_impl/asset_library_repository.py +++ b/packages/adapters/sqlalchemy_impl/asset_library_repository.py @@ -1,7 +1,7 @@ from sqlalchemy.orm import Session -from packages.domain import AssetLibrary, AssetLibraryKind from packages.adapters.sqlalchemy_impl.models import AssetLibraryModel +from packages.domain import AssetLibrary, AssetLibraryKind class SQLAlchemyAssetLibraryRepository: @@ -15,14 +15,14 @@ class SQLAlchemyAssetLibraryRepository: models = query.all() return [ AssetLibrary( - id=m.id, - workspace_id=m.workspace_id, - project_id=m.project_id, - name=m.name, - kind=AssetLibraryKind(m.kind), - created_at=m.created_at, + id=model.id, + workspace_id=model.workspace_id, + project_id=model.project_id, + name=model.name, + kind=AssetLibraryKind(model.kind), + created_at=model.created_at, ) - for m in models + for model in models ] def create(self, library: AssetLibrary) -> AssetLibrary: diff --git a/packages/adapters/sqlalchemy_impl/asset_repository.py b/packages/adapters/sqlalchemy_impl/asset_repository.py index f365f7928..e0e6ed0b3 100644 --- a/packages/adapters/sqlalchemy_impl/asset_repository.py +++ b/packages/adapters/sqlalchemy_impl/asset_repository.py @@ -1,8 +1,9 @@ import json + from sqlalchemy.orm import Session -from packages.domain import Asset from packages.adapters.sqlalchemy_impl.models import AssetModel +from packages.domain import Asset class SQLAlchemyAssetRepository: @@ -13,17 +14,17 @@ class SQLAlchemyAssetRepository: models = self.session.query(AssetModel).filter(AssetModel.library_id == library_id).all() return [ Asset( - id=m.id, - workspace_id=m.workspace_id, - project_id=m.project_id, - library_id=m.library_id, - name=m.name, - storage_key=m.storage_key, - mime_type=m.mime_type, - metadata=json.loads(m.metadata_json), - created_at=m.created_at, + id=model.id, + workspace_id=model.workspace_id, + project_id=model.project_id, + library_id=model.library_id, + name=model.name, + storage_key=model.storage_key, + mime_type=model.mime_type, + metadata=json.loads(model.metadata_json), + created_at=model.created_at, ) - for m in models + for model in models ] def create(self, asset: Asset) -> Asset: diff --git a/packages/adapters/sqlalchemy_impl/ingest_job_repository.py b/packages/adapters/sqlalchemy_impl/ingest_job_repository.py index 1ee98409c..141231f48 100644 --- a/packages/adapters/sqlalchemy_impl/ingest_job_repository.py +++ b/packages/adapters/sqlalchemy_impl/ingest_job_repository.py @@ -1,7 +1,7 @@ from sqlalchemy.orm import Session -from packages.domain import IngestJob, IngestJobStatus from packages.adapters.sqlalchemy_impl.models import IngestJobModel +from packages.domain import IngestJob, IngestJobStatus class SQLAlchemyIngestJobRepository: diff --git a/packages/adapters/sqlalchemy_impl/project_repository.py b/packages/adapters/sqlalchemy_impl/project_repository.py index 9a825928b..02c614841 100644 --- a/packages/adapters/sqlalchemy_impl/project_repository.py +++ b/packages/adapters/sqlalchemy_impl/project_repository.py @@ -1,7 +1,7 @@ from sqlalchemy.orm import Session -from packages.domain import Project from packages.adapters.sqlalchemy_impl.models import ProjectModel +from packages.domain import Project class SQLAlchemyProjectRepository: @@ -12,13 +12,13 @@ class SQLAlchemyProjectRepository: models = self.session.query(ProjectModel).filter(ProjectModel.workspace_id == workspace_id).all() return [ Project( - id=m.id, - workspace_id=m.workspace_id, - name=m.name, - description=m.description, - created_at=m.created_at, + id=model.id, + workspace_id=model.workspace_id, + name=model.name, + description=model.description, + created_at=model.created_at, ) - for m in models + for model in models ] def create(self, project: Project) -> Project: diff --git a/packages/adapters/sqlalchemy_impl/session.py b/packages/adapters/sqlalchemy_impl/session.py new file mode 100644 index 000000000..3a8614b8c --- /dev/null +++ b/packages/adapters/sqlalchemy_impl/session.py @@ -0,0 +1,46 @@ +from __future__ import annotations + +from sqlalchemy import create_engine +from sqlalchemy.orm import Session, sessionmaker + +from packages.adapters.sqlalchemy_impl.models import Base + + +def build_engine( + database_url: str, + *, + pool_size: int = 20, + max_overflow: int = 40, + pool_timeout: int = 30, + pool_recycle: int = 3600, +): + return create_engine( + database_url, + pool_size=pool_size, + max_overflow=max_overflow, + pool_timeout=pool_timeout, + pool_recycle=pool_recycle, + ) + + +def build_session_factory( + database_url: str, + *, + pool_size: int = 20, + max_overflow: int = 40, + pool_timeout: int = 30, + pool_recycle: int = 3600, +): + engine = build_engine( + database_url, + pool_size=pool_size, + max_overflow=max_overflow, + pool_timeout=pool_timeout, + pool_recycle=pool_recycle, + ) + session_factory = sessionmaker(autocommit=False, autoflush=False, bind=engine) + return engine, session_factory + + +def initialize_database(engine) -> None: + Base.metadata.create_all(bind=engine)