use crate::client::admin::FlussAdmin;
use crate::error::{Error, Result};
use crate::metadata::{Schema, SchemaInfo, TablePath};
use parking_lot::RwLock;
use std::collections::HashMap;
use std::sync::Arc;
pub(crate) struct ClientSchemaGetter {
table_path: TablePath,
admin: Arc<FlussAdmin>,
cache: RwLock<HashMap<i32, Arc<Schema>>>,
}
impl ClientSchemaGetter {
pub fn new(table_path: TablePath, admin: Arc<FlussAdmin>, latest: SchemaInfo) -> Self {
let mut map = HashMap::new();
let (schema, schema_id) = latest.into_parts();
map.insert(schema_id, Arc::new(schema));
Self {
table_path,
admin,
cache: RwLock::new(map),
}
}
pub async fn get_schema(&self, schema_id: i32) -> Result<Arc<Schema>> {
if let Some(schema) = self.cache.read().get(&schema_id).cloned() {
return Ok(schema);
}
let info = self
.admin
.get_table_schema(&self.table_path, Some(schema_id))
.await?;
let (schema, fetched_id) = info.into_parts();
if fetched_id != schema_id {
return Err(Error::UnexpectedError {
message: format!(
"Requested schema id {schema_id}, but server returned schema id {fetched_id}"
),
source: None,
});
}
let schema = Arc::new(schema);
self.cache.write().insert(schema_id, Arc::clone(&schema));
Ok(schema)
}
}