diff --git a/alembic/versions/20260812_ne_collect_trigger.py b/alembic/versions/20260812_ne_collect_trigger.py new file mode 100644 index 0000000..9f6ab66 --- /dev/null +++ b/alembic/versions/20260812_ne_collect_trigger.py @@ -0,0 +1,36 @@ +"""Add ne_collection_job.trigger_mode for scheduled batch collect. + +Revision ID: 20260812_ne_collect_trigger +Revises: 20260811_fabric_level +Create Date: 2026-08-12 +""" + +from __future__ import annotations + +from typing import Sequence, Union + +from alembic import op + +revision: str = "20260812_ne_collect_trigger" +down_revision: Union[str, Sequence[str], None] = "20260811_fabric_level" +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + from netx_api.schema_patches import apply_collection_schema_safety_net + + apply_collection_schema_safety_net(op.get_bind()) + + +def downgrade() -> None: + bind = op.get_bind() + dialect = str(getattr(bind.dialect, "name", "") or "").lower() + if dialect.startswith("postgres"): + op.execute("DROP INDEX IF EXISTS ix_ne_collection_job_trigger_mode") + op.execute("ALTER TABLE ne_collection_job DROP COLUMN IF EXISTS trigger_mode") + else: + try: + op.drop_column("ne_collection_job", "trigger_mode") + except Exception: + pass diff --git a/netx_api/app_startup.py b/netx_api/app_startup.py index 4552557..9e103ef 100644 --- a/netx_api/app_startup.py +++ b/netx_api/app_startup.py @@ -10,6 +10,7 @@ from .db import Base, SessionLocal, engine from .schema_patches import ( apply_all_legacy_startup_ddl, apply_auth_schema_patches, + apply_collection_schema_safety_net, apply_topology_schema_safety_net, run_alembic_upgrade_to_head, ) @@ -58,8 +59,9 @@ def run_api_startup() -> None: # Critical topology columns even when full legacy DDL is skipped # (e.g. alembic stamped head without applying domain patches). apply_topology_schema_safety_net(conn) + apply_collection_schema_safety_net(conn) except Exception: - _log.exception("startup: auth/topology schema safety patches failed") + _log.exception("startup: auth/topology/collection schema safety patches failed") if skip_ddl and alembic_ok: _log.info("startup: schema via Alembic (legacy inline DDL skipped)") else: diff --git a/netx_api/schema_patches.py b/netx_api/schema_patches.py index 74d5ac2..147fffe 100644 --- a/netx_api/schema_patches.py +++ b/netx_api/schema_patches.py @@ -179,6 +179,26 @@ def apply_topology_schema_safety_net(conn: Connection) -> None: ) +def apply_collection_schema_safety_net(conn: Connection) -> None: + """Always-on NE collection columns (create_all will not ALTER existing tables).""" + _run_sql( + conn, + "ALTER TABLE ne_collection_job ADD COLUMN IF NOT EXISTS last_run_at TIMESTAMP", + ) + _run_sql( + conn, + "ALTER TABLE ne_collection_job ADD COLUMN IF NOT EXISTS trigger_mode VARCHAR(32) DEFAULT 'manual'", + ) + _run_sql( + conn, + "CREATE INDEX IF NOT EXISTS ix_ne_collection_job_trigger_mode ON ne_collection_job (trigger_mode)", + ) + _run_sql( + conn, + "ALTER TABLE ne_collection_run ADD COLUMN IF NOT EXISTS ne_source VARCHAR(16) DEFAULT 'managed'", + ) + + def apply_key_alert_schema_patches( engine: Engine | None = None, *,