1792 lines
64 KiB
Rust
1792 lines
64 KiB
Rust
use crate::runtime::analytics_date_dimension::analytics_date_dimension;
|
|
use crate::runtime::creation_entry_config::{creation_entry_config, creation_entry_type_config};
|
|
use crate::runtime::feature_gate_config::feature_gate_config;
|
|
use crate::*;
|
|
use serde::{Deserialize, Serialize};
|
|
use sha2::{Digest, Sha256};
|
|
use spacetimedb::sats::de::serde::DeserializeWrapper;
|
|
use spacetimedb::sats::ser::serde::SerializeWrapper;
|
|
use std::collections::HashSet;
|
|
|
|
use crate::bark_battle::tables::{
|
|
bark_battle_draft_config, bark_battle_leaderboard_entry, bark_battle_personal_best_projection,
|
|
bark_battle_published_config, bark_battle_runtime_run, bark_battle_score_record,
|
|
bark_battle_work_stats_projection,
|
|
};
|
|
use crate::big_fish::big_fish_runtime_run;
|
|
use crate::jump_hop::tables::{
|
|
jump_hop_agent_session, jump_hop_event, jump_hop_leaderboard_entry, jump_hop_runtime_run,
|
|
jump_hop_work_profile,
|
|
};
|
|
use crate::match3d::tables::{
|
|
match_3_d_work_profile, match3d_agent_message, match3d_agent_session, match3d_runtime_run,
|
|
};
|
|
use crate::puzzle::{
|
|
puzzle_agent_message, puzzle_agent_session, puzzle_background_compile_task, puzzle_event,
|
|
puzzle_leaderboard_entry, puzzle_runtime_run, puzzle_work_profile,
|
|
};
|
|
use crate::puzzle_clear::tables::{
|
|
puzzle_clear_agent_session, puzzle_clear_event, puzzle_clear_runtime_run,
|
|
puzzle_clear_work_profile,
|
|
};
|
|
use crate::square_hole::tables::{
|
|
square_hole_agent_message, square_hole_agent_session, square_hole_runtime_run,
|
|
square_hole_work_profile,
|
|
};
|
|
use crate::wooden_fish::tables::{
|
|
wooden_fish_agent_session, wooden_fish_event, wooden_fish_runtime_run, wooden_fish_work_profile,
|
|
};
|
|
use crate::{
|
|
visual_novel_agent_message, visual_novel_agent_session, visual_novel_runtime_event,
|
|
visual_novel_runtime_history_entry, visual_novel_runtime_run, visual_novel_work_profile,
|
|
};
|
|
|
|
const MIGRATION_SCHEMA_VERSION: u32 = 1;
|
|
const MIGRATION_MAX_TABLE_NAME_LEN: usize = 96;
|
|
const MIGRATION_MAX_IMPORT_UPLOAD_ID_LEN: usize = 128;
|
|
const MIGRATION_MAX_IMPORT_CHUNK_BYTES: usize = 1024 * 1024;
|
|
const MIGRATION_MAX_OPERATOR_NOTE_CHARS: usize = 160;
|
|
const MIGRATION_BOOTSTRAP_SECRET_HEX_LEN: usize = 64;
|
|
const MIGRATION_BOOTSTRAP_SECRET_SHA256: Option<&str> =
|
|
option_env!("GENARRATIVE_SPACETIME_MIGRATION_BOOTSTRAP_SECRET_SHA256");
|
|
|
|
#[spacetimedb::table(accessor = database_migration_operator)]
|
|
pub struct DatabaseMigrationOperator {
|
|
#[primary_key]
|
|
pub operator_identity: Identity,
|
|
pub created_at: Timestamp,
|
|
pub created_by: Identity,
|
|
pub note: String,
|
|
}
|
|
|
|
#[spacetimedb::table(
|
|
accessor = database_migration_import_chunk,
|
|
index(accessor = by_database_migration_import_upload, btree(columns = [upload_id]))
|
|
)]
|
|
pub struct DatabaseMigrationImportChunk {
|
|
#[primary_key]
|
|
pub chunk_key: String,
|
|
pub upload_id: String,
|
|
pub chunk_index: u32,
|
|
pub chunk_count: u32,
|
|
pub operator_identity: Identity,
|
|
pub created_at: Timestamp,
|
|
pub chunk: String,
|
|
}
|
|
|
|
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
|
|
pub struct DatabaseMigrationExportInput {
|
|
pub include_tables: Vec<String>,
|
|
}
|
|
|
|
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
|
|
pub struct DatabaseMigrationImportInput {
|
|
pub migration_json: String,
|
|
pub include_tables: Vec<String>,
|
|
pub replace_existing: bool,
|
|
pub dry_run: bool,
|
|
}
|
|
|
|
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
|
|
pub struct DatabaseMigrationImportChunkInput {
|
|
pub upload_id: String,
|
|
pub chunk_index: u32,
|
|
pub chunk_count: u32,
|
|
pub chunk: String,
|
|
}
|
|
|
|
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
|
|
pub struct DatabaseMigrationImportChunksInput {
|
|
pub upload_id: String,
|
|
pub include_tables: Vec<String>,
|
|
pub replace_existing: bool,
|
|
pub dry_run: bool,
|
|
}
|
|
|
|
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
|
|
pub struct DatabaseMigrationImportChunksClearInput {
|
|
pub upload_id: String,
|
|
}
|
|
|
|
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
|
|
enum DatabaseMigrationImportMode {
|
|
Strict,
|
|
Incremental,
|
|
}
|
|
|
|
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
|
|
pub struct DatabaseMigrationAuthorizeOperatorInput {
|
|
pub bootstrap_secret: String,
|
|
pub operator_identity_hex: String,
|
|
pub note: String,
|
|
}
|
|
|
|
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
|
|
pub struct DatabaseMigrationRevokeOperatorInput {
|
|
pub operator_identity_hex: String,
|
|
}
|
|
|
|
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
|
|
pub struct DatabaseMigrationTableStat {
|
|
pub table_name: String,
|
|
pub exported_row_count: u64,
|
|
pub imported_row_count: u64,
|
|
pub skipped_row_count: u64,
|
|
}
|
|
|
|
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
|
|
pub struct DatabaseMigrationWarning {
|
|
pub table_name: String,
|
|
pub warning_kind: String,
|
|
pub message: String,
|
|
}
|
|
|
|
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
|
|
pub struct DatabaseMigrationProcedureResult {
|
|
pub ok: bool,
|
|
pub schema_version: u32,
|
|
pub migration_json: Option<String>,
|
|
pub table_stats: Vec<DatabaseMigrationTableStat>,
|
|
pub warnings: Vec<DatabaseMigrationWarning>,
|
|
pub error_message: Option<String>,
|
|
}
|
|
|
|
#[derive(Clone, Debug, PartialEq, Eq, SpacetimeType)]
|
|
pub struct DatabaseMigrationOperatorProcedureResult {
|
|
pub ok: bool,
|
|
pub operator_identity_hex: Option<String>,
|
|
pub error_message: Option<String>,
|
|
}
|
|
|
|
#[derive(Serialize, Deserialize)]
|
|
struct MigrationFile {
|
|
schema_version: u32,
|
|
exported_at_micros: i64,
|
|
tables: Vec<MigrationTable>,
|
|
}
|
|
|
|
#[derive(Serialize, Deserialize)]
|
|
struct MigrationTable {
|
|
name: String,
|
|
rows: Vec<serde_json::Value>,
|
|
}
|
|
|
|
macro_rules! migration_tables {
|
|
($macro_name:ident $(, $arg:expr)* $(,)?) => {
|
|
$macro_name! {
|
|
$($arg,)*
|
|
auth_store_projection_meta,
|
|
admin_account,
|
|
user_account,
|
|
auth_identity,
|
|
refresh_session,
|
|
ai_task,
|
|
ai_task_stage,
|
|
ai_text_chunk,
|
|
ai_result_reference,
|
|
ai_task_event,
|
|
external_generation_job,
|
|
external_generation_job_summary,
|
|
external_generation_job_event,
|
|
runtime_snapshot,
|
|
runtime_setting,
|
|
creation_entry_config,
|
|
creation_entry_type_config,
|
|
feature_gate_config,
|
|
user_browse_history,
|
|
profile_dashboard_state,
|
|
profile_daily_free_points,
|
|
profile_wallet_ledger,
|
|
profile_wallet_consumption_total,
|
|
asset_operation_wallet_settlement,
|
|
profile_wallet_config,
|
|
analytics_date_dimension,
|
|
tracking_event,
|
|
tracking_daily_stat,
|
|
profile_task_config,
|
|
profile_task_progress,
|
|
profile_task_reward_claim,
|
|
profile_redeem_code,
|
|
profile_redeem_code_usage,
|
|
profile_code_operation,
|
|
profile_invite_code,
|
|
profile_referral_relation,
|
|
profile_played_world,
|
|
public_work_play_daily_stat,
|
|
public_work_like,
|
|
// profile_membership / profile_recharge_product_config carry membership-cycle point fields.
|
|
profile_membership,
|
|
profile_recharge_product_config,
|
|
profile_recharge_order,
|
|
profile_recharge_refund,
|
|
profile_recharge_refund_observation,
|
|
profile_recharge_order_refund_settlement,
|
|
profile_recharge_refund_hold,
|
|
profile_wallet_manual_restriction,
|
|
profile_recharge_refund_bill_checkpoint,
|
|
profile_recharge_order_expiration_schedule,
|
|
profile_recharge_order_expiration_timer,
|
|
profile_feedback_submission,
|
|
profile_save_archive,
|
|
player_progression,
|
|
chapter_progression,
|
|
npc_state,
|
|
story_session,
|
|
story_event,
|
|
inventory_slot,
|
|
battle_state,
|
|
treasure_record,
|
|
quest_record,
|
|
quest_log,
|
|
custom_world_profile,
|
|
custom_world_session,
|
|
custom_world_agent_session,
|
|
custom_world_agent_message,
|
|
custom_world_agent_operation,
|
|
custom_world_draft_card,
|
|
custom_world_gallery_entry,
|
|
asset_object,
|
|
asset_entity_binding,
|
|
asset_event,
|
|
external_api_key,
|
|
editor_agent_conversation,
|
|
editor_project,
|
|
editor_canvas,
|
|
editor_canvas_layer,
|
|
editor_canvas_generation_dialog,
|
|
editor_canvas_layout_migration,
|
|
editor_project_resource,
|
|
editor_asset_folder,
|
|
editor_asset,
|
|
editor_asset_group_source_provenance,
|
|
editor_asset_group_cohort,
|
|
editor_showcase_asset,
|
|
editor_showcase_asset_like,
|
|
editor_showcase_campaign_config,
|
|
editor_generation_pricing_config,
|
|
editor_generation_runtime_identity_rotation,
|
|
puzzle_agent_session,
|
|
puzzle_background_compile_task,
|
|
puzzle_agent_message,
|
|
puzzle_work_profile,
|
|
puzzle_event,
|
|
puzzle_runtime_run,
|
|
puzzle_leaderboard_entry,
|
|
puzzle_clear_agent_session,
|
|
puzzle_clear_work_profile,
|
|
puzzle_clear_runtime_run,
|
|
puzzle_clear_event,
|
|
bark_battle_draft_config,
|
|
bark_battle_published_config,
|
|
bark_battle_runtime_run,
|
|
bark_battle_score_record,
|
|
bark_battle_leaderboard_entry,
|
|
bark_battle_work_stats_projection,
|
|
bark_battle_personal_best_projection,
|
|
match3d_agent_session,
|
|
match3d_agent_message,
|
|
match_3_d_work_profile,
|
|
match3d_runtime_run,
|
|
jump_hop_agent_session,
|
|
jump_hop_work_profile,
|
|
jump_hop_runtime_run,
|
|
jump_hop_event,
|
|
jump_hop_leaderboard_entry,
|
|
wooden_fish_agent_session,
|
|
wooden_fish_work_profile,
|
|
wooden_fish_runtime_run,
|
|
wooden_fish_event,
|
|
square_hole_agent_session,
|
|
square_hole_agent_message,
|
|
square_hole_work_profile,
|
|
square_hole_runtime_run,
|
|
visual_novel_agent_session,
|
|
visual_novel_agent_message,
|
|
visual_novel_work_profile,
|
|
visual_novel_runtime_run,
|
|
visual_novel_runtime_history_entry,
|
|
visual_novel_runtime_event,
|
|
big_fish_creation_session,
|
|
big_fish_agent_message,
|
|
big_fish_asset_slot,
|
|
big_fish_runtime_run,
|
|
big_fish_event
|
|
}
|
|
};
|
|
}
|
|
|
|
macro_rules! collect_all_migration_tables {
|
|
($ctx:expr, $include_tables:expr, $tables:expr) => {
|
|
migration_tables!(collect_migration_table, $ctx, $include_tables, $tables);
|
|
};
|
|
}
|
|
|
|
macro_rules! collect_migration_table {
|
|
($ctx:expr, $include_tables:expr, $tables:expr, $($table:ident),+ $(,)?) => {
|
|
$(
|
|
if should_include_table($include_tables, stringify!($table)) {
|
|
let rows = $ctx
|
|
.db
|
|
.$table()
|
|
.iter()
|
|
.map(|row| row_to_json(&row))
|
|
.collect::<Result<Vec<_>, _>>()?;
|
|
$tables.push(MigrationTable {
|
|
name: stringify!($table).to_string(),
|
|
rows,
|
|
});
|
|
}
|
|
)+
|
|
};
|
|
}
|
|
|
|
macro_rules! clear_all_migration_tables {
|
|
($ctx:expr, $include_tables:expr) => {
|
|
migration_tables!(clear_migration_table, $ctx, $include_tables);
|
|
};
|
|
}
|
|
|
|
macro_rules! clear_migration_table {
|
|
($ctx:expr, $include_tables:expr, $($table:ident),+ $(,)?) => {
|
|
$(
|
|
if should_include_table($include_tables, stringify!($table)) {
|
|
for row in $ctx.db.$table().iter().collect::<Vec<_>>() {
|
|
$ctx.db.$table().delete(row);
|
|
}
|
|
}
|
|
)+
|
|
};
|
|
}
|
|
|
|
// 迁移权限独立存表,避免把 private 表导出能力开放给任意登录身份。
|
|
#[spacetimedb::procedure]
|
|
pub fn authorize_database_migration_operator(
|
|
ctx: &mut ProcedureContext,
|
|
input: DatabaseMigrationAuthorizeOperatorInput,
|
|
) -> DatabaseMigrationOperatorProcedureResult {
|
|
match authorize_database_migration_operator_inner(ctx, input) {
|
|
Ok(operator_identity_hex) => DatabaseMigrationOperatorProcedureResult {
|
|
ok: true,
|
|
operator_identity_hex: Some(operator_identity_hex),
|
|
error_message: None,
|
|
},
|
|
Err(error) => DatabaseMigrationOperatorProcedureResult {
|
|
ok: false,
|
|
operator_identity_hex: None,
|
|
error_message: Some(error),
|
|
},
|
|
}
|
|
}
|
|
|
|
#[spacetimedb::procedure]
|
|
pub fn revoke_database_migration_operator(
|
|
ctx: &mut ProcedureContext,
|
|
input: DatabaseMigrationRevokeOperatorInput,
|
|
) -> DatabaseMigrationOperatorProcedureResult {
|
|
match revoke_database_migration_operator_inner(ctx, input) {
|
|
Ok(operator_identity_hex) => DatabaseMigrationOperatorProcedureResult {
|
|
ok: true,
|
|
operator_identity_hex: Some(operator_identity_hex),
|
|
error_message: None,
|
|
},
|
|
Err(error) => DatabaseMigrationOperatorProcedureResult {
|
|
ok: false,
|
|
operator_identity_hex: None,
|
|
error_message: Some(error),
|
|
},
|
|
}
|
|
}
|
|
|
|
// 迁移导出走 procedure 返回 JSON 字符串,避免 reducer 无返回值且不能读取 private 表给外部。
|
|
#[spacetimedb::procedure]
|
|
pub fn export_database_migration_to_file(
|
|
ctx: &mut ProcedureContext,
|
|
input: DatabaseMigrationExportInput,
|
|
) -> DatabaseMigrationProcedureResult {
|
|
match export_database_migration_to_file_inner(ctx, input) {
|
|
Ok((migration_json, stats)) => DatabaseMigrationProcedureResult {
|
|
ok: true,
|
|
schema_version: MIGRATION_SCHEMA_VERSION,
|
|
migration_json: Some(migration_json),
|
|
table_stats: stats,
|
|
warnings: Vec::new(),
|
|
error_message: None,
|
|
},
|
|
Err(error) => DatabaseMigrationProcedureResult {
|
|
ok: false,
|
|
schema_version: MIGRATION_SCHEMA_VERSION,
|
|
migration_json: None,
|
|
table_stats: Vec::new(),
|
|
warnings: Vec::new(),
|
|
error_message: Some(error),
|
|
},
|
|
}
|
|
}
|
|
|
|
// 迁移导入由 Node 侧读文件后把 JSON 字符串传入,procedure 只负责校验和写表事务。
|
|
#[spacetimedb::procedure]
|
|
pub fn import_database_migration_from_file(
|
|
ctx: &mut ProcedureContext,
|
|
input: DatabaseMigrationImportInput,
|
|
) -> DatabaseMigrationProcedureResult {
|
|
match import_database_migration_from_file_inner(ctx, input, DatabaseMigrationImportMode::Strict)
|
|
{
|
|
Ok((stats, warnings)) => DatabaseMigrationProcedureResult {
|
|
ok: true,
|
|
schema_version: MIGRATION_SCHEMA_VERSION,
|
|
migration_json: None,
|
|
table_stats: stats,
|
|
warnings,
|
|
error_message: None,
|
|
},
|
|
Err(error) => DatabaseMigrationProcedureResult {
|
|
ok: false,
|
|
schema_version: MIGRATION_SCHEMA_VERSION,
|
|
migration_json: None,
|
|
table_stats: Vec::new(),
|
|
warnings: Vec::new(),
|
|
error_message: Some(error),
|
|
},
|
|
}
|
|
}
|
|
|
|
// 增量导入只插入目标库缺失的行;主键或唯一约束冲突的行会跳过,不更新已有数据。
|
|
#[spacetimedb::procedure]
|
|
pub fn import_database_migration_incremental_from_file(
|
|
ctx: &mut ProcedureContext,
|
|
input: DatabaseMigrationImportInput,
|
|
) -> DatabaseMigrationProcedureResult {
|
|
match import_database_migration_from_file_inner(
|
|
ctx,
|
|
input,
|
|
DatabaseMigrationImportMode::Incremental,
|
|
) {
|
|
Ok((stats, warnings)) => DatabaseMigrationProcedureResult {
|
|
ok: true,
|
|
schema_version: MIGRATION_SCHEMA_VERSION,
|
|
migration_json: None,
|
|
table_stats: stats,
|
|
warnings,
|
|
error_message: None,
|
|
},
|
|
Err(error) => DatabaseMigrationProcedureResult {
|
|
ok: false,
|
|
schema_version: MIGRATION_SCHEMA_VERSION,
|
|
migration_json: None,
|
|
table_stats: Vec::new(),
|
|
warnings: Vec::new(),
|
|
error_message: Some(error),
|
|
},
|
|
}
|
|
}
|
|
|
|
// 大迁移 JSON 先按分片写入私有临时表,避免单次 HTTP request body 触发 SpacetimeDB 413。
|
|
#[spacetimedb::procedure]
|
|
pub fn put_database_migration_import_chunk(
|
|
ctx: &mut ProcedureContext,
|
|
input: DatabaseMigrationImportChunkInput,
|
|
) -> DatabaseMigrationProcedureResult {
|
|
match put_database_migration_import_chunk_inner(ctx, input) {
|
|
Ok(()) => empty_database_migration_result(true, None),
|
|
Err(error) => empty_database_migration_result(false, Some(error)),
|
|
}
|
|
}
|
|
|
|
// 分片提交保持与直接导入相同的严格追加语义;提交成功后清理临时分片。
|
|
#[spacetimedb::procedure]
|
|
pub fn import_database_migration_from_chunks(
|
|
ctx: &mut ProcedureContext,
|
|
input: DatabaseMigrationImportChunksInput,
|
|
) -> DatabaseMigrationProcedureResult {
|
|
match import_database_migration_from_chunks_inner(
|
|
ctx,
|
|
input,
|
|
DatabaseMigrationImportMode::Strict,
|
|
) {
|
|
Ok((stats, warnings)) => DatabaseMigrationProcedureResult {
|
|
ok: true,
|
|
schema_version: MIGRATION_SCHEMA_VERSION,
|
|
migration_json: None,
|
|
table_stats: stats,
|
|
warnings,
|
|
error_message: None,
|
|
},
|
|
Err(error) => empty_database_migration_result(false, Some(error)),
|
|
}
|
|
}
|
|
|
|
// 分片增量提交只插入目标库缺失的行;主键或唯一约束冲突的行会跳过。
|
|
#[spacetimedb::procedure]
|
|
pub fn import_database_migration_incremental_from_chunks(
|
|
ctx: &mut ProcedureContext,
|
|
input: DatabaseMigrationImportChunksInput,
|
|
) -> DatabaseMigrationProcedureResult {
|
|
match import_database_migration_from_chunks_inner(
|
|
ctx,
|
|
input,
|
|
DatabaseMigrationImportMode::Incremental,
|
|
) {
|
|
Ok((stats, warnings)) => DatabaseMigrationProcedureResult {
|
|
ok: true,
|
|
schema_version: MIGRATION_SCHEMA_VERSION,
|
|
migration_json: None,
|
|
table_stats: stats,
|
|
warnings,
|
|
error_message: None,
|
|
},
|
|
Err(error) => empty_database_migration_result(false, Some(error)),
|
|
}
|
|
}
|
|
|
|
// 调用方上传失败或提交失败时可显式清理同一 upload_id 的临时分片。
|
|
#[spacetimedb::procedure]
|
|
pub fn clear_database_migration_import_chunks(
|
|
ctx: &mut ProcedureContext,
|
|
input: DatabaseMigrationImportChunksClearInput,
|
|
) -> DatabaseMigrationProcedureResult {
|
|
match clear_database_migration_import_chunks_inner(ctx, input) {
|
|
Ok(()) => empty_database_migration_result(true, None),
|
|
Err(error) => empty_database_migration_result(false, Some(error)),
|
|
}
|
|
}
|
|
|
|
fn export_database_migration_to_file_inner(
|
|
ctx: &mut ProcedureContext,
|
|
input: DatabaseMigrationExportInput,
|
|
) -> Result<(String, Vec<DatabaseMigrationTableStat>), String> {
|
|
let caller = ctx.sender();
|
|
let included_tables = normalize_include_tables(&input.include_tables)?;
|
|
let exported_at_micros = ctx.timestamp.to_micros_since_unix_epoch();
|
|
|
|
let migration_file = ctx.try_with_tx(|tx| {
|
|
require_migration_operator(tx, caller)?;
|
|
build_migration_file(tx, exported_at_micros, included_tables.as_ref())
|
|
})?;
|
|
let stats = build_export_stats(&migration_file.tables);
|
|
let content = serde_json::to_string_pretty(&migration_file)
|
|
.map_err(|error| format!("迁移文件序列化失败: {error}"))?;
|
|
|
|
Ok((content, stats))
|
|
}
|
|
|
|
fn import_database_migration_from_file_inner(
|
|
ctx: &mut ProcedureContext,
|
|
input: DatabaseMigrationImportInput,
|
|
import_mode: DatabaseMigrationImportMode,
|
|
) -> Result<
|
|
(
|
|
Vec<DatabaseMigrationTableStat>,
|
|
Vec<DatabaseMigrationWarning>,
|
|
),
|
|
String,
|
|
> {
|
|
let caller = ctx.sender();
|
|
let included_tables = normalize_include_tables(&input.include_tables)?;
|
|
if import_mode == DatabaseMigrationImportMode::Incremental && input.replace_existing {
|
|
return Err("增量导入不能同时启用 replace_existing".to_string());
|
|
}
|
|
if input.migration_json.trim().is_empty() {
|
|
return Err("migration_json 不能为空".to_string());
|
|
}
|
|
ctx.try_with_tx(|tx| require_migration_operator(tx, caller))?;
|
|
|
|
let migration_file = parse_migration_file(&input.migration_json)?;
|
|
|
|
let (stats, warnings) = if input.dry_run {
|
|
build_import_dry_run_stats(&migration_file.tables, included_tables.as_ref())?
|
|
} else {
|
|
ctx.try_with_tx(|tx| {
|
|
require_migration_operator(tx, caller)?;
|
|
apply_migration_file(
|
|
tx,
|
|
&migration_file,
|
|
included_tables.as_ref(),
|
|
input.replace_existing,
|
|
import_mode,
|
|
)
|
|
})?
|
|
};
|
|
|
|
Ok((stats, warnings))
|
|
}
|
|
|
|
fn put_database_migration_import_chunk_inner(
|
|
ctx: &mut ProcedureContext,
|
|
input: DatabaseMigrationImportChunkInput,
|
|
) -> Result<(), String> {
|
|
let caller = ctx.sender();
|
|
let upload_id = normalize_import_upload_id(&input.upload_id)?;
|
|
if input.chunk_count == 0 {
|
|
return Err("分片总数必须大于 0".to_string());
|
|
}
|
|
if input.chunk_index >= input.chunk_count {
|
|
return Err(format!(
|
|
"分片序号越界: {} / {}",
|
|
input.chunk_index, input.chunk_count
|
|
));
|
|
}
|
|
if input.chunk.is_empty() {
|
|
return Err("迁移 JSON 分片不能为空".to_string());
|
|
}
|
|
if input.chunk.len() > MIGRATION_MAX_IMPORT_CHUNK_BYTES {
|
|
return Err(format!(
|
|
"迁移 JSON 分片过大,单片最多 {} bytes",
|
|
MIGRATION_MAX_IMPORT_CHUNK_BYTES
|
|
));
|
|
}
|
|
|
|
let chunk_key = build_import_chunk_key(&upload_id, input.chunk_index);
|
|
ctx.try_with_tx(|tx| {
|
|
require_migration_operator(tx, caller)?;
|
|
if let Some(existing) = tx
|
|
.db
|
|
.database_migration_import_chunk()
|
|
.chunk_key()
|
|
.find(&chunk_key)
|
|
{
|
|
if existing.operator_identity != caller {
|
|
return Err("同名迁移分片已由其他 identity 上传,已拒绝覆盖".to_string());
|
|
}
|
|
tx.db
|
|
.database_migration_import_chunk()
|
|
.chunk_key()
|
|
.delete(&chunk_key);
|
|
}
|
|
tx.db
|
|
.database_migration_import_chunk()
|
|
.insert(DatabaseMigrationImportChunk {
|
|
chunk_key: chunk_key.clone(),
|
|
upload_id: upload_id.clone(),
|
|
chunk_index: input.chunk_index,
|
|
chunk_count: input.chunk_count,
|
|
operator_identity: caller,
|
|
created_at: tx.timestamp,
|
|
chunk: input.chunk.clone(),
|
|
});
|
|
Ok(())
|
|
})?;
|
|
|
|
Ok(())
|
|
}
|
|
|
|
fn import_database_migration_from_chunks_inner(
|
|
ctx: &mut ProcedureContext,
|
|
input: DatabaseMigrationImportChunksInput,
|
|
import_mode: DatabaseMigrationImportMode,
|
|
) -> Result<
|
|
(
|
|
Vec<DatabaseMigrationTableStat>,
|
|
Vec<DatabaseMigrationWarning>,
|
|
),
|
|
String,
|
|
> {
|
|
let caller = ctx.sender();
|
|
let upload_id = normalize_import_upload_id(&input.upload_id)?;
|
|
let included_tables = normalize_include_tables(&input.include_tables)?;
|
|
if import_mode == DatabaseMigrationImportMode::Incremental && input.replace_existing {
|
|
return Err("增量导入不能同时启用 replace_existing".to_string());
|
|
}
|
|
|
|
let migration_json = ctx.try_with_tx(|tx| {
|
|
require_migration_operator(tx, caller)?;
|
|
read_database_migration_import_chunks(tx, &upload_id, caller)
|
|
})?;
|
|
let migration_file = parse_migration_file(&migration_json)?;
|
|
|
|
let (stats, warnings) = if input.dry_run {
|
|
build_import_dry_run_stats(&migration_file.tables, included_tables.as_ref())?
|
|
} else {
|
|
ctx.try_with_tx(|tx| {
|
|
require_migration_operator(tx, caller)?;
|
|
apply_migration_file(
|
|
tx,
|
|
&migration_file,
|
|
included_tables.as_ref(),
|
|
input.replace_existing,
|
|
import_mode,
|
|
)
|
|
})?
|
|
};
|
|
|
|
ctx.try_with_tx(|tx| {
|
|
require_migration_operator(tx, caller)?;
|
|
clear_database_migration_import_chunks_tx(tx, &upload_id);
|
|
Ok::<(), String>(())
|
|
})?;
|
|
|
|
Ok((stats, warnings))
|
|
}
|
|
|
|
fn clear_database_migration_import_chunks_inner(
|
|
ctx: &mut ProcedureContext,
|
|
input: DatabaseMigrationImportChunksClearInput,
|
|
) -> Result<(), String> {
|
|
let caller = ctx.sender();
|
|
let upload_id = normalize_import_upload_id(&input.upload_id)?;
|
|
ctx.try_with_tx(|tx| {
|
|
require_migration_operator(tx, caller)?;
|
|
clear_database_migration_import_chunks_tx(tx, &upload_id);
|
|
Ok::<(), String>(())
|
|
})?;
|
|
Ok(())
|
|
}
|
|
|
|
fn empty_database_migration_result(
|
|
ok: bool,
|
|
error_message: Option<String>,
|
|
) -> DatabaseMigrationProcedureResult {
|
|
DatabaseMigrationProcedureResult {
|
|
ok,
|
|
schema_version: MIGRATION_SCHEMA_VERSION,
|
|
migration_json: None,
|
|
table_stats: Vec::new(),
|
|
warnings: Vec::new(),
|
|
error_message,
|
|
}
|
|
}
|
|
|
|
fn parse_migration_file(migration_json: &str) -> Result<MigrationFile, String> {
|
|
if migration_json.trim().is_empty() {
|
|
return Err("migration_json 不能为空".to_string());
|
|
}
|
|
|
|
let migration_file = serde_json::from_str::<MigrationFile>(migration_json)
|
|
.map_err(|error| format!("迁移文件 JSON 解析失败: {error}"))?;
|
|
if migration_file.schema_version != MIGRATION_SCHEMA_VERSION {
|
|
return Err(format!(
|
|
"迁移文件 schema_version 不匹配,期望 {},实际 {}",
|
|
MIGRATION_SCHEMA_VERSION, migration_file.schema_version
|
|
));
|
|
}
|
|
|
|
Ok(migration_file)
|
|
}
|
|
|
|
fn authorize_database_migration_operator_inner(
|
|
ctx: &mut ProcedureContext,
|
|
input: DatabaseMigrationAuthorizeOperatorInput,
|
|
) -> Result<String, String> {
|
|
let caller = ctx.sender();
|
|
let operator_identity = parse_migration_operator_identity(&input.operator_identity_hex)?;
|
|
let note = normalize_migration_operator_note(&input.note)?;
|
|
let bootstrap_secret = input.bootstrap_secret.trim().to_string();
|
|
|
|
ctx.try_with_tx(|tx| {
|
|
authorize_database_migration_operator_tx(
|
|
tx,
|
|
caller,
|
|
operator_identity,
|
|
&bootstrap_secret,
|
|
note.clone(),
|
|
)
|
|
})?;
|
|
|
|
Ok(operator_identity.to_hex().to_string())
|
|
}
|
|
|
|
fn revoke_database_migration_operator_inner(
|
|
ctx: &mut ProcedureContext,
|
|
input: DatabaseMigrationRevokeOperatorInput,
|
|
) -> Result<String, String> {
|
|
let caller = ctx.sender();
|
|
let operator_identity = parse_migration_operator_identity(&input.operator_identity_hex)?;
|
|
|
|
ctx.try_with_tx(|tx| {
|
|
require_migration_operator(tx, caller)?;
|
|
if tx
|
|
.db
|
|
.database_migration_operator()
|
|
.operator_identity()
|
|
.find(&operator_identity)
|
|
.is_none()
|
|
{
|
|
return Err("迁移操作员不存在".to_string());
|
|
}
|
|
tx.db
|
|
.database_migration_operator()
|
|
.operator_identity()
|
|
.delete(&operator_identity);
|
|
Ok(())
|
|
})?;
|
|
|
|
Ok(operator_identity.to_hex().to_string())
|
|
}
|
|
|
|
fn authorize_database_migration_operator_tx(
|
|
ctx: &ReducerContext,
|
|
caller: Identity,
|
|
operator_identity: Identity,
|
|
bootstrap_secret: &str,
|
|
note: String,
|
|
) -> Result<(), String> {
|
|
let has_operator = ctx.db.database_migration_operator().iter().next().is_some();
|
|
if has_operator {
|
|
require_migration_operator(ctx, caller)?;
|
|
} else {
|
|
require_migration_bootstrap_secret(bootstrap_secret)?;
|
|
}
|
|
|
|
require_migration_operator_candidate(
|
|
crate::editor_project_storage::is_editor_generation_runtime_service_identity(
|
|
ctx,
|
|
operator_identity,
|
|
),
|
|
)?;
|
|
|
|
if ctx
|
|
.db
|
|
.database_migration_operator()
|
|
.operator_identity()
|
|
.find(&operator_identity)
|
|
.is_some()
|
|
{
|
|
ctx.db
|
|
.database_migration_operator()
|
|
.operator_identity()
|
|
.delete(&operator_identity);
|
|
}
|
|
|
|
ctx.db
|
|
.database_migration_operator()
|
|
.insert(DatabaseMigrationOperator {
|
|
operator_identity,
|
|
created_at: ctx.timestamp,
|
|
created_by: caller,
|
|
note,
|
|
});
|
|
|
|
Ok(())
|
|
}
|
|
|
|
pub(crate) fn require_migration_operator(
|
|
ctx: &ReducerContext,
|
|
caller: Identity,
|
|
) -> Result<(), String> {
|
|
if is_database_migration_operator(ctx, caller) {
|
|
Ok(())
|
|
} else {
|
|
Err("当前 identity 未被授权执行数据库迁移".to_string())
|
|
}
|
|
}
|
|
|
|
pub(crate) fn is_database_migration_operator(ctx: &ReducerContext, caller: Identity) -> bool {
|
|
ctx.db
|
|
.database_migration_operator()
|
|
.operator_identity()
|
|
.find(&caller)
|
|
.is_some()
|
|
}
|
|
|
|
fn require_migration_operator_candidate(candidate_is_runtime_writer: bool) -> Result<(), String> {
|
|
if candidate_is_runtime_writer {
|
|
Err("模型生成运行时服务 identity 不能被授权为数据库迁移操作员".to_string())
|
|
} else {
|
|
Ok(())
|
|
}
|
|
}
|
|
|
|
pub(crate) fn require_migration_bootstrap_secret(input: &str) -> Result<(), String> {
|
|
let configured_hash = MIGRATION_BOOTSTRAP_SECRET_SHA256
|
|
.map(str::trim)
|
|
.filter(|hash| !hash.is_empty())
|
|
.ok_or_else(|| "迁移引导密钥未配置,无法创建首个操作员".to_string())?;
|
|
|
|
if configured_hash.len() != 64 || !configured_hash.bytes().all(|byte| byte.is_ascii_hexdigit())
|
|
{
|
|
return Err("迁移引导密钥摘要配置不合法".to_string());
|
|
}
|
|
let normalized_input = input.trim();
|
|
if !is_valid_migration_bootstrap_secret(normalized_input) {
|
|
return Err("迁移引导密钥必须是 64 位十六进制高熵值".to_string());
|
|
}
|
|
let input_hash = migration_bootstrap_secret_sha256(normalized_input);
|
|
let normalized_configured_hash = configured_hash.to_ascii_lowercase();
|
|
if !constant_time_ascii_eq(input_hash.as_bytes(), normalized_configured_hash.as_bytes()) {
|
|
return Err("迁移引导密钥不正确".to_string());
|
|
}
|
|
|
|
Ok(())
|
|
}
|
|
|
|
fn is_valid_migration_bootstrap_secret(secret: &str) -> bool {
|
|
secret.len() == MIGRATION_BOOTSTRAP_SECRET_HEX_LEN
|
|
&& secret.bytes().all(|byte| byte.is_ascii_hexdigit())
|
|
}
|
|
|
|
fn migration_bootstrap_secret_sha256(secret: &str) -> String {
|
|
let digest = Sha256::digest(secret.as_bytes());
|
|
digest.iter().map(|byte| format!("{byte:02x}")).collect()
|
|
}
|
|
|
|
fn constant_time_ascii_eq(left: &[u8], right: &[u8]) -> bool {
|
|
if left.len() != right.len() {
|
|
return false;
|
|
}
|
|
let mut difference = 0u8;
|
|
for (left_byte, right_byte) in left.iter().zip(right.iter()) {
|
|
difference |= left_byte ^ right_byte;
|
|
}
|
|
difference == 0
|
|
}
|
|
|
|
fn parse_migration_operator_identity(input: &str) -> Result<Identity, String> {
|
|
let identity_hex = input.trim().trim_start_matches("0x");
|
|
if identity_hex.len() != 64 {
|
|
return Err("operator_identity_hex 必须是 64 位十六进制 identity".to_string());
|
|
}
|
|
|
|
Identity::from_hex(identity_hex)
|
|
.map_err(|error| format!("operator_identity_hex 格式不合法: {error}"))
|
|
}
|
|
|
|
fn normalize_migration_operator_note(input: &str) -> Result<String, String> {
|
|
let note = input.trim();
|
|
if note.chars().count() > MIGRATION_MAX_OPERATOR_NOTE_CHARS {
|
|
return Err(format!(
|
|
"迁移操作员备注过长,最多 {} 个字符",
|
|
MIGRATION_MAX_OPERATOR_NOTE_CHARS
|
|
));
|
|
}
|
|
|
|
Ok(note.to_string())
|
|
}
|
|
|
|
fn normalize_import_upload_id(input: &str) -> Result<String, String> {
|
|
let upload_id = input.trim();
|
|
if upload_id.is_empty() {
|
|
return Err("upload_id 不能为空".to_string());
|
|
}
|
|
if upload_id.len() > MIGRATION_MAX_IMPORT_UPLOAD_ID_LEN {
|
|
return Err(format!(
|
|
"upload_id 过长,最多 {} bytes",
|
|
MIGRATION_MAX_IMPORT_UPLOAD_ID_LEN
|
|
));
|
|
}
|
|
if !upload_id
|
|
.chars()
|
|
.all(|character| character.is_ascii_alphanumeric() || matches!(character, '-' | '_'))
|
|
{
|
|
return Err("upload_id 只能使用 ASCII 字母、数字、短横线或下划线".to_string());
|
|
}
|
|
Ok(upload_id.to_string())
|
|
}
|
|
|
|
fn build_import_chunk_key(upload_id: &str, chunk_index: u32) -> String {
|
|
format!("{upload_id}:{chunk_index:010}")
|
|
}
|
|
|
|
fn read_database_migration_import_chunks(
|
|
ctx: &ReducerContext,
|
|
upload_id: &str,
|
|
caller: Identity,
|
|
) -> Result<String, String> {
|
|
let mut chunks = ctx
|
|
.db
|
|
.database_migration_import_chunk()
|
|
.by_database_migration_import_upload()
|
|
.filter(upload_id)
|
|
.collect::<Vec<_>>();
|
|
if chunks.is_empty() {
|
|
return Err(format!("未找到迁移 JSON 分片: {upload_id}"));
|
|
}
|
|
if chunks.iter().any(|chunk| chunk.operator_identity != caller) {
|
|
return Err("迁移 JSON 分片包含其他 identity 上传的片段,已拒绝提交".to_string());
|
|
}
|
|
|
|
let chunk_count = chunks[0].chunk_count;
|
|
if chunk_count == 0 {
|
|
return Err("迁移 JSON 分片总数不合法".to_string());
|
|
}
|
|
if chunks
|
|
.iter()
|
|
.any(|chunk| chunk.chunk_count != chunk_count || chunk.upload_id != upload_id)
|
|
{
|
|
return Err("迁移 JSON 分片总数不一致".to_string());
|
|
}
|
|
if chunks.len() != chunk_count as usize {
|
|
return Err(format!(
|
|
"迁移 JSON 分片未上传完整,已收到 {} / {}",
|
|
chunks.len(),
|
|
chunk_count
|
|
));
|
|
}
|
|
|
|
chunks.sort_by_key(|chunk| chunk.chunk_index);
|
|
let mut expected_index = 0u32;
|
|
let mut migration_json = String::new();
|
|
for chunk in chunks {
|
|
if chunk.chunk_index != expected_index {
|
|
return Err(format!("迁移 JSON 分片缺失序号: {expected_index}"));
|
|
}
|
|
migration_json.push_str(&chunk.chunk);
|
|
expected_index = expected_index.saturating_add(1);
|
|
}
|
|
|
|
Ok(migration_json)
|
|
}
|
|
|
|
fn clear_database_migration_import_chunks_tx(ctx: &ReducerContext, upload_id: &str) {
|
|
let chunk_keys = ctx
|
|
.db
|
|
.database_migration_import_chunk()
|
|
.by_database_migration_import_upload()
|
|
.filter(upload_id)
|
|
.map(|chunk| chunk.chunk_key)
|
|
.collect::<Vec<_>>();
|
|
for chunk_key in chunk_keys {
|
|
ctx.db
|
|
.database_migration_import_chunk()
|
|
.chunk_key()
|
|
.delete(&chunk_key);
|
|
}
|
|
}
|
|
|
|
fn normalize_include_tables(input: &[String]) -> Result<Option<HashSet<String>>, String> {
|
|
if input.is_empty() {
|
|
return Ok(None);
|
|
}
|
|
|
|
let mut tables = HashSet::new();
|
|
for raw_name in input {
|
|
let name = raw_name.trim();
|
|
if name.is_empty() {
|
|
continue;
|
|
}
|
|
if name.len() > MIGRATION_MAX_TABLE_NAME_LEN {
|
|
return Err(format!("迁移表名过长: {name}"));
|
|
}
|
|
if !is_supported_migration_table(name) {
|
|
return Err(format!("迁移表不在白名单内: {name}"));
|
|
}
|
|
tables.insert(name.to_string());
|
|
}
|
|
Ok(Some(tables))
|
|
}
|
|
|
|
fn should_include_table(include_tables: Option<&HashSet<String>>, table_name: &str) -> bool {
|
|
include_tables
|
|
.map(|tables| tables.contains(table_name))
|
|
.unwrap_or(true)
|
|
}
|
|
|
|
fn build_migration_file(
|
|
ctx: &ReducerContext,
|
|
exported_at_micros: i64,
|
|
include_tables: Option<&HashSet<String>>,
|
|
) -> Result<MigrationFile, String> {
|
|
let mut tables = Vec::new();
|
|
collect_all_migration_tables!(ctx, include_tables, tables);
|
|
|
|
Ok(MigrationFile {
|
|
schema_version: MIGRATION_SCHEMA_VERSION,
|
|
exported_at_micros,
|
|
tables,
|
|
})
|
|
}
|
|
|
|
fn build_export_stats(tables: &[MigrationTable]) -> Vec<DatabaseMigrationTableStat> {
|
|
tables
|
|
.iter()
|
|
.map(|table| DatabaseMigrationTableStat {
|
|
table_name: table.name.clone(),
|
|
exported_row_count: table.rows.len() as u64,
|
|
imported_row_count: 0,
|
|
skipped_row_count: 0,
|
|
})
|
|
.collect()
|
|
}
|
|
|
|
fn build_import_dry_run_stats(
|
|
tables: &[MigrationTable],
|
|
include_tables: Option<&HashSet<String>>,
|
|
) -> Result<
|
|
(
|
|
Vec<DatabaseMigrationTableStat>,
|
|
Vec<DatabaseMigrationWarning>,
|
|
),
|
|
String,
|
|
> {
|
|
let mut stats = Vec::new();
|
|
let mut warnings = Vec::new();
|
|
for table in tables {
|
|
if !is_supported_migration_table(&table.name) {
|
|
warnings.push(build_dropped_table_warning(table));
|
|
stats.push(DatabaseMigrationTableStat {
|
|
table_name: table.name.clone(),
|
|
exported_row_count: 0,
|
|
imported_row_count: 0,
|
|
skipped_row_count: table.rows.len() as u64,
|
|
});
|
|
continue;
|
|
}
|
|
if should_include_table(include_tables, &table.name) {
|
|
stats.push(DatabaseMigrationTableStat {
|
|
table_name: table.name.clone(),
|
|
exported_row_count: 0,
|
|
imported_row_count: table.rows.len() as u64,
|
|
skipped_row_count: 0,
|
|
});
|
|
} else {
|
|
stats.push(DatabaseMigrationTableStat {
|
|
table_name: table.name.clone(),
|
|
exported_row_count: 0,
|
|
imported_row_count: 0,
|
|
skipped_row_count: table.rows.len() as u64,
|
|
});
|
|
}
|
|
}
|
|
Ok((stats, warnings))
|
|
}
|
|
|
|
fn apply_migration_file(
|
|
ctx: &ReducerContext,
|
|
migration_file: &MigrationFile,
|
|
include_tables: Option<&HashSet<String>>,
|
|
replace_existing: bool,
|
|
import_mode: DatabaseMigrationImportMode,
|
|
) -> Result<
|
|
(
|
|
Vec<DatabaseMigrationTableStat>,
|
|
Vec<DatabaseMigrationWarning>,
|
|
),
|
|
String,
|
|
> {
|
|
let mut stats = Vec::new();
|
|
let mut warnings = Vec::new();
|
|
|
|
let import_table_names = build_import_table_name_set(migration_file, include_tables);
|
|
if replace_existing {
|
|
// replace_existing 只覆盖本次迁移文件实际会导入的表,避免分批导入时误清空其它迁移白名单表。
|
|
clear_all_migration_tables!(ctx, Some(&import_table_names));
|
|
}
|
|
|
|
for table in &migration_file.tables {
|
|
if !is_supported_migration_table(&table.name) {
|
|
warnings.push(build_dropped_table_warning(table));
|
|
stats.push(DatabaseMigrationTableStat {
|
|
table_name: table.name.clone(),
|
|
exported_row_count: 0,
|
|
imported_row_count: 0,
|
|
skipped_row_count: table.rows.len() as u64,
|
|
});
|
|
continue;
|
|
}
|
|
|
|
if !should_include_table(include_tables, &table.name) {
|
|
stats.push(DatabaseMigrationTableStat {
|
|
table_name: table.name.clone(),
|
|
exported_row_count: 0,
|
|
imported_row_count: 0,
|
|
skipped_row_count: table.rows.len() as u64,
|
|
});
|
|
continue;
|
|
}
|
|
|
|
let (imported_row_count, skipped_row_count) =
|
|
insert_migration_table_rows(ctx, table, import_mode, &mut warnings)?;
|
|
stats.push(DatabaseMigrationTableStat {
|
|
table_name: table.name.clone(),
|
|
exported_row_count: 0,
|
|
imported_row_count,
|
|
skipped_row_count,
|
|
});
|
|
}
|
|
|
|
Ok((stats, warnings))
|
|
}
|
|
|
|
fn build_import_table_name_set(
|
|
migration_file: &MigrationFile,
|
|
include_tables: Option<&HashSet<String>>,
|
|
) -> HashSet<String> {
|
|
migration_file
|
|
.tables
|
|
.iter()
|
|
.filter(|table| should_include_table(include_tables, &table.name))
|
|
.map(|table| table.name.clone())
|
|
.collect()
|
|
}
|
|
|
|
fn build_dropped_table_warning(table: &MigrationTable) -> DatabaseMigrationWarning {
|
|
DatabaseMigrationWarning {
|
|
table_name: table.name.clone(),
|
|
warning_kind: "dropped_table".to_string(),
|
|
message: format!(
|
|
"迁移文件包含当前模块已删除或未加入白名单的表 {},已跳过 {} 行",
|
|
table.name,
|
|
table.rows.len()
|
|
),
|
|
}
|
|
}
|
|
|
|
fn build_dropped_field_warning(table_name: &str, field_name: &str) -> DatabaseMigrationWarning {
|
|
DatabaseMigrationWarning {
|
|
table_name: table_name.to_string(),
|
|
warning_kind: "dropped_field".to_string(),
|
|
message: format!("表 {table_name} 的旧字段 {field_name} 当前已不存在,已在导入时丢弃"),
|
|
}
|
|
}
|
|
|
|
fn row_to_json<T: spacetimedb::Serialize>(row: &T) -> Result<serde_json::Value, String> {
|
|
serde_json::to_value(SerializeWrapper::from_ref(row))
|
|
.map_err(|error| format!("迁移行序列化失败: {error}"))
|
|
}
|
|
|
|
fn row_from_json<T>(
|
|
table_name: &str,
|
|
value: &serde_json::Value,
|
|
warnings: &mut Vec<DatabaseMigrationWarning>,
|
|
) -> Result<T, String>
|
|
where
|
|
T: for<'de> spacetimedb::Deserialize<'de>,
|
|
{
|
|
let wrapped = match serde_json::from_value::<DeserializeWrapper<T>>(value.clone()) {
|
|
Ok(row) => row,
|
|
Err(original_error) => recover_row_with_deleted_fields::<T>(
|
|
table_name,
|
|
value,
|
|
&original_error.to_string(),
|
|
warnings,
|
|
)
|
|
.ok_or_else(|| format!("迁移行反序列化失败,且无法通过丢弃旧字段恢复: {original_error}"))?,
|
|
};
|
|
Ok(wrapped.0)
|
|
}
|
|
|
|
fn normalize_migration_row(table_name: &str, value: &serde_json::Value) -> serde_json::Value {
|
|
let mut next_value = value.clone();
|
|
if table_name == "profile_wallet_config" {
|
|
if let Some(object) = next_value.as_object_mut() {
|
|
// 中文注释:旧迁移包没有每日免费额度字段,导入时保持原有每日 20 泥点语义。
|
|
object
|
|
.entry("daily_free_points_per_day".to_string())
|
|
.or_insert_with(|| serde_json::Value::from(20));
|
|
}
|
|
}
|
|
if table_name == "creation_entry_config" {
|
|
if let Some(object) = next_value.as_object_mut() {
|
|
// 中文注释:入口活动横幅字段晚于创作入口配置表加入,旧迁移包按运行态默认横幅兼容。
|
|
object
|
|
.entry("event_title".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
object
|
|
.entry("event_description".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
object
|
|
.entry("event_cover_image_src".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
object
|
|
.entry("event_prize_pool_mud_points".to_string())
|
|
.or_insert_with(|| serde_json::Value::from(58_000));
|
|
object
|
|
.entry("event_starts_at_text".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
object
|
|
.entry("event_ends_at_text".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
object
|
|
.entry("event_banners_json".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
object
|
|
.entry("public_work_interactions_json".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
}
|
|
}
|
|
if table_name == "creation_entry_type_config" {
|
|
if let Some(object) = next_value.as_object_mut() {
|
|
// 中文注释:入口分类和统一创作契约字段晚于入口类型配置表加入,旧迁移包按空配置兼容。
|
|
object
|
|
.entry("category_id".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
object
|
|
.entry("category_label".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
object
|
|
.entry("category_sort_order".to_string())
|
|
.or_insert_with(|| serde_json::Value::from(0));
|
|
object
|
|
.entry("unified_creation_spec_json".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
}
|
|
}
|
|
if table_name == "user_account" {
|
|
if let Some(object) = next_value.as_object_mut() {
|
|
// 中文注释:头像字段晚于认证拆表加入,旧迁移包按未设置头像兼容。
|
|
object
|
|
.entry("avatar_url".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
// 中文注释:账号标签字段晚于认证表加入,旧迁移包默认无标签。
|
|
object
|
|
.entry("user_tags".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
}
|
|
}
|
|
if table_name == "profile_invite_code" {
|
|
if let Some(object) = next_value.as_object_mut() {
|
|
// 中文注释:邀请码 metadata 晚于邀请表加入,旧迁移包按空对象兼容。
|
|
object
|
|
.entry("metadata_json".to_string())
|
|
.or_insert_with(|| serde_json::Value::String("{}".to_string()));
|
|
}
|
|
}
|
|
if table_name == "profile_redeem_code" {
|
|
if let Some(object) = next_value.as_object_mut() {
|
|
// 中文注释:兑换码生效时间晚于兑换码表加入,旧迁移包按无时间边界兼容。
|
|
object
|
|
.entry("starts_at".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
object
|
|
.entry("expires_at".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
}
|
|
}
|
|
if table_name == "profile_recharge_order" {
|
|
if let Some(object) = next_value.as_object_mut() {
|
|
// 中文注释:真实微信支付接入后才有平台交易号,旧迁移包按未回填处理。
|
|
object
|
|
.entry("provider_transaction_id".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
object
|
|
.entry("expired_at".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
object
|
|
.entry("expiration_checked_at".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
object
|
|
.entry("expiration_provider_state".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
object
|
|
.entry("expiration_last_error".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
}
|
|
}
|
|
if table_name == "external_generation_job_summary" {
|
|
if let Some(object) = next_value.as_object_mut() {
|
|
// 中文注释:非阻断完成告警晚于轻量任务摘要表加入,旧迁移包按无告警兼容。
|
|
object
|
|
.entry("warning_message".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
}
|
|
}
|
|
if table_name == "editor_showcase_asset" {
|
|
if let Some(object) = next_value.as_object_mut() {
|
|
// 中文注释:精选分类字段晚于审核表加入,旧迁移包按未设置分类兼容。
|
|
object
|
|
.entry("showcase_category".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
}
|
|
}
|
|
if table_name == "editor_canvas" {
|
|
if let Some(object) = next_value.as_object_mut() {
|
|
// 中文注释:结构化画布上线前的迁移包仍以 layers_json 为唯一真相。
|
|
object
|
|
.entry("revision".to_string())
|
|
.or_insert_with(|| serde_json::Value::from(0));
|
|
object
|
|
.entry("layout_storage_version".to_string())
|
|
.or_insert_with(|| serde_json::Value::from(0));
|
|
object
|
|
.entry("background_color".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
}
|
|
}
|
|
if table_name == "editor_canvas_layer" {
|
|
if let Some(object) = next_value.as_object_mut() {
|
|
// 中文注释:图层素材类型覆盖晚于结构化画布图层表加入,旧迁移包默认继承资源类型。
|
|
object
|
|
.entry("asset_kind_override".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
}
|
|
}
|
|
if table_name == "editor_showcase_campaign_config" {
|
|
if let Some(object) = next_value.as_object_mut() {
|
|
// 中文注释:活动卡图片 Object Key 晚于活动卡表加入,旧迁移包按仅有图片地址兼容。
|
|
object
|
|
.entry("image_object_key".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
object
|
|
.entry("image_width".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
object
|
|
.entry("image_height".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
}
|
|
}
|
|
if table_name == "external_generation_job" || table_name == "external_generation_job_summary" {
|
|
if let Some(object) = next_value.as_object_mut() {
|
|
// 中文注释:执行阶段晚于外部生成主表和摘要投影加入,旧迁移包按未知阶段兼容;
|
|
// BFF 会把 running + phase=null 视为 generating。
|
|
object
|
|
.entry("phase".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
}
|
|
}
|
|
if table_name == "big_fish_creation_session" {
|
|
if let Some(object) = next_value.as_object_mut() {
|
|
// 中文注释:旧迁移包没有公开游玩次数字段,导入时按新建作品默认 0 兼容。
|
|
object
|
|
.entry("play_count".to_string())
|
|
.or_insert_with(|| serde_json::Value::from(0));
|
|
object
|
|
.entry("remix_count".to_string())
|
|
.or_insert_with(|| serde_json::Value::from(0));
|
|
object
|
|
.entry("like_count".to_string())
|
|
.or_insert_with(|| serde_json::Value::from(0));
|
|
object
|
|
.entry("published_at".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
}
|
|
}
|
|
if table_name == "custom_world_profile" || table_name == "custom_world_gallery_entry" {
|
|
if let Some(object) = next_value.as_object_mut() {
|
|
// 中文注释:作品可见性字段晚于首版作品表加入,旧迁移包保留历史公开默认。
|
|
object
|
|
.entry("visible".to_string())
|
|
.or_insert_with(|| serde_json::Value::Bool(true));
|
|
// 中文注释:自定义世界公开互动计数字段晚于基础作品表加入,旧迁移包按 0 兼容。
|
|
object
|
|
.entry("play_count".to_string())
|
|
.or_insert_with(|| serde_json::Value::from(0));
|
|
object
|
|
.entry("remix_count".to_string())
|
|
.or_insert_with(|| serde_json::Value::from(0));
|
|
object
|
|
.entry("like_count".to_string())
|
|
.or_insert_with(|| serde_json::Value::from(0));
|
|
}
|
|
}
|
|
if table_name == "puzzle_work_profile" {
|
|
if let Some(object) = next_value.as_object_mut() {
|
|
// 中文注释:作品可见性字段晚于首版作品表加入,旧迁移包保留历史公开默认。
|
|
object
|
|
.entry("visible".to_string())
|
|
.or_insert_with(|| serde_json::Value::Bool(true));
|
|
// 中文注释:拼图公开互动计数晚于基础作品表加入,旧迁移包按 0 兼容。
|
|
object
|
|
.entry("play_count".to_string())
|
|
.or_insert_with(|| serde_json::Value::from(0));
|
|
object
|
|
.entry("remix_count".to_string())
|
|
.or_insert_with(|| serde_json::Value::from(0));
|
|
object
|
|
.entry("like_count".to_string())
|
|
.or_insert_with(|| serde_json::Value::from(0));
|
|
object
|
|
.entry("point_incentive_total_half_points".to_string())
|
|
.or_insert_with(|| serde_json::Value::from(0));
|
|
object
|
|
.entry("point_incentive_claimed_points".to_string())
|
|
.or_insert_with(|| serde_json::Value::from(0));
|
|
// 中文注释:拼图多关卡字段晚于旧作品表加入,旧迁移包留空并由读取层补出首关。
|
|
object
|
|
.entry("levels_json".to_string())
|
|
.or_insert_with(|| serde_json::Value::from(""));
|
|
// 中文注释:作品名称/描述从旧关卡名/画面摘要拆出,旧行保留旧值做兼容回填。
|
|
let fallback_title = object
|
|
.get("level_name")
|
|
.cloned()
|
|
.unwrap_or_else(|| serde_json::Value::from(""));
|
|
object
|
|
.entry("work_title".to_string())
|
|
.or_insert(fallback_title);
|
|
let fallback_description = object
|
|
.get("summary")
|
|
.cloned()
|
|
.unwrap_or_else(|| serde_json::Value::from(""));
|
|
object
|
|
.entry("work_description".to_string())
|
|
.or_insert(fallback_description);
|
|
}
|
|
}
|
|
if table_name == "big_fish_creation_session" {
|
|
if let Some(object) = next_value.as_object_mut() {
|
|
// 中文注释:作品可见性字段晚于大鱼吃小鱼创作会话表加入,旧迁移包保留历史公开默认。
|
|
object
|
|
.entry("visible".to_string())
|
|
.or_insert_with(|| serde_json::Value::Bool(true));
|
|
}
|
|
}
|
|
if matches!(
|
|
table_name,
|
|
"jump_hop_work_profile"
|
|
| "puzzle_clear_work_profile"
|
|
| "square_hole_work_profile"
|
|
| "visual_novel_work_profile"
|
|
| "bark_battle_published_config"
|
|
) {
|
|
if let Some(object) = next_value.as_object_mut() {
|
|
// 中文注释:作品可见性字段晚于首版作品表加入,旧迁移包保留历史公开默认。
|
|
object
|
|
.entry("visible".to_string())
|
|
.or_insert_with(|| serde_json::Value::Bool(true));
|
|
if table_name == "puzzle_clear_work_profile" {
|
|
// 中文注释:拼消消底图提示词字段晚于作品表加入,旧迁移包按空提示词兼容。
|
|
object
|
|
.entry("board_background_prompt".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
}
|
|
if table_name == "jump_hop_work_profile" {
|
|
// 中文注释:跳一跳主题返回按钮资产晚于首版作品表加入,旧迁移包按未生成按钮兼容。
|
|
object
|
|
.entry("back_button_asset_json".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
}
|
|
}
|
|
}
|
|
if table_name == "match_3_d_work_profile" || table_name == "match3d_work_profile" {
|
|
if let Some(object) = next_value.as_object_mut() {
|
|
// 中文注释:作品可见性字段晚于首版作品表加入,旧迁移包保留历史公开默认。
|
|
object
|
|
.entry("visible".to_string())
|
|
.or_insert_with(|| serde_json::Value::Bool(true));
|
|
// 中文注释:抓大鹅生成素材字段晚于基础作品表加入,旧迁移包按未生成素材兼容。
|
|
object
|
|
.entry("generated_item_assets_json".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
}
|
|
}
|
|
if table_name == "wooden_fish_work_profile" {
|
|
if let Some(object) = next_value.as_object_mut() {
|
|
// 中文注释:作品可见性字段晚于首版作品表加入,旧迁移包保留历史公开默认。
|
|
object
|
|
.entry("visible".to_string())
|
|
.or_insert_with(|| serde_json::Value::Bool(true));
|
|
// 中文注释:敲木鱼背景环境图晚于首版作品表加入,旧迁移包按未生成背景兼容。
|
|
object
|
|
.entry("background_asset_json".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
// 中文注释:敲木鱼返回按钮图晚于首版作品表加入,旧迁移包按未生成返回按钮兼容。
|
|
object
|
|
.entry("back_button_asset_json".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
}
|
|
}
|
|
if table_name == "editor_project_resource" {
|
|
if let Some(object) = next_value.as_object_mut() {
|
|
// 中文注释:精选公开开关晚于画布资源表加入,旧生成资源按默认公开兼容。
|
|
object
|
|
.entry("public_showcase_enabled".to_string())
|
|
.or_insert_with(|| serde_json::Value::Bool(true));
|
|
}
|
|
}
|
|
if table_name == "editor_asset" {
|
|
if let Some(object) = next_value.as_object_mut() {
|
|
// 中文注释:账号素材与项目资源绑定晚于素材库表加入,旧素材按未绑定兼容。
|
|
object
|
|
.entry("source_resource_id".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
object
|
|
.entry("thumbnail_src".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
// 中文注释:后台任务归组晚于账号素材表加入,旧素材按未固化归组兼容。
|
|
object
|
|
.entry("group_task_id".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
object
|
|
.entry("group_task_expected_asset_count".to_string())
|
|
.or_insert(serde_json::Value::Null);
|
|
}
|
|
}
|
|
next_value
|
|
}
|
|
|
|
fn recover_row_with_deleted_fields<T>(
|
|
table_name: &str,
|
|
value: &serde_json::Value,
|
|
error_message: &str,
|
|
warnings: &mut Vec<DatabaseMigrationWarning>,
|
|
) -> Option<DeserializeWrapper<T>>
|
|
where
|
|
T: for<'de> spacetimedb::Deserialize<'de>,
|
|
{
|
|
let mut candidate = value.as_object()?.clone();
|
|
let mut next_error = error_message.to_string();
|
|
|
|
loop {
|
|
let field_name = extract_unknown_field_name(&next_error)?;
|
|
candidate.remove(&field_name)?;
|
|
warnings.push(build_dropped_field_warning(table_name, &field_name));
|
|
|
|
match serde_json::from_value::<DeserializeWrapper<T>>(serde_json::Value::Object(
|
|
candidate.clone(),
|
|
)) {
|
|
Ok(row) => return Some(row),
|
|
Err(error) => next_error = error.to_string(),
|
|
}
|
|
}
|
|
}
|
|
|
|
fn extract_unknown_field_name(error_message: &str) -> Option<String> {
|
|
let marker = "unknown field";
|
|
let marker_index = error_message.find(marker)?;
|
|
let after_marker = error_message[marker_index + marker.len()..].trim_start();
|
|
|
|
for quote in ['`', '"', '\''] {
|
|
if let Some(rest) = after_marker.strip_prefix(quote) {
|
|
let end_index = rest.find(quote)?;
|
|
return Some(rest[..end_index].to_string());
|
|
}
|
|
}
|
|
|
|
after_marker
|
|
.split(|character: char| !character.is_ascii_alphanumeric() && character != '_')
|
|
.find(|value| !value.is_empty())
|
|
.map(str::to_string)
|
|
}
|
|
|
|
fn insert_migration_table_rows(
|
|
ctx: &ReducerContext,
|
|
table: &MigrationTable,
|
|
import_mode: DatabaseMigrationImportMode,
|
|
warnings: &mut Vec<DatabaseMigrationWarning>,
|
|
) -> Result<(u64, u64), String> {
|
|
macro_rules! insert_table_match_arm {
|
|
($($table:ident),+ $(,)?) => {
|
|
match table.name.as_str() {
|
|
$(
|
|
stringify!($table) => {
|
|
let mut imported = 0u64;
|
|
let mut skipped = 0u64;
|
|
for value in &table.rows {
|
|
let normalized_value = normalize_migration_row(stringify!($table), value);
|
|
let row = row_from_json(stringify!($table), &normalized_value, warnings)
|
|
.map_err(|error| format!("{}: {error}", stringify!($table)))?;
|
|
let insert_result = ctx.db
|
|
.$table()
|
|
.try_insert(row);
|
|
match insert_result {
|
|
Ok(_) => imported = imported.saturating_add(1),
|
|
Err(error) => {
|
|
if import_mode == DatabaseMigrationImportMode::Incremental {
|
|
skipped = skipped.saturating_add(1);
|
|
} else {
|
|
return Err(format!("{} 导入失败: {error}", stringify!($table)));
|
|
}
|
|
}
|
|
}
|
|
}
|
|
Ok((imported, skipped))
|
|
}
|
|
)+
|
|
_ => Err(format!("迁移表不在白名单内: {}", table.name)),
|
|
}
|
|
};
|
|
}
|
|
|
|
migration_tables!(insert_table_match_arm)
|
|
}
|
|
|
|
fn is_supported_migration_table(table_name: &str) -> bool {
|
|
macro_rules! supported_table_match {
|
|
($($table:ident),+ $(,)?) => {
|
|
matches!(
|
|
table_name,
|
|
$(stringify!($table))|+
|
|
)
|
|
};
|
|
}
|
|
|
|
migration_tables!(supported_table_match)
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod migration_bootstrap_secret_tests {
|
|
use super::*;
|
|
|
|
#[test]
|
|
fn old_profile_redeem_code_rows_default_to_open_validity_window() {
|
|
let normalized = normalize_migration_row(
|
|
"profile_redeem_code",
|
|
&serde_json::json!({ "code": "GIFT" }),
|
|
);
|
|
|
|
assert_eq!(normalized["starts_at"], serde_json::Value::Null);
|
|
assert_eq!(normalized["expires_at"], serde_json::Value::Null);
|
|
}
|
|
|
|
#[test]
|
|
fn old_profile_wallet_config_rows_default_to_twenty_daily_free_points() {
|
|
let normalized = normalize_migration_row(
|
|
"profile_wallet_config",
|
|
&serde_json::json!({ "config_id": "profile_wallet" }),
|
|
);
|
|
|
|
assert_eq!(
|
|
normalized["daily_free_points_per_day"],
|
|
serde_json::json!(20)
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn old_external_generation_summary_rows_default_to_no_warning() {
|
|
let normalized = normalize_migration_row(
|
|
"external_generation_job_summary",
|
|
&serde_json::json!({ "job_id": "task-1" }),
|
|
);
|
|
|
|
assert_eq!(normalized["warning_message"], serde_json::Value::Null);
|
|
}
|
|
|
|
#[test]
|
|
fn old_editor_canvas_rows_default_to_legacy_storage() {
|
|
let normalized = normalize_migration_row(
|
|
"editor_canvas",
|
|
&serde_json::json!({ "canvas_id": "canvas-1", "layers_json": "[]" }),
|
|
);
|
|
|
|
assert_eq!(normalized["revision"], serde_json::json!(0));
|
|
assert_eq!(normalized["layout_storage_version"], serde_json::json!(0));
|
|
assert_eq!(normalized["background_color"], serde_json::Value::Null);
|
|
}
|
|
|
|
#[test]
|
|
fn old_editor_canvas_layer_rows_default_to_no_asset_kind_override() {
|
|
let normalized = normalize_migration_row(
|
|
"editor_canvas_layer",
|
|
&serde_json::json!({ "layer_key": "layer-1" }),
|
|
);
|
|
|
|
assert_eq!(normalized["asset_kind_override"], serde_json::Value::Null);
|
|
}
|
|
|
|
#[test]
|
|
fn bootstrap_secret_sha256_is_stable_and_does_not_contain_plaintext() {
|
|
let secret = "0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef";
|
|
let digest = migration_bootstrap_secret_sha256(secret);
|
|
|
|
assert_eq!(digest.len(), 64);
|
|
assert!(digest.bytes().all(|byte| byte.is_ascii_hexdigit()));
|
|
assert!(!digest.contains(secret));
|
|
assert!(constant_time_ascii_eq(digest.as_bytes(), digest.as_bytes()));
|
|
assert!(!constant_time_ascii_eq(
|
|
digest.as_bytes(),
|
|
migration_bootstrap_secret_sha256("another-secret").as_bytes(),
|
|
));
|
|
}
|
|
|
|
#[test]
|
|
fn bootstrap_secret_requires_exactly_64_hex_characters() {
|
|
assert!(is_valid_migration_bootstrap_secret(
|
|
"0123456789ABCDEF0123456789abcdef0123456789ABCDEF0123456789abcdef"
|
|
));
|
|
assert!(!is_valid_migration_bootstrap_secret("0123456789abcdef"));
|
|
assert!(!is_valid_migration_bootstrap_secret(
|
|
"g123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef"
|
|
));
|
|
assert!(!is_valid_migration_bootstrap_secret(
|
|
"0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef00"
|
|
));
|
|
}
|
|
|
|
#[test]
|
|
fn migration_operator_candidate_rejects_runtime_writer() {
|
|
assert!(require_migration_operator_candidate(false).is_ok());
|
|
assert!(
|
|
require_migration_operator_candidate(true)
|
|
.expect_err("runtime writer must not become migration operator")
|
|
.contains("不能被授权")
|
|
);
|
|
}
|
|
}
|