Skip to main content

platform_core/runtime_config/
postgres.rs

1use crate::db::DbPool;
2use crate::error::AppResult;
3use crate::runtime_config::descriptor::RuntimeConfigRegistry;
4use crate::runtime_config::generated::initialize_generated_values;
5use crate::runtime_config::provider::{RuntimeConfigCell, RuntimeConfigProvider};
6use crate::runtime_config::snapshot::{RuntimeConfigSnapshot, RuntimeConfigSource};
7use crate::runtime_config::store::load_all_values;
8use serde_json::Value;
9use sqlx::postgres::PgListener;
10use std::collections::BTreeMap;
11use std::sync::Arc;
12
13/// The channel name used for cross-instance config-change notifications.
14pub const CONFIG_NOTIFY_CHANNEL: &str = "config_changed";
15const GENERATED_CONFIG_ACTOR: &str = "runtime-config";
16
17async fn load_reconciled_values(
18    pool: &DbPool,
19    registry: &RuntimeConfigRegistry,
20    service_key: &str,
21) -> AppResult<BTreeMap<(String, String), Value>> {
22    let stored = load_all_values(pool).await?;
23    let generated = initialize_generated_values(
24        pool,
25        registry,
26        service_key,
27        &stored,
28        Some(GENERATED_CONFIG_ACTOR),
29    )
30    .await?;
31    if generated.is_empty() {
32        Ok(stored)
33    } else {
34        load_all_values(pool).await
35    }
36}
37
38/// Capture the startup-resolved (value, source) for every restart-only
39/// descriptor applicable to this service, so later refreshes can revert them.
40fn freeze_restart_only(
41    registry: &RuntimeConfigRegistry,
42    snapshot: &RuntimeConfigSnapshot,
43) -> BTreeMap<String, (Value, RuntimeConfigSource)> {
44    let mut frozen = BTreeMap::new();
45    for descriptor in registry.iter() {
46        if !descriptor.restart_only {
47            continue;
48        }
49        if let (Some(value), Some(source)) = (
50            snapshot.raw(&descriptor.key),
51            snapshot.source(&descriptor.key),
52        ) {
53            frozen.insert(descriptor.key.to_owned(), (value.clone(), source));
54        }
55    }
56    frozen
57}
58
59/// Database-backed runtime-config provider. Holds an atomically swappable snapshot
60/// resolved from the registry plus stored overrides for one running service.
61#[derive(Debug)]
62pub struct PostgresRuntimeConfigProvider {
63    pool: DbPool,
64    registry: Arc<RuntimeConfigRegistry>,
65    service_key: String,
66    cell: Arc<RuntimeConfigCell>,
67    /// Restart-only keys frozen at their startup-resolved (value, source).
68    /// Re-applied on every refresh so running instances keep the startup value
69    /// until the process restarts.
70    restart_only_frozen: BTreeMap<String, (Value, RuntimeConfigSource)>,
71}
72
73impl PostgresRuntimeConfigProvider {
74    /// Construct the provider and load the initial snapshot from the store.
75    pub async fn connect(
76        pool: DbPool,
77        registry: Arc<RuntimeConfigRegistry>,
78        service_key: impl Into<String>,
79    ) -> AppResult<Arc<Self>> {
80        let service_key = service_key.into();
81        let stored = load_reconciled_values(&pool, &registry, &service_key).await?;
82        let snapshot = RuntimeConfigSnapshot::resolve(&registry, &service_key, &stored);
83        let restart_only_frozen = freeze_restart_only(&registry, &snapshot);
84        let cell = Arc::new(RuntimeConfigCell::new(snapshot));
85        Ok(Arc::new(Self {
86            pool,
87            registry,
88            service_key,
89            cell,
90            restart_only_frozen,
91        }))
92    }
93
94    /// Reload all stored values and swap in a fresh snapshot.
95    ///
96    /// Restart-only keys are reverted to their frozen startup values so they do
97    /// not take effect until the process restarts.
98    pub async fn refresh(&self) -> AppResult<()> {
99        let stored = load_reconciled_values(&self.pool, &self.registry, &self.service_key).await?;
100        let snapshot = RuntimeConfigSnapshot::resolve(&self.registry, &self.service_key, &stored)
101            .with_overrides(&self.restart_only_frozen);
102        self.cell.store(snapshot);
103        Ok(())
104    }
105
106    /// Spawn the background LISTEN task. Refreshes on every notification and
107    /// fully reloads on (re)connect, so missed notifications self-heal.
108    pub fn spawn_listener(self: &Arc<Self>) {
109        let provider = Arc::clone(self);
110        tokio::spawn(async move {
111            loop {
112                match PgListener::connect_with(&provider.pool).await {
113                    Ok(mut listener) => {
114                        if let Err(error) = listener.listen(CONFIG_NOTIFY_CHANNEL).await {
115                            tracing::warn!(error = ?error, "config listener failed to subscribe");
116                            tokio::time::sleep(std::time::Duration::from_secs(2)).await;
117                            continue;
118                        }
119                        // Reconcile after (re)subscribing in case we missed events.
120                        if let Err(error) = provider.refresh().await {
121                            tracing::warn!(error = ?error, "config refresh after subscribe failed");
122                        }
123                        loop {
124                            match listener.recv().await {
125                                Ok(notification) => {
126                                    tracing::debug!(
127                                        payload = %notification.payload(),
128                                        "config change notification received"
129                                    );
130                                    if let Err(error) = provider.refresh().await {
131                                        tracing::warn!(error = ?error, "config refresh failed");
132                                    }
133                                }
134                                Err(error) => {
135                                    tracing::warn!(error = ?error, "config listener disconnected");
136                                    break;
137                                }
138                            }
139                        }
140                    }
141                    Err(error) => {
142                        tracing::warn!(error = ?error, "config listener connect failed");
143                    }
144                }
145                tokio::time::sleep(std::time::Duration::from_secs(2)).await;
146            }
147        });
148    }
149}
150
151impl RuntimeConfigProvider for PostgresRuntimeConfigProvider {
152    fn snapshot(&self) -> Arc<RuntimeConfigSnapshot> {
153        self.cell.load()
154    }
155}