use std::sync::Arc;
use corium_authz::source::{PolicySource, SourceError};
use corium_db::Db;
use crate::Connection;
pub struct ConnectionPolicySource {
connection: Arc<Connection>,
}
impl ConnectionPolicySource {
#[must_use]
pub fn new(connection: Arc<Connection>) -> Self {
Self { connection }
}
}
#[tonic::async_trait]
impl PolicySource for ConnectionPolicySource {
fn name(&self) -> &str {
self.connection.db_name()
}
async fn snapshot(&self) -> Result<Db, SourceError> {
Ok(self.connection.db())
}
async fn changed(&self, basis_t: u64) {
let mut reports = self.connection.tx_reports();
if self.connection.basis_t() > basis_t {
return;
}
loop {
match reports.recv().await {
Ok(report) if report.t > basis_t => return,
Ok(_) => {}
Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => return,
Err(tokio::sync::broadcast::error::RecvError::Closed) => {
tokio::time::sleep(corium_authz::source::DEFAULT_POLL_INTERVAL).await;
return;
}
}
}
}
}