use std::sync::Arc;
use sqlx::PgPool;
use tonic::{Request, Response, Status};
use crate::metrics::{MetricsRecorder, NoopMetrics};
use crate::proto::udb::core::vault::services::v1 as vault_pb;
use crate::proto::udb::core::vault::services::v1::vault_service_server::VaultService;
use crate::runtime::DataBrokerRuntime;
use crate::runtime::channels::ChannelManager;
pub use crate::proto::udb::core::vault::services::v1::vault_service_server::VaultServiceServer;
use super::DataBrokerService;
mod config;
mod crypto;
mod dynamic;
mod errors;
mod events;
mod handlers;
mod model;
mod store;
#[cfg(test)]
mod tests;
mod workers;
use config::{SEAL_PROBE, VAULT_MASTER_KEY_UNAVAILABLE_MESSAGE, VAULT_RUNTIME_REQUIRED_MESSAGE};
use errors::{vault_capability_status, vault_master_key_unavailable_status};
use model::{
StoredSecret, StoredTransitKey, stored_secret_from_json, stored_transit_key_from_json,
};
use store::{secret_path_read, transit_key_read};
pub use config::{VAULT_DB_LEASE_REAPER_BATCH, vault_db_lease_reaper_interval};
pub use workers::run_vault_db_lease_reaper_once;
pub struct VaultServiceImpl {
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) seal_override: Option<bool>,
}
impl VaultServiceImpl {
pub fn new() -> Self {
Self {
pg_pool: None,
runtime: None,
outbox_relation: None,
channels: None,
metrics: Arc::new(NoopMetrics),
seal_override: None,
}
}
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
}
#[cfg(test)]
pub(crate) fn with_seal_override(mut self, sealed: bool) -> Self {
self.seal_override = Some(sealed);
self
}
pub(crate) fn require_runtime(&self) -> Result<&DataBrokerRuntime, Status> {
self.runtime.as_deref().ok_or_else(|| {
vault_capability_status(
"native_entity_dispatch",
"runtime_native_entity_dispatch",
VAULT_RUNTIME_REQUIRED_MESSAGE,
)
})
}
pub(crate) fn check_seal(&self) -> Result<(), Status> {
match self.seal_override {
Some(true) => {
return Err(vault_master_key_unavailable_status(
VAULT_MASTER_KEY_UNAVAILABLE_MESSAGE,
));
}
Some(false) => return Ok(()),
None => {}
}
let runtime = self.runtime.as_deref().ok_or_else(|| {
vault_master_key_unavailable_status("vault is sealed: no runtime / master key wired")
})?;
runtime
.encrypt_secret_at_rest(SEAL_PROBE)
.map(|_| ())
.map_err(|err| {
vault_master_key_unavailable_status(format!(
"vault is sealed: master key unavailable ({err})"
))
})
}
pub(crate) fn is_sealed(&self) -> bool {
self.check_seal().is_err()
}
pub(crate) fn kek_configured(&self) -> bool {
if self.seal_override.is_some() {
return false;
}
match self.runtime.as_deref() {
Some(runtime) => runtime
.encrypt_secret_at_rest(SEAL_PROBE)
.map(|env| env.starts_with("udb-aead:"))
.unwrap_or(false),
None => false,
}
}
pub(crate) async fn read_secret_versions(
&self,
runtime: &DataBrokerRuntime,
context: &crate::RequestContext,
tenant_id: &str,
secret_path: &str,
) -> Result<Vec<StoredSecret>, Status> {
let rows = runtime
.native_entity_read_for_service(
"vault",
context,
secret_path_read(tenant_id, secret_path),
)
.await?;
Ok(rows.iter().map(stored_secret_from_json).collect())
}
pub(crate) async fn read_transit_versions(
&self,
runtime: &DataBrokerRuntime,
context: &crate::RequestContext,
tenant_id: &str,
key_name: &str,
) -> Result<Vec<StoredTransitKey>, Status> {
let rows = runtime
.native_entity_read_for_service("vault", context, transit_key_read(tenant_id, key_name))
.await?;
Ok(rows.iter().map(stored_transit_key_from_json).collect())
}
}
impl Default for VaultServiceImpl {
fn default() -> Self {
Self::new()
}
}
#[tonic::async_trait]
impl VaultService for VaultServiceImpl {
async fn put_secret(
&self,
request: Request<vault_pb::PutSecretRequest>,
) -> Result<Response<vault_pb::PutSecretResponse>, Status> {
handlers::put_secret(self, request).await
}
async fn get_secret(
&self,
request: Request<vault_pb::GetSecretRequest>,
) -> Result<Response<vault_pb::GetSecretResponse>, Status> {
handlers::get_secret(self, request).await
}
async fn list_secrets(
&self,
request: Request<vault_pb::ListSecretsRequest>,
) -> Result<Response<vault_pb::ListSecretsResponse>, Status> {
handlers::list_secrets(self, request).await
}
async fn delete_secret(
&self,
request: Request<vault_pb::DeleteSecretRequest>,
) -> Result<Response<vault_pb::DeleteSecretResponse>, Status> {
handlers::delete_secret(self, request).await
}
async fn undelete_secret(
&self,
request: Request<vault_pb::UndeleteSecretRequest>,
) -> Result<Response<vault_pb::UndeleteSecretResponse>, Status> {
handlers::undelete_secret(self, request).await
}
async fn destroy_secret(
&self,
request: Request<vault_pb::DestroySecretRequest>,
) -> Result<Response<vault_pb::DestroySecretResponse>, Status> {
handlers::destroy_secret(self, request).await
}
async fn create_transit_key(
&self,
request: Request<vault_pb::CreateTransitKeyRequest>,
) -> Result<Response<vault_pb::CreateTransitKeyResponse>, Status> {
handlers::create_transit_key(self, request).await
}
async fn rotate_transit_key(
&self,
request: Request<vault_pb::RotateTransitKeyRequest>,
) -> Result<Response<vault_pb::RotateTransitKeyResponse>, Status> {
handlers::rotate_transit_key(self, request).await
}
async fn encrypt(
&self,
request: Request<vault_pb::EncryptRequest>,
) -> Result<Response<vault_pb::EncryptResponse>, Status> {
handlers::encrypt(self, request).await
}
async fn decrypt(
&self,
request: Request<vault_pb::DecryptRequest>,
) -> Result<Response<vault_pb::DecryptResponse>, Status> {
handlers::decrypt(self, request).await
}
async fn batch_encrypt(
&self,
request: Request<vault_pb::BatchEncryptRequest>,
) -> Result<Response<vault_pb::BatchEncryptResponse>, Status> {
handlers::batch_encrypt(self, request).await
}
async fn batch_decrypt(
&self,
request: Request<vault_pb::BatchDecryptRequest>,
) -> Result<Response<vault_pb::BatchDecryptResponse>, Status> {
handlers::batch_decrypt(self, request).await
}
async fn generate_data_key(
&self,
request: Request<vault_pb::GenerateDataKeyRequest>,
) -> Result<Response<vault_pb::GenerateDataKeyResponse>, Status> {
handlers::generate_data_key(self, request).await
}
async fn rewrap(
&self,
request: Request<vault_pb::RewrapRequest>,
) -> Result<Response<vault_pb::RewrapResponse>, Status> {
handlers::rewrap(self, request).await
}
async fn get_transit_public_key(
&self,
request: Request<vault_pb::GetTransitPublicKeyRequest>,
) -> Result<Response<vault_pb::GetTransitPublicKeyResponse>, Status> {
handlers::get_transit_public_key(self, request).await
}
async fn sign(
&self,
request: Request<vault_pb::SignRequest>,
) -> Result<Response<vault_pb::SignResponse>, Status> {
handlers::sign(self, request).await
}
async fn verify(
&self,
request: Request<vault_pb::VerifyRequest>,
) -> Result<Response<vault_pb::VerifyResponse>, Status> {
handlers::verify(self, request).await
}
async fn hmac(
&self,
request: Request<vault_pb::HmacRequest>,
) -> Result<Response<vault_pb::HmacResponse>, Status> {
handlers::hmac(self, request).await
}
async fn seal_status(
&self,
request: Request<vault_pb::SealStatusRequest>,
) -> Result<Response<vault_pb::SealStatusResponse>, Status> {
handlers::seal_status(self, request).await
}
async fn generate_database_credentials(
&self,
request: Request<vault_pb::GenerateDatabaseCredentialsRequest>,
) -> Result<Response<vault_pb::GenerateDatabaseCredentialsResponse>, Status> {
handlers::generate_database_credentials(self, request).await
}
}
impl DataBrokerService {
pub(crate) fn build_vault_service(&self) -> VaultServiceImpl {
let runtime = self.runtime.load_full();
let pg_pool = runtime
.native_store_pool_for_service("vault", true, "")
.ok();
let outbox = runtime.config().cdc.outbox_relation();
let channels = Some(runtime.channels().clone());
VaultServiceImpl::new()
.with_postgres(pg_pool)
.with_runtime(Some(runtime))
.with_outbox(Some(outbox))
.with_channels(channels)
.with_metrics(self.metrics.clone())
}
}