lenso-platform-core 0.1.0

Core runtime primitives for the Lenso backend framework.
Documentation
use crate::db::DbPool;
use crate::error::AppResult;
use crate::runtime_config::descriptor::RuntimeConfigRegistry;
use crate::runtime_config::generated::initialize_generated_values;
use crate::runtime_config::provider::{RuntimeConfigCell, RuntimeConfigProvider};
use crate::runtime_config::snapshot::{RuntimeConfigSnapshot, RuntimeConfigSource};
use crate::runtime_config::store::load_all_values;
use serde_json::Value;
use sqlx::postgres::PgListener;
use std::collections::BTreeMap;
use std::sync::Arc;

/// The channel name used for cross-instance config-change notifications.
pub const CONFIG_NOTIFY_CHANNEL: &str = "config_changed";
const GENERATED_CONFIG_ACTOR: &str = "runtime-config";

async fn load_reconciled_values(
    pool: &DbPool,
    registry: &RuntimeConfigRegistry,
    service_key: &str,
) -> AppResult<BTreeMap<(String, String), Value>> {
    let stored = load_all_values(pool).await?;
    let generated = initialize_generated_values(
        pool,
        registry,
        service_key,
        &stored,
        Some(GENERATED_CONFIG_ACTOR),
    )
    .await?;
    if generated.is_empty() {
        Ok(stored)
    } else {
        load_all_values(pool).await
    }
}

/// Capture the startup-resolved (value, source) for every restart-only
/// descriptor applicable to this service, so later refreshes can revert them.
fn freeze_restart_only(
    registry: &RuntimeConfigRegistry,
    snapshot: &RuntimeConfigSnapshot,
) -> BTreeMap<String, (Value, RuntimeConfigSource)> {
    let mut frozen = BTreeMap::new();
    for descriptor in registry.iter() {
        if !descriptor.restart_only {
            continue;
        }
        if let (Some(value), Some(source)) = (
            snapshot.raw(&descriptor.key),
            snapshot.source(&descriptor.key),
        ) {
            frozen.insert(descriptor.key.to_owned(), (value.clone(), source));
        }
    }
    frozen
}

/// Database-backed runtime-config provider. Holds an atomically swappable snapshot
/// resolved from the registry plus stored overrides for one running service.
#[derive(Debug)]
pub struct PostgresRuntimeConfigProvider {
    pool: DbPool,
    registry: Arc<RuntimeConfigRegistry>,
    service_key: String,
    cell: Arc<RuntimeConfigCell>,
    /// Restart-only keys frozen at their startup-resolved (value, source).
    /// Re-applied on every refresh so running instances keep the startup value
    /// until the process restarts.
    restart_only_frozen: BTreeMap<String, (Value, RuntimeConfigSource)>,
}

impl PostgresRuntimeConfigProvider {
    /// Construct the provider and load the initial snapshot from the store.
    pub async fn connect(
        pool: DbPool,
        registry: Arc<RuntimeConfigRegistry>,
        service_key: impl Into<String>,
    ) -> AppResult<Arc<Self>> {
        let service_key = service_key.into();
        let stored = load_reconciled_values(&pool, &registry, &service_key).await?;
        let snapshot = RuntimeConfigSnapshot::resolve(&registry, &service_key, &stored);
        let restart_only_frozen = freeze_restart_only(&registry, &snapshot);
        let cell = Arc::new(RuntimeConfigCell::new(snapshot));
        Ok(Arc::new(Self {
            pool,
            registry,
            service_key,
            cell,
            restart_only_frozen,
        }))
    }

    /// Reload all stored values and swap in a fresh snapshot.
    ///
    /// Restart-only keys are reverted to their frozen startup values so they do
    /// not take effect until the process restarts.
    pub async fn refresh(&self) -> AppResult<()> {
        let stored = load_reconciled_values(&self.pool, &self.registry, &self.service_key).await?;
        let snapshot = RuntimeConfigSnapshot::resolve(&self.registry, &self.service_key, &stored)
            .with_overrides(&self.restart_only_frozen);
        self.cell.store(snapshot);
        Ok(())
    }

    /// Spawn the background LISTEN task. Refreshes on every notification and
    /// fully reloads on (re)connect, so missed notifications self-heal.
    pub fn spawn_listener(self: &Arc<Self>) {
        let provider = Arc::clone(self);
        tokio::spawn(async move {
            loop {
                match PgListener::connect_with(&provider.pool).await {
                    Ok(mut listener) => {
                        if let Err(error) = listener.listen(CONFIG_NOTIFY_CHANNEL).await {
                            tracing::warn!(error = ?error, "config listener failed to subscribe");
                            tokio::time::sleep(std::time::Duration::from_secs(2)).await;
                            continue;
                        }
                        // Reconcile after (re)subscribing in case we missed events.
                        if let Err(error) = provider.refresh().await {
                            tracing::warn!(error = ?error, "config refresh after subscribe failed");
                        }
                        loop {
                            match listener.recv().await {
                                Ok(notification) => {
                                    tracing::debug!(
                                        payload = %notification.payload(),
                                        "config change notification received"
                                    );
                                    if let Err(error) = provider.refresh().await {
                                        tracing::warn!(error = ?error, "config refresh failed");
                                    }
                                }
                                Err(error) => {
                                    tracing::warn!(error = ?error, "config listener disconnected");
                                    break;
                                }
                            }
                        }
                    }
                    Err(error) => {
                        tracing::warn!(error = ?error, "config listener connect failed");
                    }
                }
                tokio::time::sleep(std::time::Duration::from_secs(2)).await;
            }
        });
    }
}

impl RuntimeConfigProvider for PostgresRuntimeConfigProvider {
    fn snapshot(&self) -> Arc<RuntimeConfigSnapshot> {
        self.cell.load()
    }
}