use crate::__bypass::RawAccessExt as _;
use crate::auth::AuthContext;
use crate::cache::DjogiDeltaSyncMeta;
use crate::pg::accumulator::{SqlAccumulator, as_params};
use crate::pg::decode::FromPgRow;
use crate::pg::pool::DjogiPool;
use crate::query::portable::SqlEmitContext;
use heeranjid::{HeerId, HeerIdDesc};
use sassi::{BasicPredicate, DeltaPunnuFetcher, DeltaQuery, DeltaResult, FetchError, PunnuEvent};
use std::any::TypeId;
use std::collections::HashSet;
use std::marker::PhantomData;
use std::sync::Mutex;
use std::sync::atomic::{AtomicBool, Ordering};
use time::OffsetDateTime;
use tokio::sync::broadcast;
use tokio_postgres::types::ToSql;
const FIRST_TICK_WATERMARK_SAFETY_WINDOW: time::Duration = time::Duration::seconds(60);
fn t_id_decodes_from_outbox_bigint<TId: 'static>() -> bool {
TypeId::of::<TId>() == TypeId::of::<HeerId>()
|| TypeId::of::<TId>() == TypeId::of::<HeerIdDesc>()
}
fn cast_row_id_to_t_id<TId: 'static>(raw: i64) -> Option<TId> {
if TypeId::of::<TId>() == TypeId::of::<HeerId>() {
let h = HeerId::from_i64(raw).ok()?;
Some(unsafe { std::mem::transmute_copy::<HeerId, TId>(&h) })
} else if TypeId::of::<TId>() == TypeId::of::<HeerIdDesc>() {
let h = HeerIdDesc::from_i64(raw).ok()?;
Some(unsafe { std::mem::transmute_copy::<HeerIdDesc, TId>(&h) })
} else {
None
}
}
pub(crate) struct DjogiDeltaFetcher<T: sassi::DeltaSyncCacheable> {
pub(crate) pool: DjogiPool,
pub(crate) auth: AuthContext,
pub(crate) empty: bool,
pub(crate) filter: Option<BasicPredicate<T>>,
pub(crate) lru_warn_issued: AtomicBool,
pub(crate) events_rx: Mutex<broadcast::Receiver<PunnuEvent<T>>>,
pub(crate) outbox_watermark: Mutex<Option<OffsetDateTime>>,
pub(crate) _model: PhantomData<T>,
}
#[async_trait::async_trait]
impl<T> DeltaPunnuFetcher<T> for DjogiDeltaFetcher<T>
where
T: sassi::DeltaSyncCacheable
+ FromPgRow
+ crate::model::Model
+ DjogiDeltaSyncMeta
+ Send
+ Sync
+ 'static,
T::Watermark: ToSql + tokio_postgres::types::FromSqlOwned + Sync,
T::Id: ToSql + Sync,
{
async fn fetch_delta(
&self,
query: DeltaQuery<T>,
) -> Result<DeltaResult<T, T::Watermark>, FetchError> {
if self.empty {
return Ok(DeltaResult::new(Vec::new(), HashSet::new()));
}
if !self.lru_warn_issued.load(Ordering::Acquire)
&& let Ok(mut rx) = self.events_rx.try_lock()
{
'drain: loop {
match rx.try_recv() {
Ok(PunnuEvent::Invalidate {
reason: sassi::EventReason::LruEvict { .. },
..
}) => {
if !self.lru_warn_issued.swap(true, Ordering::AcqRel) {
tracing::warn!(
target: "djogi::cache",
model = std::any::type_name::<T>(),
"Punnu LRU eviction detected — `lru_size` may be \
undersized for this subscription's working set. \
Tune via `PunnuConfig::lru_size` if eviction \
collisions become frequent.",
);
}
break 'drain; }
Ok(_) => continue 'drain, Err(broadcast::error::TryRecvError::Empty) => break 'drain,
Err(broadcast::error::TryRecvError::Closed) => break 'drain,
Err(broadcast::error::TryRecvError::Lagged(_)) => {
if !self.lru_warn_issued.swap(true, Ordering::AcqRel) {
tracing::warn!(
target: "djogi::cache",
model = std::any::type_name::<T>(),
"Punnu event stream lagged — LRU eviction events \
may have been dropped. `lru_size` may be \
undersized for this subscription's working set. \
Tune via `PunnuConfig::lru_size` if eviction \
collisions become frequent.",
);
}
break 'drain;
}
}
}
}
let auth = self.auth.clone();
let since = query.since.clone();
let recover_ids = query.recover_ids.clone();
let filter = self.filter.clone();
let watermark_col = <T as DjogiDeltaSyncMeta>::WATERMARK_COLUMN;
let table_name = <T as crate::model::Model>::table_name();
let column_list = <T as FromPgRow>::COLUMN_LIST;
let outbox_enabled =
T::descriptor().has_outbox && t_id_decodes_from_outbox_bigint::<T::Id>();
let outbox_watermark_snapshot: Option<OffsetDateTime> = if outbox_enabled {
*self
.outbox_watermark
.lock()
.expect("outbox_watermark mutex poisoned")
} else {
None
};
let (items, outbox_tombstones, first_tick_server_now, source_high_watermark): (
Vec<T>,
Vec<(i64, OffsetDateTime)>,
Option<OffsetDateTime>,
Option<T::Watermark>,
) = crate::transaction::atomic(&self.pool, move |ctx| {
Box::pin(async move {
ctx.set_auth(auth);
crate::query::terminal::auto_set_tenant::<T>(ctx).await?;
let first_tick_server_now: Option<OffsetDateTime> =
if outbox_enabled && outbox_watermark_snapshot.is_none() {
let row = ctx.query_one("SELECT NOW()", &[]).await?;
Some(row.try_get::<_, OffsetDateTime>(0).map_err(|e| {
crate::DjogiError::Db(crate::error::DbError::other(format!(
"first-tick watermark: SELECT NOW() decode: {e}"
)))
})?)
} else {
None
};
let push_filter = filter.is_some() && since.is_none() && recover_ids.is_empty();
let source_high_watermark_from_max: Option<T::Watermark> = if push_filter {
let max_sql = format!("SELECT MAX({watermark_col}) FROM {table_name}");
let row = ctx.query_one(&max_sql, &[]).await?;
Some(row.try_get::<_, Option<T::Watermark>>(0).map_err(|e| {
crate::DjogiError::Db(crate::error::DbError::other(format!(
"delta refresh: decode MAX({watermark_col}): {e}"
)))
})?)
.flatten()
} else {
None
};
let mut acc = SqlAccumulator::new("SELECT ");
acc.push_sql(column_list);
acc.push_sql(" FROM ");
acc.push_sql(table_name);
match (push_filter, since.as_ref(), recover_ids.is_empty()) {
(true, None, true) => {
acc.push_sql(" WHERE ");
crate::query::portable::emit_basic_predicate::<T>(
&mut acc,
filter.as_ref().expect("push_filter implies filter"),
SqlEmitContext::root(),
crate::query::portable::JsonTrust::Trusted,
)?;
}
(false, None, true) => {}
(false, Some(s), true) => {
acc.push_sql(" WHERE ");
acc.push_sql(watermark_col);
acc.push_sql(" >= ");
acc.push_bind(s.clone());
}
(false, None, false) => {
acc.push_sql(" WHERE id IN (");
acc.push_list_binds(recover_ids.iter().cloned());
acc.push_sql(")");
}
(false, Some(s), false) => {
acc.push_sql(" WHERE (");
acc.push_sql(watermark_col);
acc.push_sql(" >= ");
acc.push_bind(s.clone());
acc.push_sql(") OR (id IN (");
acc.push_list_binds(recover_ids.iter().cloned());
acc.push_sql("))");
}
(true, _, _) => {
unreachable!("filter pushdown only occurs on full baseline ticks")
}
}
acc.push_sql(" ORDER BY ");
acc.push_sql(watermark_col);
let (sql, binds) = acc.into_parts();
let params_refs = as_params(&binds);
let items: Vec<T> = ctx.raw_query::<T>(&sql, ¶ms_refs).await?;
let source_high_watermark = source_high_watermark_from_max;
let outbox_tombstones: Vec<(i64, OffsetDateTime)> = if outbox_enabled
&& let Some(watermark) = outbox_watermark_snapshot
{
let outbox_table = format!("{table_name}_outbox");
crate::ident::check_plain_ident(&outbox_table, false).map_err(|e| {
crate::DjogiError::Db(crate::error::DbError::other(format!(
"outbox poll: invalid outbox table name {outbox_table:?}: {e:?}"
)))
})?;
let outbox_sql = format!(
"SELECT row_id, created_at FROM {outbox_table} \
WHERE action = 'delete' AND created_at >= $1 \
ORDER BY created_at"
);
let rows = ctx
.query_all(&outbox_sql, &[&watermark as &(dyn ToSql + Sync)])
.await?;
let mut decoded: Vec<(i64, OffsetDateTime)> = Vec::with_capacity(rows.len());
for row in rows {
let raw: i64 = row.try_get(0).map_err(|e| {
crate::DjogiError::Db(crate::error::DbError::other(format!(
"outbox poll: decode row_id i64: {e}"
)))
})?;
let ts: OffsetDateTime = row.try_get(1).map_err(|e| {
crate::DjogiError::Db(crate::error::DbError::other(format!(
"outbox poll: decode created_at: {e}"
)))
})?;
decoded.push((raw, ts));
}
decoded
} else {
Vec::new()
};
Ok::<_, crate::DjogiError>((
items,
outbox_tombstones,
first_tick_server_now,
source_high_watermark,
))
})
})
.await
.map_err(|e| FetchError::Custom(Box::new(e)))?;
let mut live_items = Vec::with_capacity(items.len());
let mut tombstones: HashSet<T::Id> = HashSet::new();
for item in items {
if item.__delta_should_tombstone() {
tombstones.insert(<T as sassi::Cacheable>::id(&item));
} else {
live_items.push(item);
}
}
if !outbox_tombstones.is_empty() {
let mut max_seen: Option<OffsetDateTime> = None;
for (raw, ts) in &outbox_tombstones {
if let Some(t_id) = cast_row_id_to_t_id::<T::Id>(*raw) {
tombstones.insert(t_id);
} else {
debug_assert!(
false,
"outbox poll: cast_row_id_to_t_id returned None despite TypeId gate"
);
}
max_seen = Some(match max_seen {
None => *ts,
Some(prev) if *ts > prev => *ts,
Some(prev) => prev,
});
}
if let Some(new_watermark) = max_seen {
let mut guard = self
.outbox_watermark
.lock()
.expect("outbox_watermark mutex poisoned");
let advance = match *guard {
None => true,
Some(prev) => new_watermark > prev,
};
if advance {
*guard = Some(new_watermark);
}
}
} else if let Some(server_now) = first_tick_server_now {
let initial = server_now.saturating_sub(FIRST_TICK_WATERMARK_SAFETY_WINDOW);
let mut guard = self
.outbox_watermark
.lock()
.expect("outbox_watermark mutex poisoned");
if guard.is_none() {
*guard = Some(initial);
}
}
match source_high_watermark {
Some(high_watermark) => Ok(DeltaResult::with_high_watermark(
live_items,
tombstones,
high_watermark,
)),
None => Ok(DeltaResult::new(live_items, tombstones)),
}
}
}
const _: fn() = || {
fn _assert_send_sync_static<T: Send + Sync + 'static>() {}
fn _check_fetcher<T: sassi::DeltaSyncCacheable + Send + Sync + 'static>() {
_assert_send_sync_static::<DjogiDeltaFetcher<T>>();
}
};