feat: persist ingest flow in postgres
This commit is contained in:
@@ -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(
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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:
|
||||
"""
|
||||
|
||||
+7
-6
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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)
|
||||
@@ -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()
|
||||
|
||||
@@ -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",
|
||||
]
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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)
|
||||
Reference in New Issue
Block a user