use std::str::FromStr;
use std::sync::Arc;
use std::time::Duration;
use miette::IntoDiagnostic;
use ockam::abac::expr::{eq, ident, str};
use ockam::abac::PolicyExpression::FullExpression;
use ockam::abac::SUBJECT_KEY;
use ockam::transport::HostnamePort;
use ockam_api::address::get_free_address;
use ockam_api::authenticator::direct::{
OCKAM_ROLE_ATTRIBUTE_ENROLLER_VALUE, OCKAM_ROLE_ATTRIBUTE_KEY,
};
use ockam_api::nodes::service::tcp_inlets::Inlets;
use ockam_api::ConnectionStatus;
use ockam_core::api::Reply;
use ockam_multiaddr::MultiAddr;
use tracing::{debug, error, info, warn};
use crate::background_node::BackgroundNodeClientTrait;
use crate::incoming_services::state::{IncomingService, Port};
use crate::state::AppState;
impl AppState {
pub(crate) async fn refresh_inlets(&self) -> crate::Result<()> {
info!("Refreshing inlets");
let services_arc = self.incoming_services();
let services = {
services_arc.read().await.services.clone()
};
if services.is_empty() {
debug!("No incoming services, skipping inlets refresh");
return Ok(());
}
let background_node_client = self.background_node_client().await;
for service in services {
let result = self
.refresh_inlet(background_node_client.clone(), &service)
.await;
{
let mut guard = services_arc.write().await;
match result {
Ok(port) => {
if let Some(service) = guard.find_mut_by_id(service.id()) {
if let Some(port) = port {
service.set_port(port);
}
service.set_connected(true);
}
}
Err(err) => {
warn!(%err, "Failed to refresh TCP inlet for accepted invitation");
if let Some(service) = guard.find_mut_by_id(service.id()) {
service.set_connected(false);
}
}
}
}
if service.removed() {
let mut guard = services_arc.write().await;
guard.remove_by_id(service.id());
}
self.publish_state().await;
}
info!("Inlets refreshed");
Ok(())
}
async fn refresh_inlet(
&self,
background_node_client: Arc<dyn BackgroundNodeClientTrait>,
service: &IncomingService,
) -> crate::Result<Option<Port>> {
let inlet_node_name = &service.local_node_name();
debug!(node = %inlet_node_name, "Checking node status");
if !service.enabled() {
debug!(node = %inlet_node_name, "TCP inlet is disabled by the user, deleting the node");
let _ = self.delete_background_node(inlet_node_name).await;
return Ok(None);
}
if self.is_connected(service, inlet_node_name).await {
Ok(None)
} else {
self.create_inlet(background_node_client, service)
.await
.map(Some)
}
}
async fn is_connected(&self, service: &IncomingService, inlet_node_name: &str) -> bool {
if self.state().await.get_node(inlet_node_name).await.is_ok() {
if let Ok(mut inlet_node) = self.background_node(inlet_node_name).await {
inlet_node.set_timeout_mut(Duration::from_secs(5));
if let Ok(Reply::Successful(inlet)) = inlet_node
.show_inlet(&self.context(), service.inlet_name())
.await
{
if inlet.status == ConnectionStatus::Up {
debug!(node = %inlet_node_name, alias = %inlet.alias, "TCP inlet is already up");
return true;
}
}
}
}
false
}
async fn create_inlet(
&self,
background_node_client: Arc<dyn BackgroundNodeClientTrait>,
service: &IncomingService,
) -> crate::Result<Port> {
debug!(
service_name = service.name(),
"Creating TCP inlet for accepted invitation"
);
let local_node_name = service.local_node_name();
let project_name = match self.match_owned_projects(service).await? {
Some(project_name) => project_name,
None => {
let project_name = service.enrollment_ticket().project()?.name;
background_node_client
.projects()
.enroll(&local_node_name, service.enrollment_ticket().clone())
.await?;
project_name
}
};
debug!(node = %local_node_name, "Creating node to host TCP inlet");
let res = self.backup_logs(&local_node_name).await;
if let Some(err) = res.err() {
error!(
"Error backing up logs for {}. Err={}",
&local_node_name, err
);
}
let res = self.delete_background_node(&local_node_name).await;
if let Some(err) = res.err() {
error!(
"Error deleting background node {}. Err={}",
&local_node_name, err
);
}
background_node_client
.nodes()
.create(&local_node_name, &project_name)
.await?;
tokio::time::sleep(Duration::from_millis(250)).await;
let mut inlet_node = self.background_node(&local_node_name).await?;
inlet_node.set_timeout_mut(Duration::from_secs(5));
let bind_address = match service.address() {
Some(address) => address,
None => get_free_address()?,
};
let inlet_alias = service.inlet_name().to_string();
let expr = eq([
ident(format!("{}.{}", SUBJECT_KEY, OCKAM_ROLE_ATTRIBUTE_KEY)),
str(OCKAM_ROLE_ATTRIBUTE_ENROLLER_VALUE),
]);
inlet_node
.create_inlet(
&self.context(),
&HostnamePort::from(bind_address),
&MultiAddr::from_str(&service.service_route(Some(project_name.as_str())))
.into_diagnostic()?,
&inlet_alias,
&None,
&Some(FullExpression(expr)),
Duration::from_secs(5),
true,
&None,
false,
false,
false,
&None,
)
.await
.map_err(|err| {
warn!(
"Failed to create TCP inlet for accepted invitation: {}",
err
);
err
})?;
Ok(bind_address.port())
}
async fn match_owned_projects(
&self,
service: &IncomingService,
) -> crate::Result<Option<String>> {
let ticket_project = service.enrollment_ticket().project()?;
let state = self.state().await;
let user = state.get_default_user().await?;
if let Some(my_project) = state
.projects()
.get_projects()
.await?
.iter()
.filter(|p|
p.is_admin(&user))
.find(|p| p.project_id() == ticket_project.id)
{
debug!(
"Skipping enrollment, the project {} is owned by the user",
my_project.name()
);
Ok(Some(my_project.name().to_string()))
} else {
Ok(None)
}
}
pub(crate) async fn disable_tcp_inlet(&self, invitation_id: &str) -> crate::Result<()> {
let service = {
let incoming_services_arc = self.incoming_services();
let mut writer = incoming_services_arc.write().await;
let mut service = writer.find_mut_by_id(invitation_id);
if let Some(service) = service.as_mut() {
if !service.enabled() {
debug!(node = %service.local_node_name(), alias = %service.name(), "TCP inlet was already disconnected");
return Ok(());
}
service.disable();
service.set_connected(false);
service.clone()
} else {
return Ok(());
}
};
self.publish_state().await;
self.model_mut(|model| {
let service = model.upsert_incoming_service(invitation_id);
service.enabled = false;
})
.await?;
self.background_node(&service.local_node_name())
.await?
.delete_inlet(&self.context(), service.inlet_name())
.await?;
Ok(())
}
pub(crate) async fn enable_tcp_inlet(&self, invitation_id: &str) -> crate::Result<()> {
let changed = {
let incoming_services_arc = self.incoming_services();
let mut writer = incoming_services_arc.write().await;
let mut service = writer.find_mut_by_id(invitation_id);
if let Some(service) = service.as_mut() {
if service.enabled() {
debug!(node = %service.local_node_name(), alias = %service.name(), "TCP inlet was already enabled");
return Ok(());
}
service.enable();
info!(node = %service.local_node_name(), alias = %service.name(), "Enabled TCP inlet");
true
} else {
false
}
};
if changed {
self.publish_state().await;
self.model_mut(|model| {
let service = model.upsert_incoming_service(invitation_id);
service.enabled = true;
})
.await?;
}
Ok(())
}
}