use std::sync::Arc;
use sqlx::PgPool;
use tonic::{Request, Response, Status};
use uuid::Uuid;
use crate::metrics::{MetricsRecorder, NoopMetrics};
use crate::proto::udb::core::workflow::services::v1 as workflow_pb;
use crate::proto::udb::core::workflow::services::v1::workflow_service_server::WorkflowService;
use crate::runtime::DataBrokerRuntime;
use crate::runtime::canonical_store::SystemStores;
use crate::runtime::canonical_store::system_store::{CompensationStatus, SagaStatus, SagaStore};
use crate::runtime::channels::ChannelManager;
pub use crate::proto::udb::core::workflow::services::v1::workflow_service_server::WorkflowServiceServer;
use super::DataBrokerService;
mod config;
mod errors;
mod events;
mod handlers;
mod model;
mod store;
#[cfg(test)]
mod tests;
mod tick;
pub(crate) use config::WORKFLOW_TICK_BATCH;
pub(crate) use tick::run_workflow_tick_once;
pub struct WorkflowServiceImpl {
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>,
}
impl WorkflowServiceImpl {
pub fn new() -> Self {
Self {
pg_pool: None,
runtime: None,
outbox_relation: None,
channels: None,
metrics: Arc::new(NoopMetrics),
}
}
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 require_pool(&self) -> Result<&PgPool, Status> {
self.pg_pool.as_ref().ok_or_else(|| {
errors::workflow_capability_status(
"postgres_store",
"postgres_store",
"workflow service requires a Postgres-backed store (no PG pool configured)",
)
})
}
pub(crate) fn system_stores(&self) -> Option<Arc<dyn SystemStores>> {
self.runtime.as_ref()?.default_system_stores()
}
pub(crate) async fn settle_orphan_saga(&self, saga_id: &str) {
if saga_id.is_empty() {
return;
}
let Some(store) = self.system_stores() else {
return;
};
let Ok(saga_uuid) = saga_id.parse::<Uuid>() else {
return;
};
if let Err(err) = SagaStore::update_saga_status(
store.as_ref(),
saga_uuid,
SagaStatus::Compensated,
CompensationStatus::Completed,
)
.await
{
tracing::warn!(error = %err, saga_id, "workflow start: orphan saga settle failed");
}
}
}
impl Default for WorkflowServiceImpl {
fn default() -> Self {
Self::new()
}
}
#[tonic::async_trait]
impl WorkflowService for WorkflowServiceImpl {
async fn start_workflow(
&self,
request: Request<workflow_pb::StartWorkflowRequest>,
) -> Result<Response<workflow_pb::StartWorkflowResponse>, Status> {
handlers::start_workflow(self, request).await
}
async fn get_workflow(
&self,
request: Request<workflow_pb::GetWorkflowRequest>,
) -> Result<Response<workflow_pb::GetWorkflowResponse>, Status> {
handlers::get_workflow(self, request).await
}
async fn list_workflows(
&self,
request: Request<workflow_pb::ListWorkflowsRequest>,
) -> Result<Response<workflow_pb::ListWorkflowsResponse>, Status> {
handlers::list_workflows(self, request).await
}
async fn cancel_workflow(
&self,
request: Request<workflow_pb::CancelWorkflowRequest>,
) -> Result<Response<workflow_pb::CancelWorkflowResponse>, Status> {
handlers::cancel_workflow(self, request).await
}
async fn signal_workflow(
&self,
request: Request<workflow_pb::SignalWorkflowRequest>,
) -> Result<Response<workflow_pb::SignalWorkflowResponse>, Status> {
handlers::signal_workflow(self, request).await
}
}
impl DataBrokerService {
pub(crate) fn build_workflow_service(&self) -> WorkflowServiceImpl {
let runtime = self.runtime.load_full();
let pg_pool = runtime
.native_store_pool_for_service("workflow", true, "")
.ok();
let outbox = runtime.config().cdc.outbox_relation();
let channels = Some(runtime.channels().clone());
WorkflowServiceImpl::new()
.with_postgres(pg_pool)
.with_runtime(Some(runtime))
.with_outbox(Some(outbox))
.with_channels(channels)
.with_metrics(self.metrics.clone())
}
}