Skip to main content

wowlab_cli/commands/snapshot/
sync.rs

1//! Transforms DBC data and writes it to Supabase Postgres.
2
3use std::{any::Any, pin::Pin};
4
5use anyhow::{Context, Result};
6use sqlx::PgPool;
7use wowlab_common::{
8    output::{self, fmt_duration, fmt_integer},
9    time::Instant,
10};
11use wowlab_parsers::{
12    DbcData, transform_all_armor_location, transform_all_challenge_mode_item_bonus_overrides,
13    transform_all_classes, transform_all_content_tuning_x_difficulty,
14    transform_all_content_tuning_x_expected, transform_all_content_tunings,
15    transform_all_cooldown_sets, transform_all_crafted_products,
16    transform_all_creature_difficulties, transform_all_creatures, transform_all_currencies,
17    transform_all_currency_categories, transform_all_curve_points, transform_all_curves,
18    transform_all_delves_seasons, transform_all_difficulties, transform_all_enchantments,
19    transform_all_expansion_trait_trees, transform_all_expected_stat_mods,
20    transform_all_expected_stats, transform_all_gem_properties, transform_all_global_colors,
21    transform_all_global_strings, transform_all_item_armor_quality,
22    transform_all_item_armor_shield, transform_all_item_armor_total,
23    transform_all_item_bonus_list_groups, transform_all_item_bonus_sequence_spells,
24    transform_all_item_bonuses, transform_all_item_conversions, transform_all_item_damage_scaling,
25    transform_all_item_drop_scaling, transform_all_item_group_ilvl_scaling_entries,
26    transform_all_item_level_selector_qualities, transform_all_item_level_selector_quality_sets,
27    transform_all_item_level_selectors, transform_all_item_offset_curves,
28    transform_all_item_scaling_configs, transform_all_item_squish_eras, transform_all_items,
29    transform_all_journal_instances, transform_all_journal_tiers,
30    transform_all_mythic_plus_seasons, transform_all_power_types, transform_all_professions,
31    transform_all_pvp_seasons, transform_all_pvp_tiers, transform_all_racial_spells,
32    transform_all_rand_prop_points, transform_all_specialization_spells, transform_all_specs,
33    transform_all_spells, transform_all_trait_trees, transform_armor_mitigation_by_lvl,
34    transform_base_mp, transform_base_profession_ratings, transform_challenge_mode_damage,
35    transform_challenge_mode_health, transform_combat_ratings,
36    transform_combat_ratings_mult_by_ilvl, transform_hp_per_sta,
37    transform_item_socket_cost_per_level, transform_npc_damage_by_class, transform_npc_total_hp,
38    transform_profession_ratings, transform_spell_scaling, transform_stamina_mult_by_ilvl,
39};
40#[cfg(test)]
41use wowlab_types::table_registry::PUBLISHED_SNAPSHOT_TABLES;
42use wowlab_types::{
43    data::{
44        ArmorLocationFlat, ArmorMitigationByLvlFlat, BaseMpFlat, BaseProfessionRatingsFlat,
45        ChallengeModeDamageFlat, ChallengeModeHealthFlat, ChallengeModeItemBonusOverride,
46        ClassDataFlat, CombatRatingsFlat, CombatRatingsMultByIlvlFlat, ContentTuningFlat,
47        ContentTuningXDifficultyFlat, ContentTuningXExpectedFlat, CooldownSetFlat,
48        CraftedProductFlat, CreatureDifficultyFlat, CreatureFlat, Currency, CurrencyCategory,
49        CurveFlat, CurvePointFlat, DelvesSeason, Difficulty, EnchantmentFlat,
50        ExpansionTraitTreeFlat, ExpectedStatFlat, ExpectedStatModFlat, GemPropertiesFlat,
51        GlobalColorFlat, GlobalStringFlat, HpPerStaFlat, ItemArmorQualityFlat, ItemArmorShieldFlat,
52        ItemArmorTotalFlat, ItemBonusFlat, ItemBonusListGroup, ItemBonusSequenceSpell,
53        ItemConversionFlat, ItemDamageScalingFlat, ItemDataFlat, ItemDropScalingFlat,
54        ItemGroupIlvlScalingEntry, ItemLevelSelector, ItemLevelSelectorQuality,
55        ItemLevelSelectorQualitySet, ItemOffsetCurveFlat, ItemScalingConfigFlat,
56        ItemSocketCostPerLevelFlat, ItemSquishEraFlat, JournalInstanceFlat, JournalTierFlat,
57        MythicPlusSeasonFlat, NpcDamageByClassFlat, NpcTotalHpFlat, PowerTypeFlat, Profession,
58        ProfessionRatingsFlat, PvpSeasonFlat, PvpTier, RacialSpellFlat, RandPropPointsFlat,
59        SpecDataFlat, SpecializationSpellFlat, SpellDataFlat, SpellScalingFlat,
60        StaminaMultByIlvlFlat, TraitTreeFlat,
61    },
62    sensitive::Sensitive,
63    sim::FastMap,
64    table_registry::GameDataTable,
65};
66
67use super::{SyncArgs, db};
68
69struct TransformedRows {
70    data: Box<dyn Any + Send>,
71    count: usize,
72}
73
74type InsertFn =
75    fn(PgPool, Box<dyn Any + Send>, String) -> Pin<Box<dyn Future<Output = Result<()>> + Send>>;
76
77struct TableSyncDescriptor {
78    table: GameDataTable,
79    database_name: &'static str,
80    transform: fn(&DbcData) -> TransformedRows,
81    insert: InsertFn,
82    to_ndjson: fn(&(dyn Any + Send)) -> Result<Vec<u8>>,
83}
84
85inventory::collect!(TableSyncDescriptor);
86
87pub(crate) struct SyncCredentials {
88    pub database_url: Sensitive<String>,
89    pub supabase_url: String,
90    pub service_role_key: Sensitive<String>,
91}
92
93impl std::fmt::Debug for SyncCredentials {
94    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
95        f.debug_struct("SyncCredentials")
96            .field("database_url", &self.database_url)
97            .field("supabase_url", &self.supabase_url)
98            .field("service_role_key", &self.service_role_key)
99            .finish()
100    }
101}
102
103macro_rules! register_table {
104    (@custom $table:expr, $flat_type:ty, $transform_fn:expr, $insert_fn:expr) => {
105        inventory::submit! {
106            TableSyncDescriptor {
107                table: $table,
108                database_name: $table.database_name(),
109                transform: |dbc| {
110                    let rows: Vec<$flat_type> = $transform_fn(dbc);
111                    TransformedRows { count: rows.len(), data: Box::new(rows) }
112                },
113                insert: |pool, data, patch| {
114                    async fn do_insert(pool: PgPool, data: Box<dyn Any + Send>, patch: String) -> Result<()> {
115                        let rows = data.downcast::<Vec<$flat_type>>()
116                            .map_err(|_| anyhow::anyhow!("sync registry row type mismatch: {}", $table.name()))?;
117                        $insert_fn(&pool, &rows, &patch).await?;
118                        Ok(())
119                    }
120                    Box::pin(do_insert(pool, data, patch))
121                },
122                to_ndjson: serialize_rows::<$flat_type>,
123            }
124        }
125    };
126    ($table:expr, $flat_type:ty, $transform_fn:expr, $insert_fn:expr) => {
127        inventory::submit! {
128            TableSyncDescriptor {
129                table: $table,
130                database_name: $table.database_name(),
131                transform: |dbc| {
132                    let rows: Vec<$flat_type> = $transform_fn(dbc);
133                    TransformedRows { count: rows.len(), data: Box::new(rows) }
134                },
135                insert: |pool, data, patch| {
136                    async fn do_insert(pool: PgPool, data: Box<dyn Any + Send>, patch: String) -> Result<()> {
137                        let rows = data.downcast::<Vec<$flat_type>>()
138                            .map_err(|_| anyhow::anyhow!("sync registry row type mismatch: {}", $table.name()))?;
139                        db::bulk_copy(&pool, $table.database_name(), $table.name(), &rows, &patch).await?;
140                        Ok(())
141                    }
142
143                    Box::pin(do_insert(pool, data, patch))
144                },
145                to_ndjson: serialize_rows::<$flat_type>,
146            }
147        }
148    };
149}
150
151fn serialize_rows<T>(data: &(dyn Any + Send)) -> Result<Vec<u8>>
152where
153    T: serde::Serialize + 'static,
154{
155    let rows = data
156        .downcast_ref::<Vec<T>>()
157        .context("sync registry row type does not match its descriptor")?;
158    let mut output = Vec::new();
159
160    for row in rows {
161        serde_json::to_writer(&mut output, row).context("failed to serialize snapshot row")?;
162        output.push(b'\n');
163    }
164
165    Ok(output)
166}
167
168#[rustfmt::skip]
169mod _registrations {
170    use super::{TableSyncDescriptor, GameDataTable, SpellDataFlat, transform_all_spells, TransformedRows, PgPool, Any, Result, db, serialize_rows, TraitTreeFlat, transform_all_trait_trees, ExpansionTraitTreeFlat, transform_all_expansion_trait_trees, ItemDataFlat, transform_all_items, SpecDataFlat, transform_all_specs, ClassDataFlat, transform_all_classes, GlobalColorFlat, transform_all_global_colors, GlobalStringFlat, transform_all_global_strings, ItemBonusFlat, transform_all_item_bonuses, CurveFlat, transform_all_curves, CurvePointFlat, transform_all_curve_points, RandPropPointsFlat, transform_all_rand_prop_points, ItemScalingConfigFlat, transform_all_item_scaling_configs, ItemOffsetCurveFlat, transform_all_item_offset_curves, ItemSquishEraFlat, transform_all_item_squish_eras, EnchantmentFlat, transform_all_enchantments, GemPropertiesFlat, transform_all_gem_properties, ItemDamageScalingFlat, transform_all_item_damage_scaling, CooldownSetFlat, transform_all_cooldown_sets, ExpectedStatFlat, transform_all_expected_stats, ExpectedStatModFlat, transform_all_expected_stat_mods, SpecializationSpellFlat, transform_all_specialization_spells, RacialSpellFlat, transform_all_racial_spells, PowerTypeFlat, transform_all_power_types, ItemArmorQualityFlat, transform_all_item_armor_quality, ItemArmorShieldFlat, transform_all_item_armor_shield, ItemArmorTotalFlat, transform_all_item_armor_total, ArmorLocationFlat, transform_all_armor_location, JournalTierFlat, transform_all_journal_tiers, JournalInstanceFlat, transform_all_journal_instances, CraftedProductFlat, transform_all_crafted_products, Profession, transform_all_professions, PvpSeasonFlat, transform_all_pvp_seasons, PvpTier, transform_all_pvp_tiers, MythicPlusSeasonFlat, transform_all_mythic_plus_seasons, Difficulty, transform_all_difficulties, ItemDropScalingFlat, transform_all_item_drop_scaling, Currency, transform_all_currencies, CurrencyCategory, transform_all_currency_categories, ItemBonusListGroup, transform_all_item_bonus_list_groups, ItemGroupIlvlScalingEntry, transform_all_item_group_ilvl_scaling_entries, ItemLevelSelector, transform_all_item_level_selectors, ItemLevelSelectorQualitySet, transform_all_item_level_selector_quality_sets, ItemLevelSelectorQuality, transform_all_item_level_selector_qualities, ItemConversionFlat, transform_all_item_conversions, ChallengeModeItemBonusOverride, transform_all_challenge_mode_item_bonus_overrides, ItemBonusSequenceSpell, transform_all_item_bonus_sequence_spells, DelvesSeason, transform_all_delves_seasons, HpPerStaFlat, transform_hp_per_sta, CombatRatingsFlat, transform_combat_ratings, CombatRatingsMultByIlvlFlat, transform_combat_ratings_mult_by_ilvl, StaminaMultByIlvlFlat, transform_stamina_mult_by_ilvl, BaseMpFlat, transform_base_mp, SpellScalingFlat, transform_spell_scaling, ArmorMitigationByLvlFlat, transform_armor_mitigation_by_lvl, NpcTotalHpFlat, transform_npc_total_hp, NpcDamageByClassFlat, transform_npc_damage_by_class, ChallengeModeHealthFlat, transform_challenge_mode_health, ChallengeModeDamageFlat, transform_challenge_mode_damage, ItemSocketCostPerLevelFlat, transform_item_socket_cost_per_level, ProfessionRatingsFlat, transform_profession_ratings, BaseProfessionRatingsFlat, transform_base_profession_ratings, CreatureFlat, transform_all_creatures, CreatureDifficultyFlat, transform_all_creature_difficulties, ContentTuningFlat, transform_all_content_tunings, ContentTuningXDifficultyFlat, transform_all_content_tuning_x_difficulty, ContentTuningXExpectedFlat, transform_all_content_tuning_x_expected};
171    register_table!(GameDataTable::Spells, SpellDataFlat, transform_all_spells, db::insert_spells);
172    register_table!(GameDataTable::Traits, TraitTreeFlat, transform_all_trait_trees, db::insert_traits);
173    register_table!(GameDataTable::ExpansionTraits, ExpansionTraitTreeFlat, transform_all_expansion_trait_trees, db::insert_expansion_traits);
174    register_table!(GameDataTable::Items, ItemDataFlat, transform_all_items, db::insert_items);
175    register_table!(@custom GameDataTable::Specs, SpecDataFlat, transform_all_specs, db::inserts::insert_specs);
176    register_table!(GameDataTable::Classes, ClassDataFlat, transform_all_classes, db::insert_classes);
177    register_table!(GameDataTable::GlobalColors, GlobalColorFlat, transform_all_global_colors, db::insert_global_colors);
178    register_table!(GameDataTable::GlobalStrings, GlobalStringFlat, transform_all_global_strings, db::insert_global_strings);
179    register_table!(GameDataTable::ItemBonuses, ItemBonusFlat, transform_all_item_bonuses, db::insert_item_bonuses);
180    register_table!(GameDataTable::Curves, CurveFlat, transform_all_curves, db::insert_curves);
181    register_table!(GameDataTable::CurvePoints, CurvePointFlat, transform_all_curve_points, db::insert_curve_points);
182    register_table!(GameDataTable::RandPropPoints, RandPropPointsFlat, transform_all_rand_prop_points, db::insert_rand_prop_points);
183    register_table!(GameDataTable::ItemScalingConfigs, ItemScalingConfigFlat, transform_all_item_scaling_configs, db::insert_item_scaling_configs);
184    register_table!(GameDataTable::ItemOffsetCurves, ItemOffsetCurveFlat, transform_all_item_offset_curves, db::insert_item_offset_curves);
185    register_table!(GameDataTable::ItemSquishEras, ItemSquishEraFlat, transform_all_item_squish_eras, db::insert_item_squish_eras);
186    register_table!(GameDataTable::Enchantments, EnchantmentFlat, transform_all_enchantments, db::insert_enchantments);
187    register_table!(GameDataTable::GemProperties, GemPropertiesFlat, transform_all_gem_properties, db::insert_gem_properties);
188    register_table!(GameDataTable::ItemDamageScaling, ItemDamageScalingFlat, transform_all_item_damage_scaling, db::insert_item_damage_scaling);
189    register_table!(GameDataTable::CooldownSets, CooldownSetFlat, transform_all_cooldown_sets, db::insert_cooldown_sets);
190    register_table!(GameDataTable::ExpectedStats, ExpectedStatFlat, transform_all_expected_stats, db::insert_expected_stats);
191    register_table!(GameDataTable::ExpectedStatMods, ExpectedStatModFlat, transform_all_expected_stat_mods, db::insert_expected_stat_mods);
192    register_table!(GameDataTable::SpecializationSpells, SpecializationSpellFlat, transform_all_specialization_spells, db::insert_specialization_spells);
193    register_table!(GameDataTable::RacialSpells, RacialSpellFlat, transform_all_racial_spells, db::insert_racial_spells);
194    register_table!(GameDataTable::PowerTypes, PowerTypeFlat, transform_all_power_types, db::insert_power_types);
195    register_table!(GameDataTable::ItemArmorQuality, ItemArmorQualityFlat, transform_all_item_armor_quality, db::insert_item_armor_quality);
196    register_table!(GameDataTable::ItemArmorShield, ItemArmorShieldFlat, transform_all_item_armor_shield, db::insert_item_armor_shield);
197    register_table!(GameDataTable::ItemArmorTotal, ItemArmorTotalFlat, transform_all_item_armor_total, db::insert_item_armor_total);
198    register_table!(GameDataTable::ArmorLocation, ArmorLocationFlat, transform_all_armor_location, db::insert_armor_location);
199    register_table!(GameDataTable::JournalTiers, JournalTierFlat, transform_all_journal_tiers, db::insert_journal_tiers);
200    register_table!(GameDataTable::JournalInstances, JournalInstanceFlat, transform_all_journal_instances, db::insert_journal_instances);
201    register_table!(GameDataTable::CraftedProducts, CraftedProductFlat, transform_all_crafted_products, db::insert_crafted_products);
202    register_table!(GameDataTable::Professions, Profession, transform_all_professions, db::insert_professions);
203    register_table!(GameDataTable::PvpSeasons, PvpSeasonFlat, transform_all_pvp_seasons, db::insert_pvp_seasons);
204    register_table!(GameDataTable::PvpTiers, PvpTier, transform_all_pvp_tiers, db::insert_pvp_tiers);
205    register_table!(@custom GameDataTable::MythicPlusSeasons, MythicPlusSeasonFlat, transform_all_mythic_plus_seasons, db::inserts::insert_mythic_plus_seasons);
206    register_table!(GameDataTable::Difficulties, Difficulty, transform_all_difficulties, db::insert_difficulties);
207    register_table!(GameDataTable::ItemDropScaling, ItemDropScalingFlat, transform_all_item_drop_scaling, db::insert_item_drop_scaling);
208    register_table!(GameDataTable::Currencies, Currency, transform_all_currencies, db::insert_currencies);
209    register_table!(GameDataTable::CurrencyCategories, CurrencyCategory, transform_all_currency_categories, db::insert_currency_categories);
210    register_table!(GameDataTable::ItemBonusListGroups, ItemBonusListGroup, transform_all_item_bonus_list_groups, db::insert_item_bonus_list_groups);
211    register_table!(GameDataTable::ItemGroupIlvlScalingEntries, ItemGroupIlvlScalingEntry, transform_all_item_group_ilvl_scaling_entries, db::insert_item_group_ilvl_scaling_entries);
212    register_table!(GameDataTable::ItemLevelSelectors, ItemLevelSelector, transform_all_item_level_selectors, db::insert_item_level_selectors);
213    register_table!(GameDataTable::ItemLevelSelectorQualitySets, ItemLevelSelectorQualitySet, transform_all_item_level_selector_quality_sets, db::insert_item_level_selector_quality_sets);
214    register_table!(GameDataTable::ItemLevelSelectorQualities, ItemLevelSelectorQuality, transform_all_item_level_selector_qualities, db::insert_item_level_selector_qualities);
215    register_table!(GameDataTable::ItemConversions, ItemConversionFlat, transform_all_item_conversions, db::insert_item_conversions);
216    register_table!(GameDataTable::ChallengeModeItemBonusOverrides, ChallengeModeItemBonusOverride, transform_all_challenge_mode_item_bonus_overrides, db::insert_challenge_mode_item_bonus_overrides);
217    register_table!(GameDataTable::ItemBonusSequenceSpells, ItemBonusSequenceSpell, transform_all_item_bonus_sequence_spells, db::insert_item_bonus_sequence_spells);
218    register_table!(GameDataTable::DelvesSeasons, DelvesSeason, transform_all_delves_seasons, db::insert_delves_seasons);
219    register_table!(GameDataTable::HpPerSta, HpPerStaFlat, transform_hp_per_sta, db::insert_hp_per_sta);
220    register_table!(GameDataTable::CombatRatings, CombatRatingsFlat, transform_combat_ratings, db::insert_combat_ratings);
221    register_table!(GameDataTable::CombatRatingsMultByIlvl, CombatRatingsMultByIlvlFlat, transform_combat_ratings_mult_by_ilvl, db::insert_combat_ratings_mult_by_ilvl);
222    register_table!(GameDataTable::StaminaMultByIlvl, StaminaMultByIlvlFlat, transform_stamina_mult_by_ilvl, db::insert_stamina_mult_by_ilvl);
223    register_table!(GameDataTable::BaseMp, BaseMpFlat, transform_base_mp, db::insert_base_mp);
224    register_table!(GameDataTable::SpellScaling, SpellScalingFlat, transform_spell_scaling, db::insert_spell_scaling);
225    register_table!(GameDataTable::ArmorMitigationByLvl, ArmorMitigationByLvlFlat, transform_armor_mitigation_by_lvl, db::insert_armor_mitigation_by_lvl);
226    register_table!(GameDataTable::NpcTotalHp, NpcTotalHpFlat, transform_npc_total_hp, db::insert_npc_total_hp);
227    register_table!(GameDataTable::NpcDamageByClass, NpcDamageByClassFlat, transform_npc_damage_by_class, db::insert_npc_damage_by_class);
228    register_table!(GameDataTable::ChallengeModeHealth, ChallengeModeHealthFlat, transform_challenge_mode_health, db::insert_challenge_mode_health);
229    register_table!(GameDataTable::ChallengeModeDamage, ChallengeModeDamageFlat, transform_challenge_mode_damage, db::insert_challenge_mode_damage);
230    register_table!(GameDataTable::ItemSocketCostPerLevel, ItemSocketCostPerLevelFlat, transform_item_socket_cost_per_level, db::insert_item_socket_cost_per_level);
231    register_table!(GameDataTable::ProfessionRatings, ProfessionRatingsFlat, transform_profession_ratings, db::insert_profession_ratings);
232    register_table!(GameDataTable::BaseProfessionRatings, BaseProfessionRatingsFlat, transform_base_profession_ratings, db::insert_base_profession_ratings);
233    register_table!(GameDataTable::Creatures, CreatureFlat, transform_all_creatures, db::insert_creatures);
234    register_table!(GameDataTable::CreatureDifficulties, CreatureDifficultyFlat, transform_all_creature_difficulties, db::insert_creature_difficulties);
235    register_table!(GameDataTable::ContentTunings, ContentTuningFlat, transform_all_content_tunings, db::insert_content_tunings);
236    register_table!(GameDataTable::ContentTuningXDifficulty, ContentTuningXDifficultyFlat, transform_all_content_tuning_x_difficulty, db::insert_content_tuning_x_difficulty);
237    register_table!(GameDataTable::ContentTuningXExpected, ContentTuningXExpectedFlat, transform_all_content_tuning_x_expected, db::insert_content_tuning_x_expected);
238}
239
240pub(super) async fn run_sync(args: SyncArgs) -> Result<()> {
241    run_sync_with_optional_credentials(args, None).await
242}
243
244pub(crate) async fn run_sync_with_credentials(
245    args: SyncArgs,
246    credentials: &SyncCredentials,
247) -> Result<()> {
248    run_sync_with_optional_credentials(args, Some(credentials)).await
249}
250
251// #t(fn: rust_max_fn_lines) snapshot sync is one ordered load, transform, publish, and reporting transaction
252async fn run_sync_with_optional_credentials(
253    args: SyncArgs,
254    credentials: Option<&SyncCredentials>,
255) -> Result<()> {
256    let total_start = Instant::now();
257
258    output::header(&format!("Snapshot sync — patch {}", args.patch));
259
260    if args.tables.is_empty() {
261        output::kv("Tables", "all");
262    } else {
263        let names: Vec<_> = args.tables.iter().map(ToString::to_string).collect();
264
265        output::kv("Tables", &names.join(", "));
266    }
267
268    output::kv("Data dir", &args.data_dir.display().to_string());
269
270    if args.dry_run {
271        output::warning("Dry run mode — no database writes");
272    }
273
274    output::blank();
275
276    let start = Instant::now();
277    let dbc = DbcData::load_all(&args.data_dir)?;
278
279    output::success(&format!(
280        "Loaded DBC data in {}",
281        fmt_duration(start.elapsed().as_secs_f64())
282    ));
283
284    let transform_start = Instant::now();
285    let mut transformed: Vec<(&TableSyncDescriptor, TransformedRows)> = Vec::new();
286    let mut total_records: usize = 0;
287
288    for desc in inventory::iter::<TableSyncDescriptor> {
289        if !db::should_sync(&args.tables, desc.table) {
290            continue;
291        }
292
293        let rows = (desc.transform)(&dbc);
294
295        total_records += rows.count;
296
297        if rows.count > 0 {
298            output::detail(&format!(
299                "{}: {} records",
300                desc.table.name(),
301                fmt_integer(rows.count as u64)
302            ));
303        }
304
305        transformed.push((desc, rows));
306        tokio::task::yield_now().await;
307    }
308
309    transformed.sort_by_key(|(desc, _)| desc.table.name());
310
311    output::success(&format!(
312        "Transformed {} records across {} tables in {}",
313        fmt_integer(total_records as u64),
314        transformed.len(),
315        fmt_duration(transform_start.elapsed().as_secs_f64()),
316    ));
317
318    if args.dry_run {
319        output::blank();
320        output::success(&format!(
321            "Dry run complete in {}",
322            fmt_duration(total_start.elapsed().as_secs_f64())
323        ));
324
325        return Ok(());
326    }
327
328    if args.json_only {
329        let mut snapshot_files: Vec<(&'static str, Vec<u8>)> = Vec::new();
330
331        for (desc, rows) in &transformed {
332            if desc.table.is_snapshot_exposed() {
333                let ndjson = (desc.to_ndjson)(rows.data.as_ref())?;
334
335                snapshot_files.push((desc.database_name, gzip(&ndjson)?));
336            }
337
338            tokio::task::yield_now().await;
339        }
340
341        anyhow::ensure!(
342            !snapshot_files.is_empty(),
343            "no snapshot tables in the selection — nothing to upload"
344        );
345
346        let pool = connect(credentials).await?;
347
348        publish_snapshot(&pool, &args.patch, &snapshot_files, credentials).await?;
349        output::blank();
350        output::success(&format!(
351            "Snapshot-only run complete in {}",
352            fmt_duration(total_start.elapsed().as_secs_f64())
353        ));
354
355        return Ok(());
356    }
357
358    let start = Instant::now();
359    let pool = connect(credentials).await?;
360
361    output::success(&format!(
362        "Connected to database in {}",
363        fmt_duration(start.elapsed().as_secs_f64())
364    ));
365
366    let db_tables: Vec<&str> = transformed
367        .iter()
368        .map(|(desc, _)| desc.database_name)
369        .collect();
370    let before_counts = db::count_rows(&pool, &db_tables).await?;
371
372    output::blank();
373    output::subheader("Syncing");
374    let mut total_inserted: usize = 0;
375    let mut snapshot_files: Vec<(&'static str, Vec<u8>)> = Vec::new();
376    let insert_start = Instant::now();
377
378    for (desc, rows) in transformed {
379        let start = Instant::now();
380        let count = rows.count;
381
382        if desc.table.is_snapshot_exposed() {
383            let ndjson = (desc.to_ndjson)(rows.data.as_ref())?;
384
385            snapshot_files.push((desc.database_name, gzip(&ndjson)?));
386        }
387
388        (desc.insert)(pool.clone(), rows.data, args.patch.clone()).await?;
389        let elapsed = start.elapsed().as_secs_f64();
390
391        total_inserted += count;
392
393        if count > 0 {
394            report_insert_rate(desc.table.name(), count, elapsed)?;
395        }
396    }
397
398    let insert_elapsed = insert_start.elapsed().as_secs_f64();
399
400    let after_counts = db::count_rows(&pool, &db_tables).await?;
401
402    output::blank();
403    output::subheader("Delta");
404    print_delta(&db_tables, &before_counts, &after_counts)?;
405
406    if !snapshot_files.is_empty() {
407        publish_snapshot(&pool, &args.patch, &snapshot_files, credentials).await?;
408    }
409
410    let total_elapsed = total_start.elapsed().as_secs_f64();
411    let total_inserted = u32::try_from(total_inserted).context("total row count exceeds u32")?;
412    let total_rate = f64::from(total_inserted) / insert_elapsed;
413
414    output::blank();
415    output::separator();
416    output::success(&format!(
417        "Sync complete — {} rows in {} ({} rows/s)",
418        fmt_integer(u64::from(total_inserted)),
419        fmt_duration(total_elapsed),
420        format_args!("{total_rate:.0}"),
421    ));
422
423    Ok(())
424}
425
426fn report_insert_rate(label: &str, count: usize, elapsed: f64) -> Result<()> {
427    let count = u32::try_from(count).context("table row count exceeds u32")?;
428    let rate = f64::from(count) / elapsed;
429
430    output::detail(&format!(
431        "{label}: {} rows in {} ({rate:.0} rows/s)",
432        fmt_integer(u64::from(count)),
433        fmt_duration(elapsed),
434    ));
435
436    Ok(())
437}
438
439fn gzip(bytes: &[u8]) -> Result<Vec<u8>> {
440    use std::io::Write;
441
442    let mut encoder = flate2::write::GzEncoder::new(Vec::new(), flate2::Compression::default());
443
444    encoder.write_all(bytes)?;
445
446    Ok(encoder.finish()?)
447}
448
449async fn publish_snapshot(
450    pool: &PgPool,
451    patch: &str,
452    snapshot_files: &[(&'static str, Vec<u8>)],
453    credentials: Option<&SyncCredentials>,
454) -> Result<()> {
455    output::blank();
456    output::subheader("Snapshot");
457
458    let client = match credentials {
459        Some(credentials) => wowlab_supabase::SupabaseClient::new(
460            &credentials.supabase_url,
461            credentials.service_role_key.expose(),
462        ),
463        None => wowlab_supabase::SupabaseClient::from_env_service_role(),
464    }
465    .context("snapshot upload client")?;
466    let retry = wowlab_supabase::RetryConfig::default();
467
468    let current = db::read_meta_revision(pool).await?;
469    let target = current + 1;
470
471    for (db_table, bytes) in snapshot_files {
472        let table = db_table.strip_prefix("game.").unwrap_or(db_table);
473        let object_path = format!("{target}/{table}.ndjson.gz");
474
475        wowlab_supabase::with_retry(&retry, || {
476            client.upload_object(
477                "game-snapshot",
478                &object_path,
479                bytes.clone(),
480                "application/gzip",
481                "max-age=31536000, immutable",
482            )
483        })
484        .await
485        .with_context(|| format!("snapshot upload failed for {table}"))?;
486        output::detail(&format!(
487            "{}: uploaded {} gzipped",
488            table,
489            fmt_integer(bytes.len() as u64)
490        ));
491    }
492
493    let affected = db::write_meta(
494        pool,
495        patch,
496        current,
497        target,
498        // TODO(meta): populate per-table revisions
499        serde_json::json!({}),
500    )
501    .await?;
502
503    anyhow::ensure!(
504        affected == 1,
505        "game.meta revision moved past {current} during the sync (concurrent sync?); \
506         the uploaded objects at {target}/ are harmless — rerun the sync"
507    );
508
509    output::success(&format!("Published snapshot revision {target}"));
510
511    Ok(())
512}
513
514async fn connect(credentials: Option<&SyncCredentials>) -> Result<PgPool> {
515    let database_url = match credentials {
516        Some(credentials) => credentials.database_url.expose().clone(),
517        None => std::env::var("DATABASE_URL").context(
518            "DATABASE_URL not set. Get it from Supabase Dashboard -> Settings -> Database -> Connection string",
519        )?,
520    };
521
522    Ok(db::connect(&database_url).await?)
523}
524
525fn print_delta(
526    tables: &[&str],
527    before: &FastMap<String, i64>,
528    after: &FastMap<String, i64>,
529) -> Result<()> {
530    for &table in tables {
531        let b = *before.get(table).unwrap_or(&0);
532        let a = *after.get(table).unwrap_or(&0);
533        let d = a - b;
534
535        if d == 0 && a == 0 {
536            continue;
537        }
538
539        let label = table.strip_prefix("game.").unwrap_or(table);
540        let delta_str = match d.cmp(&0) {
541            std::cmp::Ordering::Greater => format!("+{}", fmt_integer(d.unsigned_abs())),
542            std::cmp::Ordering::Less => format!("-{}", fmt_integer(d.unsigned_abs())),
543            std::cmp::Ordering::Equal => "0".to_string(),
544        };
545        let before_count = u64::try_from(b).context("database returned a negative row count")?;
546        let after_count = u64::try_from(a).context("database returned a negative row count")?;
547
548        output::kv(
549            label,
550            &format!(
551                "{} -> {} ({})",
552                fmt_integer(before_count),
553                fmt_integer(after_count),
554                delta_str
555            ),
556        );
557    }
558
559    Ok(())
560}
561
562#[cfg(test)]
563mod tests {
564    use std::{collections::BTreeSet, io::Read};
565
566    use googletest::{Result as GtestResult, prelude::*};
567    use serde::Serialize;
568
569    use super::*;
570
571    #[gtest]
572    fn every_sync_table_has_one_consistent_registry_descriptor() -> GtestResult<()> {
573        let descriptors = inventory::iter::<TableSyncDescriptor>
574            .into_iter()
575            .collect::<Vec<_>>();
576        let all_tables = GameDataTable::iter().collect::<Vec<_>>();
577
578        verify_that!(descriptors.len(), eq(all_tables.len()))?;
579        let mut labels = BTreeSet::new();
580        let mut db_tables = BTreeSet::new();
581
582        for table in all_tables {
583            let matching = descriptors
584                .iter()
585                .filter(|descriptor| descriptor.table == table)
586                .collect::<Vec<_>>();
587
588            verify_that!(matching.len(), eq(1))?;
589            let descriptor = *matching.first().or_fail()?;
590
591            verify_true!(labels.insert(descriptor.table.name()))?;
592            verify_true!(db_tables.insert(descriptor.database_name))?;
593            verify_that!(descriptor.database_name, eq(table.database_name()))?;
594        }
595
596        Ok(())
597    }
598
599    #[gtest]
600    fn published_snapshot_table_set_and_order_are_stable() -> GtestResult<()> {
601        let names = PUBLISHED_SNAPSHOT_TABLES
602            .iter()
603            .map(|table| table.database_name())
604            .collect::<Vec<_>>();
605
606        verify_that!(
607            names,
608            eq(&vec![
609                "game.curve_points",
610                "game.curves",
611                "game.item_bonuses",
612                "game.rand_prop_points",
613                "game.item_scaling_configs",
614                "game.item_offset_curves",
615                "game.item_squish_eras",
616                "game.combat_ratings",
617                "game.combat_ratings_mult_by_ilvl",
618                "game.hp_per_sta",
619                "game.spell_scaling",
620                "game.item_drop_scaling",
621                "game.specs_traits",
622                "game.expansion_traits",
623                "game.journal_instances",
624                "game.global_colors",
625                "game.specs",
626                "game.power_types",
627                "game.classes",
628            ])
629        )
630    }
631
632    #[gtest]
633    fn ndjson_and_gzip_bytes_are_deterministic() -> GtestResult<()> {
634        #[derive(Serialize)]
635        struct Row {
636            id: u32,
637            name: &'static str,
638        }
639
640        let rows = vec![
641            Row {
642                id: 1,
643                name: "alpha",
644            },
645            Row {
646                id: 2,
647                name: "beta",
648            },
649        ];
650        let ndjson = serialize_rows::<Row>(&rows).or_fail()?;
651
652        verify_that!(
653            ndjson.as_slice(),
654            eq(b"{\"id\":1,\"name\":\"alpha\"}\n{\"id\":2,\"name\":\"beta\"}\n")
655        )?;
656
657        let first = gzip(&ndjson).or_fail()?;
658        let second = gzip(&ndjson).or_fail()?;
659
660        verify_that!(first, eq(&second))?;
661
662        let mut decoded = Vec::new();
663
664        flate2::read::GzDecoder::new(first.as_slice())
665            .read_to_end(&mut decoded)
666            .or_fail()?;
667
668        verify_that!(decoded, eq(&ndjson))
669    }
670
671    #[gtest]
672    fn registry_type_mismatch_is_a_regular_error() -> GtestResult<()> {
673        let rows = vec!["wrong row type".to_string()];
674        let error = serialize_rows::<u32>(&rows).err().or_fail()?;
675
676        verify_that!(
677            error.to_string(),
678            eq("sync registry row type does not match its descriptor")
679        )
680    }
681
682    #[gtest]
683    fn credential_debug_output_redacts_database_password_and_service_key() -> GtestResult<()> {
684        let credentials = SyncCredentials {
685            database_url: Sensitive::new(
686                "postgresql://user:database-password@example.invalid/db".to_string(),
687            ),
688            supabase_url: "https://example.invalid".to_string(),
689            service_role_key: Sensitive::new("service-role-secret".to_string()),
690        };
691        let debug = format!("{credentials:?}");
692
693        verify_that!(debug.as_str(), not(contains_substring("database-password")))?;
694        verify_that!(
695            debug.as_str(),
696            not(contains_substring("service-role-secret"))
697        )?;
698        verify_that!(
699            debug.as_str(),
700            contains_substring("https://example.invalid")
701        )?;
702
703        verify_that!(debug.matches("[REDACTED]").count(), eq(2))
704    }
705}