diff --git a/Cargo.lock b/Cargo.lock index 95eafb6..dabe10e 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2507,9 +2507,9 @@ dependencies = [ [[package]] name = "tucana" -version = "0.0.76" +version = "0.0.77" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "473db4c6dff2f83d8275b92abe90873a5e68b5ca5cf02e2aa0154fdd69bc925a" +checksum = "7d717ad031739316d921424288d56d594f9deb283a660bf372590b35ab65b931" dependencies = [ "pbjson", "pbjson-build", diff --git a/Cargo.toml b/Cargo.toml index fe61486..a63a126 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -11,7 +11,7 @@ opentelemetry = { version = "0.32.0", features = ["metrics"] } tracing = { version = "0.1.41", features = ["log"] } prost = "0.14.1" tonic = "0.14.1" -tucana = { version = "0.0.76", features = ["all"] } +tucana = { version = "0.0.77", features = ["aquila", "sagittarius_gateway"] } code0-flow = { version = "0.0.42", features = ["flow_config", "flow_health", "flow_telemetry"] } serde_json = "1.0.140" async-nats = "0.50.0" diff --git a/src/authorization/mod.rs b/src/authorization/mod.rs index a112574..de39045 100644 --- a/src/authorization/mod.rs +++ b/src/authorization/mod.rs @@ -1,5 +1,5 @@ //! Bearer-token helpers shared by every gRPC client and server in Aquila: -//! [`authorization::get_authorization_metadata`] to attach a token to an +//! [`authorization::get_authentication_metadata`] to attach a token to an //! outgoing request, [`authorization::extract_token`] to read one back off //! an incoming request. @@ -10,21 +10,21 @@ pub mod authorization { metadata::{MetadataMap, MetadataValue}, }; - /// get_authorization_metadata + /// get_authentication_metadata /// - /// Creates a `MetadataMap` that contains the defined token as a value of the `authorization` key - /// Used for setting the runtime_token to authorize Sagittarius request + /// Creates a `MetadataMap` that contains the defined token as a value of the `authentication` key. + /// Used for setting the runtime_token to authenticate Aquila against the Sagittarius gateway. /// /// # Examples /// /// ``` - /// use aquila_grpc::get_authorization_metadata; + /// use aquila_grpc::get_authentication_metadata; /// let token = String::from("token"); - /// let metadata = get_authorization_metadata(&token); - /// assert!(metadata.get("authorization").is_some()); - /// assert_eq!(metadata.get("authorization").unwrap(), "token"); + /// let metadata = get_authentication_metadata(&token); + /// assert!(metadata.get("authentication").is_some()); + /// assert_eq!(metadata.get("authentication").unwrap(), "token"); /// ``` - pub fn get_authorization_metadata(token: &str) -> MetadataMap { + pub fn get_authentication_metadata(token: &str) -> MetadataMap { let metadata_value = MetadataValue::from_str(token).unwrap_or_else(|error| { panic!( "An error occurred trying to convert runtime_token into metadata: {}", @@ -33,7 +33,7 @@ pub mod authorization { }); let mut map = MetadataMap::new(); - map.insert("authorization", metadata_value); + map.insert("authentication", metadata_value); map } diff --git a/src/sagittarius/flow_service_client_impl/mod.rs b/src/sagittarius/flow_service_client_impl/mod.rs index 1fe790f..8e3997e 100644 --- a/src/sagittarius/flow_service_client_impl/mod.rs +++ b/src/sagittarius/flow_service_client_impl/mod.rs @@ -1,7 +1,7 @@ //! Client for Sagittarius' flow synchronization stream: keeps Aquila's local //! flow KV store in sync with whatever Sagittarius considers the current -//! set, and rebroadcasts action module configuration updates it receives -//! along the way. +//! set. Module configuration updates arrive over their own stream now — see +//! [`super::module_configuration_client_impl`]. //! //! - [`flow_store`] applies the sync operations (delete/replace/update) to the KV store. //! - [`dev_export`] mirrors the synced flows to a local JSON file, development only, @@ -14,24 +14,12 @@ use std::sync::Arc; use std::sync::atomic::{AtomicBool, Ordering}; use futures::StreamExt; -use tokio::sync::broadcast; use tonic::{Extensions, Request, transport::Channel}; -use tucana::sagittarius::{ +use tucana::sagittarius_gateway::{ FlowLogonRequest, FlowResponse, flow_response::Data, flow_service_client::FlowServiceClient, }; -use crate::{authorization::authorization::get_authorization_metadata, telemetry::metrics}; - -fn module_config_stats(configs: &tucana::shared::ModuleConfigurations) -> (usize, usize) { - let project_count = configs.module_configurations.len(); - let config_count = configs - .module_configurations - .iter() - .map(|project_cfg| project_cfg.module_configurations.len()) - .sum(); - - (project_count, config_count) -} +use crate::{authorization::authorization::get_authentication_metadata, telemetry::metrics}; #[derive(Clone)] pub struct SagittariusFlowClient { @@ -46,7 +34,6 @@ pub struct SagittariusFlowClient { /// components can hold off on work that depends on Sagittarius state /// actually being loaded. sagittarius_ready: Arc, - action_config_tx: broadcast::Sender, } impl SagittariusFlowClient { @@ -57,7 +44,6 @@ impl SagittariusFlowClient { flow_export_path: String, channel: Channel, sagittarius_ready: Arc, - action_config_tx: broadcast::Sender, ) -> SagittariusFlowClient { let client = FlowServiceClient::new(channel); @@ -68,7 +54,6 @@ impl SagittariusFlowClient { token, flow_export_path, sagittarius_ready, - action_config_tx, } } @@ -157,25 +142,6 @@ impl SagittariusFlowClient { received_count.saturating_sub(stored_count) as u64, ); } - Data::ModuleConfigurations(action_configurations) => { - let (project_count, config_count) = module_config_stats(&action_configurations); - log::debug!( - "Received module configurations from flow stream module_identifier={} project_count={} config_count={}", - action_configurations.module_identifier, - project_count, - config_count - ); - - match self.action_config_tx.send(action_configurations) { - Ok(receiver_count) => log::debug!( - "Broadcasted module configurations to action forwarders receiver_count={}", - receiver_count - ), - Err(err) => { - log::warn!("No action configuration receivers available: {:?}", err); - } - } - } } } @@ -185,7 +151,7 @@ impl SagittariusFlowClient { self.sagittarius_ready.store(false, Ordering::SeqCst); let request = Request::from_parts( - get_authorization_metadata(&self.token), + get_authentication_metadata(&self.token), Extensions::new(), FlowLogonRequest {}, ); diff --git a/src/sagittarius/mod.rs b/src/sagittarius/mod.rs index ae86b39..d46c698 100644 --- a/src/sagittarius/mod.rs +++ b/src/sagittarius/mod.rs @@ -1,9 +1,11 @@ //! Clients Aquila uses to talk *to* Sagittarius: flow synchronization, -//! module registration, runtime status forwarding, and the test/live -//! execution stream. See [`retry`] for the shared reconnect-with-backoff -//! logic every long-lived stream client here is built on. +//! module registration, module configuration sync, runtime status +//! forwarding, and the test/live execution stream. See [`retry`] for the +//! shared reconnect-with-backoff logic every long-lived stream client here +//! is built on. pub mod flow_service_client_impl; +pub mod module_configuration_client_impl; pub mod module_service_client_impl; pub mod retry; pub mod runtime_status_service_client_impl; diff --git a/src/sagittarius/module_configuration_client_impl.rs b/src/sagittarius/module_configuration_client_impl.rs new file mode 100644 index 0000000..4b9b59e --- /dev/null +++ b/src/sagittarius/module_configuration_client_impl.rs @@ -0,0 +1,114 @@ +//! Client for Sagittarius' module configuration stream: forwards module +//! configuration updates onto `action_config_tx` for the action forwarders to +//! pick up. This used to arrive embedded in the flow synchronization stream +//! (see [`super::flow_service_client_impl`]) but now has its own dedicated +//! stream. + +use futures::StreamExt; +use tokio::sync::broadcast; +use tonic::{Extensions, Request, transport::Channel}; +use tucana::sagittarius_gateway::{ + ModuleConfigurationRequest, module_service_client::ModuleServiceClient, +}; + +use crate::authorization::authorization::get_authentication_metadata; + +fn module_config_stats(configs: &tucana::shared::ModuleConfigurations) -> (usize, usize) { + let project_count = configs.module_configurations.len(); + let config_count = configs + .module_configurations + .iter() + .map(|project_cfg| project_cfg.module_configurations.len()) + .sum(); + + (project_count, config_count) +} + +#[derive(Clone)] +pub struct SagittariusModuleConfigurationClient { + client: ModuleServiceClient, + token: String, + action_config_tx: broadcast::Sender, +} + +impl SagittariusModuleConfigurationClient { + pub fn new( + channel: Channel, + token: String, + action_config_tx: broadcast::Sender, + ) -> Self { + Self { + client: ModuleServiceClient::new(channel), + token, + action_config_tx, + } + } + + fn handle_response(&self, module_configurations: tucana::shared::ModuleConfigurations) { + let (project_count, config_count) = module_config_stats(&module_configurations); + log::debug!( + "Received module configurations module_identifier={} project_count={} config_count={}", + module_configurations.module_identifier, + project_count, + config_count + ); + + match self.action_config_tx.send(module_configurations) { + Ok(receiver_count) => log::debug!( + "Broadcasted module configurations to action forwarders receiver_count={}", + receiver_count + ), + Err(err) => { + log::warn!("No action configuration receivers available: {:?}", err); + } + } + } + + /// Opens the module configuration stream and services it until it ends + /// or errors, at which point the caller is expected to reconnect. + pub async fn init_configuration_stream(&mut self) -> Result<(), tonic::Status> { + let request = Request::from_parts( + get_authentication_metadata(&self.token), + Extensions::new(), + ModuleConfigurationRequest {}, + ); + + let response = match self.client.configurations(request).await { + Ok(res) => { + log::info!("Sagittarius module configuration stream established"); + res + } + Err(status) => { + log::warn!( + "Sagittarius module configuration stream connection failed status={:?}", + status + ); + return Err(status); + } + }; + + let mut stream = response.into_inner(); + + while let Some(result) = stream.next().await { + match result { + Ok(res) => { + if let Some(module_configurations) = res.module_configurations { + self.handle_response(module_configurations); + } else { + log::warn!("Received empty Sagittarius module configuration response"); + } + } + Err(status) => { + log::warn!( + "Sagittarius module configuration stream failed; reconnecting status={:?}", + status + ); + return Err(status); + } + } + } + + log::warn!("Sagittarius closed the module configuration stream; reconnecting"); + Err(tonic::Status::unavailable("module configuration stream ended")) + } +} diff --git a/src/sagittarius/module_service_client_impl.rs b/src/sagittarius/module_service_client_impl.rs index c2b930f..3e9f308 100644 --- a/src/sagittarius/module_service_client_impl.rs +++ b/src/sagittarius/module_service_client_impl.rs @@ -2,33 +2,26 @@ //! runtimes) to Sagittarius, so Sagittarius knows what's available to //! reference from a flow. -use crate::configuration::service::ServiceConfiguration; -use crate::{authorization::authorization::get_authorization_metadata, telemetry::errors}; +use crate::{authorization::authorization::get_authentication_metadata, telemetry::errors}; use std::time::Duration; use tonic::transport::Channel; use tonic::{Extensions, Request}; pub struct SagittariusModuleServiceClient { - service: ServiceConfiguration, - client: tucana::sagittarius::module_service_client::ModuleServiceClient, + client: tucana::sagittarius_gateway::module_service_client::ModuleServiceClient, token: String, unary_rpc_timeout: Duration, } impl SagittariusModuleServiceClient { - pub fn new( - channel: Channel, - token: String, - unary_rpc_timeout: Duration, - service_configuration: ServiceConfiguration, - ) -> Self { - let client = tucana::sagittarius::module_service_client::ModuleServiceClient::new(channel); + pub fn new(channel: Channel, token: String, unary_rpc_timeout: Duration) -> Self { + let client = + tucana::sagittarius_gateway::module_service_client::ModuleServiceClient::new(channel); Self { client, token, unary_rpc_timeout, - service: service_configuration, } } @@ -37,9 +30,7 @@ impl SagittariusModuleServiceClient { skip_all, fields(rpc.system = "grpc", rpc.service = "ModuleService", rpc.method = "Update") )] - /// Forwards a module update to Sagittarius, alongside the full list of - /// definition sources this Aquila instance currently has available — - /// Sagittarius needs that list to validate references in flows. + /// Forwards a module update to Sagittarius. pub async fn update_modules( &mut self, modules_update_request: tucana::aquila::ModuleUpdateRequest, @@ -51,11 +42,10 @@ impl SagittariusModuleServiceClient { ); let mut request = Request::from_parts( - get_authorization_metadata(&self.token), + get_authentication_metadata(&self.token), Extensions::new(), - tucana::sagittarius::ModuleUpdateRequest { + tucana::sagittarius_gateway::ModuleUpdateRequest { modules: modules_update_request.modules, - available_defintition_soruces: self.service.collect_modules(), }, ); request.set_timeout(self.unary_rpc_timeout); diff --git a/src/sagittarius/runtime_status_service_client_impl.rs b/src/sagittarius/runtime_status_service_client_impl.rs index 03e4766..ff439c9 100644 --- a/src/sagittarius/runtime_status_service_client_impl.rs +++ b/src/sagittarius/runtime_status_service_client_impl.rs @@ -3,10 +3,10 @@ //! `server::runtime_status_service_server_impl`, tracked in #360 — but kept //! ready for when that's re-enabled. -use crate::{authorization::authorization::get_authorization_metadata, telemetry::errors}; +use crate::{authorization::authorization::get_authentication_metadata, telemetry::errors}; use std::time::Duration; use tonic::{Extensions, Request, transport::Channel}; -use tucana::sagittarius::runtime_status_service_client::RuntimeStatusServiceClient; +use tucana::sagittarius_gateway::runtime_status_service_client::RuntimeStatusServiceClient; pub struct SagittariusRuntimeStatusServiceClient { client: RuntimeStatusServiceClient, @@ -30,9 +30,9 @@ impl SagittariusRuntimeStatusServiceClient { ) -> tucana::aquila::RuntimeStatusUpdateResponse { log::debug!("Forwarding runtime status update to Sagittarius"); let mut request = Request::from_parts( - get_authorization_metadata(&self.token), + get_authentication_metadata(&self.token), Extensions::new(), - tucana::sagittarius::RuntimeStatusUpdateRequest { + tucana::sagittarius_gateway::RuntimeStatusUpdateRequest { status: runtime_status_request.status, }, ); diff --git a/src/sagittarius/test_execution_client_impl/mod.rs b/src/sagittarius/test_execution_client_impl/mod.rs index aae3de9..e2175c1 100644 --- a/src/sagittarius/test_execution_client_impl/mod.rs +++ b/src/sagittarius/test_execution_client_impl/mod.rs @@ -23,12 +23,12 @@ use prost::Message; use tokio_stream::wrappers::ReceiverStream; use tonic::transport::Channel; use tonic::{Extensions, Request}; -use tucana::sagittarius::execution_logon_request::Data; -use tucana::sagittarius::execution_service_client::ExecutionServiceClient; -use tucana::sagittarius::{ExecutionLogonRequest, Logon}; +use tucana::sagittarius_gateway::execution_logon_request::Data; +use tucana::sagittarius_gateway::execution_service_client::ExecutionServiceClient; +use tucana::sagittarius_gateway::{ExecutionLogonRequest, Logon}; use tucana::shared::{ExecutionFlow, ValidationFlow}; -use crate::{authorization::authorization::get_authorization_metadata, flow::key_has_flow_id}; +use crate::{authorization::authorization::get_authentication_metadata, flow::key_has_flow_id}; pub struct SagittariusTestExecutionServiceClient { nats_client: async_nats::Client, @@ -165,13 +165,13 @@ impl SagittariusTestExecutionServiceClient { let ack = ReceiverStream::new(rx); let request = Request::from_parts( - get_authorization_metadata(&self.token), + get_authentication_metadata(&self.token), Extensions::new(), ack, ); log::debug!("Opening Sagittarius execution stream"); - let mut test_execution_stream = match self.client.test(request).await { + let mut test_execution_stream = match self.client.update(request).await { Ok(response) => { log::info!("Sagittarius execution stream established"); response.into_inner() diff --git a/src/sagittarius/test_execution_client_impl/response_sender.rs b/src/sagittarius/test_execution_client_impl/response_sender.rs index b5e4eb6..b05de7d 100644 --- a/src/sagittarius/test_execution_client_impl/response_sender.rs +++ b/src/sagittarius/test_execution_client_impl/response_sender.rs @@ -6,8 +6,8 @@ use std::sync::Arc; use tokio::sync::Mutex; use tonic::Status; -use tucana::sagittarius::execution_logon_request::Data; -use tucana::sagittarius::ExecutionLogonRequest; +use tucana::sagittarius_gateway::execution_logon_request::Data; +use tucana::sagittarius_gateway::ExecutionLogonRequest; use tucana::shared::ExecutionResult; use super::flow_id_cache::ExecutionFlowIdCache; diff --git a/src/server/dynamic_server.rs b/src/server/dynamic_server.rs index 48eea54..44bcd8b 100644 --- a/src/server/dynamic_server.rs +++ b/src/server/dynamic_server.rs @@ -101,7 +101,6 @@ impl AquilaDynamicServer { self.channel.clone(), self.token.clone(), self.sagittarius_unary_rpc_timeout, - self.service_configuration.clone(), ))); info!("ModuleService started"); diff --git a/src/startup/dynamic_mode.rs b/src/startup/dynamic_mode.rs index e20bf63..474fb70 100644 --- a/src/startup/dynamic_mode.rs +++ b/src/startup/dynamic_mode.rs @@ -1,8 +1,8 @@ -//! Dynamic mode wiring: the gRPC server and two independent, self-healing -//! Sagittarius streams (flow sync, test execution) all run as separate -//! tasks supervised by a single `select!` — if any one of them exits or -//! panics, the others are aborted and Aquila shuts down rather than -//! continuing in a partially working state. +//! Dynamic mode wiring: the gRPC server and three independent, self-healing +//! Sagittarius streams (flow sync, module configuration sync, test +//! execution) all run as separate tasks supervised by a single `select!` — +//! if any one of them exits or panics, the others are aborted and Aquila +//! shuts down rather than continuing in a partially working state. use async_nats::Client; @@ -12,6 +12,7 @@ use crate::{ }, sagittarius::{ flow_service_client_impl::SagittariusFlowClient, + module_configuration_client_impl::SagittariusModuleConfigurationClient, retry::create_channel_with_retry, test_execution_client_impl::{ SagittariusExecutionResponseSender, SagittariusTestExecutionServiceClient, @@ -83,6 +84,11 @@ pub async fn run( let flow_export_path_for_flow = config.static_config.flow_path.clone(); let sagittarius_ready_for_flow = app_readiness.sagittarius_ready.clone(); + let backend_url_for_module_configuration = config.dynamic_config.backend_url.clone(); + let runtime_token_for_module_configuration = config.dynamic_config.backend_token.clone(); + let sagittarius_ready_for_module_configuration = app_readiness.sagittarius_ready.clone(); + let action_config_tx_for_module_configuration = action_config_tx.clone(); + let env = match config.environment { crate::configuration::env::Environment::Development => String::from("DEVELOPMENT"), crate::configuration::env::Environment::Staging => String::from("STAGING"), @@ -150,7 +156,6 @@ pub async fn run( flow_export_path_for_flow.clone(), ch, sagittarius_ready_for_flow.clone(), - action_config_tx.clone(), ); match flow_client.init_flow_stream().await { @@ -176,6 +181,51 @@ pub async fn run( } }); + let mut module_configuration_task = tokio::spawn(async move { + let mut backoff = Duration::from_millis(200); + let max_backoff = Duration::from_secs(10); + + loop { + log::debug!( + "Attempting to initialize Sagittarius module configuration stream backoff_ms={}", + backoff.as_millis() + ); + let ch = create_channel_with_retry( + "Sagittarius Module Configuration Stream", + backend_url_for_module_configuration.clone(), + sagittarius_ready_for_module_configuration.clone(), + ) + .await; + + let mut module_configuration_client = SagittariusModuleConfigurationClient::new( + ch, + runtime_token_for_module_configuration.clone(), + action_config_tx_for_module_configuration.clone(), + ); + + match module_configuration_client.init_configuration_stream().await { + Ok(_) => { + log::warn!( + "Sagittarius module configuration stream ended normally; reconnecting" + ); + } + Err(e) => { + log::warn!( + "Sagittarius module configuration stream dropped; reconnecting error={:?}", + e + ); + } + } + + tokio::time::sleep(backoff).await; + backoff = std::cmp::min(backoff * 2, max_backoff); + log::debug!( + "Next module configuration stream reconnect backoff_ms={}", + backoff.as_millis() + ); + } + }); + #[cfg(unix)] let sigterm = async { use tokio::signal::unix::{SignalKind, signal}; @@ -195,6 +245,7 @@ pub async fn run( Err(err) => errors::record("task", "grpc.task", &err, "mode=dynamic"), } flow_task.abort(); + module_configuration_task.abort(); test_execution_task.abort(); } result = &mut test_execution_task => { @@ -205,6 +256,7 @@ pub async fn run( } server_task.abort(); flow_task.abort(); + module_configuration_task.abort(); } result = &mut flow_task => { match result { @@ -213,18 +265,31 @@ pub async fn run( Err(err) => errors::record("task", "flow_stream.task", &err, "mode=dynamic"), } server_task.abort(); + module_configuration_task.abort(); + test_execution_task.abort(); + } + result = &mut module_configuration_task => { + match result { + Ok(()) => log::warn!("Module configuration stream task exited unexpectedly; shutting down"), + Err(err) if err.is_panic() => {} + Err(err) => errors::record("task", "module_configuration_stream.task", &err, "mode=dynamic"), + } + server_task.abort(); + flow_task.abort(); test_execution_task.abort(); } _ = tokio::signal::ctrl_c() => { log::info!("Ctrl+C/Exit signal received, shutting down"); server_task.abort(); flow_task.abort(); + module_configuration_task.abort(); test_execution_task.abort(); } _ = sigterm => { log::info!("SIGTERM received, shutting down"); server_task.abort(); flow_task.abort(); + module_configuration_task.abort(); test_execution_task.abort(); } }