From ceca49910b9a499deee457ffaa7f6a9eaff25679 Mon Sep 17 00:00:00 2001 From: anilb Date: Wed, 22 Jul 2026 20:44:13 +0200 Subject: [PATCH 1/5] feat: add sequin date-based backfill script and packages sort column indexes Signed-off-by: anilb --- ...9__sequin_backfill_sort_column_indexes.sql | 11 + scripts/sequin-backfill.py | 471 ++++++++++++++++++ 2 files changed, 482 insertions(+) create mode 100644 backend/src/osspckgs/migrations/V1784745489__sequin_backfill_sort_column_indexes.sql create mode 100755 scripts/sequin-backfill.py diff --git a/backend/src/osspckgs/migrations/V1784745489__sequin_backfill_sort_column_indexes.sql b/backend/src/osspckgs/migrations/V1784745489__sequin_backfill_sort_column_indexes.sql new file mode 100644 index 0000000000..724a394f5b --- /dev/null +++ b/backend/src/osspckgs/migrations/V1784745489__sequin_backfill_sort_column_indexes.sql @@ -0,0 +1,11 @@ +-- Sequin backfills paginate with a keyset cursor of (sort column, primary key), +-- e.g. WHERE ("last_synced_at", "id", "package_id") >= ($1, $2, $3) ORDER BY ... LIMIT n. +-- Without these indexes every page is a full sort (28M+ rows on versions, +-- 44M+ on package_dependencies) and hits Sequin's per-query timeout. +-- CONCURRENTLY relies on flyway's -mixed=true to run outside a transaction. + +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_last_synced_at_id_package_id_idx + ON versions (last_synced_at, id, package_id); + +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_updated_at_id_depends_on_id_idx + ON package_dependencies (updated_at, id, depends_on_id); diff --git a/scripts/sequin-backfill.py b/scripts/sequin-backfill.py new file mode 100755 index 0000000000..117f22fec7 --- /dev/null +++ b/scripts/sequin-backfill.py @@ -0,0 +1,471 @@ +#!/usr/bin/env python3 +"""Trigger date-based Sequin backfills for all sinks in a database. + +Interactive flow: + 1. Pick a Sequin environment (from `sequin` CLI contexts in ~/.sequin/contexts) + 2. Pick a database + 3. Enter a start date (UTC) — backfills replay rows with sort column >= date + 4. Review the table checklist (all selected by default, deselect any) + 5. Confirm — one backfill is created per (sink, table) pair + +Reads (databases, sinks) go through the Sequin Management API using the CLI +context's host + API token. Backfill creation goes through `bin/sequin rpc` +on the running node, because the HTTP API in v0.13.x only supports full-table +backfills — a date-bounded start position requires +`Sequin.Consumers.create_backfills_for_form/3` with a partial-backfill config. + +Non-interactive example: + scripts/sequin-backfill.py --context lfx-prod --database CM2 \ + --date 2026-07-01 --tables public.members,public.organizations --yes + +Only remote environments with an exec mapping in CONTEXT_EXEC (or via +SEQUIN_RPC_EXEC) are supported. + +Overrides: + SEQUIN_RPC_EXEC command prefix that reaches the sequin container + (e.g. "kubectl -n default exec -i deploy/sequin --") + SEQUIN_RELEASE_BIN path of the release script inside the container + (default: /home/app/prod/rel/sequin/bin/sequin) +""" + +import argparse +import json +import os +import re +import shlex +import ssl +import subprocess +import sys +import urllib.error +import urllib.request +from datetime import datetime, timezone +from pathlib import Path + +CONTEXTS_DIR = Path.home() / ".sequin" / "contexts" +SORT_COLUMN_CANDIDATES = [ + "updatedAt", "updated_at", + "lastSyncedAt", "last_synced_at", + "snapshotAt", "snapshot_at", + "verifiedAt", "verified_at", +] +RELEASE_BIN = os.environ.get("SEQUIN_RELEASE_BIN", "/home/app/prod/rel/sequin/bin/sequin") + +# How to reach `bin/sequin rpc` for each CLI context. Contexts not listed here +# require SEQUIN_RPC_EXEC. +CONTEXT_EXEC = { + "lfx-prod": ["kubectl", "--context", "cm-prod-oracle", "-n", "default", + "exec", "-i", "deploy/sequin", "--"], +} + +ELIXIR_TEMPLATE = r''' +payload = Jason.decode!("__PAYLOAD__") +sort_cols = payload["sort_columns"] +date = payload["date"] + +out = + Enum.map(payload["sinks"], fn %{"id" => sid, "tables" => trefs} -> + case Sequin.Repo.get(Sequin.Consumers.SinkConsumer, sid) do + nil -> + %{sink: sid, created: [], failed: ["sink consumer not found"], skipped: []} + + c -> + c = Sequin.Repo.preload(c, :postgres_database) + + {cfgs, skipped, oid_map} = + Enum.reduce(trefs, {[], [], %{}}, fn tref, {cfgs, skipped, oid_map} -> + [schema, tname] = String.split(tref, ".", parts: 2) + table = Enum.find(c.postgres_database.tables, fn t -> t.name == tname and t.schema == schema end) + col = + table && + Enum.find_value(sort_cols, fn name -> + Enum.find(table.columns, fn col -> col.name == name end) + end) + + cond do + is_nil(table) -> + {cfgs, [%{table: tref, reason: "table not found in database"} | skipped], oid_map} + + is_nil(col) -> + reason = "no sort column found (tried: " <> Enum.join(sort_cols, ", ") <> ")" + {cfgs, [%{table: tref, reason: reason} | skipped], oid_map} + + true -> + cfg = %{ + "oid" => table.oid, + "type" => "partial", + "sortColumnAttnum" => col.attnum, + "initialMinCursor" => date + } + + {[cfg | cfgs], skipped, Map.put(oid_map, table.oid, %{table: tref, sort_column: col.name})} + end + end) + + res = + if cfgs == [] do + %{} + else + Sequin.Consumers.create_backfills_for_form(c.account_id, c, cfgs) + end + + created = + res + |> Map.get(:created, []) + |> Enum.map(fn b -> + info = Map.get(oid_map, b.table_oid, %{}) + %{id: b.id, table: info[:table], sort_column: info[:sort_column]} + end) + + failed = + res + |> Map.get(:failed, []) + |> Enum.map(fn cs -> + Enum.map_join(cs.errors, "; ", fn {field, {msg, _}} -> "#{field}: #{msg}" end) + end) + + if created != [] do + Sequin.Runtime.Supervisor.maybe_start_table_readers( + Sequin.Repo.preload(c, :active_backfills, force: true) + ) + end + + %{sink: c.name, created: created, failed: failed, skipped: skipped} + end + end) + +IO.puts("SEQUIN_BACKFILL_RESULT:" <> Jason.encode!(out)) +''' + + +def die(msg, code=1): + print(f"error: {msg}", file=sys.stderr) + sys.exit(code) + + +# ---------------------------------------------------------------- CLI contexts + +def load_contexts(): + if not CONTEXTS_DIR.is_dir(): + die(f"no sequin CLI contexts found in {CONTEXTS_DIR} (run `sequin context add` first)") + contexts = {} + for path in sorted(CONTEXTS_DIR.glob("*.json")): + try: + data = json.loads(path.read_text()) + contexts[data["name"]] = data + except (json.JSONDecodeError, KeyError): + print(f"warning: skipping unreadable context file {path}", file=sys.stderr) + if not contexts: + die(f"no usable contexts in {CONTEXTS_DIR}") + return contexts + + +# ------------------------------------------------------------- Management API + +class SequinApi: + def __init__(self, ctx): + self.token = ctx["api_token"] + self.host = ctx["hostname"] + self.base = None # resolved on first request + + def _try(self, base, path): + req = urllib.request.Request( + f"{base}{path}", headers={"Authorization": f"Bearer {self.token}"} + ) + ssl_ctx = ssl.create_default_context() + return urllib.request.urlopen(req, timeout=15, context=ssl_ctx) + + def get(self, path): + if self.base: + with self._try(self.base, path) as resp: + return json.load(resp) + last_err = None + for scheme in ("https", "http"): + base = f"{scheme}://{self.host}" + try: + with self._try(base, path) as resp: + data = json.load(resp) + self.base = base + return data + except (urllib.error.URLError, OSError, json.JSONDecodeError) as e: + last_err = e + die(f"cannot reach Sequin API at {self.host}: {last_err}") + + def databases(self): + for path in ("/api/postgres_databases", "/api/databases"): + try: + return self.get(path)["data"] + except urllib.error.HTTPError as e: + if e.code != 404: + raise + die("no database listing endpoint on this Sequin version") + + def sinks(self): + try: + return self.get("/api/sinks")["data"] + except urllib.error.HTTPError as e: + if e.code == 404: + die("this Sequin version has no /api/sinks endpoint — upgrade required") + raise + + +# ------------------------------------------------------------------ prompting + +def choose(title, labels): + print(f"\n{title}") + for i, label in enumerate(labels, 1): + print(f" {i}. {label}") + while True: + raw = input("> ").strip() + if raw.isdigit() and 1 <= int(raw) <= len(labels): + return int(raw) - 1 + print(f"enter a number between 1 and {len(labels)}") + + +def parse_toggles(raw, count): + """'1 3-5,8' -> indices; None if malformed.""" + indices = set() + for part in re.split(r"[,\s]+", raw.strip()): + if not part: + continue + m = re.fullmatch(r"(\d+)(?:-(\d+))?", part) + if not m: + return None + lo = int(m.group(1)) + hi = int(m.group(2) or lo) + if not (1 <= lo <= hi <= count): + return None + indices.update(range(lo - 1, hi)) + return indices + + +def choose_tables(tables, sinks_by_table): + selected = [True] * len(tables) + while True: + print("\nTables to backfill (all sinks on a deselected table are skipped):") + for i, table in enumerate(tables, 1): + mark = "x" if selected[i - 1] else " " + sinks = ", ".join(s["name"] for s in sinks_by_table[table]) + print(f" [{mark}] {i:2}. {table} ({sinks})") + print("toggle: numbers/ranges (e.g. '2 4-6') · a=all & confirm · n=none · Enter=confirm · q=abort") + raw = input("> ").strip().lower() + if raw == "": + if any(selected): + return [t for t, keep in zip(tables, selected) if keep] + print("nothing selected — select at least one table or 'q' to abort") + elif raw == "a": + return list(tables) + elif raw == "n": + selected = [False] * len(tables) + elif raw == "q": + die("aborted", 130) + else: + toggles = parse_toggles(raw, len(tables)) + if toggles is None: + print("unrecognized input") + else: + for i in toggles: + selected[i] = not selected[i] + + +def normalize_date(raw): + """Accept YYYY-MM-DD[ HH:MM[:SS]] or ISO-8601; return UTC ISO string.""" + raw = raw.strip() + candidate = raw.replace(" ", "T", 1) + if re.fullmatch(r"\d{4}-\d{2}-\d{2}", candidate): + candidate += "T00:00:00" + elif re.fullmatch(r"\d{4}-\d{2}-\d{2}T\d{2}:\d{2}", candidate): + candidate += ":00" + try: + parsed = datetime.fromisoformat(candidate.replace("Z", "+00:00")) + except ValueError: + return None + if parsed.tzinfo is None: + parsed = parsed.replace(tzinfo=timezone.utc) + return parsed.astimezone(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") + + +# ------------------------------------------------------------------ rpc layer + +def resolve_exec(ctx): + override = os.environ.get("SEQUIN_RPC_EXEC") + if override: + return shlex.split(override) + if ctx["name"] in CONTEXT_EXEC: + return CONTEXT_EXEC[ctx["name"]] + die(f"no exec mapping for context '{ctx['name']}' — set SEQUIN_RPC_EXEC") + + +def build_elixir(plan, date, sort_columns): + payload = json.dumps({ + "date": date, + "sort_columns": sort_columns, + "sinks": [{"id": s["id"], "tables": tables} for s, tables in plan], + }) + escaped = ( + payload.replace("\\", "\\\\").replace('"', '\\"').replace("#{", "\\#{") + ) + return ELIXIR_TEMPLATE.replace("__PAYLOAD__", escaped) + + +def run_rpc(exec_prefix, code): + cmd = exec_prefix + [RELEASE_BIN, "rpc", code] + proc = subprocess.run(cmd, capture_output=True, text=True) + marker = "SEQUIN_BACKFILL_RESULT:" + for line in proc.stdout.splitlines(): + if line.startswith(marker): + return json.loads(line[len(marker):]) + print(proc.stdout, file=sys.stderr) + print(proc.stderr, file=sys.stderr) + die(f"rpc did not return a result (exit {proc.returncode})") + + +# ----------------------------------------------------------------------- main + +def main(): + parser = argparse.ArgumentParser(description="Trigger date-based Sequin backfills") + parser.add_argument("--context", help="sequin CLI context name") + parser.add_argument("--database", help="Sequin database name") + parser.add_argument("--date", help="start date, UTC (YYYY-MM-DD or ISO-8601)") + parser.add_argument("--tables", help="comma-separated schema.table list (default: all)") + parser.add_argument("--sort-column", + help="force a specific sort column; default auto-detects " + f"per table ({', '.join(SORT_COLUMN_CANDIDATES)})") + parser.add_argument("--dry-run", action="store_true", + help="print the plan and generated code without executing") + parser.add_argument("--yes", action="store_true", help="skip confirmation") + args = parser.parse_args() + + contexts = load_contexts() + if not os.environ.get("SEQUIN_RPC_EXEC"): + contexts = {n: c for n, c in contexts.items() if n in CONTEXT_EXEC} + if not contexts: + die("no contexts with an exec mapping — add one to CONTEXT_EXEC or set SEQUIN_RPC_EXEC") + if args.context: + if args.context not in contexts: + die(f"unknown or unsupported context '{args.context}' (available: {', '.join(contexts)})") + ctx = contexts[args.context] + elif len(contexts) == 1: + ctx = next(iter(contexts.values())) + print(f"Using context '{ctx['name']}' ({ctx['hostname']})") + else: + names = list(contexts) + labels = [f"{n} ({contexts[n]['hostname']})" for n in names] + ctx = contexts[names[choose("Select Sequin environment:", labels)]] + + api = SequinApi(ctx) + databases = api.databases() + if not databases: + die("no databases found in this Sequin instance") + if args.database: + db = next((d for d in databases if d["name"] == args.database), None) + if not db: + die(f"unknown database '{args.database}' " + f"(available: {', '.join(d['name'] for d in databases)})") + else: + db = databases[choose("Select database:", [d["name"] for d in databases])] + + all_sinks = api.sinks() + sinks = [s for s in all_sinks if s.get("database") == db["name"]] + if not sinks: + die(f"no sinks found for database '{db['name']}'") + for sink in sinks: + if not (sink.get("source") or {}).get("include_tables"): + print(f"warning: sink '{sink['name']}' has no include_tables " + f"(schema-wide source) — skipping it", file=sys.stderr) + sinks = [s for s in sinks if (s.get("source") or {}).get("include_tables")] + if not sinks: + die("no sinks with explicit source tables to backfill") + + sinks_by_table = {} + for sink in sinks: + for table in sink["source"]["include_tables"]: + sinks_by_table.setdefault(table, []).append(sink) + tables = sorted(sinks_by_table) + + if args.date: + date = normalize_date(args.date) + if not date: + die(f"cannot parse date '{args.date}'") + else: + while True: + date = normalize_date(input("\nBackfill start date, UTC (YYYY-MM-DD or ISO-8601): ")) + if date: + break + print("cannot parse that date, try again") + + if args.tables: + wanted = [t.strip() for t in args.tables.split(",") if t.strip()] + wanted = [t if "." in t else f"public.{t}" for t in wanted] + unknown = [t for t in wanted if t not in sinks_by_table] + if unknown: + die(f"tables not covered by any sink in '{db['name']}': {', '.join(unknown)}") + selected_tables = wanted + elif sys.stdin.isatty(): + selected_tables = choose_tables(tables, sinks_by_table) + else: + selected_tables = tables + + # plan: one backfill per (sink, table-of-that-sink-still-selected) + plan = [] + for sink in sinks: + sink_tables = [t for t in sink["source"]["include_tables"] if t in selected_tables] + if sink_tables: + plan.append((sink, sink_tables)) + if not plan: + die("selection matches no sinks — nothing to do") + + sort_columns = [args.sort_column] if args.sort_column else SORT_COLUMN_CANDIDATES + sort_label = args.sort_column or f"auto ({' > '.join(SORT_COLUMN_CANDIDATES)})" + print(f"\nPlan — backfill from {date} (sort column: {sort_label}) " + f"on '{db['name']}' via context '{ctx['name']}':") + for sink, sink_tables in plan: + status = sink.get("status", "?") + active = len(sink.get("active_backfills") or []) + note = f" [status: {status}]" if status != "active" else "" + note += f" [!! {active} active backfill(s) already]" if active else "" + print(f" {sink['name']}{note}: {', '.join(sink_tables)}") + total = sum(len(t) for _, t in plan) + print(f" -> {total} backfill(s) across {len(plan)} sink(s)") + + code = build_elixir(plan, date, sort_columns) + if args.dry_run: + print("\n--- dry run: generated Elixir (executed via `bin/sequin rpc`) ---") + print(code) + return + + exec_prefix = resolve_exec(ctx) + if not args.yes: + prod = "prod" in ctx["name"] or "production" in ctx["hostname"] + if prod: + print("\n*** PRODUCTION environment ***") + answer = input(f"\nType 'yes' to create {total} backfill(s): ").strip().lower() + if answer != "yes": + die("aborted", 130) + + print("\nTriggering backfills via rpc...") + results = run_rpc(exec_prefix, code) + + print() + failures = 0 + for res in results: + for entry in res.get("created", []): + print(f" OK {res['sink']} / {entry['table']} " + f"(sort column: {entry.get('sort_column')}, backfill {entry['id']})") + for msg in res.get("failed", []): + failures += 1 + print(f" FAIL {res['sink']}: {msg}") + for skip in res.get("skipped", []): + failures += 1 + print(f" SKIP {res['sink']} / {skip['table']}: {skip['reason']}") + created_total = sum(len(r.get("created", [])) for r in results) + print(f"\n{created_total}/{total} backfill(s) created." + + (" Check failures above." if failures else "")) + sys.exit(1 if failures else 0) + + +if __name__ == "__main__": + try: + main() + except KeyboardInterrupt: + print() + sys.exit(130) From b170758fc6ac95c1b8290a8718bb8d97c737099d Mon Sep 17 00:00:00 2001 From: anilb Date: Wed, 22 Jul 2026 21:31:22 +0200 Subject: [PATCH 2/5] chore: remove backfill script Signed-off-by: anilb --- scripts/sequin-backfill.py | 471 ------------------------------------- 1 file changed, 471 deletions(-) delete mode 100755 scripts/sequin-backfill.py diff --git a/scripts/sequin-backfill.py b/scripts/sequin-backfill.py deleted file mode 100755 index 117f22fec7..0000000000 --- a/scripts/sequin-backfill.py +++ /dev/null @@ -1,471 +0,0 @@ -#!/usr/bin/env python3 -"""Trigger date-based Sequin backfills for all sinks in a database. - -Interactive flow: - 1. Pick a Sequin environment (from `sequin` CLI contexts in ~/.sequin/contexts) - 2. Pick a database - 3. Enter a start date (UTC) — backfills replay rows with sort column >= date - 4. Review the table checklist (all selected by default, deselect any) - 5. Confirm — one backfill is created per (sink, table) pair - -Reads (databases, sinks) go through the Sequin Management API using the CLI -context's host + API token. Backfill creation goes through `bin/sequin rpc` -on the running node, because the HTTP API in v0.13.x only supports full-table -backfills — a date-bounded start position requires -`Sequin.Consumers.create_backfills_for_form/3` with a partial-backfill config. - -Non-interactive example: - scripts/sequin-backfill.py --context lfx-prod --database CM2 \ - --date 2026-07-01 --tables public.members,public.organizations --yes - -Only remote environments with an exec mapping in CONTEXT_EXEC (or via -SEQUIN_RPC_EXEC) are supported. - -Overrides: - SEQUIN_RPC_EXEC command prefix that reaches the sequin container - (e.g. "kubectl -n default exec -i deploy/sequin --") - SEQUIN_RELEASE_BIN path of the release script inside the container - (default: /home/app/prod/rel/sequin/bin/sequin) -""" - -import argparse -import json -import os -import re -import shlex -import ssl -import subprocess -import sys -import urllib.error -import urllib.request -from datetime import datetime, timezone -from pathlib import Path - -CONTEXTS_DIR = Path.home() / ".sequin" / "contexts" -SORT_COLUMN_CANDIDATES = [ - "updatedAt", "updated_at", - "lastSyncedAt", "last_synced_at", - "snapshotAt", "snapshot_at", - "verifiedAt", "verified_at", -] -RELEASE_BIN = os.environ.get("SEQUIN_RELEASE_BIN", "/home/app/prod/rel/sequin/bin/sequin") - -# How to reach `bin/sequin rpc` for each CLI context. Contexts not listed here -# require SEQUIN_RPC_EXEC. -CONTEXT_EXEC = { - "lfx-prod": ["kubectl", "--context", "cm-prod-oracle", "-n", "default", - "exec", "-i", "deploy/sequin", "--"], -} - -ELIXIR_TEMPLATE = r''' -payload = Jason.decode!("__PAYLOAD__") -sort_cols = payload["sort_columns"] -date = payload["date"] - -out = - Enum.map(payload["sinks"], fn %{"id" => sid, "tables" => trefs} -> - case Sequin.Repo.get(Sequin.Consumers.SinkConsumer, sid) do - nil -> - %{sink: sid, created: [], failed: ["sink consumer not found"], skipped: []} - - c -> - c = Sequin.Repo.preload(c, :postgres_database) - - {cfgs, skipped, oid_map} = - Enum.reduce(trefs, {[], [], %{}}, fn tref, {cfgs, skipped, oid_map} -> - [schema, tname] = String.split(tref, ".", parts: 2) - table = Enum.find(c.postgres_database.tables, fn t -> t.name == tname and t.schema == schema end) - col = - table && - Enum.find_value(sort_cols, fn name -> - Enum.find(table.columns, fn col -> col.name == name end) - end) - - cond do - is_nil(table) -> - {cfgs, [%{table: tref, reason: "table not found in database"} | skipped], oid_map} - - is_nil(col) -> - reason = "no sort column found (tried: " <> Enum.join(sort_cols, ", ") <> ")" - {cfgs, [%{table: tref, reason: reason} | skipped], oid_map} - - true -> - cfg = %{ - "oid" => table.oid, - "type" => "partial", - "sortColumnAttnum" => col.attnum, - "initialMinCursor" => date - } - - {[cfg | cfgs], skipped, Map.put(oid_map, table.oid, %{table: tref, sort_column: col.name})} - end - end) - - res = - if cfgs == [] do - %{} - else - Sequin.Consumers.create_backfills_for_form(c.account_id, c, cfgs) - end - - created = - res - |> Map.get(:created, []) - |> Enum.map(fn b -> - info = Map.get(oid_map, b.table_oid, %{}) - %{id: b.id, table: info[:table], sort_column: info[:sort_column]} - end) - - failed = - res - |> Map.get(:failed, []) - |> Enum.map(fn cs -> - Enum.map_join(cs.errors, "; ", fn {field, {msg, _}} -> "#{field}: #{msg}" end) - end) - - if created != [] do - Sequin.Runtime.Supervisor.maybe_start_table_readers( - Sequin.Repo.preload(c, :active_backfills, force: true) - ) - end - - %{sink: c.name, created: created, failed: failed, skipped: skipped} - end - end) - -IO.puts("SEQUIN_BACKFILL_RESULT:" <> Jason.encode!(out)) -''' - - -def die(msg, code=1): - print(f"error: {msg}", file=sys.stderr) - sys.exit(code) - - -# ---------------------------------------------------------------- CLI contexts - -def load_contexts(): - if not CONTEXTS_DIR.is_dir(): - die(f"no sequin CLI contexts found in {CONTEXTS_DIR} (run `sequin context add` first)") - contexts = {} - for path in sorted(CONTEXTS_DIR.glob("*.json")): - try: - data = json.loads(path.read_text()) - contexts[data["name"]] = data - except (json.JSONDecodeError, KeyError): - print(f"warning: skipping unreadable context file {path}", file=sys.stderr) - if not contexts: - die(f"no usable contexts in {CONTEXTS_DIR}") - return contexts - - -# ------------------------------------------------------------- Management API - -class SequinApi: - def __init__(self, ctx): - self.token = ctx["api_token"] - self.host = ctx["hostname"] - self.base = None # resolved on first request - - def _try(self, base, path): - req = urllib.request.Request( - f"{base}{path}", headers={"Authorization": f"Bearer {self.token}"} - ) - ssl_ctx = ssl.create_default_context() - return urllib.request.urlopen(req, timeout=15, context=ssl_ctx) - - def get(self, path): - if self.base: - with self._try(self.base, path) as resp: - return json.load(resp) - last_err = None - for scheme in ("https", "http"): - base = f"{scheme}://{self.host}" - try: - with self._try(base, path) as resp: - data = json.load(resp) - self.base = base - return data - except (urllib.error.URLError, OSError, json.JSONDecodeError) as e: - last_err = e - die(f"cannot reach Sequin API at {self.host}: {last_err}") - - def databases(self): - for path in ("/api/postgres_databases", "/api/databases"): - try: - return self.get(path)["data"] - except urllib.error.HTTPError as e: - if e.code != 404: - raise - die("no database listing endpoint on this Sequin version") - - def sinks(self): - try: - return self.get("/api/sinks")["data"] - except urllib.error.HTTPError as e: - if e.code == 404: - die("this Sequin version has no /api/sinks endpoint — upgrade required") - raise - - -# ------------------------------------------------------------------ prompting - -def choose(title, labels): - print(f"\n{title}") - for i, label in enumerate(labels, 1): - print(f" {i}. {label}") - while True: - raw = input("> ").strip() - if raw.isdigit() and 1 <= int(raw) <= len(labels): - return int(raw) - 1 - print(f"enter a number between 1 and {len(labels)}") - - -def parse_toggles(raw, count): - """'1 3-5,8' -> indices; None if malformed.""" - indices = set() - for part in re.split(r"[,\s]+", raw.strip()): - if not part: - continue - m = re.fullmatch(r"(\d+)(?:-(\d+))?", part) - if not m: - return None - lo = int(m.group(1)) - hi = int(m.group(2) or lo) - if not (1 <= lo <= hi <= count): - return None - indices.update(range(lo - 1, hi)) - return indices - - -def choose_tables(tables, sinks_by_table): - selected = [True] * len(tables) - while True: - print("\nTables to backfill (all sinks on a deselected table are skipped):") - for i, table in enumerate(tables, 1): - mark = "x" if selected[i - 1] else " " - sinks = ", ".join(s["name"] for s in sinks_by_table[table]) - print(f" [{mark}] {i:2}. {table} ({sinks})") - print("toggle: numbers/ranges (e.g. '2 4-6') · a=all & confirm · n=none · Enter=confirm · q=abort") - raw = input("> ").strip().lower() - if raw == "": - if any(selected): - return [t for t, keep in zip(tables, selected) if keep] - print("nothing selected — select at least one table or 'q' to abort") - elif raw == "a": - return list(tables) - elif raw == "n": - selected = [False] * len(tables) - elif raw == "q": - die("aborted", 130) - else: - toggles = parse_toggles(raw, len(tables)) - if toggles is None: - print("unrecognized input") - else: - for i in toggles: - selected[i] = not selected[i] - - -def normalize_date(raw): - """Accept YYYY-MM-DD[ HH:MM[:SS]] or ISO-8601; return UTC ISO string.""" - raw = raw.strip() - candidate = raw.replace(" ", "T", 1) - if re.fullmatch(r"\d{4}-\d{2}-\d{2}", candidate): - candidate += "T00:00:00" - elif re.fullmatch(r"\d{4}-\d{2}-\d{2}T\d{2}:\d{2}", candidate): - candidate += ":00" - try: - parsed = datetime.fromisoformat(candidate.replace("Z", "+00:00")) - except ValueError: - return None - if parsed.tzinfo is None: - parsed = parsed.replace(tzinfo=timezone.utc) - return parsed.astimezone(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") - - -# ------------------------------------------------------------------ rpc layer - -def resolve_exec(ctx): - override = os.environ.get("SEQUIN_RPC_EXEC") - if override: - return shlex.split(override) - if ctx["name"] in CONTEXT_EXEC: - return CONTEXT_EXEC[ctx["name"]] - die(f"no exec mapping for context '{ctx['name']}' — set SEQUIN_RPC_EXEC") - - -def build_elixir(plan, date, sort_columns): - payload = json.dumps({ - "date": date, - "sort_columns": sort_columns, - "sinks": [{"id": s["id"], "tables": tables} for s, tables in plan], - }) - escaped = ( - payload.replace("\\", "\\\\").replace('"', '\\"').replace("#{", "\\#{") - ) - return ELIXIR_TEMPLATE.replace("__PAYLOAD__", escaped) - - -def run_rpc(exec_prefix, code): - cmd = exec_prefix + [RELEASE_BIN, "rpc", code] - proc = subprocess.run(cmd, capture_output=True, text=True) - marker = "SEQUIN_BACKFILL_RESULT:" - for line in proc.stdout.splitlines(): - if line.startswith(marker): - return json.loads(line[len(marker):]) - print(proc.stdout, file=sys.stderr) - print(proc.stderr, file=sys.stderr) - die(f"rpc did not return a result (exit {proc.returncode})") - - -# ----------------------------------------------------------------------- main - -def main(): - parser = argparse.ArgumentParser(description="Trigger date-based Sequin backfills") - parser.add_argument("--context", help="sequin CLI context name") - parser.add_argument("--database", help="Sequin database name") - parser.add_argument("--date", help="start date, UTC (YYYY-MM-DD or ISO-8601)") - parser.add_argument("--tables", help="comma-separated schema.table list (default: all)") - parser.add_argument("--sort-column", - help="force a specific sort column; default auto-detects " - f"per table ({', '.join(SORT_COLUMN_CANDIDATES)})") - parser.add_argument("--dry-run", action="store_true", - help="print the plan and generated code without executing") - parser.add_argument("--yes", action="store_true", help="skip confirmation") - args = parser.parse_args() - - contexts = load_contexts() - if not os.environ.get("SEQUIN_RPC_EXEC"): - contexts = {n: c for n, c in contexts.items() if n in CONTEXT_EXEC} - if not contexts: - die("no contexts with an exec mapping — add one to CONTEXT_EXEC or set SEQUIN_RPC_EXEC") - if args.context: - if args.context not in contexts: - die(f"unknown or unsupported context '{args.context}' (available: {', '.join(contexts)})") - ctx = contexts[args.context] - elif len(contexts) == 1: - ctx = next(iter(contexts.values())) - print(f"Using context '{ctx['name']}' ({ctx['hostname']})") - else: - names = list(contexts) - labels = [f"{n} ({contexts[n]['hostname']})" for n in names] - ctx = contexts[names[choose("Select Sequin environment:", labels)]] - - api = SequinApi(ctx) - databases = api.databases() - if not databases: - die("no databases found in this Sequin instance") - if args.database: - db = next((d for d in databases if d["name"] == args.database), None) - if not db: - die(f"unknown database '{args.database}' " - f"(available: {', '.join(d['name'] for d in databases)})") - else: - db = databases[choose("Select database:", [d["name"] for d in databases])] - - all_sinks = api.sinks() - sinks = [s for s in all_sinks if s.get("database") == db["name"]] - if not sinks: - die(f"no sinks found for database '{db['name']}'") - for sink in sinks: - if not (sink.get("source") or {}).get("include_tables"): - print(f"warning: sink '{sink['name']}' has no include_tables " - f"(schema-wide source) — skipping it", file=sys.stderr) - sinks = [s for s in sinks if (s.get("source") or {}).get("include_tables")] - if not sinks: - die("no sinks with explicit source tables to backfill") - - sinks_by_table = {} - for sink in sinks: - for table in sink["source"]["include_tables"]: - sinks_by_table.setdefault(table, []).append(sink) - tables = sorted(sinks_by_table) - - if args.date: - date = normalize_date(args.date) - if not date: - die(f"cannot parse date '{args.date}'") - else: - while True: - date = normalize_date(input("\nBackfill start date, UTC (YYYY-MM-DD or ISO-8601): ")) - if date: - break - print("cannot parse that date, try again") - - if args.tables: - wanted = [t.strip() for t in args.tables.split(",") if t.strip()] - wanted = [t if "." in t else f"public.{t}" for t in wanted] - unknown = [t for t in wanted if t not in sinks_by_table] - if unknown: - die(f"tables not covered by any sink in '{db['name']}': {', '.join(unknown)}") - selected_tables = wanted - elif sys.stdin.isatty(): - selected_tables = choose_tables(tables, sinks_by_table) - else: - selected_tables = tables - - # plan: one backfill per (sink, table-of-that-sink-still-selected) - plan = [] - for sink in sinks: - sink_tables = [t for t in sink["source"]["include_tables"] if t in selected_tables] - if sink_tables: - plan.append((sink, sink_tables)) - if not plan: - die("selection matches no sinks — nothing to do") - - sort_columns = [args.sort_column] if args.sort_column else SORT_COLUMN_CANDIDATES - sort_label = args.sort_column or f"auto ({' > '.join(SORT_COLUMN_CANDIDATES)})" - print(f"\nPlan — backfill from {date} (sort column: {sort_label}) " - f"on '{db['name']}' via context '{ctx['name']}':") - for sink, sink_tables in plan: - status = sink.get("status", "?") - active = len(sink.get("active_backfills") or []) - note = f" [status: {status}]" if status != "active" else "" - note += f" [!! {active} active backfill(s) already]" if active else "" - print(f" {sink['name']}{note}: {', '.join(sink_tables)}") - total = sum(len(t) for _, t in plan) - print(f" -> {total} backfill(s) across {len(plan)} sink(s)") - - code = build_elixir(plan, date, sort_columns) - if args.dry_run: - print("\n--- dry run: generated Elixir (executed via `bin/sequin rpc`) ---") - print(code) - return - - exec_prefix = resolve_exec(ctx) - if not args.yes: - prod = "prod" in ctx["name"] or "production" in ctx["hostname"] - if prod: - print("\n*** PRODUCTION environment ***") - answer = input(f"\nType 'yes' to create {total} backfill(s): ").strip().lower() - if answer != "yes": - die("aborted", 130) - - print("\nTriggering backfills via rpc...") - results = run_rpc(exec_prefix, code) - - print() - failures = 0 - for res in results: - for entry in res.get("created", []): - print(f" OK {res['sink']} / {entry['table']} " - f"(sort column: {entry.get('sort_column')}, backfill {entry['id']})") - for msg in res.get("failed", []): - failures += 1 - print(f" FAIL {res['sink']}: {msg}") - for skip in res.get("skipped", []): - failures += 1 - print(f" SKIP {res['sink']} / {skip['table']}: {skip['reason']}") - created_total = sum(len(r.get("created", [])) for r in results) - print(f"\n{created_total}/{total} backfill(s) created." - + (" Check failures above." if failures else "")) - sys.exit(1 if failures else 0) - - -if __name__ == "__main__": - try: - main() - except KeyboardInterrupt: - print() - sys.exit(130) From 81937e0ab93ceaf9d31bf408b3fe4ba6a855247a Mon Sep 17 00:00:00 2001 From: anilb Date: Wed, 22 Jul 2026 21:36:22 +0200 Subject: [PATCH 3/5] chore: trim migration comments Signed-off-by: anilb --- .../V1784745489__sequin_backfill_sort_column_indexes.sql | 6 ------ 1 file changed, 6 deletions(-) diff --git a/backend/src/osspckgs/migrations/V1784745489__sequin_backfill_sort_column_indexes.sql b/backend/src/osspckgs/migrations/V1784745489__sequin_backfill_sort_column_indexes.sql index 724a394f5b..423e1fe562 100644 --- a/backend/src/osspckgs/migrations/V1784745489__sequin_backfill_sort_column_indexes.sql +++ b/backend/src/osspckgs/migrations/V1784745489__sequin_backfill_sort_column_indexes.sql @@ -1,9 +1,3 @@ --- Sequin backfills paginate with a keyset cursor of (sort column, primary key), --- e.g. WHERE ("last_synced_at", "id", "package_id") >= ($1, $2, $3) ORDER BY ... LIMIT n. --- Without these indexes every page is a full sort (28M+ rows on versions, --- 44M+ on package_dependencies) and hits Sequin's per-query timeout. --- CONCURRENTLY relies on flyway's -mixed=true to run outside a transaction. - CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_last_synced_at_id_package_id_idx ON versions (last_synced_at, id, package_id); From edbea78772abe591f357d88c280a2e6d2c9977ca Mon Sep 17 00:00:00 2001 From: anilb Date: Thu, 23 Jul 2026 11:04:02 +0200 Subject: [PATCH 4/5] fix: partition-aware backfill index creation Signed-off-by: anilb --- ...9__sequin_backfill_sort_column_indexes.sql | 234 +++++++++++++++++- 1 file changed, 230 insertions(+), 4 deletions(-) diff --git a/backend/src/osspckgs/migrations/V1784745489__sequin_backfill_sort_column_indexes.sql b/backend/src/osspckgs/migrations/V1784745489__sequin_backfill_sort_column_indexes.sql index 423e1fe562..db7b8ffd24 100644 --- a/backend/src/osspckgs/migrations/V1784745489__sequin_backfill_sort_column_indexes.sql +++ b/backend/src/osspckgs/migrations/V1784745489__sequin_backfill_sort_column_indexes.sql @@ -1,5 +1,231 @@ -CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_last_synced_at_id_package_id_idx - ON versions (last_synced_at, id, package_id); +CREATE INDEX IF NOT EXISTS versions_last_synced_at_id_package_id_idx + ON ONLY versions (last_synced_at, id, package_id); -CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_updated_at_id_depends_on_id_idx - ON package_dependencies (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p0_last_synced_at_id_package_id_idx + ON versions_p0 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p1_last_synced_at_id_package_id_idx + ON versions_p1 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p2_last_synced_at_id_package_id_idx + ON versions_p2 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p3_last_synced_at_id_package_id_idx + ON versions_p3 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p4_last_synced_at_id_package_id_idx + ON versions_p4 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p5_last_synced_at_id_package_id_idx + ON versions_p5 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p6_last_synced_at_id_package_id_idx + ON versions_p6 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p7_last_synced_at_id_package_id_idx + ON versions_p7 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p8_last_synced_at_id_package_id_idx + ON versions_p8 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p9_last_synced_at_id_package_id_idx + ON versions_p9 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p10_last_synced_at_id_package_id_idx + ON versions_p10 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p11_last_synced_at_id_package_id_idx + ON versions_p11 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p12_last_synced_at_id_package_id_idx + ON versions_p12 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p13_last_synced_at_id_package_id_idx + ON versions_p13 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p14_last_synced_at_id_package_id_idx + ON versions_p14 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p15_last_synced_at_id_package_id_idx + ON versions_p15 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p16_last_synced_at_id_package_id_idx + ON versions_p16 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p17_last_synced_at_id_package_id_idx + ON versions_p17 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p18_last_synced_at_id_package_id_idx + ON versions_p18 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p19_last_synced_at_id_package_id_idx + ON versions_p19 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p20_last_synced_at_id_package_id_idx + ON versions_p20 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p21_last_synced_at_id_package_id_idx + ON versions_p21 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p22_last_synced_at_id_package_id_idx + ON versions_p22 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p23_last_synced_at_id_package_id_idx + ON versions_p23 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p24_last_synced_at_id_package_id_idx + ON versions_p24 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p25_last_synced_at_id_package_id_idx + ON versions_p25 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p26_last_synced_at_id_package_id_idx + ON versions_p26 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p27_last_synced_at_id_package_id_idx + ON versions_p27 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p28_last_synced_at_id_package_id_idx + ON versions_p28 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p29_last_synced_at_id_package_id_idx + ON versions_p29 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p30_last_synced_at_id_package_id_idx + ON versions_p30 (last_synced_at, id, package_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p31_last_synced_at_id_package_id_idx + ON versions_p31 (last_synced_at, id, package_id); + +DO $$ +DECLARE part text; +BEGIN + FOR part IN + SELECT c.relname FROM pg_inherits i JOIN pg_class c ON c.oid = i.inhrelid + WHERE i.inhparent = 'versions'::regclass + LOOP + IF NOT EXISTS ( + SELECT 1 FROM pg_inherits + WHERE inhrelid = (part || '_last_synced_at_id_package_id_idx')::regclass + ) THEN + EXECUTE format('ALTER INDEX versions_last_synced_at_id_package_id_idx ATTACH PARTITION %I', part || '_last_synced_at_id_package_id_idx'); + END IF; + END LOOP; +END$$; + +CREATE INDEX IF NOT EXISTS package_dependencies_updated_at_id_depends_on_id_idx + ON ONLY package_dependencies (updated_at, id, depends_on_id); + +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p0_updated_at_id_depends_on_id_idx + ON package_dependencies_p0 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p1_updated_at_id_depends_on_id_idx + ON package_dependencies_p1 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p2_updated_at_id_depends_on_id_idx + ON package_dependencies_p2 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p3_updated_at_id_depends_on_id_idx + ON package_dependencies_p3 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p4_updated_at_id_depends_on_id_idx + ON package_dependencies_p4 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p5_updated_at_id_depends_on_id_idx + ON package_dependencies_p5 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p6_updated_at_id_depends_on_id_idx + ON package_dependencies_p6 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p7_updated_at_id_depends_on_id_idx + ON package_dependencies_p7 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p8_updated_at_id_depends_on_id_idx + ON package_dependencies_p8 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p9_updated_at_id_depends_on_id_idx + ON package_dependencies_p9 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p10_updated_at_id_depends_on_id_idx + ON package_dependencies_p10 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p11_updated_at_id_depends_on_id_idx + ON package_dependencies_p11 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p12_updated_at_id_depends_on_id_idx + ON package_dependencies_p12 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p13_updated_at_id_depends_on_id_idx + ON package_dependencies_p13 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p14_updated_at_id_depends_on_id_idx + ON package_dependencies_p14 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p15_updated_at_id_depends_on_id_idx + ON package_dependencies_p15 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p16_updated_at_id_depends_on_id_idx + ON package_dependencies_p16 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p17_updated_at_id_depends_on_id_idx + ON package_dependencies_p17 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p18_updated_at_id_depends_on_id_idx + ON package_dependencies_p18 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p19_updated_at_id_depends_on_id_idx + ON package_dependencies_p19 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p20_updated_at_id_depends_on_id_idx + ON package_dependencies_p20 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p21_updated_at_id_depends_on_id_idx + ON package_dependencies_p21 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p22_updated_at_id_depends_on_id_idx + ON package_dependencies_p22 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p23_updated_at_id_depends_on_id_idx + ON package_dependencies_p23 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p24_updated_at_id_depends_on_id_idx + ON package_dependencies_p24 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p25_updated_at_id_depends_on_id_idx + ON package_dependencies_p25 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p26_updated_at_id_depends_on_id_idx + ON package_dependencies_p26 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p27_updated_at_id_depends_on_id_idx + ON package_dependencies_p27 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p28_updated_at_id_depends_on_id_idx + ON package_dependencies_p28 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p29_updated_at_id_depends_on_id_idx + ON package_dependencies_p29 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p30_updated_at_id_depends_on_id_idx + ON package_dependencies_p30 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p31_updated_at_id_depends_on_id_idx + ON package_dependencies_p31 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p32_updated_at_id_depends_on_id_idx + ON package_dependencies_p32 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p33_updated_at_id_depends_on_id_idx + ON package_dependencies_p33 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p34_updated_at_id_depends_on_id_idx + ON package_dependencies_p34 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p35_updated_at_id_depends_on_id_idx + ON package_dependencies_p35 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p36_updated_at_id_depends_on_id_idx + ON package_dependencies_p36 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p37_updated_at_id_depends_on_id_idx + ON package_dependencies_p37 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p38_updated_at_id_depends_on_id_idx + ON package_dependencies_p38 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p39_updated_at_id_depends_on_id_idx + ON package_dependencies_p39 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p40_updated_at_id_depends_on_id_idx + ON package_dependencies_p40 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p41_updated_at_id_depends_on_id_idx + ON package_dependencies_p41 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p42_updated_at_id_depends_on_id_idx + ON package_dependencies_p42 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p43_updated_at_id_depends_on_id_idx + ON package_dependencies_p43 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p44_updated_at_id_depends_on_id_idx + ON package_dependencies_p44 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p45_updated_at_id_depends_on_id_idx + ON package_dependencies_p45 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p46_updated_at_id_depends_on_id_idx + ON package_dependencies_p46 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p47_updated_at_id_depends_on_id_idx + ON package_dependencies_p47 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p48_updated_at_id_depends_on_id_idx + ON package_dependencies_p48 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p49_updated_at_id_depends_on_id_idx + ON package_dependencies_p49 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p50_updated_at_id_depends_on_id_idx + ON package_dependencies_p50 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p51_updated_at_id_depends_on_id_idx + ON package_dependencies_p51 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p52_updated_at_id_depends_on_id_idx + ON package_dependencies_p52 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p53_updated_at_id_depends_on_id_idx + ON package_dependencies_p53 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p54_updated_at_id_depends_on_id_idx + ON package_dependencies_p54 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p55_updated_at_id_depends_on_id_idx + ON package_dependencies_p55 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p56_updated_at_id_depends_on_id_idx + ON package_dependencies_p56 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p57_updated_at_id_depends_on_id_idx + ON package_dependencies_p57 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p58_updated_at_id_depends_on_id_idx + ON package_dependencies_p58 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p59_updated_at_id_depends_on_id_idx + ON package_dependencies_p59 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p60_updated_at_id_depends_on_id_idx + ON package_dependencies_p60 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p61_updated_at_id_depends_on_id_idx + ON package_dependencies_p61 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p62_updated_at_id_depends_on_id_idx + ON package_dependencies_p62 (updated_at, id, depends_on_id); +CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p63_updated_at_id_depends_on_id_idx + ON package_dependencies_p63 (updated_at, id, depends_on_id); + +DO $$ +DECLARE part text; +BEGIN + FOR part IN + SELECT c.relname FROM pg_inherits i JOIN pg_class c ON c.oid = i.inhrelid + WHERE i.inhparent = 'package_dependencies'::regclass + LOOP + IF NOT EXISTS ( + SELECT 1 FROM pg_inherits + WHERE inhrelid = (part || '_updated_at_id_depends_on_id_idx')::regclass + ) THEN + EXECUTE format('ALTER INDEX package_dependencies_updated_at_id_depends_on_id_idx ATTACH PARTITION %I', part || '_updated_at_id_depends_on_id_idx'); + END IF; + END LOOP; +END$$; From 9d60260dae60e079ecdffc39d91108310d32b588 Mon Sep 17 00:00:00 2001 From: anilb Date: Thu, 23 Jul 2026 11:14:40 +0200 Subject: [PATCH 5/5] fix: self-heal invalid indexes on migration retry Signed-off-by: anilb --- ...9__sequin_backfill_sort_column_indexes.sql | 22 +++++++++++++++++++ 1 file changed, 22 insertions(+) diff --git a/backend/src/osspckgs/migrations/V1784745489__sequin_backfill_sort_column_indexes.sql b/backend/src/osspckgs/migrations/V1784745489__sequin_backfill_sort_column_indexes.sql index db7b8ffd24..659869ead0 100644 --- a/backend/src/osspckgs/migrations/V1784745489__sequin_backfill_sort_column_indexes.sql +++ b/backend/src/osspckgs/migrations/V1784745489__sequin_backfill_sort_column_indexes.sql @@ -1,6 +1,17 @@ CREATE INDEX IF NOT EXISTS versions_last_synced_at_id_package_id_idx ON ONLY versions (last_synced_at, id, package_id); +DO $$ +DECLARE idx text; +BEGIN + FOR idx IN + SELECT c.relname FROM pg_index i JOIN pg_class c ON c.oid = i.indexrelid + WHERE NOT i.indisvalid AND c.relname ~ '^versions_p[0-9]+_last_synced_at_id_package_id_idx$' + LOOP + EXECUTE format('DROP INDEX IF EXISTS %I', idx); + END LOOP; +END$$; + CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p0_last_synced_at_id_package_id_idx ON versions_p0 (last_synced_at, id, package_id); CREATE INDEX CONCURRENTLY IF NOT EXISTS versions_p1_last_synced_at_id_package_id_idx @@ -85,6 +96,17 @@ END$$; CREATE INDEX IF NOT EXISTS package_dependencies_updated_at_id_depends_on_id_idx ON ONLY package_dependencies (updated_at, id, depends_on_id); +DO $$ +DECLARE idx text; +BEGIN + FOR idx IN + SELECT c.relname FROM pg_index i JOIN pg_class c ON c.oid = i.indexrelid + WHERE NOT i.indisvalid AND c.relname ~ '^package_dependencies_p[0-9]+_updated_at_id_depends_on_id_idx$' + LOOP + EXECUTE format('DROP INDEX IF EXISTS %I', idx); + END LOOP; +END$$; + CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p0_updated_at_id_depends_on_id_idx ON package_dependencies_p0 (updated_at, id, depends_on_id); CREATE INDEX CONCURRENTLY IF NOT EXISTS package_dependencies_p1_updated_at_id_depends_on_id_idx