Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
20 commits
Select commit Hold shift + click to select a range
2dcba6c
Fix SQLite migration compatibility and idempotency
Bornunique911 Aug 5, 2026
187dff7
migrations: add downgrades, use sa.inspect, and verify unique constra…
Bornunique911 Aug 5, 2026
57284a0
fixed black formatting issue
Bornunique911 Aug 5, 2026
d4bea6f
migrations: improve downgrade safety and constraint validation
Bornunique911 Aug 5, 2026
4ea91bd
migrations: remove document_metadata migration for cre and node
Bornunique911 Aug 7, 2026
50e69a7
migrations: enhance idempotency for document_metadata column addition…
Bornunique911 Aug 7, 2026
1493bfa
migrations: update revision ID for embedding_vec migration
Bornunique911 Aug 7, 2026
c5b1a56
migrations: remove unused import of inspect from sqlalchemy
Bornunique911 Aug 7, 2026
f18e2ec
migrations: change document_metadata column type to JSON for cre and …
Bornunique911 Aug 7, 2026
dcfafbb
migrations: enhance idempotency for document_metadata column tracking…
Bornunique911 Aug 7, 2026
2c714dc
migrations: ensure tracking record removal for document_metadata colu…
Bornunique911 Aug 7, 2026
0ba1848
migrations: refine tracking record checks and cleanup in downgrade
Bornunique911 Aug 7, 2026
c392850
migrations: ensure tracking record removal for node's document_metada…
Bornunique911 Aug 7, 2026
7395cf8
migrations: remove unnecessary revision identifiers comment
Bornunique911 Aug 7, 2026
1e31f25
migrations: enhance idempotency for tracking records in upgrade and d…
Bornunique911 Aug 7, 2026
3dbbdc8
migrations: update down_revision to correct previous migration reference
Bornunique911 Aug 7, 2026
adfa5a7
migrations: clean up upgrade function by removing print statements
Bornunique911 Aug 7, 2026
debedcf
migrations: add missing docstrings to improve docstring coverage
Bornunique911 Aug 7, 2026
891f105
migrations: add missing docstrings to improve docstring coverage
Bornunique911 Aug 7, 2026
4bc02ef
migrations: enhance docstring coverage and improve upgrade/downgrade …
Bornunique911 Aug 7, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,38 @@
"""add embedding_vec to embeddings for SQLite

Revision ID: 967016ee10fa
Revises: b5ac48010165
Create Date: 2026-08-06 00:13:49.191527

"""

from alembic import op
import sqlalchemy as sa

revision = "967016ee10fa"
down_revision = "b5ac48010165"
branch_labels = None
depends_on = None


def column_exists(table, column):
"""Check if a column exists in the given table."""
conn = op.get_bind()
inspector = sa.inspect(conn)
columns = [c["name"] for c in inspector.get_columns(table)]
return column in columns


def upgrade():
"""Add embedding_vec as a TEXT column for SQLite if it doesn't already exist."""
if not column_exists("embeddings", "embedding_vec"):
op.add_column(
"embeddings", sa.Column("embedding_vec", sa.Text(), nullable=True)
)


def downgrade():
"""Remove embedding_vec from embeddings if it was added by this migration."""
with op.batch_alter_table("embeddings") as batch_op:
if column_exists("embeddings", "embedding_vec"):
batch_op.drop_column("embedding_vec")
222 changes: 165 additions & 57 deletions migrations/versions/9f1a2b3c4d5e_add_artifact_ingest_persistence.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,74 +9,182 @@
from alembic import op
import sqlalchemy as sa


revision = "9f1a2b3c4d5e"
down_revision = "a1b2c3d4e5f6"
branch_labels = None
depends_on = None


def upgrade():
op.create_table(
"artifact_ingest_event",
sa.Column("id", sa.String(), primary_key=True),
sa.Column("run_id", sa.String(), nullable=False),
sa.Column("artifact_id", sa.String(), nullable=False),
sa.Column("harvest_mode", sa.String(), nullable=False),
sa.Column("event_type", sa.String(), nullable=False),
sa.Column("source_json", sa.Text(), nullable=False),
sa.Column("locator_json", sa.Text(), nullable=False),
sa.Column("artifact_json", sa.Text(), nullable=False),
sa.Column("harvest_json", sa.Text(), nullable=False),
sa.Column("observed_at", sa.DateTime(), nullable=False),
sa.Column("created_at", sa.DateTime(), nullable=False),
sa.ForeignKeyConstraint(
["run_id"],
["import_run.id"],
onupdate="CASCADE",
ondelete="CASCADE",
),
)
op.create_unique_constraint(
"uq_artifact_ingest_event_run_artifact",
"artifact_ingest_event",
["run_id", "artifact_id"],
def table_exists(table_name):
"""Check if a table exists in the current database."""
conn = op.get_bind()
inspector = sa.inspect(conn)
return table_name in inspector.get_table_names()


def has_unique_on_columns(table, columns):
"""Return True if any unique constraint (named or unnamed) exists on exactly these columns."""
conn = op.get_bind()
inspector = sa.inspect(conn)
constraints = inspector.get_unique_constraints(table)
for c in constraints:
if c["column_names"] == columns:
return True
return False


def constraint_name_exists(table, constraint_name):
"""Return True if a unique constraint with the given name exists on the table."""
conn = op.get_bind()
inspector = sa.inspect(conn)
return any(
c["name"] == constraint_name for c in inspector.get_unique_constraints(table)
)

op.create_table(
"ingest_chunk",
sa.Column("id", sa.String(), primary_key=True),
sa.Column("artifact_event_id", sa.String(), nullable=False),
sa.Column("chunk_id", sa.String(), nullable=False),
sa.Column("text", sa.Text(), nullable=False),
sa.Column("char_count", sa.Integer(), nullable=False),
sa.Column("span_json", sa.Text(), nullable=False),
sa.Column("delta_json", sa.Text(), nullable=True),
sa.Column("created_at", sa.DateTime(), nullable=False),
sa.ForeignKeyConstraint(
["artifact_event_id"],
["artifact_ingest_event.id"],
onupdate="CASCADE",
ondelete="CASCADE",

def tracking_record_exists(revision, table_name):
"""Check if this migration created the given table."""
conn = op.get_bind()
if not table_exists("_migration_tracking"):
return False
result = conn.execute(
sa.text(
"SELECT 1 FROM _migration_tracking WHERE revision = :rev AND table_name = :tbl"
),
{"rev": revision, "tbl": table_name},
)
op.create_unique_constraint(
"uq_ingest_chunk_artifact_chunk",
"ingest_chunk",
["artifact_event_id", "chunk_id"],
return result.first() is not None


def add_tracking_record(revision, table_name):
"""Record that this migration created a table."""
if not table_exists("_migration_tracking"):
op.create_table(
"_migration_tracking",
sa.Column("revision", sa.String, primary_key=True),
sa.Column("table_name", sa.String, primary_key=True),
)
op.execute(
sa.text(
"INSERT INTO _migration_tracking (revision, table_name) VALUES (:rev, :tbl)"
),
{"rev": revision, "tbl": table_name},
)


def remove_tracking_record(revision, table_name):
"""Remove the tracking record for a table created by this migration."""
if table_exists("_migration_tracking"):
op.execute(
sa.text(
"DELETE FROM _migration_tracking WHERE revision = :rev AND table_name = :tbl"
),
{"rev": revision, "tbl": table_name},
)


def upgrade():
"""Create artifact_ingest_event and ingest_chunk tables with proper constraints.

If the tables already exist, ensure the required unique constraints are present,
handling SQLite batch rewrites with foreign key enforcement disabled.
"""
# Create artifact_ingest_event only if it doesn't exist
if not table_exists("artifact_ingest_event"):
op.create_table(
"artifact_ingest_event",
sa.Column("id", sa.String(), primary_key=True),
sa.Column("run_id", sa.String(), nullable=False),
sa.Column("artifact_id", sa.String(), nullable=False),
sa.Column("harvest_mode", sa.String(), nullable=False),
sa.Column("event_type", sa.String(), nullable=False),
sa.Column("source_json", sa.Text(), nullable=False),
sa.Column("locator_json", sa.Text(), nullable=False),
sa.Column("artifact_json", sa.Text(), nullable=False),
sa.Column("harvest_json", sa.Text(), nullable=False),
sa.Column("observed_at", sa.DateTime(), nullable=False),
sa.Column("created_at", sa.DateTime(), nullable=False),
sa.ForeignKeyConstraint(
["run_id"],
["import_run.id"],
onupdate="CASCADE",
ondelete="CASCADE",
),
sa.UniqueConstraint(
"run_id", "artifact_id", name="uq_artifact_ingest_event_run_artifact"
),
)
add_tracking_record(revision, "artifact_ingest_event")
else:
# Ensure the unique constraint exists
if not has_unique_on_columns(
"artifact_ingest_event", ["run_id", "artifact_id"]
):
# Temporarily disable foreign keys for SQLite because we may drop the parent table
op.execute("PRAGMA foreign_keys=OFF")
try:
with op.batch_alter_table("artifact_ingest_event") as batch_op:
if constraint_name_exists(
"artifact_ingest_event",
"uq_artifact_ingest_event_run_artifact",
):
batch_op.drop_constraint(
"uq_artifact_ingest_event_run_artifact", type_="unique"
)
batch_op.create_unique_constraint(
"uq_artifact_ingest_event_run_artifact",
["run_id", "artifact_id"],
)
finally:
op.execute("PRAGMA foreign_keys=ON")

# Create ingest_chunk only if it doesn't exist
if not table_exists("ingest_chunk"):
op.create_table(
"ingest_chunk",
sa.Column("id", sa.String(), primary_key=True),
sa.Column("artifact_event_id", sa.String(), nullable=False),
sa.Column("chunk_id", sa.String(), nullable=False),
sa.Column("text", sa.Text(), nullable=False),
sa.Column("char_count", sa.Integer(), nullable=False),
sa.Column("span_json", sa.Text(), nullable=False),
sa.Column("delta_json", sa.Text(), nullable=True),
sa.Column("created_at", sa.DateTime(), nullable=False),
sa.ForeignKeyConstraint(
["artifact_event_id"],
["artifact_ingest_event.id"],
onupdate="CASCADE",
ondelete="CASCADE",
),
sa.UniqueConstraint(
"artifact_event_id",
"chunk_id",
name="uq_ingest_chunk_artifact_chunk",
),
)
add_tracking_record(revision, "ingest_chunk")
else:
if not has_unique_on_columns("ingest_chunk", ["artifact_event_id", "chunk_id"]):
with op.batch_alter_table("ingest_chunk") as batch_op:
if constraint_name_exists(
"ingest_chunk", "uq_ingest_chunk_artifact_chunk"
):
batch_op.drop_constraint(
"uq_ingest_chunk_artifact_chunk", type_="unique"
)
batch_op.create_unique_constraint(
"uq_ingest_chunk_artifact_chunk",
["artifact_event_id", "chunk_id"],
)


def downgrade():
op.drop_constraint(
"uq_ingest_chunk_artifact_chunk",
"ingest_chunk",
type_="unique",
)
op.drop_table("ingest_chunk")
op.drop_constraint(
"uq_artifact_ingest_event_run_artifact",
"artifact_ingest_event",
type_="unique",
)
op.drop_table("artifact_ingest_event")
"""Drop only tables that were created by this migration."""
# Drop only tables that were created by this migration
if tracking_record_exists(revision, "ingest_chunk"):
op.drop_table("ingest_chunk")
remove_tracking_record(revision, "ingest_chunk")

if tracking_record_exists(revision, "artifact_ingest_event"):
op.drop_table("artifact_ingest_event")
remove_tracking_record(revision, "artifact_ingest_event")
Loading
Loading