use std::time::Duration;
use std::collections::HashSet;
use super::knative::unique_k8s_name;
use super::{Bus, TransportError};
use super::{Message, MessageKind, SubscriptionPlan};
const SEND_TIMEOUT: Duration = Duration::from_secs(10);
fn kind_str(kind: MessageKind) -> &'static str {
match kind {
MessageKind::Command => "command",
MessageKind::Event => "event",
}
}
#[derive(Clone)]
pub struct KnativeBus {
client: reqwest::Client,
ingress_base: String,
namespace: String,
source: String,
commands_broker: String,
events_broker: String,
publishes_events: bool,
local: Option<String>,
}
impl KnativeBus {
pub fn new(
ingress_base: impl Into<String>,
namespace: impl Into<String>,
source: impl Into<String>,
commands_broker: impl Into<String>,
events_broker: impl Into<String>,
) -> Self {
Self {
client: reqwest::Client::new(),
ingress_base: ingress_base.into(),
namespace: namespace.into(),
source: source.into(),
commands_broker: commands_broker.into(),
events_broker: events_broker.into(),
publishes_events: true,
local: None,
}
}
pub fn publishes_events(mut self, publishes: bool) -> Self {
self.publishes_events = publishes;
self
}
pub fn local(mut self, addr: impl Into<String>) -> Self {
self.local = Some(addr.into());
self
}
fn ingress_url(&self, broker: &str) -> String {
let mut url = self.ingress_base.trim_end_matches('/').to_string();
for segment in [self.namespace.as_str(), broker] {
if !segment.is_empty() {
url.push('/');
url.push_str(segment);
}
}
url
}
async fn post_cloud_event(
&self,
broker: &str,
message: &Message,
) -> Result<(), TransportError> {
let Some(id) = message.id() else {
return Err(TransportError::permanent(
"knative produce requires a message id (CloudEvents `id` is mandatory)",
));
};
let mut request = self
.client
.post(self.ingress_url(broker))
.timeout(SEND_TIMEOUT)
.header("ce-specversion", "1.0")
.header("ce-id", id)
.header("ce-type", message.name())
.header("ce-source", self.source.as_str())
.header("ce-sourcedkind", kind_str(message.kind))
.header("content-type", message.content_type.as_str());
for (key, value) in &message.metadata {
request = request.header(format!("ce-{key}"), value);
}
let response = request
.body(message.payload.clone())
.send()
.await
.map_err(|err| TransportError::retryable(format!("knative POST: {err}")))?;
if !response.status().is_success() {
return Err(TransportError::retryable(format!(
"knative broker-ingress returned {}",
response.status()
)));
}
Ok(())
}
pub fn manifests(&self, plan: &SubscriptionPlan, subscriptions: &[(&str, &str)]) -> String {
let mut out = String::new();
let mut used = HashSet::new();
if !plan.commands.is_empty() {
let own_commands = format!("{}-commands", self.source);
out.push_str(&self.broker_yaml(&own_commands));
for command in &plan.commands {
out.push_str(&self.trigger_yaml(&own_commands, command, &mut used));
}
}
if self.publishes_events {
out.push_str(&self.broker_yaml(&self.events_broker));
}
for event in &plan.events {
let broker = subscriptions
.iter()
.find(|(name, _)| *name == event.as_str())
.map(|(_, broker)| *broker)
.unwrap_or("UNMAPPED-events");
out.push_str(&self.trigger_yaml(broker, event, &mut used));
}
out
}
fn broker_yaml(&self, name: &str) -> String {
format!(
"apiVersion: eventing.knative.dev/v1\n\
kind: Broker\n\
metadata:\n\
\x20 name: {name}\n\
\x20 namespace: {ns}\n\
---\n",
ns = self.namespace,
)
}
fn trigger_yaml(&self, broker: &str, event: &str, used: &mut HashSet<String>) -> String {
let trigger_name = unique_k8s_name(&format!("{}-{}", self.source, event), used);
let subscriber = match &self.local {
Some(addr) => format!(
"\x20 subscriber:\n\
\x20 uri: http://{addr}/cloudevent/{event}\n"
),
None => format!(
"\x20 subscriber:\n\
\x20 ref:\n\
\x20 apiVersion: serving.knative.dev/v1\n\
\x20 kind: Service\n\
\x20 name: {source}\n\
\x20 uri: /cloudevent/{event}\n",
source = self.source,
),
};
format!(
"apiVersion: eventing.knative.dev/v1\n\
kind: Trigger\n\
metadata:\n\
\x20 name: {trigger_name}\n\
\x20 namespace: {ns}\n\
spec:\n\
\x20 broker: {broker}\n\
\x20 filter:\n\
\x20 attributes:\n\
\x20 type: {event}\n\
{subscriber}\
---\n",
ns = self.namespace,
)
}
}
impl Bus for KnativeBus {
async fn send(&self, name: &str, payload: Vec<u8>) -> Result<(), TransportError> {
self.send_message(Message::new(name, MessageKind::Command, payload))
.await
}
async fn publish(&self, name: &str, payload: Vec<u8>) -> Result<(), TransportError> {
self.publish_message(Message::new(name, MessageKind::Event, payload))
.await
}
async fn send_message(&self, message: Message) -> Result<(), TransportError> {
self.post_cloud_event(&self.commands_broker, &message).await
}
async fn publish_message(&self, message: Message) -> Result<(), TransportError> {
self.post_cloud_event(&self.events_broker, &message).await
}
}