mod causal;
pub mod cell_host;
mod context;
mod descriptor;
mod dependencies;
mod error;
mod message_router;
mod projector;
mod runtime;
mod service;
mod session;
#[cfg(any(
feature = "http",
feature = "grpc",
feature = "postgres",
feature = "sqlite",
feature = "nats",
feature = "rabbitmq",
feature = "kafka",
))]
mod workers;
pub use crate::bus::{Message, MessageKind, PayloadDecodeError, SubscriptionPlan};
pub use causal::AggregateCheckout;
pub use context::Context;
pub use dependencies::{
CausalProjectionRouteDependencies, CausalProjectionStore, CausalRepositoryBackend,
CausalRouteDependencies, ConfigurableOutboxPublisher, HasOutboxStore, HasReadModelStore,
HasRepo, ReadModelStoreDependencies, RepoDependencies, RepoReadModelDependencies,
};
pub use error::HandlerError;
pub use descriptor::{
MessageEndpointDescriptor, MetricsEndpointDescriptor, ServiceDescriptor,
ServiceObservabilityDescriptor, TraceExportMode, TracePropagationMode, TracingDescriptor,
TransportDescriptor,
};
pub use projector::{
CausalProjectorContext, CausalProjectorRouteBuilder, LoadedProjection, ProjectionRepairHandle,
ProjectionRepairHandleParseError,
};
#[cfg(feature = "graphql")]
pub use projector::{ModeledProjection, ModeledProjectorRouteBuilder};
pub use runtime::{DEFAULT_MAX_PUBLISH_ATTEMPTS, DEFAULT_PUBLISH_LEASE};
#[cfg(all(feature = "graphql", test))]
pub(crate) use service::CausalCommandProjectionEvidence;
#[cfg(feature = "graphql")]
pub use service::GraphqlServiceBindError;
#[cfg(any(
feature = "http",
feature = "grpc",
feature = "postgres",
feature = "sqlite",
feature = "nats",
feature = "rabbitmq",
feature = "kafka",
))]
pub use workers::{
spawn_outbox_publish_loop, spawn_service_consumer_loop, CONSUMER_IDLE_POLL,
};
pub use service::{
direct_read_model, invoke_transition, require_loaded, CausalCommandContext, CausalCommitBuilder,
CausalRepository, CommandRequest, CommandResponse, DeliveryKind, DirectReadModelProjection,
HandlerNames, HandlerSpec, PortableCommand, PreparedCausalCommit, PreparedCommandHandler,
RouteBuilder, Routes, Service, ThinCommandBuilder, ThinCommandInvoked, ThinCommandLoaded,
TypedRouteBuilder,
};
#[cfg(feature = "graphql")]
pub use service::{CausalCommandPublicStatus, CausalDispatchError, CausalDispatchResult};
#[cfg(feature = "graphql")]
pub(crate) use service::{
CausalCommandProjectionObligation, CausalCommandPublicState, CausalCommandReceiptSource,
CausalProjectionEvidenceState,
};
#[cfg(feature = "graphql")]
pub(crate) mod wait_path;
pub use session::{Session, ROLE_KEY, USER_ID_KEY};
#[cfg(feature = "http")]
pub const MAX_HTTP_BODY_BYTES: usize = 1024 * 1024;
pub(crate) fn lifecycle_mutations_open() -> bool {
let Some(root) = std::env::var_os("DISTRIBUTED_LIFECYCLE_DIR") else {
return true;
};
let Some(generation) = std::env::var_os("DISTRIBUTED_GENERATION_ID") else {
return false;
};
let root = std::path::PathBuf::from(root);
let Some(generation) = generation.to_str() else {
return false;
};
if !root.is_absolute() {
return false;
}
let state = root.join("dev.json");
let Ok(metadata) = std::fs::symlink_metadata(&state) else {
return false;
};
if metadata.file_type().is_symlink()
|| !metadata.is_file()
|| metadata.len() > 1024 * 1024
{
return false;
}
std::fs::read(&state)
.ok()
.is_some_and(|source| lifecycle_state_allows_mutations(&source, generation))
}
fn lifecycle_state_allows_mutations(source: &[u8], generation: &str) -> bool {
let Ok(state) = serde_json::from_slice::<serde_json::Value>(source) else {
return false;
};
state.get("phase").and_then(serde_json::Value::as_str) == Some("active")
&& state
.get("active")
.and_then(|active| active.get("generationId"))
.and_then(serde_json::Value::as_str)
== Some(generation)
}
#[cfg(test)]
mod lifecycle_gate_tests {
use super::lifecycle_state_allows_mutations;
#[test]
fn mutations_open_only_for_the_exact_active_generation() {
let active = br#"{"phase":"active","active":{"generationId":"sha256:one"}}"#;
let preparing =
br#"{"phase":"preparing","active":{"generationId":"sha256:one"}}"#;
assert!(lifecycle_state_allows_mutations(active, "sha256:one"));
assert!(!lifecycle_state_allows_mutations(active, "sha256:two"));
assert!(!lifecycle_state_allows_mutations(preparing, "sha256:one"));
assert!(!lifecycle_state_allows_mutations(b"not-json", "sha256:one"));
}
}
#[cfg(feature = "http")]
mod http;
#[cfg(feature = "http")]
#[allow(unused_imports)]
pub(crate) use http::session_from_headers;
#[cfg(feature = "http")]
pub use http::{router, serve};
#[cfg(feature = "http")]
mod knative_ingress;
#[cfg(feature = "http")]
pub use knative_ingress::cloud_events_router;
#[cfg(feature = "grpc")]
pub mod grpc;
#[cfg(feature = "grpc")]
pub use grpc::{grpc_server, serve_grpc, GrpcServeError};
#[macro_export]
macro_rules! routes {
($routes:expr $(,)?) => {
$routes
};
($routes:expr, $($rest:tt)+) => {
$crate::__routes!($routes, $($rest)+)
};
}
#[doc(hidden)]
#[macro_export]
macro_rules! __routes {
($routes:expr, command $($seg:ident)::+ $(, $($rest:tt)*)?) => {
$crate::__routes_continue!(
$routes.command($($seg)::+::COMMAND).guarded(
$($seg)::+::guard,
$($seg)::+::handle,
)
$(, $($rest)*)?
)
};
($routes:expr, event $($seg:ident)::+ $(, $($rest:tt)*)?) => {
$crate::__routes_continue!(
$routes.event($($seg)::+::EVENT).guarded(
$($seg)::+::guard,
$($seg)::+::handle,
)
$(, $($rest)*)?
)
};
($routes:expr, events $($seg:ident)::+ $(, $($rest:tt)*)?) => {
$crate::__routes_continue!(
$routes.events($($seg)::+::EVENTS).guarded(
$($seg)::+::guard,
$($seg)::+::handle,
)
$(, $($rest)*)?
)
};
($routes:expr, $($seg:ident)::+ $(, $($rest:tt)*)?) => {
compile_error!(
"routes! entries must be prefixed with `command`, `event`, or `events`"
)
};
}
#[doc(hidden)]
#[macro_export]
macro_rules! __routes_continue {
($routes:expr) => {
$routes
};
($routes:expr,) => {
$routes
};
($routes:expr, $($rest:tt)+) => {
$crate::__routes!($routes, $($rest)+)
};
}