use std::sync::Arc;
use sqlx::PgPool;
use tonic::{Request, Response, Status};
use crate::generation::CatalogManifest;
use crate::metrics::{MetricsRecorder, NoopMetrics};
use crate::proto::udb::core::backup::services::v1 as backup_pb;
use crate::proto::udb::core::backup::services::v1::backup_service_server::BackupService;
use crate::runtime::DataBrokerRuntime;
use crate::runtime::channels::ChannelManager;
pub use crate::proto::udb::core::backup::services::v1::backup_service_server::BackupServiceServer;
use super::DataBrokerService;
use super::native_helpers::{
DEFAULT_OBJECT_BACKEND, DEFAULT_OBJECT_BUCKET, storage_object_defaults,
};
mod config;
mod errors;
mod events;
mod export;
mod handlers;
mod import;
mod model;
mod store;
#[cfg(test)]
mod tests;
pub struct BackupServiceImpl {
pub(crate) pg_pool: Option<PgPool>,
pub(crate) runtime: Option<Arc<DataBrokerRuntime>>,
pub(crate) outbox_relation: Option<String>,
pub(crate) channels: Option<ChannelManager>,
pub(crate) metrics: Arc<dyn MetricsRecorder>,
pub(crate) manifest: Option<CatalogManifest>,
pub(crate) object_backend: String,
pub(crate) object_bucket: String,
}
impl BackupServiceImpl {
pub fn new() -> Self {
Self {
pg_pool: None,
runtime: None,
outbox_relation: None,
channels: None,
metrics: Arc::new(NoopMetrics),
manifest: None,
object_backend: DEFAULT_OBJECT_BACKEND.to_string(),
object_bucket: DEFAULT_OBJECT_BUCKET.to_string(),
}
}
pub fn with_postgres(mut self, pool: Option<PgPool>) -> Self {
self.pg_pool = pool;
self
}
pub(crate) fn with_runtime(mut self, runtime: Option<Arc<DataBrokerRuntime>>) -> Self {
self.runtime = runtime;
self
}
pub(crate) fn with_outbox(mut self, relation: Option<String>) -> Self {
self.outbox_relation = relation;
self
}
pub(crate) fn with_channels(mut self, channels: Option<ChannelManager>) -> Self {
self.channels = channels;
self
}
pub(crate) fn with_metrics(mut self, metrics: Arc<dyn MetricsRecorder>) -> Self {
self.metrics = metrics;
self
}
pub(crate) fn with_manifest(mut self, manifest: Option<CatalogManifest>) -> Self {
self.manifest = manifest;
self
}
pub(crate) fn with_object(mut self, backend: String, bucket: String) -> Self {
if !backend.trim().is_empty() {
self.object_backend = backend;
}
if !bucket.trim().is_empty() {
self.object_bucket = bucket;
}
self
}
pub(crate) fn require_runtime(&self) -> Result<&DataBrokerRuntime, Status> {
self.runtime.as_deref().ok_or_else(|| {
errors::backup_capability_status(
"native_entity_dispatch",
"runtime_native_entity_dispatch",
"backup service requires runtime native-entity dispatch (no runtime configured)",
)
})
}
pub(crate) fn require_pool(&self) -> Result<&PgPool, Status> {
self.pg_pool.as_ref().ok_or_else(|| {
errors::backup_capability_status(
"postgres_store",
"postgres_store",
"backup service requires a Postgres-backed store (no PG pool configured)",
)
})
}
pub(crate) fn require_manifest(&self) -> Result<&CatalogManifest, Status> {
self.manifest.as_ref().ok_or_else(|| {
errors::backup_capability_status(
"tenant_table_enumeration",
"catalog_manifest",
"backup service requires the catalog manifest to enumerate tenant tables",
)
})
}
}
impl Default for BackupServiceImpl {
fn default() -> Self {
Self::new()
}
}
#[tonic::async_trait]
impl BackupService for BackupServiceImpl {
async fn start_tenant_backup(
&self,
request: Request<backup_pb::StartTenantBackupRequest>,
) -> Result<Response<backup_pb::StartTenantBackupResponse>, Status> {
export::start_tenant_backup(self, request).await
}
async fn restore_tenant(
&self,
request: Request<backup_pb::RestoreTenantRequest>,
) -> Result<Response<backup_pb::RestoreTenantResponse>, Status> {
import::restore_tenant(self, request).await
}
async fn list_backups(
&self,
request: Request<backup_pb::ListBackupsRequest>,
) -> Result<Response<backup_pb::ListBackupsResponse>, Status> {
handlers::list_backups(self, request).await
}
async fn get_backup(
&self,
request: Request<backup_pb::GetBackupRequest>,
) -> Result<Response<backup_pb::GetBackupResponse>, Status> {
handlers::get_backup(self, request).await
}
async fn put_backup_policy(
&self,
request: Request<backup_pb::PutBackupPolicyRequest>,
) -> Result<Response<backup_pb::PutBackupPolicyResponse>, Status> {
handlers::put_backup_policy(self, request).await
}
async fn get_backup_policy(
&self,
request: Request<backup_pb::GetBackupPolicyRequest>,
) -> Result<Response<backup_pb::GetBackupPolicyResponse>, Status> {
handlers::get_backup_policy(self, request).await
}
async fn list_backup_policies(
&self,
request: Request<backup_pb::ListBackupPoliciesRequest>,
) -> Result<Response<backup_pb::ListBackupPoliciesResponse>, Status> {
handlers::list_backup_policies(self, request).await
}
async fn delete_backup_policy(
&self,
request: Request<backup_pb::DeleteBackupPolicyRequest>,
) -> Result<Response<backup_pb::DeleteBackupPolicyResponse>, Status> {
handlers::delete_backup_policy(self, request).await
}
}
impl DataBrokerService {
pub(crate) fn build_backup_service(&self) -> BackupServiceImpl {
let runtime = self.runtime.load_full();
let pg_pool = runtime
.native_store_pool_for_service("backup", true, "")
.ok();
let outbox = runtime.config().cdc.outbox_relation();
let channels = Some(runtime.channels().clone());
let (object_backend, object_bucket) = storage_object_defaults(
std::env::var("UDB_STORAGE_OBJECT_BACKEND").ok(),
std::env::var("UDB_STORAGE_BUCKET").ok(),
);
BackupServiceImpl::new()
.with_postgres(pg_pool)
.with_runtime(Some(runtime))
.with_outbox(Some(outbox))
.with_channels(channels)
.with_metrics(self.metrics.clone())
.with_manifest(Some(self.manifest.clone()))
.with_object(object_backend, object_bucket)
}
}