use std::sync::Arc;
use sqlx::PgPool;
use tonic::{Request, Response, Status};
use crate::metrics::{MetricsRecorder, NoopMetrics};
use crate::proto::udb::core::asset::services::v1 as asset_pb;
use crate::proto::udb::core::asset::services::v1::asset_service_server::AssetService;
use crate::runtime::DataBrokerRuntime;
use crate::runtime::channels::{ChannelManager, ChannelPermit, OperationChannel};
pub use crate::proto::udb::core::asset::services::v1::asset_service_server::AssetServiceServer;
use super::DataBrokerService;
use super::native_helpers::admit_on as native_admit_on;
mod config;
mod consumer;
mod embedding;
mod errors;
mod execution;
mod handlers;
mod model;
mod steps;
mod store;
#[cfg(test)]
mod tests;
use config::DEFAULT_VECTOR_COLLECTION;
use errors::{
asset_capability_status, native_state_decryption_failed_status,
native_state_encryption_failed_status,
};
pub struct AssetServiceImpl {
pub(crate) pg_pool: Option<PgPool>,
pub(crate) outbox_relation: Option<String>,
pub(crate) runtime: Option<Arc<DataBrokerRuntime>>,
pub(crate) channels: Option<ChannelManager>,
pub(crate) vector_collection: String,
pub(crate) metrics: Arc<dyn MetricsRecorder>,
}
impl AssetServiceImpl {
pub fn new() -> Self {
Self {
pg_pool: None,
outbox_relation: None,
runtime: None,
channels: None,
vector_collection: DEFAULT_VECTOR_COLLECTION.to_string(),
metrics: Arc::new(NoopMetrics),
}
}
pub fn with_postgres(mut self, pool: Option<PgPool>) -> Self {
self.pg_pool = pool;
self
}
pub(crate) fn with_metrics(mut self, metrics: Arc<dyn MetricsRecorder>) -> Self {
self.metrics = metrics;
self
}
pub(crate) fn with_vector(
mut self,
runtime: Option<Arc<DataBrokerRuntime>>,
collection: String,
) -> Self {
self.channels = runtime.as_ref().map(|rt| rt.channels().clone());
self.runtime = runtime;
if !collection.trim().is_empty() {
self.vector_collection = collection;
}
self
}
fn encrypt_native_json_state(&self, raw_json: &str) -> Result<String, Status> {
match self.runtime.as_ref() {
Some(runtime) => runtime
.encrypt_native_json_state_at_rest(raw_json)
.map_err(native_state_encryption_failed_status),
None => Ok(raw_json.to_string()),
}
}
fn decrypt_native_json_state(&self, stored_json: &str) -> Result<String, Status> {
if stored_json.trim().is_empty() {
return Ok(String::new());
}
match self.runtime.as_ref() {
Some(runtime) => runtime
.decrypt_native_json_state_at_rest(stored_json)
.map_err(native_state_decryption_failed_status),
None => Ok(stored_json.to_string()),
}
}
fn require_runtime(&self) -> Result<&DataBrokerRuntime, Status> {
self.runtime.as_deref().ok_or_else(|| {
asset_capability_status(
"native_entity_dispatch",
"runtime_native_entity_dispatch",
"asset service requires runtime native entity dispatch",
)
})
}
async fn admit(&self, tenant: &str, project: &str) -> Result<Option<ChannelPermit>, Status> {
native_admit_on(
self.channels.as_ref(),
&self.metrics,
"asset",
OperationChannel::Object,
tenant,
Some(project),
)
.await
}
async fn admit_read(&self, tenant: &str) -> Result<Option<ChannelPermit>, Status> {
native_admit_on(
self.channels.as_ref(),
&self.metrics,
"asset",
OperationChannel::Read,
tenant,
Some(""),
)
.await
}
pub(crate) fn with_outbox(mut self, relation: Option<String>) -> Self {
self.outbox_relation = relation;
self
}
fn require_pool(&self) -> Result<&PgPool, Status> {
self.pg_pool.as_ref().ok_or_else(|| {
asset_capability_status(
"postgres_store",
"postgres_store",
"asset service requires a Postgres-backed store (no PG pool configured)",
)
})
}
}
impl Default for AssetServiceImpl {
fn default() -> Self {
Self::new()
}
}
#[tonic::async_trait]
impl AssetService for AssetServiceImpl {
async fn create_pipeline_definition(
&self,
request: Request<asset_pb::CreatePipelineDefinitionRequest>,
) -> Result<Response<asset_pb::CreatePipelineDefinitionResponse>, Status> {
handlers::create_pipeline_definition(self, request).await
}
async fn get_pipeline_definition(
&self,
request: Request<asset_pb::GetPipelineDefinitionRequest>,
) -> Result<Response<asset_pb::GetPipelineDefinitionResponse>, Status> {
handlers::get_pipeline_definition(self, request).await
}
async fn register_asset(
&self,
request: Request<asset_pb::RegisterAssetRequest>,
) -> Result<Response<asset_pb::RegisterAssetResponse>, Status> {
handlers::register_asset(self, request).await
}
async fn start_pipeline(
&self,
request: Request<asset_pb::StartPipelineRequest>,
) -> Result<Response<asset_pb::StartPipelineResponse>, Status> {
handlers::start_pipeline(self, request).await
}
async fn get_pipeline(
&self,
request: Request<asset_pb::GetPipelineRequest>,
) -> Result<Response<asset_pb::GetPipelineResponse>, Status> {
handlers::get_pipeline(self, request).await
}
async fn complete_step(
&self,
request: Request<asset_pb::CompleteStepRequest>,
) -> Result<Response<asset_pb::CompleteStepResponse>, Status> {
handlers::complete_step(self, request).await
}
async fn list_assets(
&self,
request: Request<asset_pb::ListAssetsRequest>,
) -> Result<Response<asset_pb::ListAssetsResponse>, Status> {
handlers::list_assets(self, request).await
}
async fn get_asset(
&self,
request: Request<asset_pb::GetAssetRequest>,
) -> Result<Response<asset_pb::GetAssetResponse>, Status> {
handlers::get_asset(self, request).await
}
}
impl DataBrokerService {
pub(crate) fn build_asset_service(&self) -> AssetServiceImpl {
let runtime = self.runtime.load_full();
let pg_pool = runtime
.native_store_pool_for_service("asset", true, "")
.ok();
let outbox = runtime.config().cdc.outbox_relation();
let collection = std::env::var("UDB_ASSET_VECTOR_COLLECTION")
.unwrap_or_else(|_| DEFAULT_VECTOR_COLLECTION.to_string());
AssetServiceImpl::new()
.with_postgres(pg_pool)
.with_outbox(Some(outbox))
.with_metrics(self.metrics.clone())
.with_vector(Some(runtime.clone()), collection)
}
}