platform_core/runtime_config/
postgres.rs1use 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
13pub 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
38fn 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#[derive(Debug)]
62pub struct PostgresRuntimeConfigProvider {
63 pool: DbPool,
64 registry: Arc<RuntimeConfigRegistry>,
65 service_key: String,
66 cell: Arc<RuntimeConfigCell>,
67 restart_only_frozen: BTreeMap<String, (Value, RuntimeConfigSource)>,
71}
72
73impl PostgresRuntimeConfigProvider {
74 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, ®istry, &service_key).await?;
82 let snapshot = RuntimeConfigSnapshot::resolve(®istry, &service_key, &stored);
83 let restart_only_frozen = freeze_restart_only(®istry, &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 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 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 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}