use std::time::{Duration, Instant};
use async_trait::async_trait;
use serde_json::Value as Json;
use uuid::Uuid;
use super::{CanonicalStore, DurabilityToken};
use crate::runtime::executors::clickhouse::ClickHouseExecutor;
const OUTBOX_SEQ_ID: &str = "outbox_seq";
const OUTBOX_SEQ_LOCK: &str = "__udb_clickhouse_outbox_seq";
const CLICKHOUSE_COUNTER_LOCK_TTL: Duration = Duration::from_secs(30);
const CLICKHOUSE_COUNTER_LOCK_WAIT: Duration = Duration::from_secs(30);
const CLICKHOUSE_COUNTER_LOCK_POLL: Duration = Duration::from_millis(50);
const CLICKHOUSE_SYSTEM_MUTATION_LOCK_TTL: Duration = Duration::from_secs(30);
const CLICKHOUSE_SYSTEM_MUTATION_LOCK_WAIT: Duration = Duration::from_secs(30);
const CLICKHOUSE_SYSTEM_MUTATION_LOCK_POLL: Duration = Duration::from_millis(50);
const KEEPER_LEASE_TABLE: &str = "udb_keeper_advisory_leases";
const KEEPER_LEASE_LIMIT: usize = 100_000;
const KEEPER_STRICT_SETTING: &str = "SETTINGS keeper_map_strict_mode=1";
fn safe_ident(id: &str) -> Result<(), String> {
if id.is_empty() || id.len() > 64 {
return Err(format!(
"ClickHouse identifier '{id}' is invalid: must be 1-64 characters"
));
}
let first = id.chars().next().unwrap();
if !first.is_ascii_alphabetic() && first != '_' {
return Err(format!(
"ClickHouse identifier '{id}' must start with a letter or underscore"
));
}
if !id.chars().all(|c| c.is_ascii_alphanumeric() || c == '_') {
return Err(format!(
"ClickHouse identifier '{id}' contains invalid characters; \
only ASCII letters, digits, and underscores are allowed"
));
}
Ok(())
}
fn safe_keeper_path_component(raw: &str) -> String {
let mut out = String::with_capacity(raw.len().max(1));
for ch in raw.chars().take(64) {
if ch.is_ascii_alphanumeric() || matches!(ch, '_' | '-' | '.') {
out.push(ch);
} else {
out.push('_');
}
}
if out.is_empty() {
"default".to_string()
} else {
out
}
}
pub(super) fn sql_lit(s: &str) -> String {
format!("'{}'", s.replace('\'', "''"))
}
pub struct ClickHouseCanonicalStore {
pub(super) executor: ClickHouseExecutor,
pub(super) instance_name: String,
pub(super) database: String,
}
impl ClickHouseCanonicalStore {
pub fn new(
executor: ClickHouseExecutor,
instance_name: impl Into<String>,
database: impl Into<String>,
) -> Self {
Self {
executor,
instance_name: instance_name.into(),
database: database.into(),
}
}
pub(super) fn executor(&self) -> &ClickHouseExecutor {
&self.executor
}
pub(super) fn qualified(&self, table: &str) -> Result<String, String> {
safe_ident(&self.database)?;
Ok(format!("`{}`.`{}`", self.database, table))
}
fn keeper_lease_table(&self) -> Result<String, String> {
self.qualified(KEEPER_LEASE_TABLE)
}
fn keeper_lease_path(&self) -> String {
format!(
"/udb/{}/{}/advisory_leases",
safe_keeper_path_component(&self.database),
safe_keeper_path_component(&self.instance_name)
)
}
pub(super) fn now_unix_ms() -> i64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as i64)
.unwrap_or(0)
}
async fn current_counter_seq(&self) -> Result<i64, String> {
let counters = self.qualified("udb_counters")?;
let sql = format!(
"SELECT seq FROM {counters} FINAL WHERE id = {id}",
id = sql_lit(OUTBOX_SEQ_ID)
);
let rows = self.executor.select_rows(&sql).await?;
Ok(cell_i64(rows.first(), "seq"))
}
async fn acquire_counter_sequence_lease(&self, owner_id: &str) -> Result<(), String> {
let started = Instant::now();
loop {
if self
.try_acquire_advisory_lease(OUTBOX_SEQ_LOCK, owner_id, CLICKHOUSE_COUNTER_LOCK_TTL)
.await?
{
return Ok(());
}
if started.elapsed() >= CLICKHOUSE_COUNTER_LOCK_WAIT {
return Err(format!(
"clickhouse outbox sequence lock was not acquired within {:?}",
CLICKHOUSE_COUNTER_LOCK_WAIT
));
}
tokio::time::sleep(CLICKHOUSE_COUNTER_LOCK_POLL).await;
}
}
pub(super) async fn acquire_system_mutation_lease(
&self,
lease_name: &str,
owner_id: &str,
op: &str,
) -> Result<(), String> {
let started = Instant::now();
loop {
if self
.try_acquire_advisory_lease(
lease_name,
owner_id,
CLICKHOUSE_SYSTEM_MUTATION_LOCK_TTL,
)
.await?
{
return Ok(());
}
if started.elapsed() >= CLICKHOUSE_SYSTEM_MUTATION_LOCK_WAIT {
return Err(format!(
"clickhouse {op} mutation lock '{lease_name}' was not acquired within {:?}",
CLICKHOUSE_SYSTEM_MUTATION_LOCK_WAIT
));
}
tokio::time::sleep(CLICKHOUSE_SYSTEM_MUTATION_LOCK_POLL).await;
}
}
pub(super) async fn release_system_mutation_lease(
&self,
lease_name: &str,
owner_id: &str,
op: &str,
) -> Result<(), String> {
self.release_advisory_lease(lease_name, owner_id)
.await
.map_err(|err| {
format!("clickhouse {op} mutation lock '{lease_name}' release failed: {err}")
})
}
fn keeper_insert_lease_sql(
table: &str,
lease_name: &str,
owner_id: &str,
expires_at: i64,
token: &str,
) -> String {
format!(
"INSERT INTO {table} (lease_name, owner_id, expires_at, token) \
{KEEPER_STRICT_SETTING} VALUES ({name}, {owner}, {expires}, {token})",
name = sql_lit(lease_name),
owner = sql_lit(owner_id),
expires = expires_at,
token = sql_lit(token),
)
}
fn keeper_update_lease_sql(
table: &str,
lease_name: &str,
owner_id: &str,
expires_at: i64,
token: &str,
previous_token: &str,
) -> String {
format!(
"ALTER TABLE {table} UPDATE owner_id = {owner}, expires_at = {expires}, token = {token} \
WHERE lease_name = {name} AND token = {prev_token} {KEEPER_STRICT_SETTING}",
owner = sql_lit(owner_id),
expires = expires_at,
token = sql_lit(token),
name = sql_lit(lease_name),
prev_token = sql_lit(previous_token),
)
}
fn keeper_select_lease_sql(table: &str, lease_name: &str) -> String {
format!(
"SELECT owner_id, expires_at, token FROM {table} WHERE lease_name = {name}",
name = sql_lit(lease_name)
)
}
}
#[async_trait]
impl CanonicalStore for ClickHouseCanonicalStore {
fn backend_label(&self) -> &'static str {
"clickhouse"
}
fn instance_name(&self) -> &str {
&self.instance_name
}
async fn ensure_system_tables(&self) -> Result<(), String> {
let outbox = self.qualified("udb_outbox_events")?;
self.executor
.execute_ddl(&format!(
"CREATE TABLE IF NOT EXISTS {outbox} (\
event_seq Int64, \
event_id String, \
topic String, \
partition_key String, \
payload String, \
created_at DateTime64(3)\
) ENGINE = MergeTree ORDER BY event_seq"
))
.await
.map_err(|e| format!("ensure_system_tables (clickhouse outbox) failed: {e}"))?;
let counters = self.qualified("udb_counters")?;
self.executor
.execute_ddl(&format!(
"CREATE TABLE IF NOT EXISTS {counters} (\
id String, \
seq Int64, \
version UInt64\
) ENGINE = ReplacingMergeTree(version) ORDER BY id"
))
.await
.map_err(|e| format!("ensure_system_tables (clickhouse counters) failed: {e}"))?;
self.ensure_advisory_lease_table().await?;
Ok(())
}
async fn enqueue_outbox_event(
&self,
event_id: &str,
topic: &str,
partition_key: &str,
payload: &serde_json::Value,
) -> Result<i64, String> {
let sequence_owner = format!("outbox:{event_id}:{}", Uuid::new_v4());
self.acquire_counter_sequence_lease(&sequence_owner).await?;
let allocation = async {
let current = self.current_counter_seq().await?;
let next = current + 1;
let counters = self.qualified("udb_counters")?;
self.executor
.execute_ddl(&format!(
"INSERT INTO {counters} (id, seq, version) VALUES ({id}, {next}, {next})",
id = sql_lit(OUTBOX_SEQ_ID)
))
.await
.map_err(|e| format!("outbox seq counter insert failed: {e}"))?;
let confirmed = self.current_counter_seq().await?;
if confirmed < next {
return Err(format!(
"outbox seq allocation lost a race: inserted {next}, re-read {confirmed} \
(Keeper sequence lock invariant violated)"
));
}
Ok(next)
}
.await;
let release = self
.release_advisory_lease(OUTBOX_SEQ_LOCK, &sequence_owner)
.await;
if let Err(err) = release {
return Err(format!(
"clickhouse outbox sequence lock release failed: {err}"
));
}
let next = allocation?;
let outbox = self.qualified("udb_outbox_events")?;
let payload_text = serde_json::to_string(payload)
.map_err(|e| format!("outbox payload serialise failed: {e}"))?;
self.executor
.execute_ddl(&format!(
"INSERT INTO {outbox} \
(event_seq, event_id, topic, partition_key, payload, created_at) \
VALUES ({next}, {eid}, {topic}, {pk}, {payload}, now64(3))",
eid = sql_lit(event_id),
topic = sql_lit(topic),
pk = sql_lit(partition_key),
payload = sql_lit(&payload_text),
))
.await
.map_err(|e| format!("outbox event insert failed: {e}"))?;
Ok(next)
}
async fn outbox_max_seq(&self) -> Result<i64, String> {
self.current_counter_seq().await
}
async fn current_durability_token(&self) -> Result<DurabilityToken, String> {
let seq = self.outbox_max_seq().await?;
Ok(DurabilityToken::new("clickhouse", seq.to_string()))
}
async fn wait_for_token(
&self,
token: &DurabilityToken,
timeout: Duration,
) -> Result<bool, String> {
if !token.is_for("clickhouse") {
return Err(format!(
"ClickHouseCanonicalStore cannot wait on a '{}' token",
token.backend_label
));
}
let target: i64 = token.value.parse().map_err(|e| {
format!(
"malformed clickhouse durability token '{}': {e}",
token.value
)
})?;
let started = Instant::now();
let poll = super::durability_poll_interval(timeout, super::CLICKHOUSE_DURABILITY_POLL_MS);
loop {
if self.outbox_max_seq().await? >= target {
return Ok(true);
}
if started.elapsed() >= timeout {
return Ok(false);
}
tokio::time::sleep(poll).await;
}
}
async fn ensure_advisory_lease_table(&self) -> Result<(), String> {
let leases = self.keeper_lease_table()?;
let keeper_path = self.keeper_lease_path();
self.executor
.execute_ddl(&format!(
"CREATE TABLE IF NOT EXISTS {leases} (\
lease_name String, \
owner_id String, \
expires_at Int64, \
token String\
) ENGINE = KeeperMap({path}, {limit}) PRIMARY KEY lease_name",
path = sql_lit(&keeper_path),
limit = KEEPER_LEASE_LIMIT,
))
.await
.map_err(|e| {
format!(
"ensure_advisory_lease_table (clickhouse KeeperMap) failed: {e}; \
configure ClickHouse Keeper/KeeperMap for HA-canonical ClickHouse"
)
})?;
Ok(())
}
async fn try_acquire_advisory_lease(
&self,
lease_name: &str,
owner_id: &str,
ttl: std::time::Duration,
) -> Result<bool, String> {
let leases = self.keeper_lease_table()?;
let now = Self::now_unix_ms();
let new_expires = now + (ttl.as_millis() as i64);
let token = format!("{}:{now}:{}", owner_id, Uuid::new_v4());
let insert_sql =
Self::keeper_insert_lease_sql(&leases, lease_name, owner_id, new_expires, &token);
match self.executor.execute_ddl(&insert_sql).await {
Ok(()) => return Ok(true),
Err(insert_err) => {
let read_sql = Self::keeper_select_lease_sql(&leases, lease_name);
let rows = self
.executor
.select_rows(&read_sql)
.await
.map_err(|read_err| {
format!(
"try_acquire_advisory_lease fresh insert failed ({insert_err}); \
current lease read also failed: {read_err}"
)
})?;
if rows.first().is_none() {
return Err(format!(
"try_acquire_advisory_lease fresh insert failed and no existing \
KeeperMap row was visible: {insert_err}"
));
}
}
}
let read_sql = Self::keeper_select_lease_sql(&leases, lease_name);
let rows = self.executor.select_rows(&read_sql).await?;
let current = rows.first();
let cur_owner = cell_str(current, "owner_id");
let cur_expires = cell_i64(current, "expires_at");
let cur_token = cell_str(current, "token");
let is_free = current.is_none() || cur_owner.is_empty() || cur_expires <= now;
let may_acquire = is_free || cur_owner == owner_id;
if !may_acquire {
return Ok(false);
}
let update_sql = Self::keeper_update_lease_sql(
&leases,
lease_name,
owner_id,
new_expires,
&token,
&cur_token,
);
self.executor.execute_ddl(&update_sql).await.map_err(|e| {
format!("try_acquire_advisory_lease KeeperMap update (clickhouse) failed: {e}")
})?;
let confirm = self.executor.select_rows(&read_sql).await?;
let won = cell_str(confirm.first(), "owner_id") == owner_id
&& cell_str(confirm.first(), "token") == token;
Ok(won)
}
async fn release_advisory_lease(&self, lease_name: &str, owner_id: &str) -> Result<(), String> {
let leases = self.keeper_lease_table()?;
let read_sql = Self::keeper_select_lease_sql(&leases, lease_name);
let rows = self.executor.select_rows(&read_sql).await?;
let current = rows.first();
let cur_owner = cell_str(current, "owner_id");
if cur_owner != owner_id {
return Ok(());
}
let cur_token = cell_str(current, "token");
let release_token = format!("released:{}:{}", Self::now_unix_ms(), Uuid::new_v4());
let release_sql =
Self::keeper_update_lease_sql(&leases, lease_name, "", 0, &release_token, &cur_token);
self.executor.execute_ddl(&release_sql).await.map_err(|e| {
format!("release_advisory_lease KeeperMap update (clickhouse) failed: {e}")
})?;
Ok(())
}
}
fn cell_i64(row: Option<&Json>, key: &str) -> i64 {
let Some(cell) = row.and_then(|r| r.get(key)) else {
return 0;
};
cell.as_i64()
.or_else(|| cell.as_str().and_then(|s| s.parse::<i64>().ok()))
.unwrap_or(0)
}
#[cfg(test)]
fn cell_u64(row: Option<&Json>, key: &str) -> u64 {
let Some(cell) = row.and_then(|r| r.get(key)) else {
return 0;
};
cell.as_str()
.and_then(|s| s.parse::<u64>().ok())
.or_else(|| cell.as_u64())
.unwrap_or(0)
}
fn cell_str(row: Option<&Json>, key: &str) -> String {
row.and_then(|r| r.get(key))
.and_then(Json::as_str)
.map(|s| s.to_string())
.unwrap_or_default()
}
#[cfg(test)]
mod tests {
use super::*;
use crate::runtime::executors::clickhouse::ClickHouseConfig;
fn dummy_store() -> ClickHouseCanonicalStore {
let exec = ClickHouseExecutor::new(ClickHouseConfig {
http_base: "http://localhost:8123".to_string(),
username: "default".to_string(),
password: String::new(),
database: "udb".to_string(),
is_cloud: false,
connect_timeout_secs: 10,
query_timeout_secs: 30,
});
ClickHouseCanonicalStore::new(exec, "primary", "udb")
}
#[test]
fn backend_label_is_pinned() {
let store = dummy_store();
assert_eq!(store.backend_label(), "clickhouse");
assert_eq!(store.instance_name(), "primary");
}
#[tokio::test]
async fn wait_for_token_rejects_foreign_backend() {
let store = dummy_store();
let foreign = DurabilityToken::new("postgres", "0/100");
let err = store
.wait_for_token(&foreign, Duration::from_millis(1))
.await
.expect_err("foreign token must be rejected");
assert!(err.contains("cannot wait on"));
}
#[tokio::test]
async fn wait_for_token_rejects_malformed_value() {
let store = dummy_store();
let bad = DurabilityToken::new("clickhouse", "not-an-int");
let err = store
.wait_for_token(&bad, Duration::from_millis(1))
.await
.expect_err("malformed token must error");
assert!(err.contains("malformed"));
}
#[test]
fn unsafe_database_is_rejected() {
let exec = ClickHouseExecutor::new(ClickHouseConfig {
http_base: "http://localhost:8123".to_string(),
username: "default".to_string(),
password: String::new(),
database: "udb".to_string(),
is_cloud: false,
connect_timeout_secs: 10,
query_timeout_secs: 30,
});
let store = ClickHouseCanonicalStore::new(exec, "primary", "evil`; DROP");
assert!(store.qualified("udb_counters").is_err());
}
#[test]
fn keeper_lease_table_and_path_are_pinned() {
let store = dummy_store();
assert_eq!(
store.keeper_lease_table().unwrap(),
"`udb`.`udb_keeper_advisory_leases`"
);
assert_eq!(
store.keeper_lease_path(),
"/udb/udb/primary/advisory_leases"
);
let exec = ClickHouseExecutor::new(ClickHouseConfig {
http_base: "http://localhost:8123".to_string(),
username: "default".to_string(),
password: String::new(),
database: "udb".to_string(),
is_cloud: false,
connect_timeout_secs: 10,
query_timeout_secs: 30,
});
let weird = ClickHouseCanonicalStore::new(exec, "primary/blue:1", "udb");
assert_eq!(
weird.keeper_lease_path(),
"/udb/udb/primary_blue_1/advisory_leases"
);
}
#[test]
fn keeper_lease_sql_uses_strict_mode_and_token_cas() {
let table = "`udb`.`udb_keeper_advisory_leases`";
let insert = ClickHouseCanonicalStore::keeper_insert_lease_sql(
table, "lease-a", "owner-a", 10, "t1",
);
assert!(insert.contains("INSERT INTO `udb`.`udb_keeper_advisory_leases`"));
assert!(insert.contains("SETTINGS keeper_map_strict_mode=1 VALUES"));
assert!(insert.contains("'lease-a'"));
assert!(insert.contains("'owner-a'"));
let update = ClickHouseCanonicalStore::keeper_update_lease_sql(
table, "lease-a", "owner-b", 20, "t2", "t1",
);
assert!(update.starts_with("ALTER TABLE `udb`.`udb_keeper_advisory_leases` UPDATE"));
assert!(update.contains("owner_id = 'owner-b'"));
assert!(update.contains("token = 't2'"));
assert!(update.contains("WHERE lease_name = 'lease-a' AND token = 't1'"));
assert!(update.ends_with("SETTINGS keeper_map_strict_mode=1"));
}
#[test]
fn cell_helpers_accept_string_and_number_forms() {
let row = serde_json::json!({
"seq": 42,
"version": "7",
"owner_id": "owner-a",
});
assert_eq!(cell_i64(Some(&row), "seq"), 42);
assert_eq!(cell_u64(Some(&row), "version"), 7);
assert_eq!(cell_str(Some(&row), "owner_id"), "owner-a");
assert_eq!(cell_i64(Some(&row), "missing"), 0);
assert_eq!(cell_u64(None, "version"), 0);
assert_eq!(cell_str(None, "owner_id"), "");
}
}