use std::collections::BTreeMap;
use std::num::NonZeroUsize;
use std::sync::{Arc, Mutex};
use lru::LruCache;
use crate::graphql_federation::{CompositionError, Supergraph};
use crate::graphql_plan::{plan, PlanError, QueryPlan};
use boatramp_core::config::HandlerGraphqlDataConfig;
use boatramp_core::kv::KvStore;
type SqlSubgraphs = BTreeMap<String, (String, HandlerGraphqlDataConfig)>;
#[derive(Clone)]
pub(crate) struct CachedGraph {
pub version: u64,
pub supergraph: Arc<Supergraph>,
pub sql_subgraphs: Arc<SqlSubgraphs>,
}
const SUPERGRAPH_CAPACITY: usize = 256;
const PLAN_CAPACITY: usize = 1024;
pub(crate) struct GraphqlCache {
supergraphs: Mutex<LruCache<String, CachedGraph>>,
plans: Mutex<LruCache<(String, u64, String), Arc<QueryPlan>>>,
}
impl Default for GraphqlCache {
fn default() -> Self {
Self {
supergraphs: Mutex::new(LruCache::new(
NonZeroUsize::new(SUPERGRAPH_CAPACITY).expect("nonzero"),
)),
plans: Mutex::new(LruCache::new(
NonZeroUsize::new(PLAN_CAPACITY).expect("nonzero"),
)),
}
}
}
impl GraphqlCache {
pub(crate) async fn supergraph(
&self,
kv: &dyn KvStore,
project: &str,
) -> Result<CachedGraph, CompositionError> {
let version = crate::graphql_registry::composition_version(kv, project).await;
let hit = self
.supergraphs
.lock()
.unwrap()
.get(project)
.filter(|c| c.version == version)
.cloned();
if let Some(hit) = hit {
return Ok(hit);
}
let supergraph = Arc::new(crate::graphql_registry::supergraph(kv, project).await?);
let sql_subgraphs = Arc::new(crate::graphql_registry::sql_subgraphs(kv, project).await);
let cached = CachedGraph {
version,
supergraph,
sql_subgraphs,
};
self.supergraphs
.lock()
.unwrap()
.put(project.to_string(), cached.clone());
Ok(cached)
}
pub(crate) fn plan(
&self,
project: &str,
version: u64,
op_hash: &str,
query: &str,
graph: &Supergraph,
) -> Result<Arc<QueryPlan>, PlanError> {
let key = (project.to_string(), version, op_hash.to_string());
let hit = self.plans.lock().unwrap().get(&key).cloned();
if let Some(hit) = hit {
return Ok(hit);
}
let planned = Arc::new(plan(query, graph)?);
self.plans.lock().unwrap().put(key, planned.clone());
Ok(planned)
}
}
#[cfg(test)]
mod tests {
use super::*;
use boatramp_core::kv::{KvStore, MemoryKv};
use std::sync::atomic::{AtomicUsize, Ordering};
const ACCOUNTS: &str = r#"
type Query { me: User }
type User @key(fields: "id") { id: ID! name: String }
"#;
struct CountingKv {
inner: MemoryKv,
lists: AtomicUsize,
}
#[async_trait::async_trait]
impl KvStore for CountingKv {
async fn get(&self, key: &str) -> Result<Option<Vec<u8>>, boatramp_core::kv::KvError> {
self.inner.get(key).await
}
async fn put(&self, key: &str, value: Vec<u8>) -> Result<(), boatramp_core::kv::KvError> {
self.inner.put(key, value).await
}
async fn delete(&self, key: &str) -> Result<(), boatramp_core::kv::KvError> {
self.inner.delete(key).await
}
async fn list_prefix(
&self,
prefix: &str,
) -> Result<Vec<String>, boatramp_core::kv::KvError> {
self.lists.fetch_add(1, Ordering::Relaxed);
self.inner.list_prefix(prefix).await
}
}
#[tokio::test]
async fn a_cache_hit_does_not_recompose() {
let kv = CountingKv {
inner: MemoryKv::new(),
lists: AtomicUsize::new(0),
};
crate::graphql_registry::publish(&kv, "acme", "accounts", ACCOUNTS)
.await
.unwrap();
let cache = GraphqlCache::default();
let first = cache.supergraph(&kv, "acme").await.unwrap();
let after_first = kv.lists.load(Ordering::Relaxed);
assert!(first.supergraph.root_query.contains_key("me"));
let _second = cache.supergraph(&kv, "acme").await.unwrap();
assert_eq!(
kv.lists.load(Ordering::Relaxed),
after_first,
"a cache hit must not re-list/recompose"
);
}
#[tokio::test]
async fn a_registry_mutation_invalidates_the_cache() {
let kv = MemoryKv::new();
crate::graphql_registry::publish(&kv, "acme", "accounts", ACCOUNTS)
.await
.unwrap();
let cache = GraphqlCache::default();
let v1 = cache.supergraph(&kv, "acme").await.unwrap().version;
crate::graphql_registry::publish(
&kv,
"acme",
"reviews",
"type Query { topReviews: [Review] } type Review { id: ID! }",
)
.await
.unwrap();
let after = cache.supergraph(&kv, "acme").await.unwrap();
assert!(after.version > v1, "version advanced after a mutation");
assert!(after.supergraph.root_query.contains_key("topReviews"));
}
#[tokio::test]
async fn plans_are_cached_per_operation_and_projects_are_isolated() {
let kv = MemoryKv::new();
crate::graphql_registry::publish(&kv, "acme", "accounts", ACCOUNTS)
.await
.unwrap();
let cache = GraphqlCache::default();
let graph = cache.supergraph(&kv, "acme").await.unwrap();
let p1 = cache
.plan(
"acme",
graph.version,
"op-a",
"{ me { id } }",
&graph.supergraph,
)
.unwrap();
let p1_again = cache
.plan(
"acme",
graph.version,
"op-a",
"{ me { id } }",
&graph.supergraph,
)
.unwrap();
assert!(Arc::ptr_eq(&p1, &p1_again), "same op → cached plan");
let p_other = cache
.plan(
"other",
graph.version,
"op-a",
"{ me { id } }",
&graph.supergraph,
)
.unwrap();
assert!(
!Arc::ptr_eq(&p1, &p_other),
"distinct projects never share a plan"
);
}
}