pub(crate) mod audit;
pub(crate) mod backups;
pub(crate) mod channels;
pub(crate) mod connectors;
pub(crate) mod engine;
pub(crate) mod functions;
pub(crate) mod packages;
pub(crate) mod trace_dlq;
pub(crate) mod workflows;
use axum::Router;
use axum::routing::{get, patch, post};
use serde::Deserialize;
use serde::Serialize;
use serde_json::json;
use axum::Extension;
use crate::engine::reload_engine;
use crate::errors::OrionError;
use crate::server::admin_auth::AdminPrincipal;
use crate::server::state::AppState;
#[derive(Debug)]
pub(crate) enum StatusAction {
Activate,
Archive,
}
impl StatusAction {
pub(crate) fn parse(
status: crate::storage::models::EntityStatus,
) -> Result<Self, crate::errors::OrionError> {
use crate::storage::models::EntityStatus;
match status {
EntityStatus::Active => Ok(Self::Activate),
EntityStatus::Archived => Ok(Self::Archive),
EntityStatus::Draft => Err(crate::errors::OrionError::validation(
"Invalid status transition to 'draft'. Use 'active' or 'archived'".to_string(),
)),
}
}
}
pub(crate) const MAX_IMPORT_ITEMS: usize = 1000;
pub(crate) fn check_import_batch_size(len: usize) -> Result<(), crate::errors::OrionError> {
if len > MAX_IMPORT_ITEMS {
return Err(crate::errors::OrionError::validation(format!(
"import accepts at most {MAX_IMPORT_ITEMS} items per request, got {len} — \
split the batch"
)));
}
Ok(())
}
pub(crate) use orion_api::{ImportAction, ImportItemError, ImportItemResult, OnConflict};
pub(crate) struct ImportOps<V, K, E, C, U> {
pub validate: V,
pub conflict_key: K,
pub exists: E,
pub create: C,
pub upsert: U,
}
#[derive(Default)]
pub(crate) struct ImportOutcome {
pub imported: u64,
pub failed: u64,
pub unchanged: u64,
pub skipped: u64,
pub errors: Vec<ImportItemError>,
pub results: Vec<ImportItemResult>,
written: Vec<String>,
}
impl ImportOutcome {
pub(crate) fn written(&self) -> impl Iterator<Item = &str> {
self.written.iter().map(String::as_str)
}
fn record(&mut self, index: usize, key: Option<&str>, action: ImportAction) {
if action.is_write() {
self.imported += 1;
if let Some(key) = key {
self.written.push(key.to_string());
}
} else if action == ImportAction::Unchanged {
self.unchanged += 1;
} else {
self.skipped += 1;
}
self.results.push(ImportItemResult {
index: index as u64,
id: key.map(str::to_string),
action: action.as_str().to_string(),
});
}
fn fail(&mut self, index: usize, error: String) {
self.failed += 1;
self.errors.push(ImportItemError {
index: index as u64,
error,
});
}
}
pub(crate) async fn import_items<T, V, K, E, EFut, C, CFut, U, UFut>(
items: Vec<serde_json::Value>,
dry_run: bool,
on_conflict: OnConflict,
ops: ImportOps<V, K, E, C, U>,
) -> ImportOutcome
where
T: serde::de::DeserializeOwned,
V: Fn(&T) -> Result<(), crate::errors::OrionError>,
K: Fn(&T) -> Option<String>,
E: Fn(String) -> EFut,
EFut: std::future::Future<Output = Result<bool, crate::errors::OrionError>>,
C: Fn(T) -> CFut,
CFut: std::future::Future<Output = Result<(), crate::errors::OrionError>>,
U: Fn(T, bool) -> UFut,
UFut: std::future::Future<Output = Result<ImportAction, crate::errors::OrionError>>,
{
let mut out = ImportOutcome::default();
let mut seen: std::collections::HashSet<String> = std::collections::HashSet::new();
for (i, item) in items.into_iter().enumerate() {
let parsed: T = match serde_json::from_value(item) {
Ok(v) => v,
Err(e) => {
out.fail(i, e.to_string());
continue;
}
};
if let Err(e) = (ops.validate)(&parsed) {
out.fail(i, e.client_message());
continue;
}
let key = (ops.conflict_key)(&parsed);
if let Some(ref key) = key {
if seen.contains(key) {
out.fail(
i,
format!(
"'{key}' appears more than once in this batch — the second \
item would conflict with the first"
),
);
continue;
}
if dry_run || on_conflict != OnConflict::Fail {
seen.insert(key.clone());
}
}
if on_conflict == OnConflict::NewVersion {
match (ops.upsert)(parsed, dry_run).await {
Ok(action) => out.record(i, key.as_deref(), action),
Err(e) => out.fail(i, e.client_message()),
}
continue;
}
if let Some(ref key) = key
&& (dry_run || on_conflict == OnConflict::Skip)
{
match (ops.exists)(key.clone()).await {
Ok(true) => {
if on_conflict == OnConflict::Skip {
out.record(i, Some(key), ImportAction::Skipped);
} else {
out.fail(i, format!("'{key}' already exists"));
}
continue;
}
Ok(false) => {}
Err(e) => {
out.fail(
i,
format!("could not check for a conflict: {}", e.client_message()),
);
continue;
}
}
}
if dry_run {
out.record(i, key.as_deref(), ImportAction::Created);
} else {
match (ops.create)(parsed).await {
Ok(()) => out.record(i, key.as_deref(), ImportAction::Created),
Err(e) => out.fail(i, e.client_message()),
}
}
}
out
}
pub(crate) fn versioned_upsert_action(status: &str, identical: bool) -> ImportAction {
use crate::storage::models::EntityStatus;
if status == EntityStatus::Draft.as_str() {
if identical {
ImportAction::Unchanged
} else {
ImportAction::UpdatedDraft
}
} else if identical && status == EntityStatus::Active.as_str() {
ImportAction::Unchanged
} else {
ImportAction::NewVersion
}
}
#[derive(Serialize, utoipa::ToSchema)]
pub(crate) struct ValidationIssue {
pub(crate) field: String,
pub(crate) message: String,
}
#[derive(Serialize, utoipa::ToSchema)]
pub(crate) struct ValidationResponse {
pub(crate) valid: bool,
pub(crate) errors: Vec<ValidationIssue>,
pub(crate) warnings: Vec<ValidationIssue>,
}
impl ValidationEnvelope {
pub(crate) fn new(errors: Vec<ValidationIssue>, warnings: Vec<ValidationIssue>) -> Self {
Self {
data: ValidationResponse {
valid: errors.is_empty(),
errors,
warnings,
},
}
}
}
#[derive(Serialize, utoipa::ToSchema)]
pub(crate) struct ValidationEnvelope {
pub(crate) data: ValidationResponse,
}
pub(crate) fn issues_from_error(err: OrionError) -> Vec<ValidationIssue> {
match err {
OrionError::Validation { details, .. } if !details.is_empty() => details
.into_iter()
.map(|d| ValidationIssue {
field: d.path,
message: d.message,
})
.collect(),
other => vec![ValidationIssue {
field: "(root)".to_string(),
message: other.client_message(),
}],
}
}
pub(crate) fn import_response(
dry_run: bool,
outcome: ImportOutcome,
) -> axum::Json<serde_json::Value> {
let report = orion_api::ImportResult {
dry_run,
imported: outcome.imported,
failed: outcome.failed,
unchanged: outcome.unchanged,
skipped: outcome.skipped,
errors: outcome.errors,
results: outcome.results,
};
axum::Json(json!({ "data": report }))
}
#[derive(Debug, Default, Deserialize, utoipa::IntoParams)]
#[into_params(parameter_in = Query)]
pub(crate) struct ImportQuery {
#[serde(default)]
pub dry_run: bool,
#[serde(default)]
pub on_conflict: OnConflict,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Deserialize, utoipa::ToSchema)]
#[serde(rename_all = "lowercase")]
pub(crate) enum ReloadMode {
#[default]
Now,
Defer,
}
#[derive(Debug, Default, Deserialize, utoipa::IntoParams)]
#[into_params(parameter_in = Query)]
pub(crate) struct StatusChangeQuery {
#[serde(default)]
pub dry_run: bool,
#[serde(default)]
pub reload: ReloadMode,
}
#[derive(Debug, Default, Deserialize, utoipa::IntoParams)]
#[into_params(parameter_in = Query)]
pub(crate) struct ReloadQuery {
#[serde(default)]
pub reload: ReloadMode,
}
pub(crate) const ANONYMOUS_PRINCIPAL: &str = "anonymous";
fn request_details() -> Option<String> {
let ctx = crate::server::request_context::current()?;
let mut details = serde_json::Map::new();
if !ctx.request_id.is_empty() {
details.insert("request_id".into(), json!(ctx.request_id));
}
if !ctx.client_ip.is_empty() {
details.insert("client_ip".into(), json!(ctx.client_ip));
}
if let Some(ua) = ctx.user_agent {
details.insert("user_agent".into(), json!(ua));
}
if let Some(cc) = ctx.change_context {
details.insert("change_context".into(), json!(cc));
}
(!details.is_empty()).then(|| serde_json::Value::Object(details).to_string())
}
fn audit_log(
queue: &crate::queue::audit_queue::AuditQueue,
principal: &Option<Extension<AdminPrincipal>>,
action: &str,
resource_type: &str,
resource_id: &str,
) {
let who = principal
.as_ref()
.map(|e| e.0.key_id.as_str())
.unwrap_or(ANONYMOUS_PRINCIPAL);
let details = request_details();
tracing::info!(
target: "audit",
principal = %who,
action = %action,
resource_type = %resource_type,
resource_id = %resource_id,
details = details.as_deref().unwrap_or("{}"),
"admin_audit_event"
);
crate::metrics::record_admin_audit(action, resource_type);
queue.submit(crate::queue::audit_queue::AuditEvent {
principal: who.to_string(),
action: action.to_string(),
resource_type: resource_type.to_string(),
resource_id: resource_id.to_string(),
details,
});
}
fn audit_log_draft_only(
queue: &crate::queue::audit_queue::AuditQueue,
principal: &Option<Extension<AdminPrincipal>>,
action: &str,
resource_type: &str,
resource_id: &str,
) {
audit_log(queue, principal, action, resource_type, resource_id);
}
async fn audit_and_reload(
state: &AppState,
principal: &Option<Extension<AdminPrincipal>>,
action: &str,
resource_type: &str,
resource_id: &str,
reload: ReloadMode,
) -> Result<(), crate::errors::OrionError> {
audit_log(
&state.audit_queue,
principal,
action,
resource_type,
resource_id,
);
if reload == ReloadMode::Defer {
return Ok(());
}
reload_engine(state).await?;
state.cluster.bump_config_epoch().await
}
pub fn admin_routes(max_body_size: usize) -> Router<AppState> {
let channel_routes = Router::new()
.route(
"/",
get(channels::list_channels).post(channels::create_channel),
)
.route("/import", post(channels::import_channels))
.route("/export", get(channels::export_channels))
.route("/validate", post(channels::validate_channel))
.route(
"/{id}",
get(channels::get_channel)
.put(channels::update_channel)
.delete(channels::delete_channel),
)
.route("/{id}/status", patch(channels::change_channel_status))
.route(
"/{id}/versions",
get(channels::list_channel_versions).post(channels::create_new_channel_version),
);
let workflow_routes = Router::new()
.route(
"/",
get(workflows::list_workflows).post(workflows::create_workflow),
)
.route("/import", post(workflows::import_workflows))
.route("/export", get(workflows::export_workflows))
.route("/validate", post(workflows::validate_workflow))
.route(
"/{id}",
get(workflows::get_workflow)
.put(workflows::update_workflow)
.delete(workflows::delete_workflow),
)
.route("/{id}/status", patch(workflows::change_workflow_status))
.route("/{id}/dependencies", get(workflows::workflow_dependencies))
.route(
"/{id}/versions",
get(workflows::list_workflow_versions).post(workflows::create_new_workflow_version),
)
.route("/{id}/rollout", patch(workflows::update_rollout))
.route("/{id}/test", post(workflows::test_workflow));
let connector_routes = Router::new()
.route(
"/",
get(connectors::list_connectors).post(connectors::create_connector),
)
.route("/import", post(connectors::import_connectors))
.route("/export", get(connectors::export_connectors))
.route("/validate", post(connectors::validate_connector))
.route(
"/{id}",
get(connectors::get_connector)
.put(connectors::update_connector)
.delete(connectors::delete_connector),
)
.route("/{id}/test", post(connectors::test_connector))
.route("/circuit-breakers", get(connectors::list_circuit_breakers))
.route(
"/circuit-breakers/{key}",
post(connectors::reset_circuit_breaker),
);
let engine_routes = Router::new()
.route("/status", get(engine::engine_status))
.route("/reload", post(engine::engine_reload));
let audit_routes = Router::new().route("/", get(audit::list_audit_logs));
let function_routes = Router::new().route("/", get(functions::list_functions));
let trace_routes = Router::new()
.route("/", get(crate::server::routes::data::traces::list_traces))
.route("/{id}", get(crate::server::routes::data::traces::get_trace));
let trace_dlq_routes = Router::new()
.route("/", get(trace_dlq::list_trace_dlq))
.route("/purge", post(trace_dlq::purge_trace_dlq))
.route("/{id}", get(trace_dlq::get_trace_dlq_entry))
.route("/{id}/requeue", post(trace_dlq::requeue_trace_dlq_entry));
let backup_routes =
Router::new().route("/", post(backups::create_backup).get(backups::list_backups));
let package_routes = Router::new()
.route("/", get(packages::list_packages))
.route(
"/{name}",
get(packages::get_package).put(packages::put_package),
);
Router::new()
.nest("/channels", channel_routes)
.nest("/workflows", workflow_routes)
.nest("/connectors", connector_routes)
.nest("/engine", engine_routes)
.nest("/functions", function_routes)
.nest("/audit-logs", audit_routes)
.nest("/traces", trace_routes)
.nest("/trace-dlq", trace_dlq_routes)
.nest("/backups", backup_routes)
.nest("/packages", package_routes)
.layer(axum::extract::DefaultBodyLimit::max(max_body_size))
}