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;
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
}
}
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
}
#[derive(Debug)]
pub struct PostgresRuntimeConfigProvider {
pool: DbPool,
registry: Arc<RuntimeConfigRegistry>,
service_key: String,
cell: Arc<RuntimeConfigCell>,
restart_only_frozen: BTreeMap<String, (Value, RuntimeConfigSource)>,
}
impl PostgresRuntimeConfigProvider {
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, ®istry, &service_key).await?;
let snapshot = RuntimeConfigSnapshot::resolve(®istry, &service_key, &stored);
let restart_only_frozen = freeze_restart_only(®istry, &snapshot);
let cell = Arc::new(RuntimeConfigCell::new(snapshot));
Ok(Arc::new(Self {
pool,
registry,
service_key,
cell,
restart_only_frozen,
}))
}
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(())
}
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;
}
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()
}
}