use chrono::Utc;
use tonic::{Request, Response, Status};
use uuid::Uuid;
use crate::proto::udb::core::backup::services::v1 as backup_pb;
use crate::runtime::channels::OperationChannel;
use crate::runtime::core::tenant_purge::plan_tenant_purge;
use crate::runtime::executor_utils::qi_runtime;
use crate::runtime::tenant_movement::{
TenantMovementOperation, TenantMovementRequest, tenant_movement_policy_status,
validate_tenant_movement_scope,
};
use super::super::native_helpers::{
admit_on as native_admit_on, native_service_context, parse_uuid, validate_request_tenant,
};
use super::BackupServiceImpl;
use super::config::{KIND_BACKUP, MANIFEST_SUFFIX, TOPIC_BACKUP_COMPLETED};
use super::errors::backup_internal_status;
use super::events::emit_event;
use super::model::{qualified_relation, sha256_hex};
use super::store::journal_run;
pub(crate) fn resolve_object_target(
svc: &BackupServiceImpl,
backend: &str,
bucket: &str,
) -> (String, String) {
let backend = if backend.trim().is_empty() {
svc.object_backend.clone()
} else {
backend.trim().to_string()
};
let bucket = if bucket.trim().is_empty() {
svc.object_bucket.clone()
} else {
bucket.trim().to_string()
};
(backend, bucket)
}
pub(crate) async fn start_tenant_backup(
svc: &BackupServiceImpl,
request: Request<backup_pb::StartTenantBackupRequest>,
) -> Result<Response<backup_pb::StartTenantBackupResponse>, Status> {
let metadata = request.metadata().clone();
let req = request.into_inner();
validate_request_tenant(&metadata, &req.tenant_id)?;
let tenant_id = req.tenant_id.trim().to_string();
if tenant_id.is_empty() {
return Err(crate::runtime::executor_utils::invalid_argument_fields(
"tenant_id is required",
[("tenant_id", "must be a non-empty tenant id")],
));
}
let movement = TenantMovementRequest {
operation: TenantMovementOperation::BackupExport,
tenant_id: &tenant_id,
target_tenant_id: None,
tenant_filter_present: true,
privileged_cross_tenant: false,
};
validate_tenant_movement_scope(&movement)
.map_err(|err| tenant_movement_policy_status(movement.operation, err))?;
let _admit = native_admit_on(
svc.channels.as_ref(),
&svc.metrics,
"backup",
OperationChannel::Admin,
&tenant_id,
None,
)
.await?;
let runtime = svc.require_runtime()?;
let pool = svc.require_pool()?;
let manifest = svc.require_manifest()?;
let _ = parse_uuid("tenant_id", &tenant_id)?;
let context = native_service_context(&metadata, &tenant_id, "");
let (object_backend, object_bucket) =
resolve_object_target(svc, &req.object_backend, &req.object_bucket);
let plan = plan_tenant_purge(manifest);
let excluded_count = plan.excluded.len() as i64;
let backup_id = Uuid::new_v4().to_string();
let started = Utc::now();
let object_prefix = format!(
"backups/{}/{}-{}/",
tenant_id,
started.format("%Y%m%dT%H%M%SZ"),
&backup_id[..8.min(backup_id.len())]
);
let mut table_entries: Vec<backup_pb::BackupTableEntry> = Vec::new();
let mut total_rows: i64 = 0;
for target in &plan.targets {
let rel = qualified_relation(&target.schema, &target.table);
let select_sql = format!(
"SELECT row_to_json(t)::text FROM {rel} t WHERE {col}::text = $1",
col = qi_runtime(&target.tenant_column),
);
let rows: Vec<String> = sqlx::query_scalar(&select_sql)
.bind(&tenant_id)
.fetch_all(pool)
.await
.map_err(|err| {
backup_internal_status(
"start_backup_read_table",
format!(
"backup failed reading {}.{}: {err}",
target.schema, target.table
),
)
})?;
let row_count = rows.len() as i64;
total_rows += row_count;
let jsonl = rows.join("\n");
let ciphertext = runtime.encrypt_secret_at_rest(&jsonl).map_err(|err| {
backup_internal_status(
"start_backup_encrypt_artifact",
format!("backup encryption failed: {err}"),
)
})?;
let bytes = ciphertext.into_bytes();
let checksum = sha256_hex(&bytes);
let object_key = format!(
"{object_prefix}{}.{}.jsonl.enc",
target.schema, target.table
);
let put_req = crate::runtime::core::setup_data::object_request_json(
"put",
&object_bucket,
&object_key,
"application/octet-stream",
);
runtime
.put_object_backend_target_for_project(
&object_backend,
None,
&context.project_id,
&put_req,
bytes,
)
.await?;
table_entries.push(backup_pb::BackupTableEntry {
schema: target.schema.clone(),
table: target.table.clone(),
tenant_column: target.tenant_column.clone(),
object_key,
row_count,
checksum_sha256: checksum,
});
}
let excluded: Vec<backup_pb::BackupExcludedTable> = plan
.excluded
.iter()
.map(|e| backup_pb::BackupExcludedTable {
schema: e.schema.clone(),
table: e.table.clone(),
reason: e.reason.clone(),
})
.collect();
let manifest_value = serde_json::json!({
"backup_id": backup_id,
"tenant_id": tenant_id,
"created_at": started.to_rfc3339(),
"object_prefix": object_prefix,
"object_backend": object_backend,
"object_bucket": object_bucket,
"encrypted": true,
"fk_ordered": plan.fk_ordered,
"tables": table_entries.iter().map(|t| serde_json::json!({
"schema": t.schema,
"table": t.table,
"tenant_column": t.tenant_column,
"object_key": t.object_key,
"row_count": t.row_count,
"checksum_sha256": t.checksum_sha256,
})).collect::<Vec<_>>(),
"excluded": excluded.iter().map(|e| serde_json::json!({
"schema": e.schema,
"table": e.table,
"reason": e.reason,
})).collect::<Vec<_>>(),
});
let manifest_bytes = serde_json::to_vec(&manifest_value).map_err(|err| {
backup_internal_status(
"start_backup_serialize_manifest",
format!("backup manifest serialize failed: {err}"),
)
})?;
let manifest_checksum = sha256_hex(&manifest_bytes);
let manifest_key = format!("{object_prefix}{MANIFEST_SUFFIX}");
let manifest_put = crate::runtime::core::setup_data::object_request_json(
"put",
&object_bucket,
&manifest_key,
"application/json",
);
runtime
.put_object_backend_target_for_project(
&object_backend,
None,
&context.project_id,
&manifest_put,
manifest_bytes,
)
.await?;
let table_count = table_entries.len() as i64;
journal_run(
runtime,
&context,
&backup_id,
&tenant_id,
KIND_BACKUP,
&object_prefix,
&manifest_checksum,
table_count,
total_rows,
excluded_count,
"",
"",
)
.await?;
emit_event(
svc,
TOPIC_BACKUP_COMPLETED,
&tenant_id,
&tenant_id,
&context.project_id,
&backup_id,
serde_json::json!({
"backup_id": backup_id,
"tenant_id": tenant_id,
"object_prefix": object_prefix,
"manifest_checksum": manifest_checksum,
"table_count": table_count,
"total_rows": total_rows,
"excluded_count": excluded_count,
}),
)
.await;
Ok(Response::new(backup_pb::StartTenantBackupResponse {
backup_id,
object_prefix,
manifest_checksum,
table_count: table_count as i32,
total_rows,
excluded_count: excluded_count as i32,
tables: table_entries,
excluded,
message: "tenant backup completed".to_string(),
error: None,
}))
}