diff --git a/Cargo.lock b/Cargo.lock index 8f99dd62..e43adfb0 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -965,7 +965,7 @@ checksum = "c8d4a3bb8b1e0c1050499d1815f5ab16d04f0959b233085fb31653fbfc9d98f9" [[package]] name = "cluster_mgr" -version = "1.12.2" +version = "1.12.3" dependencies = [ "anyhow", "async-trait", diff --git a/src/cluster_mgr/src/cli/task/backup_utils.rs b/src/cluster_mgr/src/cli/task/backup_utils.rs index bc3b5291..f0916d51 100644 --- a/src/cluster_mgr/src/cli/task/backup_utils.rs +++ b/src/cluster_mgr/src/cli/task/backup_utils.rs @@ -10,6 +10,8 @@ use anyhow::{Context, Result}; use regex::Regex; use tracing::{info, warn}; +pub const ROCKSDB_CLOUD_RESTORE_BRANCH: &str = "main"; + /// Join manifest filenames into comma-separated string /// Validates that manifest filenames don't contain commas /// For single manifest, returns as-is (backward compatible) @@ -329,7 +331,7 @@ pub fn parse_snapshot_manifest(manifest: &str) -> Result<(String, u32, String)> } /// Parse current database manifest filename to extract ng_id and epoch -/// Format: CLOUDMANIFEST-development-- +/// Format: CLOUDMANIFEST-main-- /// Returns: (ng_id, epoch) pub fn parse_database_manifest(manifest: &str) -> Result<(u32, u64)> { // Remove CLOUDMANIFEST- prefix @@ -337,8 +339,8 @@ pub fn parse_database_manifest(manifest: &str) -> Result<(u32, u64)> { .strip_prefix("CLOUDMANIFEST-") .ok_or_else(|| anyhow::anyhow!("Manifest must start with CLOUDMANIFEST-: {}", manifest))?; - // Pattern: development-- - let re = Regex::new(r"^development-(\d+)-(\d+)$").context("Failed to compile regex")?; + // Pattern: main-- + let re = Regex::new(r"^main-(\d+)-(\d+)$").context("Failed to compile regex")?; let caps = re .captures(manifest) @@ -353,7 +355,7 @@ pub fn parse_database_manifest(manifest: &str) -> Result<(u32, u64)> { /// Find maximum epoch for a given ng_id from list of manifest keys /// Returns the maximum epoch found, or 0 if none found pub fn find_max_epoch_for_ng(manifest_keys: &[String], ng_id: u32) -> Result { - let prefix = format!("CLOUDMANIFEST-development-{}-", ng_id); + let prefix = format!("CLOUDMANIFEST-{}-{}-", ROCKSDB_CLOUD_RESTORE_BRANCH, ng_id); let mut max_epoch = 0u64; for key in manifest_keys { diff --git a/src/cluster_mgr/src/cli/task/s3_restore_task.rs b/src/cluster_mgr/src/cli/task/s3_restore_task.rs index cc07b670..8d6940a8 100644 --- a/src/cluster_mgr/src/cli/task/s3_restore_task.rs +++ b/src/cluster_mgr/src/cli/task/s3_restore_task.rs @@ -1,5 +1,6 @@ use crate::cli::task::backup_utils::{ find_max_epoch_for_ng, parse_database_manifest, parse_snapshot_manifest, split_manifests, + ROCKSDB_CLOUD_RESTORE_BRANCH, }; use crate::cli::task::s3_utils::{copy_s3_object, list_s3_objects, S3ClientBuilder}; use crate::cli::task::task_base::{ExecutionValue, TaskArgValue, TaskExecutor, TaskHost, TaskId}; @@ -10,6 +11,10 @@ use async_trait::async_trait; use std::collections::{HashMap, HashSet}; use tracing::{error, info, warn}; +// TODO: Read the object path and branch name from cluster configuration, and derive +// the ds_ path for every snapshot manifest in multi-shard deployments. +const ROCKSDB_CLOUD_RESTORE_OBJECT_PATH: &str = "rocksdb_cloud/ds_0"; + #[derive(Clone, Debug)] pub struct S3RestoreTask { task_id: TaskId, @@ -100,8 +105,11 @@ impl TaskExecutor for S3RestoreTask { info!("Found {} manifest(s) to restore", manifest_list.len()); // List all current database manifests to find max epochs - let db_manifest_prefix = "rocksdb_cloud/CLOUDMANIFEST-development-"; - let all_manifests = list_s3_objects(&s3_client, &self.bucket, db_manifest_prefix) + let db_manifest_prefix = format!( + "{}/CLOUDMANIFEST-{}-", + ROCKSDB_CLOUD_RESTORE_OBJECT_PATH, ROCKSDB_CLOUD_RESTORE_BRANCH + ); + let all_manifests = list_s3_objects(&s3_client, &self.bucket, &db_manifest_prefix) .await .context("Failed to list current database manifests")?; @@ -133,7 +141,10 @@ impl TaskExecutor for S3RestoreTask { ); // Construct source S3 key - let source_key = format!("rocksdb_cloud/CLOUDMANIFEST-{}", snapshot_manifest); + let source_key = format!( + "{}/CLOUDMANIFEST-{}", + ROCKSDB_CLOUD_RESTORE_OBJECT_PATH, snapshot_manifest + ); // Check if source file exists match s3_client @@ -172,8 +183,8 @@ impl TaskExecutor for S3RestoreTask { // Construct destination S3 key let dest_key = format!( - "rocksdb_cloud/CLOUDMANIFEST-development-{}-{}", - ng_id, new_epoch + "{}/CLOUDMANIFEST-{}-{}-{}", + ROCKSDB_CLOUD_RESTORE_OBJECT_PATH, ROCKSDB_CLOUD_RESTORE_BRANCH, ng_id, new_epoch ); // Copy manifest file diff --git a/src/cluster_mgr/src/config/storage_service_config.rs b/src/cluster_mgr/src/config/storage_service_config.rs index d3bdf2fd..109fa331 100644 --- a/src/cluster_mgr/src/config/storage_service_config.rs +++ b/src/cluster_mgr/src/config/storage_service_config.rs @@ -208,7 +208,7 @@ impl StorageService { s3.aws_access_key_id.clone(), s3.aws_secret_key.clone(), s3.region.clone(), - None, + s3.endpoint.clone(), )) } RocksDB::MINIO(minio) => {