Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
20 changes: 10 additions & 10 deletions src/authorization/mod.rs
Original file line number Diff line number Diff line change
@@ -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.

Expand All @@ -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: {}",
Expand All @@ -33,7 +33,7 @@ pub mod authorization {
});

let mut map = MetadataMap::new();
map.insert("authorization", metadata_value);
map.insert("authentication", metadata_value);
map
}

Expand Down
44 changes: 5 additions & 39 deletions src/sagittarius/flow_service_client_impl/mod.rs
Original file line number Diff line number Diff line change
@@ -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,
Expand All @@ -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 {
Expand All @@ -46,7 +34,6 @@ pub struct SagittariusFlowClient {
/// components can hold off on work that depends on Sagittarius state
/// actually being loaded.
sagittarius_ready: Arc<AtomicBool>,
action_config_tx: broadcast::Sender<tucana::shared::ModuleConfigurations>,
}

impl SagittariusFlowClient {
Expand All @@ -57,7 +44,6 @@ impl SagittariusFlowClient {
flow_export_path: String,
channel: Channel,
sagittarius_ready: Arc<AtomicBool>,
action_config_tx: broadcast::Sender<tucana::shared::ModuleConfigurations>,
) -> SagittariusFlowClient {
let client = FlowServiceClient::new(channel);

Expand All @@ -68,7 +54,6 @@ impl SagittariusFlowClient {
token,
flow_export_path,
sagittarius_ready,
action_config_tx,
}
}

Expand Down Expand Up @@ -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);
}
}
}
}
}

Expand All @@ -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 {},
);
Expand Down
8 changes: 5 additions & 3 deletions src/sagittarius/mod.rs
Original file line number Diff line number Diff line change
@@ -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;
Expand Down
114 changes: 114 additions & 0 deletions src/sagittarius/module_configuration_client_impl.rs
Original file line number Diff line number Diff line change
@@ -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<Channel>,
token: String,
action_config_tx: broadcast::Sender<tucana::shared::ModuleConfigurations>,
}

impl SagittariusModuleConfigurationClient {
pub fn new(
channel: Channel,
token: String,
action_config_tx: broadcast::Sender<tucana::shared::ModuleConfigurations>,
) -> 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"))
}
}
26 changes: 8 additions & 18 deletions src/sagittarius/module_service_client_impl.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Channel>,
client: tucana::sagittarius_gateway::module_service_client::ModuleServiceClient<Channel>,
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,
}
}

Expand All @@ -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,
Expand All @@ -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);
Expand Down
8 changes: 4 additions & 4 deletions src/sagittarius/runtime_status_service_client_impl.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Channel>,
Expand All @@ -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,
},
);
Expand Down
Loading