Files
xiaoxia-saas/packages/adapters/sqlalchemy_impl/ingest_job_repository.py
T
xiaoxia f9f3e6bfb9
CI/CD Pipeline / Check if frontend-only change (push) Has been skipped
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (push) Successful in 2s
CI/CD Pipeline / Check push changed paths (pull_request) Has been skipped
CI/CD Pipeline / Dedup Check - skip PR tests when covered by push pipeline (pull_request) Successful in 3s
CI/CD Pipeline / PR Build Web Image (push) Has been skipped
CI/CD Pipeline / PR Build API Image (push) Has been skipped
CI/CD Pipeline / PR Build Worker Image (push) Has been skipped
CI/CD Pipeline / Check if frontend-only change (pull_request) Successful in 2s
CI/CD Pipeline / Frontend Lint (push) Has been skipped
CI/CD Pipeline / Build Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Check push changed paths (push) Successful in 10s
CI/CD Pipeline / Validate - Style (pull_request) Has been skipped
CI/CD Pipeline / Validate - Security (pull_request) Has been skipped
CI/CD Pipeline / Validate - Python (mypy + alembic) (pull_request) Has been skipped
CI/CD Pipeline / Unit Tests (pull_request) Has been skipped
CI/CD Pipeline / Integration Tests (pull_request) Has been skipped
CI/CD Pipeline / Frontend Lint (pull_request) Has been skipped
CI/CD Pipeline / Frontend Unit Tests (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging API Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Web Image (pull_request) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (pull_request) Has been skipped
CI/CD Pipeline / PR Build API Image (pull_request) Successful in 50s
CI/CD Pipeline / Build Staging API Image (push) Successful in 1m12s
CI/CD Pipeline / Build Production API Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Web Image (pull_request) Has been skipped
CI/CD Pipeline / Build Production Worker Image (pull_request) Has been skipped
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (pull_request) Has been skipped
CI/CD Pipeline / Staging E2E Tests (pull_request) Has been skipped
CI/CD Pipeline / Staging API Integration Tests (pull_request) Has been skipped
CI/CD Pipeline / Deploy Production (pull_request) Has been skipped
CI/CD Pipeline / ACR Image Cleanup (pull_request) Has been skipped
CI/CD Pipeline / Production Browser E2E (pull_request) Has been skipped
CI/CD Pipeline / Canary Release to Production (pull_request) Has been skipped
CI/CD Pipeline / Build Staging Worker Image (push) Successful in 58s
CI/CD Pipeline / PR Build Web Image (pull_request) Successful in 1m51s
Preview Deploy / Deploy Preview Environment (pull_request) Successful in 2m4s
CI/CD Pipeline / Build Staging Web Image (push) Successful in 1m31s
CI/CD Pipeline / Retag skipped Staging API Image (push) Has been skipped
CI/CD Pipeline / Retag skipped Staging Web Image (push) Has been skipped
CI/CD Pipeline / Retag skipped Staging Worker Image (push) Has been skipped
PR Automation / Auto Approve on CI Green (pull_request) Successful in 3m13s
PR Automation / Auto Merge on CI Green + Approved (pull_request) Has been skipped
CI/CD Pipeline / Frontend Unit Tests (push) Failing after 3m20s
CI/CD Pipeline / Deploy Staging (Watchtower auto-deploy) (push) Successful in 1m17s
CI/CD Pipeline / Integration Tests (push) Successful in 3m39s
CI/CD Pipeline / Validate - Python (mypy + alembic) (push) Successful in 4m50s
CI/CD Pipeline / ACR Image Cleanup (push) Successful in 1m51s
CI/CD Pipeline / Validate - Style (push) Successful in 6m33s
AI Code Review / AI Code Review (pull_request) Successful in 6m54s
CI/CD Pipeline / Staging API Integration Tests (push) Successful in 4m8s
CI/CD Pipeline / Staging E2E Tests (push) Failing after 5m18s
CI/CD Pipeline / PR Build Worker Image (pull_request) Failing after 9m14s
CI/CD Pipeline / CI Gate (pull_request) Failing after 1s
CI/CD Pipeline / Validate - Security (push) Successful in 9m51s
CI/CD Pipeline / Unit Tests (push) Successful in 10m42s
CI/CD Pipeline / Build Production API Image (push) Has been skipped
CI/CD Pipeline / Build Production Web Image (push) Has been skipped
CI/CD Pipeline / Build Production Worker Image (push) Has been skipped
CI/CD Pipeline / CI Gate (push) Has been skipped
CI/CD Pipeline / Canary Release to Production (push) Has been skipped
CI/CD Pipeline / Deploy Production (push) Has been skipped
CI/CD Pipeline / Production Browser E2E (push) Has been skipped
fix(worker+api): P1 封面评分时序bug + direct/complete吞ingest占位bug (#2092)
Co-authored-by: xiaoxia <dev@xiaoxiajianji.com>
Co-committed-by: xiaoxia <dev@xiaoxiajianji.com>
2026-09-29 04:24:48 +08:00

107 lines
4.3 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
from sqlalchemy.orm import Session
from packages.adapters.sqlalchemy_impl.models import IngestJobModel
from packages.domain import IngestJob, IngestJobStatus
class SQLAlchemyIngestJobRepository:
def __init__(self, session: Session):
self.session = session
def create(self, job: IngestJob) -> IngestJob:
model = IngestJobModel(
id=job.id,
project_id=job.project_id,
library_id=job.library_id,
storage_key=job.storage_key,
status=job.status.value,
error_message=job.error_message,
result_asset_id=job.result_asset_id,
file_hash=job.file_hash,
asset_id=job.asset_id or "",
celery_task_id=getattr(job, "celery_task_id", "") or "",
created_at=job.created_at,
updated_at=job.updated_at,
)
self.session.add(model)
self.session.commit()
return job
def get(self, job_id: str) -> IngestJob | None:
model = self.session.query(IngestJobModel).filter(IngestJobModel.id == job_id).first()
if model is None:
return None
return IngestJob(
id=model.id,
project_id=model.project_id,
library_id=model.library_id,
storage_key=model.storage_key,
status=IngestJobStatus(model.status),
error_message=model.error_message,
result_asset_id=model.result_asset_id,
file_hash=model.file_hash or "",
asset_id=getattr(model, "asset_id", "") or "",
celery_task_id=getattr(model, "celery_task_id", "") or "",
created_at=model.created_at,
updated_at=model.updated_at,
)
def find_by_asset_id(self, asset_id: str) -> IngestJob | None:
"""返回 asset 最近一条未失败的 ingest job(PENDING/PROCESSING/COMPLETED 均算存在,用于幂等判断)。"""
if not asset_id:
return None
# 优先返回仍在跑的 (PENDING/PROCESSING),否则返回最新一条 COMPLETED
model = (
self.session.query(IngestJobModel)
.filter(IngestJobModel.asset_id == asset_id)
.filter(IngestJobModel.status.in_([IngestJobStatus.PENDING.value, IngestJobStatus.PROCESSING.value]))
.order_by(IngestJobModel.created_at.desc())
.first()
)
if model is None:
model = (
self.session.query(IngestJobModel)
.filter(IngestJobModel.asset_id == asset_id)
.filter(IngestJobModel.status == IngestJobStatus.COMPLETED.value)
.order_by(IngestJobModel.created_at.desc())
.first()
)
if model is None:
return None
return IngestJob(
id=model.id,
project_id=model.project_id,
library_id=model.library_id,
storage_key=model.storage_key,
status=IngestJobStatus(model.status),
error_message=model.error_message,
result_asset_id=model.result_asset_id,
file_hash=model.file_hash or "",
asset_id=getattr(model, "asset_id", "") or "",
celery_task_id=getattr(model, "celery_task_id", "") or "",
created_at=model.created_at,
updated_at=model.updated_at,
)
def list_by_project(self, project_id: str) -> list[IngestJob]:
models = self.session.query(IngestJobModel).filter(IngestJobModel.project_id == project_id).all()
return [self.get(model.id) for model in models if self.get(model.id) is not None]
def update(self, job: IngestJob) -> IngestJob:
model = self.session.query(IngestJobModel).filter(IngestJobModel.id == job.id).first()
if model is None:
raise ValueError(f"IngestJob {job.id} not found")
model.status = job.status.value
model.error_message = job.error_message
model.result_asset_id = job.result_asset_id
model.file_hash = job.file_hash
model.storage_key = job.storage_key
if job.asset_id:
model.asset_id = job.asset_id
celery_tid = getattr(job, "celery_task_id", "")
if celery_tid:
model.celery_task_id = celery_tid
model.updated_at = job.updated_at
self.session.commit()
return job