feat: add asset manifest knowledge ingest
This commit is contained in:
72
app/db/migrations/versions/0002_asset_ingest_jobs.py
Normal file
72
app/db/migrations/versions/0002_asset_ingest_jobs.py
Normal file
@@ -0,0 +1,72 @@
|
||||
"""asset manifest ingest jobs
|
||||
|
||||
Revision ID: 0002_asset_ingest_jobs
|
||||
Revises: 0001_initial
|
||||
Create Date: 2026-06-14
|
||||
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
from collections.abc import Sequence
|
||||
|
||||
import sqlalchemy as sa
|
||||
from alembic import op
|
||||
from sqlalchemy.dialects import postgresql
|
||||
|
||||
revision: str = "0002_asset_ingest_jobs"
|
||||
down_revision: str | None = "0001_initial"
|
||||
branch_labels: str | Sequence[str] | None = None
|
||||
depends_on: str | Sequence[str] | None = None
|
||||
|
||||
|
||||
def upgrade() -> None:
|
||||
op.create_table(
|
||||
"asset_ingest_jobs",
|
||||
sa.Column("id", postgresql.UUID(as_uuid=True), primary_key=True),
|
||||
sa.Column(
|
||||
"document_id",
|
||||
postgresql.UUID(as_uuid=True),
|
||||
sa.ForeignKey("documents.id", ondelete="CASCADE"),
|
||||
nullable=False,
|
||||
),
|
||||
sa.Column("asset_id", sa.Text, nullable=False),
|
||||
sa.Column("manifest_version", sa.String(32), nullable=False),
|
||||
sa.Column("sha256", sa.String(64), nullable=False),
|
||||
sa.Column("owner_module", sa.Text, nullable=False),
|
||||
sa.Column("owner_record_type", sa.Text, nullable=False),
|
||||
sa.Column("owner_record_id", sa.Text, nullable=False),
|
||||
sa.Column("status", sa.String(32), nullable=False, server_default="QUEUED"),
|
||||
sa.Column("idempotency_key", sa.Text, nullable=True),
|
||||
sa.Column("storage_bucket", sa.Text, nullable=False),
|
||||
sa.Column("storage_key", sa.Text, nullable=False),
|
||||
sa.Column("storage_version_id", sa.Text, nullable=True),
|
||||
sa.Column("manifest", postgresql.JSONB, nullable=False),
|
||||
sa.Column(
|
||||
"ingest_options",
|
||||
postgresql.JSONB,
|
||||
nullable=False,
|
||||
server_default=sa.text("'{}'::jsonb"),
|
||||
),
|
||||
sa.Column("created_at", sa.DateTime(timezone=True), server_default=sa.func.now(), nullable=False),
|
||||
sa.Column("updated_at", sa.DateTime(timezone=True), server_default=sa.func.now(), nullable=False),
|
||||
sa.UniqueConstraint(
|
||||
"asset_id",
|
||||
"sha256",
|
||||
"manifest_version",
|
||||
name="uq_asset_ingest_identity",
|
||||
),
|
||||
)
|
||||
op.create_index("ix_asset_ingest_document", "asset_ingest_jobs", ["document_id"])
|
||||
op.create_index("ix_asset_ingest_asset", "asset_ingest_jobs", ["asset_id"])
|
||||
op.create_index(
|
||||
"ix_asset_ingest_owner",
|
||||
"asset_ingest_jobs",
|
||||
["owner_module", "owner_record_type", "owner_record_id"],
|
||||
)
|
||||
|
||||
|
||||
def downgrade() -> None:
|
||||
op.drop_index("ix_asset_ingest_owner", table_name="asset_ingest_jobs")
|
||||
op.drop_index("ix_asset_ingest_asset", table_name="asset_ingest_jobs")
|
||||
op.drop_index("ix_asset_ingest_document", table_name="asset_ingest_jobs")
|
||||
op.drop_table("asset_ingest_jobs")
|
||||
@@ -246,6 +246,45 @@ class IngestionRun(Base):
|
||||
)
|
||||
|
||||
|
||||
class AssetIngestJob(Base):
|
||||
__tablename__ = "asset_ingest_jobs"
|
||||
__table_args__ = (
|
||||
UniqueConstraint(
|
||||
"asset_id",
|
||||
"sha256",
|
||||
"manifest_version",
|
||||
name="uq_asset_ingest_identity",
|
||||
),
|
||||
Index("ix_asset_ingest_document", "document_id"),
|
||||
Index("ix_asset_ingest_asset", "asset_id"),
|
||||
Index("ix_asset_ingest_owner", "owner_module", "owner_record_type", "owner_record_id"),
|
||||
)
|
||||
|
||||
id: Mapped[uuid.UUID] = mapped_column(UUID(as_uuid=True), primary_key=True, default=uuid.uuid4)
|
||||
document_id: Mapped[uuid.UUID] = mapped_column(
|
||||
UUID(as_uuid=True), ForeignKey("documents.id", ondelete="CASCADE"), nullable=False
|
||||
)
|
||||
asset_id: Mapped[str] = mapped_column(Text, nullable=False)
|
||||
manifest_version: Mapped[str] = mapped_column(String(32), nullable=False)
|
||||
sha256: Mapped[str] = mapped_column(String(64), nullable=False)
|
||||
owner_module: Mapped[str] = mapped_column(Text, nullable=False)
|
||||
owner_record_type: Mapped[str] = mapped_column(Text, nullable=False)
|
||||
owner_record_id: Mapped[str] = mapped_column(Text, nullable=False)
|
||||
status: Mapped[str] = mapped_column(String(32), nullable=False, default="QUEUED")
|
||||
idempotency_key: Mapped[str | None] = mapped_column(Text, nullable=True)
|
||||
storage_bucket: Mapped[str] = mapped_column(Text, nullable=False)
|
||||
storage_key: Mapped[str] = mapped_column(Text, nullable=False)
|
||||
storage_version_id: Mapped[str | None] = mapped_column(Text, nullable=True)
|
||||
manifest_json: Mapped[dict[str, Any]] = mapped_column("manifest", JSONB, nullable=False)
|
||||
ingest_options: Mapped[dict[str, Any]] = mapped_column(JSONB, nullable=False, default=dict)
|
||||
created_at: Mapped[datetime] = mapped_column(
|
||||
DateTime(timezone=True), server_default=func.now(), nullable=False
|
||||
)
|
||||
updated_at: Mapped[datetime] = mapped_column(
|
||||
DateTime(timezone=True), server_default=func.now(), onupdate=func.now(), nullable=False
|
||||
)
|
||||
|
||||
|
||||
class ProcessingEvent(Base):
|
||||
__tablename__ = "processing_events"
|
||||
__table_args__ = (
|
||||
|
||||
Reference in New Issue
Block a user