use std::sync::Arc;
use crate::Position;
use crate::event::EventRef;
use crate::index::IndexSet;
use crate::log::set::SegmentSet;
use crate::query::{AppendCondition, Matches, Query};
use super::tips::{StagedTips, TagTips, Verdict};
use super::{AppendError, ConflictClause, ConflictSite};
struct EvalCtx<'a> {
main: &'a TagTips,
staged: &'a StagedTips,
index: &'a IndexSet,
set: &'a SegmentSet,
verify: bool,
force_scan: bool,
}
pub fn evaluate(
cond: &AppendCondition,
main: &TagTips,
staged: &StagedTips,
index: &IndexSet,
set: &SegmentSet,
verify: bool,
force_scan: bool,
) -> Result<Option<(ConflictClause, ConflictSite)>, AppendError> {
let ctx = EvalCtx {
main,
staged,
index,
set,
verify,
force_scan,
};
let clauses = [
Some((
ConflictClause::Boundary,
&cond.fail_if_events_match,
cond.after,
)),
cond.fail_if_exists
.as_ref()
.map(|query| (ConflictClause::Existence, query, Position::ZERO)),
];
for (clause, query, after) in clauses.into_iter().flatten() {
if let Some(at) = evaluate_clause(&ctx, query, after)? {
return Ok(Some((clause, at)));
}
}
Ok(None)
}
fn evaluate_clause(
ctx: &EvalCtx<'_>,
query: &Query,
after: Position,
) -> Result<Option<ConflictSite>, AppendError> {
if ctx.staged.may_conflict(query) {
return Ok(Some(ConflictSite::SameBatch));
}
match ctx.main.may_match(query, after) {
Verdict::DefinitelyNoMatch if ctx.verify => {
Ok(verified_against_scan(None, ctx.set, query, after)?.map(ConflictSite::Durable))
}
Verdict::DefinitelyNoMatch => Ok(None),
Verdict::Unknown if ctx.force_scan => {
Ok(scan_for_match(ctx.set, query, after)?.map(ConflictSite::Durable))
}
Verdict::Unknown => match ctx.index.find_match(query, after) {
Ok(found) if ctx.verify => {
Ok(verified_against_scan(found, ctx.set, query, after)?.map(ConflictSite::Durable))
}
Ok(found) => Ok(found.map(ConflictSite::Durable)),
Err(_err) => {
#[cfg(feature = "tracing")]
tracing::warn!(
"index existence check unavailable ({_err}); scanning the log for the condition range"
);
Ok(scan_for_match(ctx.set, query, after)?.map(ConflictSite::Durable))
}
},
}
}
fn verified_against_scan(
fast: Option<Position>,
set: &SegmentSet,
query: &Query,
after: Position,
) -> Result<Option<Position>, AppendError> {
let scanned = scan_for_match(set, query, after)?;
if fast != scanned {
#[cfg(feature = "tracing")]
tracing::error!(
"verify: fast-path {fast:?} disagreed with scan oracle {scanned:?} for query {query:?} after {after}"
);
debug_assert_eq!(
fast, scanned,
"verify: fast-path disagreed with the scan oracle"
);
}
Ok(scanned)
}
fn scan_for_match(
set: &SegmentSet,
query: &Query,
after: Position,
) -> Result<Option<Position>, AppendError> {
let mut scan = set.scan_after(after);
while let Some(item) = scan.next() {
let record = item.map_err(|err| AppendError::Log(Arc::new(err)))?;
let event = EventRef::from_bytes(record.data).map_err(AppendError::Corrupt)?;
if query.matches(event) {
return Ok(Some(record.position));
}
}
Ok(None)
}