osint-dashboard/alembic/versions/007_event_dedup.py
Sirius DevOps d81ea68dec fix: pipeline reliability — WS, summarizer, ingest, health
Stop the live HUD reconnect storm (nginx WS snippet + backoff), copy
intel/nous_client into the summarizer image, and make event ingest
idempotent on URL. GDELT uses the DOC API; NWS no longer sends bbox;
FIRMS is one ON CONFLICT batch; GET /api/aircraft serves last-known.
Health reports freshness without 503ing docker. EONET + CISA KEV added.
2026-08-28 21:49:05 -04:00

115 lines
3.5 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.

"""event_dedup + Timescale compression/retention
Revision ID: 007_event_dedup
Revises: 006_merge_heads
Create Date: 2026-08-28
"""
from alembic import op
revision = "007_event_dedup"
down_revision = "006_merge_heads"
branch_labels = None
depends_on = None
def upgrade() -> None:
op.execute("""
CREATE TABLE IF NOT EXISTS event_dedup (
url TEXT PRIMARY KEY,
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
)
""")
# Keep the earliest row per URL, drop the 10× USGS/camera dupes.
op.execute("""
DELETE FROM events a
USING events b
WHERE a.url IS NOT NULL AND a.url <> ''
AND a.url = b.url
AND (a.ingested_at, a.id) > (b.ingested_at, b.id)
""")
op.execute("""
INSERT INTO event_dedup (url)
SELECT DISTINCT url FROM events
WHERE url IS NOT NULL AND url <> ''
ON CONFLICT (url) DO NOTHING
""")
# Compression + retention. Policies no-op if Timescale rejects (fresh PG).
op.execute("""
DO $$
BEGIN
PERFORM add_compression_policy('events', INTERVAL '7 days', if_not_exists => TRUE);
EXCEPTION WHEN OTHERS THEN
BEGIN
ALTER TABLE events SET (
timescaledb.compress,
timescaledb.compress_segmentby = 'source_type',
timescaledb.compress_orderby = 'ingested_at DESC'
);
PERFORM add_compression_policy('events', INTERVAL '7 days', if_not_exists => TRUE);
EXCEPTION WHEN OTHERS THEN
NULL;
END;
END
$$;
""")
op.execute("""
DO $$
BEGIN
PERFORM add_retention_policy('events', INTERVAL '180 days', if_not_exists => TRUE);
EXCEPTION WHEN OTHERS THEN
NULL;
END
$$;
""")
op.execute("""
DO $$
BEGIN
ALTER TABLE fires SET (
timescaledb.compress,
timescaledb.compress_segmentby = 'satellite',
timescaledb.compress_orderby = 'acq_time DESC'
);
PERFORM add_compression_policy('fires', INTERVAL '7 days', if_not_exists => TRUE);
PERFORM add_retention_policy('fires', INTERVAL '90 days', if_not_exists => TRUE);
EXCEPTION WHEN OTHERS THEN
NULL;
END
$$;
""")
op.execute("""
DO $$
BEGIN
ALTER TABLE aircraft_positions SET (
timescaledb.compress,
timescaledb.compress_segmentby = 'hex',
timescaledb.compress_orderby = 'ts DESC'
);
PERFORM add_compression_policy('aircraft_positions', INTERVAL '1 day', if_not_exists => TRUE);
PERFORM add_retention_policy('aircraft_positions', INTERVAL '14 days', if_not_exists => TRUE);
EXCEPTION WHEN OTHERS THEN
NULL;
END
$$;
""")
op.execute("""
DO $$
BEGIN
ALTER TABLE vessel_positions SET (
timescaledb.compress,
timescaledb.compress_segmentby = 'mmsi',
timescaledb.compress_orderby = 'ts DESC'
);
PERFORM add_compression_policy('vessel_positions', INTERVAL '1 day', if_not_exists => TRUE);
PERFORM add_retention_policy('vessel_positions', INTERVAL '14 days', if_not_exists => TRUE);
EXCEPTION WHEN OTHERS THEN
NULL;
END
$$;
""")
def downgrade() -> None:
op.execute("DROP TABLE IF EXISTS event_dedup")