|
| 1 | +-- Wire packages-db into the Sequin → Kafka → Tinybird pipeline. |
| 2 | +-- |
| 3 | +-- Two related changes bundled here because both serve the same goal — making |
| 4 | +-- packages-db row changes replicate cleanly into Tinybird: |
| 5 | +-- |
| 6 | +-- 1. Publication + REPLICA IDENTITY FULL on the 11 tables the Tinybird |
| 7 | +-- datasources read from. publish_via_partition_root collapses the |
| 8 | +-- versions (32) / package_dependencies (64) partition leaves into a |
| 9 | +-- single logical topic each. REPLICA IDENTITY on a partitioned root |
| 10 | +-- does not cascade, so every leaf is set explicitly via pg_inherits. |
| 11 | +-- |
| 12 | +-- 2. rank_packages() bumps last_synced_at on every UPDATE that touches a |
| 13 | +-- DS-exported field (impact, is_critical, last_rank_pass_at). |
| 14 | +-- last_synced_at is the Tinybird ENGINE_VER for the packages datasource; |
| 15 | +-- without this bump, ReplacingMergeTree may keep an older row when |
| 16 | +-- criticality changes without any other write path touching the row. |
| 17 | + |
| 18 | +-- ─── 1. Sequin publication ────────────────────────────────────────────────── |
| 19 | + |
| 20 | +DO $$ |
| 21 | +BEGIN |
| 22 | + IF NOT EXISTS ( |
| 23 | + SELECT 1 FROM pg_publication WHERE pubname = 'sequin_pub' |
| 24 | + ) THEN |
| 25 | + CREATE PUBLICATION sequin_pub |
| 26 | + FOR TABLE |
| 27 | + packages, |
| 28 | + versions, |
| 29 | + package_dependencies, |
| 30 | + package_maintainers, |
| 31 | + package_repos, |
| 32 | + maintainers, |
| 33 | + repos, |
| 34 | + repo_scorecard_checks, |
| 35 | + advisories, |
| 36 | + advisory_packages, |
| 37 | + advisory_affected_ranges |
| 38 | + WITH (publish_via_partition_root = true); |
| 39 | + END IF; |
| 40 | +END$$; |
| 41 | + |
| 42 | +ALTER TABLE public.packages REPLICA IDENTITY FULL; |
| 43 | +ALTER TABLE public.versions REPLICA IDENTITY FULL; |
| 44 | +ALTER TABLE public.package_dependencies REPLICA IDENTITY FULL; |
| 45 | +ALTER TABLE public.package_maintainers REPLICA IDENTITY FULL; |
| 46 | +ALTER TABLE public.package_repos REPLICA IDENTITY FULL; |
| 47 | +ALTER TABLE public.maintainers REPLICA IDENTITY FULL; |
| 48 | +ALTER TABLE public.repos REPLICA IDENTITY FULL; |
| 49 | +ALTER TABLE public.repo_scorecard_checks REPLICA IDENTITY FULL; |
| 50 | +ALTER TABLE public.advisories REPLICA IDENTITY FULL; |
| 51 | +ALTER TABLE public.advisory_packages REPLICA IDENTITY FULL; |
| 52 | +ALTER TABLE public.advisory_affected_ranges REPLICA IDENTITY FULL; |
| 53 | + |
| 54 | +-- versions (32) and package_dependencies (64) are hash-partitioned. REPLICA |
| 55 | +-- IDENTITY on the partitioned root does not cascade; set it on every leaf. |
| 56 | +DO $$ |
| 57 | +DECLARE |
| 58 | + parent_table text; |
| 59 | + partition_oid regclass; |
| 60 | +BEGIN |
| 61 | + FOREACH parent_table IN ARRAY ARRAY['public.versions', 'public.package_dependencies'] |
| 62 | + LOOP |
| 63 | + FOR partition_oid IN |
| 64 | + SELECT inhrelid::regclass |
| 65 | + FROM pg_inherits |
| 66 | + WHERE inhparent = parent_table::regclass |
| 67 | + LOOP |
| 68 | + EXECUTE format('ALTER TABLE %s REPLICA IDENTITY FULL', partition_oid); |
| 69 | + END LOOP; |
| 70 | + END LOOP; |
| 71 | +END$$; |
| 72 | + |
| 73 | +-- ─── 2. rank_packages() bumps last_synced_at ──────────────────────────────── |
| 74 | + |
| 75 | +CREATE OR REPLACE FUNCTION rank_packages( |
| 76 | + weight_downloads numeric DEFAULT 0.25, |
| 77 | + weight_dependent_packages numeric DEFAULT 0.25, |
| 78 | + weight_transitive numeric DEFAULT 0.50, |
| 79 | + critical_top_n_by_ecosystem jsonb DEFAULT '{"npm":400000,"go":100000,"maven":200000,"pypi":100000,"nuget":50000,"cargo":75000}'::jsonb |
| 80 | +) |
| 81 | +RETURNS TABLE(scored_rows int, ranked_rows int) |
| 82 | +LANGUAGE plpgsql AS $$ |
| 83 | +DECLARE |
| 84 | + n_scored int; |
| 85 | + n_ranked int; |
| 86 | +BEGIN |
| 87 | + -- Step 1: score |
| 88 | + WITH percentile_scores AS ( |
| 89 | + SELECT |
| 90 | + id, |
| 91 | + ( |
| 92 | + weight_downloads * PERCENT_RANK() OVER ( |
| 93 | + PARTITION BY ecosystem ORDER BY LOG(1 + COALESCE(downloads_last_30d, 0))) |
| 94 | + |
| 95 | + + weight_dependent_packages * PERCENT_RANK() OVER ( |
| 96 | + PARTITION BY ecosystem ORDER BY LOG(1 + COALESCE(dependent_count, 0))) |
| 97 | + |
| 98 | + + weight_transitive * PERCENT_RANK() OVER ( |
| 99 | + PARTITION BY ecosystem ORDER BY LOG(1 + COALESCE(transitive_dependent_count, 0))) |
| 100 | + )::numeric(10, 4) AS new_impact |
| 101 | + FROM packages |
| 102 | + WHERE ecosystem IN (SELECT jsonb_object_keys(critical_top_n_by_ecosystem)) |
| 103 | + ) |
| 104 | + UPDATE packages p |
| 105 | + SET impact = ps.new_impact, |
| 106 | + last_synced_at = NOW() |
| 107 | + FROM percentile_scores ps |
| 108 | + WHERE p.id = ps.id |
| 109 | + AND p.impact IS DISTINCT FROM ps.new_impact; |
| 110 | + |
| 111 | + GET DIAGNOSTICS n_scored = ROW_COUNT; |
| 112 | + |
| 113 | + -- Step 2: rank + flag |
| 114 | + WITH ranked AS ( |
| 115 | + SELECT |
| 116 | + id, ecosystem, |
| 117 | + ROW_NUMBER() OVER ( |
| 118 | + PARTITION BY ecosystem |
| 119 | + ORDER BY impact DESC NULLS LAST, id |
| 120 | + ) AS r |
| 121 | + FROM packages |
| 122 | + WHERE purl IS NOT NULL |
| 123 | + AND ecosystem IN (SELECT jsonb_object_keys(critical_top_n_by_ecosystem)) |
| 124 | + ), |
| 125 | + flagged AS ( |
| 126 | + SELECT |
| 127 | + id, r, |
| 128 | + COALESCE( |
| 129 | + r <= (critical_top_n_by_ecosystem ->> ecosystem)::int, |
| 130 | + FALSE |
| 131 | + ) AS new_is_critical |
| 132 | + FROM ranked |
| 133 | + ) |
| 134 | + UPDATE packages p |
| 135 | + SET rank_in_ecosystem = f.r, |
| 136 | + is_critical = f.new_is_critical, |
| 137 | + last_synced_at = NOW() |
| 138 | + FROM flagged f |
| 139 | + WHERE p.id = f.id |
| 140 | + AND ( |
| 141 | + p.rank_in_ecosystem IS DISTINCT FROM f.r |
| 142 | + OR p.is_critical IS DISTINCT FROM f.new_is_critical |
| 143 | + ); |
| 144 | + |
| 145 | + GET DIAGNOSTICS n_ranked = ROW_COUNT; |
| 146 | + |
| 147 | + -- Step 2.5: spotlight overrides |
| 148 | + UPDATE packages p |
| 149 | + SET is_critical = TRUE, |
| 150 | + last_synced_at = NOW() |
| 151 | + FROM package_criticality_spotlight s |
| 152 | + WHERE p.ecosystem = s.ecosystem |
| 153 | + AND (p.namespace IS NOT DISTINCT FROM s.namespace) |
| 154 | + AND p.name = s.name |
| 155 | + AND p.is_critical = FALSE; |
| 156 | + |
| 157 | + -- Step 3: stamp last_rank_pass_at unconditionally |
| 158 | + UPDATE packages |
| 159 | + SET last_rank_pass_at = NOW(), |
| 160 | + last_synced_at = NOW() |
| 161 | + WHERE ecosystem IN (SELECT jsonb_object_keys(critical_top_n_by_ecosystem)); |
| 162 | + |
| 163 | + RETURN QUERY SELECT n_scored, n_ranked; |
| 164 | +END; |
| 165 | +$$; |
0 commit comments