use std::{net::SocketAddr, path::PathBuf, time::Duration};
use serde::{Deserialize, Deserializer, de};
use crate::error::ServerError;
use super::{
config_error,
defaults::{
CLUSTER_BROADCAST_CAPACITY_REQUIRED, DEFAULT_GRPC_ADDRESS, DEFAULT_HTTP_ADDRESS,
DEFAULT_MAX_IN_FLIGHT_ACTIVITIES, DEFAULT_MCP_DISCOVER_TTL_MS,
DEFAULT_MCP_TASK_POLL_INTERVAL_MS, DEFAULT_MCP_TOOLS_LIST_TTL_MS,
DEFAULT_OBSERVABILITY_MAX_EVENT_BYTES, DEFAULT_OBSERVABILITY_MAX_STREAM_EVENTS,
EVENT_BROADCAST_CAPACITY_REQUIRED, MCP_TASK_POLL_INTERVAL_REQUIRED,
OBSERVABILITY_MAX_EVENT_BYTES_REQUIRED, OBSERVABILITY_MAX_STREAM_EVENTS_REQUIRED,
},
};
#[derive(Clone, Debug, Deserialize)]
#[serde(default, deny_unknown_fields)]
pub struct ServerSection {
pub listen_address: SocketAddr,
pub grpc_address: SocketAddr,
#[serde(default)]
pub cors_allowed_origins: Vec<String>,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum StoreBackend {
Memory,
Haematite,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) enum RetiredStoreInput {
Backend,
BackendEnvironment,
Url,
Environment,
Flag,
}
#[derive(Clone, Debug)]
pub struct StoreConfig {
pub backend: StoreBackend,
pub(crate) retired_input: Option<RetiredStoreInput>,
pub owned_shards: Vec<usize>,
pub data_dir: Option<String>,
pub shard_count: usize,
pub cluster: Option<ClusterConfig>,
pub node_cache_budget: Option<haematite::NodeCacheBudget>,
pub lock_acquisition_patience_ms: Option<u64>,
pub lock_acquisition_retry_cadence_ms: Option<u64>,
}
#[derive(Default, Deserialize)]
#[serde(default, deny_unknown_fields)]
struct StoreConfigWire {
backend: Option<String>,
url: Option<String>,
owned_shards: Vec<usize>,
data_dir: Option<String>,
shard_count: Option<usize>,
cluster: Option<ClusterConfig>,
node_cache_budget: Option<haematite::NodeCacheBudget>,
lock_acquisition_patience_ms: Option<u64>,
lock_acquisition_retry_cadence_ms: Option<u64>,
}
impl<'de> Deserialize<'de> for StoreConfig {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: Deserializer<'de>,
{
let wire = StoreConfigWire::deserialize(deserializer)?;
let mut config = Self::default();
if let Some(backend) = wire.backend.as_deref() {
match backend.to_ascii_lowercase().as_str() {
"memory" => config.backend = StoreBackend::Memory,
"haematite" => config.backend = StoreBackend::Haematite,
"libsql" => config.retired_input = Some(RetiredStoreInput::Backend),
_ => {
return Err(de::Error::custom(
"store.backend must be one of: memory, haematite",
));
}
}
}
if wire.url.is_some() && config.retired_input.is_none() {
config.retired_input = Some(RetiredStoreInput::Url);
}
config.owned_shards = wire.owned_shards;
config.data_dir = wire.data_dir;
if let Some(shard_count) = wire.shard_count {
config.shard_count = shard_count;
}
config.cluster = wire.cluster;
config.node_cache_budget = wire.node_cache_budget;
config.lock_acquisition_patience_ms = wire.lock_acquisition_patience_ms;
config.lock_acquisition_retry_cadence_ms = wire.lock_acquisition_retry_cadence_ms;
Ok(config)
}
}
#[derive(Clone, Debug, Deserialize, PartialEq, Eq)]
#[serde(deny_unknown_fields)]
pub struct ClusterConfig {
pub node_id: String,
pub bind_address: SocketAddr,
#[serde(default)]
pub members: Vec<String>,
#[serde(default)]
pub peers: Vec<ClusterPeer>,
#[serde(default)]
pub failover_poll_interval_ms: Option<u64>,
#[serde(default)]
pub failover_confirmations: Option<u32>,
}
#[derive(Clone, Debug, Deserialize, PartialEq, Eq)]
#[serde(deny_unknown_fields)]
pub struct ClusterPeer {
pub name: String,
pub address: SocketAddr,
#[serde(default)]
pub grpc_address: Option<SocketAddr>,
#[serde(default)]
pub owned_shards: Vec<usize>,
}
#[derive(Clone, Debug, Deserialize)]
#[serde(default, deny_unknown_fields)]
pub struct RuntimeSection {
pub scheduler_threads: usize,
pub jit_threshold: Option<u32>,
pub query_timeout_ms: Option<u64>,
}
#[derive(Clone, Debug, Deserialize)]
#[serde(default, deny_unknown_fields)]
pub struct DrainConfig {
pub timeout_seconds: u64,
}
#[derive(Clone, Debug, Deserialize)]
#[serde(default, deny_unknown_fields)]
pub struct AuthConfig {
pub enabled: bool,
pub jwks_url: Option<String>,
pub jwks_refresh_seconds: u64,
}
#[derive(Clone, Debug, Deserialize)]
#[serde(default, deny_unknown_fields)]
pub struct MetricsConfig {
pub enabled: bool,
}
#[derive(Clone, Debug, Deserialize)]
#[serde(default, deny_unknown_fields)]
pub struct NamespacesConfig {
pub default: String,
pub auto_create: AutoCreate,
pub max_in_flight_activities: u32,
}
#[derive(Clone, Copy, Debug, Default, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
pub enum AutoCreate {
#[default]
Open,
Closed,
}
#[derive(Clone, Debug, Deserialize)]
#[serde(default, deny_unknown_fields)]
pub struct ListenConfig {
pub grpc: SocketAddr,
pub http: SocketAddr,
}
#[derive(Clone, Debug, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct TlsConfig {
pub certificate_chain_path: PathBuf,
pub private_key_path: PathBuf,
}
#[derive(Clone, Debug, Deserialize)]
#[serde(default, deny_unknown_fields)]
pub struct OpsConsoleConfig {
pub source: OpsConsoleAssetSource,
}
#[derive(Clone, Debug, Deserialize)]
pub enum OpsConsoleAssetSource {
FileSystem {
asset_path: PathBuf,
},
Embedded,
}
#[derive(Clone, Debug, Deserialize)]
#[serde(default, deny_unknown_fields)]
pub struct NamespaceConfig {
pub mode: NamespaceMode,
}
#[derive(Clone, Debug, Deserialize)]
pub enum NamespaceMode {
SharedEngine,
SingleTenant {
namespace: String,
},
}
#[derive(Clone, Debug, Deserialize)]
#[serde(default, deny_unknown_fields)]
pub struct WorkerConfig {
#[serde(with = "duration_millis")]
pub heartbeat_window: Duration,
#[serde(default)]
pub queue_service: crate::worker::queue_service::QueueServiceConfig,
}
#[derive(Clone, Debug, Deserialize)]
#[serde(default, deny_unknown_fields)]
pub struct WebSocketConfig {
pub outbound_buffer_bound: usize,
pub event_broadcast_capacity: Option<usize>,
pub cluster_broadcast_capacity: Option<usize>,
}
impl WebSocketConfig {
pub(super) fn validate(&self) -> Result<(), ServerError> {
if self.outbound_buffer_bound == 0 {
return config_error("websocket.outbound_buffer_bound must be greater than zero");
}
match self.event_broadcast_capacity {
None | Some(0) => return config_error(EVENT_BROADCAST_CAPACITY_REQUIRED),
Some(_) => {}
}
match self.cluster_broadcast_capacity {
None | Some(0) => return config_error(CLUSTER_BROADCAST_CAPACITY_REQUIRED),
Some(_) => {}
}
Ok(())
}
}
#[derive(Clone, Debug, Default, Deserialize, PartialEq, Eq)]
#[serde(default, deny_unknown_fields)]
pub struct McpConfig {
pub enabled: bool,
pub allowed_origins: Vec<String>,
pub discover_ttl_ms: Option<u64>,
pub tools_list_ttl_ms: Option<u64>,
pub task_poll_interval_ms: Option<u64>,
}
impl McpConfig {
pub(super) fn validate(&self) -> Result<(), ServerError> {
if !self.enabled {
return Ok(());
}
if self.task_poll_interval_ms == Some(0) {
return config_error(MCP_TASK_POLL_INTERVAL_REQUIRED);
}
Ok(())
}
#[must_use]
pub fn resolved(&self) -> ResolvedMcpConfig {
ResolvedMcpConfig {
enabled: self.enabled,
allowed_origins: self.allowed_origins.clone(),
discover_ttl_ms: self.discover_ttl_ms.unwrap_or(DEFAULT_MCP_DISCOVER_TTL_MS),
tools_list_ttl_ms: self
.tools_list_ttl_ms
.unwrap_or(DEFAULT_MCP_TOOLS_LIST_TTL_MS),
task_poll_interval_ms: self
.task_poll_interval_ms
.filter(|value| *value > 0)
.unwrap_or(DEFAULT_MCP_TASK_POLL_INTERVAL_MS),
}
}
}
impl Default for ResolvedMcpConfig {
fn default() -> Self {
McpConfig::default().resolved()
}
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct ResolvedMcpConfig {
pub enabled: bool,
pub allowed_origins: Vec<String>,
pub discover_ttl_ms: u64,
pub tools_list_ttl_ms: u64,
pub task_poll_interval_ms: u64,
}
#[derive(Clone, Copy, Debug, Deserialize, PartialEq, Eq)]
#[serde(default, deny_unknown_fields)]
pub struct ObservabilityConfig {
pub max_event_bytes: usize,
pub max_stream_events: u64,
pub max_batch_events: Option<usize>,
pub max_batch_hold_ms: Option<u64>,
}
impl ObservabilityConfig {
#[must_use]
pub fn with_flush_policy(max_batch_events: usize, max_batch_hold_ms: u64) -> Self {
Self {
max_batch_events: Some(max_batch_events),
max_batch_hold_ms: Some(max_batch_hold_ms),
..Self::default()
}
}
pub(super) fn validate(&self) -> Result<(), ServerError> {
if self.max_event_bytes == 0 {
return config_error(OBSERVABILITY_MAX_EVENT_BYTES_REQUIRED);
}
if self.max_stream_events == 0 {
return config_error(OBSERVABILITY_MAX_STREAM_EVENTS_REQUIRED);
}
Ok(())
}
}
impl Default for ObservabilityConfig {
fn default() -> Self {
Self {
max_event_bytes: DEFAULT_OBSERVABILITY_MAX_EVENT_BYTES,
max_stream_events: DEFAULT_OBSERVABILITY_MAX_STREAM_EVENTS,
max_batch_events: None,
max_batch_hold_ms: None,
}
}
}
#[derive(Clone, Debug, Default, Deserialize)]
#[serde(default, deny_unknown_fields)]
pub struct DeployConfig {
pub enabled: bool,
pub max_archive_bytes: Option<u64>,
pub max_inflated_bytes: Option<u64>,
}
#[derive(Clone, Debug, Default, Deserialize)]
#[serde(default, deny_unknown_fields)]
pub struct DevConfig {
pub enabled: bool,
}
#[derive(Clone, Debug, Default, Deserialize)]
#[serde(default, deny_unknown_fields)]
pub struct OutboxConfig {
pub enabled: bool,
pub poll_interval_ms: Option<u64>,
pub batch_size: Option<u32>,
pub max_attempts: Option<u32>,
pub backoff_base_ms: Option<u64>,
pub backoff_multiplier: Option<u32>,
pub backoff_max_ms: Option<u64>,
pub reconcile_interval_ms: Option<u64>,
pub reconcile_stale_after_ms: Option<u64>,
pub transport: OutboxTransport,
pub liminal_listen_address: Option<String>,
}
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum OutboxTransport {
#[cfg_attr(not(feature = "liminal-transport"), default)]
Grpc,
#[cfg_attr(feature = "liminal-transport", default)]
Liminal,
}
#[derive(Clone, Debug, Default, Deserialize)]
#[serde(default, deny_unknown_fields)]
pub struct AuthoringConfig {
pub gleam_path: Option<PathBuf>,
pub project_root: Option<PathBuf>,
pub workspace_dir: Option<PathBuf>,
}
impl Default for ServerSection {
fn default() -> Self {
Self {
listen_address: DEFAULT_HTTP_ADDRESS,
grpc_address: DEFAULT_GRPC_ADDRESS,
cors_allowed_origins: Vec::new(),
}
}
}
impl Default for StoreConfig {
fn default() -> Self {
Self {
backend: StoreBackend::Haematite,
retired_input: None,
owned_shards: Vec::new(),
data_dir: None,
node_cache_budget: None,
lock_acquisition_patience_ms: None,
lock_acquisition_retry_cadence_ms: None,
shard_count: 64,
cluster: None,
}
}
}
impl Default for RuntimeSection {
fn default() -> Self {
Self {
scheduler_threads: 1,
jit_threshold: None,
query_timeout_ms: None,
}
}
}
impl Default for DrainConfig {
fn default() -> Self {
Self {
timeout_seconds: 30,
}
}
}
impl Default for AuthConfig {
fn default() -> Self {
Self {
enabled: false,
jwks_url: None,
jwks_refresh_seconds: 300,
}
}
}
impl Default for MetricsConfig {
fn default() -> Self {
Self { enabled: true }
}
}
impl Default for NamespacesConfig {
fn default() -> Self {
Self {
default: "default".to_owned(),
auto_create: AutoCreate::default(),
max_in_flight_activities: DEFAULT_MAX_IN_FLIGHT_ACTIVITIES,
}
}
}
impl Default for ListenConfig {
fn default() -> Self {
Self {
grpc: DEFAULT_GRPC_ADDRESS,
http: DEFAULT_HTTP_ADDRESS,
}
}
}
impl Default for OpsConsoleConfig {
fn default() -> Self {
Self {
source: OpsConsoleAssetSource::Embedded,
}
}
}
impl Default for NamespaceConfig {
fn default() -> Self {
Self {
mode: NamespaceMode::SharedEngine,
}
}
}
impl Default for WorkerConfig {
fn default() -> Self {
Self {
heartbeat_window: Duration::from_secs(30),
queue_service: crate::worker::queue_service::QueueServiceConfig::default(),
}
}
}
impl Default for WebSocketConfig {
fn default() -> Self {
Self {
outbound_buffer_bound: 32,
event_broadcast_capacity: None,
cluster_broadcast_capacity: None,
}
}
}
mod duration_millis {
use std::time::Duration;
use serde::{Deserialize, Deserializer};
pub(super) fn deserialize<'de, D>(deserializer: D) -> Result<Duration, D::Error>
where
D: Deserializer<'de>,
{
let millis = u64::deserialize(deserializer)?;
Ok(Duration::from_millis(millis))
}
}