1use 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
251async 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 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}