use super::{
MemoryAccessEvent, MemoryChangeResult, MemoryChangeSet, MemoryNamespace,
MemoryNamespaceChangeToken, MemoryNamespaceSnapshot, MemoryNode, MemoryQuery,
MemoryQueryResult, MemoryRepositoryError, MemorySnapshotRequest, MemoryUsageSummary,
MAX_QUERY_LIMIT,
};
#[async_trait::async_trait]
pub trait MemoryRepository: Send + Sync {
async fn apply(
&self,
change_set: MemoryChangeSet,
) -> Result<MemoryChangeResult, MemoryRepositoryError>;
async fn get(
&self,
namespace: &MemoryNamespace,
node_id: &str,
) -> Result<Option<MemoryNode>, MemoryRepositoryError>;
async fn query(&self, query: MemoryQuery) -> Result<MemoryQueryResult, MemoryRepositoryError>;
async fn snapshot_namespace(
&self,
request: MemorySnapshotRequest,
) -> Result<MemoryNamespaceSnapshot, MemoryRepositoryError> {
request.validate()?;
let query_limit = request.max_nodes.saturating_add(1).min(MAX_QUERY_LIMIT);
let result = self
.query(
MemoryQuery::new(request.namespace.clone())
.with_statuses(request.statuses.iter().copied())
.with_limit(query_limit),
)
.await?;
if request.max_nodes >= MAX_QUERY_LIMIT && result.hits.len() == MAX_QUERY_LIMIT {
return Err(MemoryRepositoryError::LimitExceeded {
resource: "namespace snapshot query horizon".into(),
limit: MAX_QUERY_LIMIT - 1,
actual: MAX_QUERY_LIMIT,
});
}
super::snapshot::snapshot_from_nodes(
request,
result.hits.into_iter().map(|hit| hit.node).collect(),
)
}
async fn namespace_change_token(
&self,
namespace: &MemoryNamespace,
) -> Result<Option<MemoryNamespaceChangeToken>, MemoryRepositoryError> {
namespace.validate()?;
Ok(None)
}
async fn record_admission(&self, event: MemoryAccessEvent)
-> Result<(), MemoryRepositoryError>;
async fn record_use(&self, event: MemoryAccessEvent) -> Result<(), MemoryRepositoryError>;
async fn usage_summary(
&self,
namespace: &MemoryNamespace,
node_id: &str,
) -> Result<MemoryUsageSummary, MemoryRepositoryError>;
}