use std::borrow::Cow;
use std::collections::{BTreeMap, BTreeSet, HashSet};
use std::sync::Arc;
use std::sync::atomic::Ordering;
use anyhow::Result;
use chrono::Utc;
use common::time::sleep;
use reblessive::TreeStack;
use uuid::Uuid;
use web_time::Instant;
use super::builder::Building;
use super::{
Appending, BUILD_CLOSING_SLEEP, BuildGeneration, ExistingPrimaryAppending, IndexBuildPhase,
IndexBuildReportStatus, LEGACY_BATCH_ID, PrimaryAppendingTicket,
};
use crate::catalog::providers::{NodeProvider, TableProvider};
use crate::catalog::{Index, Record};
use crate::ctx::FrozenContext;
use crate::doc::{CursorDoc, Document};
use crate::exe::FlowResultExt as _;
use crate::idx::docids::{DocId, TableDocIds};
use crate::idx::ft::fulltext::FullTextIndex;
use crate::idx::index::IndexOperation;
use crate::idx::{IndexKeyBase, PreviousBuildPrimaryKey};
use crate::key::schema::{
BuildAppendKey, BuildReservationKey, DocLookupKey, DocLookupPrefix, DocPendingKey,
DocPendingPrefix, IndexAppendKey, RecordKey, RecordPrefix,
};
use crate::key::{KVKey, KVKeyDecode, KVSubspace, KVValue, Key, Resumable};
use crate::kvs::{
DatastoreError, Direction, INDEXING_BATCH_SIZE, Transaction, Val,
is_retryable_transaction_conflict,
};
use crate::legacy::analyzer_function::LegacyAnalyzerFunction;
use crate::val::{Array, Number, Object, RecordId, RecordIdKey, RecordIdentity, Value};
const DOC_ID_RECLAIM_MAX_RETRIES: usize = 10;
const PRIMARY_APPEND_REBUILD_MAX_RETRIES: usize = 10;
const LEGACY_MARKER_MAX_READINGS: usize = 27;
struct InitialIndexValue<'a> {
rid: &'a RecordId,
opt_values: Option<Vec<Value>>,
count_cond_match: Option<(bool, bool)>,
}
#[derive(Clone, Debug, Default)]
pub(super) struct CountPrimaryProgress {
cursor: Option<RecordIdKey>,
early: BTreeMap<RecordIdentity, PrimaryAppendingTicket>,
late_read: BTreeSet<RecordIdentity>,
}
impl CountPrimaryProgress {
pub(super) fn resuming_at(cursor: Option<RecordIdKey>) -> Self {
Self {
cursor,
..Self::default()
}
}
fn take_due(
&mut self,
through: Option<&RecordIdentity>,
) -> Vec<(RecordIdentity, PrimaryAppendingTicket)> {
match through {
Some(through) => self.early.extract_if(..=through, |_, _| true).collect(),
None => std::mem::take(&mut self.early).into_iter().collect(),
}
}
}
struct CountPrimaryAppendingScan<'a> {
lookup_tx: &'a Transaction,
progress: &'a mut CountPrimaryProgress,
through: Option<&'a RecordIdKey>,
live_ids: &'a BTreeSet<RecordIdentity>,
older_markers: &'a [(RecordIdentity, PreviousBuildPrimaryKey)],
initial_count: usize,
}
fn queued_no_later(a: PrimaryAppendingTicket, b: PrimaryAppendingTicket) -> bool {
(a.ticket, a.mutation_seq) <= (b.ticket, b.mutation_seq)
}
fn baselined_match(count_cond_match: Option<(bool, bool)>) -> Option<(bool, bool)> {
count_cond_match.map(|(old_matches, _)| (false, old_matches))
}
struct Marker {
ptr: PrimaryAppendingTicket,
appending: Appending,
}
impl Marker {
fn earlier(self, other: Self) -> Self {
if queued_no_later(self.ptr, other.ptr) {
self
} else {
other
}
}
}
enum Met {
Before,
In,
After,
}
impl Met {
fn place(
record: &RecordIdentity,
after: Option<&RecordIdentity>,
through: Option<&RecordIdentity>,
) -> Self {
if through.is_some_and(|through| record > through) {
Met::Before
} else if after.is_some_and(|after| record <= after) {
Met::After
} else {
Met::In
}
}
}
impl Building {
pub(super) async fn index_appending_loop(
&self,
initial_count: usize,
updates_count: &mut usize,
last_prepare_remove_check: &mut Instant,
) -> Result<()> {
let rng = self.ikb.new_ig_range()?;
let generation = self.build_generation.load(Ordering::Acquire);
let durable_rng = if generation == 0 {
None
} else {
Some(self.ikb.new_bg_range(generation)?)
};
loop {
if self.is_aborted().await {
return Ok(());
}
self.is_beyond_threshold(None)?;
self.check_prepare_remove(last_prepare_remove_check).await?;
let (keys, durable_keys) = {
let tx = self.new_read_tx().await?;
let keys = catch!(tx, tx.keys(rng.clone(), INDEXING_BATCH_SIZE, 0, None).await);
let durable_keys = if let Some(durable_rng) = &durable_rng {
catch!(tx, tx.keys(durable_rng.clone(), INDEXING_BATCH_SIZE, 0, None).await)
} else {
Vec::new()
};
tx.cancel().await?;
(keys, durable_keys)
};
let pending = keys.len() + durable_keys.len();
if keys.is_empty() && durable_keys.is_empty() {
self.mark_durable_report(
generation,
IndexBuildReportStatus::Indexing,
Some(initial_count),
Some(0),
Some(*updates_count),
)
.await?;
break;
}
self.mark_durable_report(
generation,
IndexBuildReportStatus::Indexing,
Some(initial_count),
Some(pending),
Some(*updates_count),
)
.await?;
if !keys.is_empty() {
let ctx = self.new_write_tx_ctx().await?;
let tx = ctx.tx();
let saved_updates_count = *updates_count;
let allowed = if generation == 0 {
&[][..]
} else {
&[IndexBuildPhase::Building, IndexBuildPhase::Closing][..]
};
if generation != 0
&& let Err(err) = self.maintain_build_ownership(&tx, generation, allowed).await
{
*updates_count = saved_updates_count;
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict maintaining build ownership, retrying",
)
.await
{
continue;
}
return Err(err);
}
match self.index_appending_range(&ctx, &tx, keys, updates_count).await {
Ok(()) => {}
Err(err) => {
*updates_count = saved_updates_count;
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict in appending range, retrying",
)
.await
{
continue;
}
return Err(err);
}
};
match tx.commit().await {
Ok(()) => {
self.mark_durable_report(
generation,
IndexBuildReportStatus::Indexing,
Some(initial_count),
Some(0),
Some(*updates_count),
)
.await?;
}
Err(err) => {
*updates_count = saved_updates_count;
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict on commit, retrying",
)
.await
{
continue;
}
return Err(err);
}
}
}
if !durable_keys.is_empty() {
let ctx = self.new_write_tx_ctx().await?;
let tx = ctx.tx();
let saved_updates_count = *updates_count;
let allowed = &[IndexBuildPhase::Building, IndexBuildPhase::Closing];
if let Err(err) = self.maintain_build_ownership(&tx, generation, allowed).await {
*updates_count = saved_updates_count;
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict maintaining build ownership, retrying",
)
.await
{
continue;
}
return Err(err);
}
match self
.index_durable_appending_range(&ctx, &tx, durable_keys, updates_count)
.await
{
Ok(()) => match tx.commit().await {
Ok(()) => {
self.mark_durable_report(
generation,
IndexBuildReportStatus::Indexing,
Some(initial_count),
Some(0),
Some(*updates_count),
)
.await?;
}
Err(err) => {
*updates_count = saved_updates_count;
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict on durable appending commit, retrying",
)
.await
{
continue;
}
return Err(err);
}
},
Err(err) => {
*updates_count = saved_updates_count;
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict in durable appending range, retrying",
)
.await
{
continue;
}
return Err(err);
}
}
}
}
Ok(())
}
pub(super) async fn wait_for_durable_reservations(
&self,
generation: BuildGeneration,
last_prepare_remove_check: &mut Instant,
) -> Result<()> {
let rng = self.ikb.new_br_range(generation)?;
loop {
if self.is_aborted().await {
return Ok(());
}
self.check_prepare_remove(last_prepare_remove_check).await?;
let keys = {
let tx = self.new_read_tx().await?;
let keys = catch!(tx, tx.keys(rng.clone(), INDEXING_BATCH_SIZE, 0, None).await);
tx.cancel().await?;
keys
};
if keys.is_empty() {
return Ok(());
}
let now = Utc::now();
let mut blocked = false;
let ctx = self.new_write_tx_ctx().await?;
let tx = ctx.tx();
if let Err(err) =
self.maintain_build_ownership(&tx, generation, &[IndexBuildPhase::Closing]).await
{
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict maintaining build ownership, retrying",
)
.await
{
continue;
}
return Err(err);
}
for key in keys {
let br = BuildReservationKey::decode_key(&key)?;
if let Some(reservation) = tx.get_key(&br, None).await? {
let appending_committed = !tx
.keys(self.ikb.new_bg_ticket_range(br.generation, br.ticket)?, 1, 0, None)
.await?
.is_empty();
let writer_dead = reservation.expires_at <= now
&& !self.reservation_node_is_live(&tx, reservation.node).await?;
if appending_committed || writer_dead {
tx.del_key(&br).await?;
} else {
blocked = true;
}
}
}
if let Err(err) = tx.commit().await {
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict while cleaning build reservations, retrying",
)
.await
{
continue;
}
return Err(err);
}
if blocked {
sleep(BUILD_CLOSING_SLEEP).await;
}
}
}
pub(super) async fn wait_for_prior_generation_reservations(
&self,
below: BuildGeneration,
) -> Result<()> {
let rng = self.ikb.new_br_range_below(below)?;
loop {
if self.is_aborted().await {
return Ok(());
}
let keys = {
let tx = self.new_read_tx().await?;
let keys = catch!(tx, tx.keys(rng.clone(), INDEXING_BATCH_SIZE, 0, None).await);
tx.cancel().await?;
keys
};
if keys.is_empty() {
return Ok(());
}
let now = Utc::now();
let mut blocked = false;
let ctx = self.new_write_tx_ctx().await?;
let tx = ctx.tx();
for key in keys {
let br = BuildReservationKey::decode_key(&key)?;
if let Some(reservation) = tx.get_key(&br, None).await? {
let appending_committed = !tx
.keys(self.ikb.new_bg_ticket_range(br.generation, br.ticket)?, 1, 0, None)
.await?
.is_empty();
let writer_dead = reservation.expires_at <= now
&& !self.reservation_node_is_live(&tx, reservation.node).await?;
if appending_committed || writer_dead {
tx.del_key(&br).await?;
} else {
blocked = true;
}
}
}
if let Err(err) = tx.commit().await {
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict draining prior-generation reservations, retrying",
)
.await
{
continue;
}
return Err(err);
}
if blocked {
sleep(BUILD_CLOSING_SLEEP).await;
}
}
}
async fn reservation_node_is_live(&self, tx: &Transaction, node: Uuid) -> Result<bool> {
Ok(tx.get_node(node).await?.is_some_and(|node| node.is_active()))
}
pub(super) async fn index_initial_batch(
&self,
ctx: &FrozenContext,
tx: &Transaction,
values: &[(Vec<u8>, Val)],
initial_count: usize,
v1_appending_sentinel: &mut bool,
count_progress: &mut Option<CountPrimaryProgress>,
) -> Result<usize> {
let mut rc = false;
let mut count = 0;
let mut live_ids = BTreeSet::new();
let mut last_live_id = None;
let mut older_markers = Vec::new();
let generation = self.build_generation.load(Ordering::Acquire);
let keys =
values.iter().map(|(k, _)| RecordKey::decode_key(k)).collect::<Result<Vec<_>>>()?;
let mut stack = TreeStack::new();
let fulltext_index =
IndexOperation::create_fulltext_index(ctx, self.ix_key.ns, self.ix_key.db, &self.ix)
.await?;
let lookup_tx = self.new_read_tx().await?;
let count_cond_expr: Option<crate::expr::Expr> =
self.ix.count_cond.as_ref().map(|c| c.0.clone());
let result = async {
let mut older_held = if generation != 0 {
self.older_markers_held(&lookup_tx, generation, values, &keys).await?
} else {
Vec::new()
};
for (i, ((_, v), key)) in values.iter().zip(keys).enumerate() {
yield_now!();
if self.is_aborted().await {
return Ok(count);
}
self.is_beyond_threshold(Some(initial_count + count))?;
let val: Record = revision::from_slice(v.as_slice())?;
let rid: Arc<RecordId> = RecordId {
table: key.tb.into_owned(),
key: key.id.into_owned(),
}
.into();
if count_progress.is_some() {
live_ids.insert(RecordIdentity(rid.key.clone()));
last_live_id = Some(rid.key.clone());
}
let older = match older_held.get_mut(i).and_then(Option::take) {
Some((key, ptr)) => {
let marker =
self.confirmed_marker(&lookup_tx, generation, ptr, &rid.key).await?;
if marker.is_some() && count_progress.is_some() {
older_markers.push((RecordIdentity(rid.key.clone()), key));
}
marker
}
None => None,
};
let (opt_values, count_cond_match) = if let Some(a) = self
.check_existing_primary_appending(
&lookup_tx,
&rid.key,
older,
v1_appending_sentinel,
)
.await?
{
(a.old_values, baselined_match(a.count_cond_match))
} else {
let doc = CursorDoc::new(Some(Arc::clone(&rid)), None, val);
let opt_values = stack
.enter(|stk| {
Document::build_opt_values(stk, ctx, &self.opt, &self.ix, &doc)
})
.finish()
.await?;
let count_cond_match = if let Some(expr) = &count_cond_expr {
let new_matches = stack
.enter(|stk| {
crate::legacy::expr_compute(expr, stk, ctx, &self.opt, Some(&doc))
})
.finish()
.await
.catch_return()?
.is_truthy();
Some((false, new_matches))
} else {
None
};
(opt_values, count_cond_match)
};
self.index_initial_values(
ctx,
&mut stack,
&fulltext_index,
InitialIndexValue {
rid: rid.as_ref(),
opt_values,
count_cond_match,
},
&mut rc,
)
.await?;
count += 1;
}
if let Some(progress) = count_progress.as_mut()
&& let Some(through) = last_live_id.as_ref()
{
count += self
.index_missing_count_primary_appendings(
ctx,
CountPrimaryAppendingScan {
lookup_tx: &lookup_tx,
progress,
through: Some(through),
live_ids: &live_ids,
older_markers: &older_markers,
initial_count: initial_count + count,
},
&mut stack,
&mut rc,
)
.await?;
}
self.check_index_compaction(tx, &mut rc).await?;
Ok(count)
}
.await;
let cancel_result = lookup_tx.cancel().await;
match result {
Ok(count) => {
cancel_result?;
Ok(count)
}
Err(err) => {
let _ = cancel_result;
Err(err)
}
}
}
pub(super) async fn index_remaining_count_primary_appendings(
&self,
ctx: &FrozenContext,
tx: &Transaction,
count_progress: &mut Option<CountPrimaryProgress>,
initial_count: usize,
) -> Result<usize> {
let Some(progress) = count_progress.as_mut() else {
return Ok(0);
};
let lookup_tx = self.new_read_tx().await?;
let mut stack = TreeStack::new();
let mut rc = false;
let live_ids = BTreeSet::new();
let result = self
.index_missing_count_primary_appendings(
ctx,
CountPrimaryAppendingScan {
lookup_tx: &lookup_tx,
progress,
through: None,
live_ids: &live_ids,
older_markers: &[],
initial_count,
},
&mut stack,
&mut rc,
)
.await;
let cancel_result = lookup_tx.cancel().await;
match result {
Ok(count) => {
cancel_result?;
self.check_index_compaction(tx, &mut rc).await?;
Ok(count)
}
Err(err) => {
let _ = cancel_result;
Err(err)
}
}
}
async fn index_initial_values(
&self,
ctx: &FrozenContext,
stack: &mut TreeStack,
fulltext_index: &Option<FullTextIndex>,
value: InitialIndexValue<'_>,
rc: &mut bool,
) -> Result<()> {
let InitialIndexValue {
rid,
opt_values,
count_cond_match,
} = value;
let az_fn = LegacyAnalyzerFunction::new(ctx, &self.opt);
let mut io = IndexOperation::new(
ctx,
self.ix_key.ns,
self.ix_key.db,
self.tb,
&self.ix,
None,
opt_values,
rid,
);
if let Some((old_matches, new_matches)) = count_cond_match {
io = io.with_count_cond_match(old_matches, new_matches);
}
if let Some(fulltext_index) = fulltext_index {
stack
.enter(|stk| io.compute_fulltext_with_index(stk, &az_fn, fulltext_index, rc))
.finish()
.await
} else {
stack.enter(|stk| io.compute(stk, &az_fn, rc)).finish().await
}
}
async fn index_missing_count_primary_appendings(
&self,
ctx: &FrozenContext,
scan: CountPrimaryAppendingScan<'_>,
stack: &mut TreeStack,
rc: &mut bool,
) -> Result<usize> {
let Index::Count(_) = &self.ix.index else {
return Ok(0);
};
let generation = self.build_generation.load(Ordering::Acquire);
if generation == 0 {
return Ok(0);
}
let after = scan.progress.cursor.clone().map(RecordIdentity);
let through = scan.through.map(|t| RecordIdentity(t.clone()));
let span_end =
scan.through.map(|t| self.ikb.new_bp_key(generation, t).encode_key()).transpose()?;
let past_span = |key: &[u8]| span_end.as_ref().is_some_and(|end| key > &**end);
for (record, key) in scan.older_markers {
if past_span(key.as_bytes()) {
scan.progress.late_read.insert(record.clone());
}
}
let mut count = 0;
for (record, ptr) in scan.progress.take_due(through.as_ref()) {
yield_now!();
if self.is_aborted().await {
return Ok(count);
}
self.is_beyond_threshold(Some(scan.initial_count + count))?;
if self.settled_in_span(&scan, generation, &record).await? {
continue;
}
let appending = self.queued_appending(scan.lookup_tx, generation, ptr).await?;
self.index_queued_baseline(ctx, stack, rc, appending).await?;
count += 1;
}
let range =
self.ikb.new_bp_span_range(generation, scan.progress.cursor.as_ref(), scan.through)?;
let mut next = (range.start() < range.end()).then_some(range);
while let Some(rng) = next {
if self.is_aborted().await {
return Ok(count);
}
let batch =
scan.lookup_tx.batch_keys_vals(rng.clone(), INDEXING_BATCH_SIZE, None).await?;
next = match (&batch.next, batch.result.last()) {
(Some(_), Some((k, _))) => Some(rng.resume_after(k, Direction::Forward)),
_ => None,
};
for (key, val) in batch.result {
yield_now!();
if self.is_aborted().await {
return Ok(count);
}
self.is_beyond_threshold(Some(scan.initial_count + count))?;
let ptr = PrimaryAppendingTicket::kv_decode_value(&val, ())?;
let appending = self.queued_appending(scan.lookup_tx, generation, ptr).await?;
let record = RecordIdentity(appending.id.clone());
if scan.live_ids.contains(&record) {
continue;
}
let own_key = self.ikb.new_bp_key(generation, &appending.id).encode_key()?;
let appending = if *own_key == *key {
let own = Marker {
ptr,
appending,
};
match self.older_marker(scan.lookup_tx, generation, &record.0).await? {
Some((older_key, older)) => {
if past_span(older_key.as_bytes()) {
scan.progress.late_read.insert(record);
}
own.earlier(older).appending
}
None => own.appending,
}
} else {
match Met::place(&record, after.as_ref(), through.as_ref()) {
Met::Before => {
scan.progress.early.entry(record).or_insert(ptr);
continue;
}
Met::In => {
if self.settled_in_span(&scan, generation, &record).await? {
continue;
}
}
Met::After => {
if scan.progress.late_read.remove(&record) {
continue;
}
let own =
self.own_marker(scan.lookup_tx, generation, &record.0).await?;
if own.is_some_and(|own| queued_no_later(own.ptr, ptr))
|| scan
.lookup_tx
.exists_key(&self.record_key(&record.0), None)
.await?
{
continue;
}
}
}
appending
};
self.index_queued_baseline(ctx, stack, rc, appending).await?;
count += 1;
}
}
if let Some(last) = scan.through {
scan.progress.cursor = Some(last.clone());
}
Ok(count)
}
async fn queued_appending(
&self,
tx: &Transaction,
generation: BuildGeneration,
ptr: PrimaryAppendingTicket,
) -> Result<Appending> {
let bg = self.ikb.new_bg_key(generation, ptr.ticket, ptr.mutation_seq);
let Some(appending) = tx.get_key(&bg, None).await? else {
return Err(
DatastoreError::CorruptedIndex("Durable appending record is missing").into()
);
};
Ok(appending)
}
async fn own_marker(
&self,
tx: &Transaction,
generation: BuildGeneration,
record: &RecordIdKey,
) -> Result<Option<Marker>> {
let Some(ptr) = tx.get_key(&self.ikb.new_bp_key(generation, record), None).await? else {
return Ok(None);
};
self.confirmed_marker(tx, generation, ptr, record).await
}
async fn older_marker(
&self,
tx: &Transaction,
generation: BuildGeneration,
record: &RecordIdKey,
) -> Result<Option<(PreviousBuildPrimaryKey, Marker)>> {
let Some(key) = self.ikb.new_bp_key_in_previous_spelling(generation, record)? else {
return Ok(None);
};
let Some(ptr) = tx.get_key(&key, None).await? else {
return Ok(None);
};
Ok(self.confirmed_marker(tx, generation, ptr, record).await?.map(|marker| (key, marker)))
}
async fn settled_in_span(
&self,
scan: &CountPrimaryAppendingScan<'_>,
generation: BuildGeneration,
record: &RecordIdentity,
) -> Result<bool> {
Ok(scan.live_ids.contains(record)
|| self.own_marker(scan.lookup_tx, generation, &record.0).await?.is_some())
}
async fn older_markers_held(
&self,
tx: &Transaction,
generation: BuildGeneration,
values: &[(Vec<u8>, Val)],
keys: &[RecordKey<'_>],
) -> Result<Vec<Option<(PreviousBuildPrimaryKey, PrimaryAppendingTicket)>>> {
let spelling = self.ikb.previous_bp_spelling(generation)?;
let prefix = RecordPrefix {
ns: self.ix_key.ns,
db: self.ix_key.db,
tb: Cow::Borrowed(self.ikb.table()),
}
.encode_bound()?
.len();
let mut older = Vec::new();
let mut spelled = Vec::with_capacity(values.len());
for ((k, _), key) in values.iter().zip(keys) {
let found = spelling.key(&key.id, &k[prefix..])?;
spelled.push(found.is_some());
older.extend(found);
}
if older.is_empty() {
return Ok(vec![None; values.len()]);
}
let held = tx.get_many_key(older.iter().collect(), None).await?;
let mut held = older.into_iter().zip(held).map(|(key, ptr)| ptr.map(|ptr| (key, ptr)));
Ok(spelled
.into_iter()
.map(|spelled| {
if spelled {
held.next().flatten()
} else {
None
}
})
.collect())
}
async fn confirmed_marker(
&self,
tx: &Transaction,
generation: BuildGeneration,
ptr: PrimaryAppendingTicket,
record: &RecordIdKey,
) -> Result<Option<Marker>> {
let appending = self.queued_appending(tx, generation, ptr).await?;
Ok(appending.id.addresses_same_record(record).then_some(Marker {
ptr,
appending,
}))
}
async fn index_queued_baseline(
&self,
ctx: &FrozenContext,
stack: &mut TreeStack,
rc: &mut bool,
appending: Appending,
) -> Result<()> {
let rid = RecordId {
table: self.ikb.table().clone(),
key: appending.id,
};
let count_cond_match = baselined_match(appending.count_cond_match);
self.index_initial_values(
ctx,
stack,
&None,
InitialIndexValue {
rid: &rid,
opt_values: appending.old_values,
count_cond_match,
},
rc,
)
.await
}
pub(super) async fn resume_strands_an_older_marker(
&self,
generation: BuildGeneration,
cursor: &RecordIdKey,
) -> Result<bool> {
let tx = self.new_read_tx().await?;
let result = async {
let passed = RecordIdentity(cursor.clone());
let range = self.ikb.new_bp_span_range(generation, Some(cursor), None)?;
let mut next = (range.start() < range.end()).then_some(range);
while let Some(rng) = next {
let batch = tx.batch_keys_vals(rng.clone(), INDEXING_BATCH_SIZE, None).await?;
next = match (&batch.next, batch.result.last()) {
(Some(_), Some((k, _))) => Some(rng.resume_after(k, Direction::Forward)),
_ => None,
};
let queued = batch
.result
.iter()
.map(|(_, val)| {
let ptr = PrimaryAppendingTicket::kv_decode_value(val, ())?;
Ok(self.ikb.new_bg_key(generation, ptr.ticket, ptr.mutation_seq))
})
.collect::<Result<Vec<_>>>()?;
for appending in tx.get_many_key(queued, None).await? {
let Some(appending) = appending else {
return Err(DatastoreError::CorruptedIndex(
"Durable appending record is missing",
)
.into());
};
if RecordIdentity(appending.id) <= passed {
return Ok(true);
}
}
}
Ok(false)
}
.await;
let cancel_result = tx.cancel().await;
match result {
Ok(strands) => {
cancel_result?;
Ok(strands)
}
Err(err) => {
let _ = cancel_result;
Err(err)
}
}
}
pub(super) async fn rebuild_primary_appendings(&self) -> Result<()> {
let generation = self.build_generation.load(Ordering::Acquire);
if generation == 0 {
return Ok(());
}
let mut retries = 0usize;
let mut next = Some(self.ikb.new_bp_range(generation)?);
while let Some(rng) = next.clone() {
if self.is_aborted().await {
return Ok(());
}
let ctx = self.new_write_tx_ctx().await?;
let tx = ctx.tx();
catch!(
tx,
self.maintain_build_ownership(&tx, generation, &[IndexBuildPhase::Building]).await
);
let batch = catch!(tx, tx.batch_keys(rng.clone(), INDEXING_BATCH_SIZE, None).await);
let after = batch
.next
.and(batch.result.last())
.map(|last| rng.resume_after(last, Direction::Forward));
for key in &batch.result {
catch!(tx, tx.del(Key::from(key.as_slice())).await);
}
if self
.commit_and_retryable_conflict(
&tx,
"transient conflict clearing primary-append markers, retrying",
)
.await?
{
retries += 1;
if retries > PRIMARY_APPEND_REBUILD_MAX_RETRIES {
return Err(self.primary_append_rebuild_failed());
}
continue;
}
retries = 0;
next = after;
}
let mut retries = 0usize;
let mut next = Some(self.ikb.new_bg_range(generation)?);
while let Some(rng) = next.clone() {
if self.is_aborted().await {
return Ok(());
}
let ctx = self.new_write_tx_ctx().await?;
let tx = ctx.tx();
catch!(
tx,
self.maintain_build_ownership(&tx, generation, &[IndexBuildPhase::Building]).await
);
let batch =
catch!(tx, tx.batch_keys_vals(rng.clone(), INDEXING_BATCH_SIZE, None).await);
let after = match (&batch.next, batch.result.last()) {
(Some(_), Some((k, _))) => Some(rng.resume_after(k, Direction::Forward)),
_ => None,
};
let mut seen = HashSet::new();
let mut firsts = Vec::new();
for (key, val) in &batch.result {
let bg = catch!(tx, BuildAppendKey::decode_key(key));
let appending = catch!(tx, Appending::kv_decode_value(val, ()));
if !seen.insert(RecordIdentity(appending.id.clone())) {
continue;
}
let ptr = PrimaryAppendingTicket {
ticket: bg.ticket,
mutation_seq: bg.mutation_seq,
};
firsts.push((appending.id, ptr));
}
let markers =
firsts.iter().map(|(id, _)| self.ikb.new_bp_key(generation, id)).collect();
let held = catch!(tx, tx.get_many_key(markers, None).await);
for ((id, ptr), held) in firsts.iter().zip(held) {
if held.is_some_and(|held| queued_no_later(held, *ptr)) {
continue;
}
catch!(tx, tx.set_key(&self.ikb.new_bp_key(generation, id), ptr).await);
}
if self
.commit_and_retryable_conflict(
&tx,
"transient conflict rebuilding primary-append markers, retrying",
)
.await?
{
retries += 1;
if retries > PRIMARY_APPEND_REBUILD_MAX_RETRIES {
return Err(self.primary_append_rebuild_failed());
}
continue;
}
retries = 0;
next = after;
}
Ok(())
}
fn primary_append_rebuild_failed(&self) -> anyhow::Error {
DatastoreError::QueryNotExecuted {
message: format!(
"could not rebuild the primary-append markers of index `{}` on table `{}`; rebuild the index",
self.ix.name, self.ix.table_name
),
}
.into()
}
async fn check_existing_primary_appending(
&self,
lookup_tx: &Transaction,
id_key: &RecordIdKey,
older: Option<Marker>,
v1_appending_sentinel: &mut bool,
) -> Result<Option<Appending>> {
match self.load_existing_primary_appending(lookup_tx, id_key, older).await? {
ExistingPrimaryAppending::None => Ok(None),
ExistingPrimaryAppending::Appending(appending) => Ok(Some(appending)),
ExistingPrimaryAppending::Legacy => {
self.cleanup_legacy_primary_appending(id_key).await?;
if !*v1_appending_sentinel {
*v1_appending_sentinel = true;
warn!(
"Found legacy v1 primary appending entry from an older version; legacy queued updates will be ignored. Consider rebuilding index {} on table {}.",
self.ix.name, self.ix.table_name
);
}
Ok(None)
}
}
}
async fn load_existing_primary_appending(
&self,
tx: &Transaction,
id_key: &RecordIdKey,
older: Option<Marker>,
) -> Result<ExistingPrimaryAppending> {
let generation = self.build_generation.load(Ordering::Acquire);
if generation != 0 {
let own = self.own_marker(tx, generation, id_key).await?;
if let Some(first) = own.into_iter().chain(older).reduce(Marker::earlier) {
return Ok(ExistingPrimaryAppending::Appending(first.appending));
}
}
let ip = self.ikb.new_ip_key(id_key.clone());
let Some(pa) = tx.get_key(&ip, None).await? else {
return Ok(ExistingPrimaryAppending::None);
};
if pa.1 == LEGACY_BATCH_ID {
return Ok(ExistingPrimaryAppending::Legacy);
}
let ig = self.ikb.new_ig_key(pa.0, pa.1);
let Some(appending) = tx.get_key(&ig, None).await? else {
return Err(DatastoreError::CorruptedIndex("Appending record is missing").into());
};
Ok(ExistingPrimaryAppending::Appending(appending))
}
async fn cleanup_legacy_primary_appending(&self, id_key: &RecordIdKey) -> Result<()> {
let ctx = self.new_write_tx_ctx().await?;
let tx = ctx.tx();
let ip = self.ikb.new_ip_key(id_key.clone());
let pa = catch!(tx, tx.get_key(&ip, None).await);
if matches!(pa, Some(pa) if pa.1 == LEGACY_BATCH_ID) {
catch!(tx, tx.del_key(&ip).await);
}
let res = tx.commit().await;
match res {
Ok(()) => Ok(()),
Err(err) if is_retryable_transaction_conflict(&err) => {
let _ = tx.cancel().await;
warn!(
"{}: transient conflict while cleaning legacy primary appending entry; continuing",
self.ix.name
);
Ok(())
}
Err(err) => {
let _ = tx.cancel().await;
Err(err)
}
}
}
async fn apply_appending(
&self,
ctx: &FrozenContext,
stack: &mut TreeStack,
fulltext_index: &Option<FullTextIndex>,
appending: Appending,
rc: &mut bool,
) -> Result<RecordIdKey> {
let rid_key = appending.id;
let reclaim_doc_id = appending.new_values.is_none()
&& appending.old_values.is_some()
&& self.ix.uses_shared_doc_ids();
let rid = RecordId {
table: self.ikb.table().clone(),
key: rid_key.clone(),
};
let az_fn = LegacyAnalyzerFunction::new(ctx, &self.opt);
let mut io = IndexOperation::new(
ctx,
self.ix_key.ns,
self.ix_key.db,
self.tb,
&self.ix,
appending.old_values,
appending.new_values,
&rid,
);
if let Some((old_matches, new_matches)) = appending.count_cond_match {
io = io.with_count_cond_match(old_matches, new_matches);
}
if let Some(fulltext_index) = fulltext_index {
stack
.enter(|stk| io.compute_fulltext_with_index(stk, &az_fn, fulltext_index, rc))
.finish()
.await?;
} else {
stack.enter(|stk| io.compute(stk, &az_fn, rc)).finish().await?;
}
if reclaim_doc_id {
let tx = ctx.tx();
if self.other_doc_id_index_building(ctx).await? {
self.defer_doc_id_reclaim(&tx, &rid_key).await?;
} else {
TableDocIds::new(self.ix_key.ns, self.ix_key.db, self.ikb.table().clone())
.remove(&tx, &rid_key)
.await?;
let dp = DocPendingKey::new(
self.ix_key.ns,
self.ix_key.db,
Cow::Borrowed(self.ikb.table()),
RecordIdentity(rid_key.clone()),
);
tx.del_key(&dp).await?;
}
}
Ok(rid_key)
}
async fn other_doc_id_index_building(&self, ctx: &FrozenContext) -> Result<bool> {
let tx = ctx.tx();
let indexes =
tx.all_tb_indexes(self.ix_key.ns, self.ix_key.db, self.ikb.table(), None).await?;
for other in indexes.iter() {
if other.index_id == self.ix.index_id || !other.uses_shared_doc_ids() {
continue;
}
let other_ikb = IndexKeyBase::new(
self.ix_key.ns,
self.ix_key.db,
self.ikb.table().clone(),
other.index_id,
);
if let Some(state) = tx.get_key(&other_ikb.new_bs_key(), None).await?
&& matches!(state.phase, IndexBuildPhase::Building | IndexBuildPhase::Closing)
{
return Ok(true);
}
}
Ok(false)
}
async fn defer_doc_id_reclaim(&self, tx: &Transaction, id: &RecordIdKey) -> Result<()> {
let dp = DocPendingKey::new(
self.ix_key.ns,
self.ix_key.db,
Cow::Borrowed(self.ikb.table()),
RecordIdentity(id.clone()),
);
tx.set_key(&dp, &()).await
}
async fn record_named_by_legacy_marker(
&self,
tx: &Transaction,
marker: &[u8],
) -> Result<Option<RecordIdKey>> {
let table = self.ikb.table();
let pending = DocPendingPrefix::new(self.ix_key.ns, self.ix_key.db, Cow::Borrowed(table))
.encode_bound()?;
let Some(spelling) = marker.strip_prefix(&*pending) else {
return Ok(None);
};
TableDocIds::new(self.ix_key.ns, self.ix_key.db, table.clone())
.record_spelled_as(tx, spelling)
.await
}
fn legacy_marker_readings(&self, marker: &[u8]) -> Result<Vec<RecordIdKey>> {
let table = Cow::Borrowed(self.ikb.table());
let pending =
DocPendingPrefix::new(self.ix_key.ns, self.ix_key.db, table.clone()).encode_bound()?;
let Some(spelling) = marker.strip_prefix(&*pending) else {
return Ok(Vec::new());
};
let mut erased =
DocLookupPrefix::new(self.ix_key.ns, self.ix_key.db, table).encode_bound()?.to_vec();
erased.extend_from_slice(spelling);
let Ok(di) = DocLookupKey::decode_key(&erased) else {
return Ok(Vec::new());
};
Ok(numeric_variants(&di.id, LEGACY_MARKER_MAX_READINGS))
}
pub(super) async fn reclaim_deferred_doc_ids(&self) -> Result<()> {
if !self.ix.uses_shared_doc_ids() {
return Ok(());
}
let docs = TableDocIds::new(self.ix_key.ns, self.ix_key.db, self.ikb.table().clone());
let mut retries = 0usize;
loop {
if self.is_aborted().await {
return Ok(());
}
let ctx = self.new_write_tx_ctx().await?;
let tx = ctx.tx();
if catch!(tx, self.other_doc_id_index_building(&ctx).await) {
tx.cancel().await?;
return Ok(());
}
let rng = catch!(
tx,
DocPendingPrefix::new(
self.ix_key.ns,
self.ix_key.db,
Cow::Borrowed(self.ikb.table())
)
.range()
);
let keys = catch!(tx, tx.keys(rng, INDEXING_BATCH_SIZE, 0, None).await);
if keys.is_empty() {
tx.cancel().await?;
return Ok(());
}
let mut ids: Vec<RecordIdKey> = Vec::with_capacity(keys.len());
let mut seen = HashSet::with_capacity(keys.len());
let mut take = |id: RecordIdKey, ids: &mut Vec<RecordIdKey>| {
if seen.insert(RecordIdentity(id.clone())) {
ids.push(id);
}
};
for k in &keys {
let decodes = match DocPendingKey::decode_key(k) {
Ok(dp) => {
let decoded = dp.id.0;
let has_second_reading = !decoded.hash_agrees_with_eq();
take(decoded, &mut ids);
if !has_second_reading {
continue;
}
true
}
Err(_) => false,
};
if let Some(id) = catch!(tx, self.record_named_by_legacy_marker(&tx, k).await) {
take(id, &mut ids);
}
if !decodes {
for id in catch!(tx, self.legacy_marker_readings(k)) {
take(id, &mut ids);
}
}
}
let record_keys: Vec<RecordKey> = ids.iter().map(|id| self.record_key(id)).collect();
let records = catch!(tx, tx.get_many_key(record_keys, None).await);
let mut reclaimed: Vec<(RecordIdKey, DocId)> = Vec::new();
for (id, existing) in ids.iter().zip(records) {
if existing.is_none()
&& let Some(doc_id) = catch!(tx, docs.remove(&tx, id).await)
{
reclaimed.push((id.clone(), doc_id));
}
}
for k in &keys {
catch!(tx, tx.del(Key::from(k.as_slice())).await);
}
if self
.commit_and_retryable_conflict(
&tx,
"transient conflict reclaiming deferred doc-ID mappings, retrying",
)
.await?
{
retries += 1;
if retries > DOC_ID_RECLAIM_MAX_RETRIES {
warn!(
index = %self.ix.name,
table = %self.ix.table_name,
"giving up the deferred doc-ID reclaim sweep after repeated \
conflicts; leftover markers will be reclaimed by the next \
doc-ID index build on the table"
);
return Ok(());
}
continue;
}
retries = 0;
self.repair_recreated_doc_ids(&docs, &reclaimed).await?;
}
}
fn record_key<'a>(&'a self, id: &'a RecordIdKey) -> RecordKey<'a> {
RecordKey {
ns: self.ix_key.ns,
db: self.ix_key.db,
tb: std::borrow::Cow::Borrowed(self.ikb.table()),
id: std::borrow::Cow::Borrowed(id),
}
}
async fn repair_recreated_doc_ids(
&self,
docs: &TableDocIds,
reclaimed: &[(RecordIdKey, DocId)],
) -> Result<()> {
if reclaimed.is_empty() {
return Ok(());
}
let mut retries = 0usize;
loop {
let ctx = self.new_write_tx_ctx().await?;
let tx = ctx.tx();
let record_keys: Vec<RecordKey> =
reclaimed.iter().map(|(id, _)| self.record_key(id)).collect();
let records = catch!(tx, tx.get_many_key(record_keys, None).await);
let mut restored = 0usize;
for ((id, doc_id), existing) in reclaimed.iter().zip(records) {
if existing.is_some() && catch!(tx, docs.restore(&tx, id, *doc_id).await) {
restored += 1;
}
}
if restored == 0 {
tx.cancel().await?;
return Ok(());
}
if !self
.commit_and_retryable_conflict(
&tx,
"transient conflict repairing re-created doc-ID mappings, retrying",
)
.await?
{
warn!(
index = %self.ix.name,
table = %self.ix.table_name,
restored,
"restored doc-ID mappings of records re-created concurrently \
with the deferred reclaim sweep"
);
return Ok(());
}
retries += 1;
if retries > DOC_ID_RECLAIM_MAX_RETRIES {
warn!(
index = %self.ix.name,
table = %self.ix.table_name,
"giving up doc-ID mapping repair after repeated conflicts; \
affected records re-index on their next write"
);
return Ok(());
}
}
}
async fn index_appending_range(
&self,
ctx: &FrozenContext,
tx: &Transaction,
keys: Vec<Vec<u8>>,
count: &mut usize,
) -> Result<()> {
let mut rc = false;
let mut stack = TreeStack::new();
let fulltext_index =
IndexOperation::create_fulltext_index(ctx, self.ix_key.ns, self.ix_key.db, &self.ix)
.await?;
for k in keys {
yield_now!();
if self.is_aborted().await {
return Ok(());
}
self.is_beyond_threshold(Some(*count))?;
let ig = IndexAppendKey::decode_key(&k)?;
if let Some(appending) = tx.get_key(&ig, None).await? {
let rid_key = self
.apply_appending(ctx, &mut stack, &fulltext_index, appending, &mut rc)
.await?;
tx.del_key(&ig).await?;
let ip = self.ikb.new_ip_key(rid_key);
tx.del_key(&ip).await?;
}
*count += 1;
}
self.check_index_compaction(tx, &mut rc).await?;
Ok(())
}
async fn index_durable_appending_range(
&self,
ctx: &FrozenContext,
tx: &Transaction,
keys: Vec<Vec<u8>>,
count: &mut usize,
) -> Result<()> {
let mut rc = false;
let mut stack = TreeStack::new();
let fulltext_index =
IndexOperation::create_fulltext_index(ctx, self.ix_key.ns, self.ix_key.db, &self.ix)
.await?;
for k in keys {
yield_now!();
if self.is_aborted().await {
return Ok(());
}
self.is_beyond_threshold(Some(*count))?;
let bg = BuildAppendKey::decode_key(&k)?;
if let Some(appending) = tx.get_key(&bg, None).await? {
let rid_key = self
.apply_appending(ctx, &mut stack, &fulltext_index, appending, &mut rc)
.await?;
tx.del_key(&bg).await?;
let bp = self.ikb.new_bp_key(bg.generation, &rid_key);
tx.del_key(&bp).await?;
tx.del_key(&self.ikb.new_br_key(bg.generation, bg.ticket)).await?;
}
*count += 1;
}
self.check_index_compaction(tx, &mut rc).await?;
Ok(())
}
async fn check_index_compaction(&self, tx: &Transaction, rc: &mut bool) -> Result<()> {
if !*rc {
return Ok(());
}
IndexOperation::compaction_trigger(&self.ikb, tx, self.ctx.node_id()).await?;
*rc = false;
Ok(())
}
}
fn numeric_variants(id: &RecordIdKey, cap: usize) -> Vec<RecordIdKey> {
let rebuild = |v: Value| match v {
Value::Array(a) => Some(RecordIdKey::Array(a)),
Value::Object(o) => Some(RecordIdKey::Object(o)),
_ => None,
};
match id {
RecordIdKey::Array(a) => {
value_variants(&Value::Array(a.clone()), cap).into_iter().filter_map(rebuild).collect()
}
RecordIdKey::Object(o) => {
value_variants(&Value::Object(o.clone()), cap).into_iter().filter_map(rebuild).collect()
}
other => vec![other.clone()],
}
}
fn value_variants(value: &Value, cap: usize) -> Vec<Value> {
match value {
Value::Number(n) => number_variants(*n).into_iter().map(Value::Number).take(cap).collect(),
Value::Array(a) => {
let parts = a.0.iter().map(|v| value_variants(v, cap)).collect();
capped_product(parts, cap).into_iter().map(|vs| Value::Array(Array(vs))).collect()
}
Value::Object(o) => {
let keys: Vec<_> = o.0.iter().map(|(k, _)| k.clone()).collect();
let parts = o.0.iter().map(|(_, v)| value_variants(v, cap)).collect();
capped_product(parts, cap)
.into_iter()
.map(|vs| Value::Object(Object(keys.iter().cloned().zip(vs).collect())))
.collect()
}
other => vec![other.clone()],
}
}
fn number_variants(n: Number) -> Vec<Number> {
use rust_decimal::Decimal;
use rust_decimal::prelude::{FromPrimitive, ToPrimitive};
let exact = match n {
Number::Int(i) => Decimal::from(i),
Number::Float(f) => match Decimal::from_f64(f) {
Some(d) => d,
None => return vec![n],
},
Number::Decimal(d) => d,
};
let mut out = Vec::with_capacity(3);
if exact.is_integer()
&& let Some(i) = exact.to_i64()
{
out.push(Number::Int(i));
}
if let Some(f) = exact.to_f64() {
out.push(Number::Float(f));
}
out.push(Number::Decimal(exact));
out
}
fn capped_product(parts: Vec<Vec<Value>>, cap: usize) -> Vec<Vec<Value>> {
let mut acc: Vec<Vec<Value>> = vec![Vec::new()];
for choices in parts {
let mut next = Vec::new();
'fill: for prefix in &acc {
for choice in &choices {
if next.len() == cap {
break 'fill;
}
let mut combined = prefix.clone();
combined.push(choice.clone());
next.push(combined);
}
}
acc = next;
}
acc
}