uqa-execution 0.4.9

Volcano physical operators with row-batch pipelines
//
// Unified Query Algebra
//
// Copyright (c) 2023-2026 Cognica, Inc.
//

//! Reserve sequence blocks, consume session caches and publish allocation observations.
use super::{
    positions::{Reserved, ReservedBlock},
    NextvalTarget, SequenceValueContext, SequenceValueError,
};
use crate::catalog::sequence::{
    session::{
        NontransactionalSequenceValue, SessionLastSequenceReference, SessionSequenceCache,
        SessionSequenceValue,
    },
    SequenceState,
};
use std::collections::BTreeMap;
use uqa_core::RelationIdentity;
use uqa_storage::SequenceReservationResult;
impl SequenceValueContext<'_> {
    pub(super) fn take_cached_nextval(
        target: &NextvalTarget,
        caches: &mut BTreeMap<RelationIdentity, SessionSequenceCache>,
    ) -> Result<Option<(i64, bool)>, SequenceValueError> {
        let cache_relation = caches
            .contains_key(&target.relation)
            .then(|| target.relation.clone())
            .or_else(|| {
                caches.iter().find_map(|(relation, cache)| {
                    (cache.object_id == target.object_id).then(|| relation.clone())
                })
            });
        let Some(cache_relation) = cache_relation else {
            return Ok(None);
        };
        let cache = caches
            .remove(&cache_relation)
            .expect("selected sequence cache entry must exist");
        if cache.object_id != target.object_id
            || cache.definition_generation != target.state.definition_generation
        {
            return Ok(None);
        }
        let current = cache.next_value;
        if cache.remaining > 1 {
            let next_value = current.checked_add(target.state.increment).ok_or_else(|| {
                SequenceValueError::Internal(format!(
                    "cached sequence `{}` value overflow",
                    target.name
                ))
            })?;
            caches.insert(
                target.relation.clone(),
                SessionSequenceCache {
                    next_value,
                    remaining: cache.remaining - 1,
                    ..cache
                },
            );
        }
        Ok(Some((current, cache.autonomous)))
    }
    pub(super) fn reserve_nextval_block(
        &self,
        target: &NextvalTarget,
    ) -> Result<Reserved, SequenceValueError> {
        let private = self.sequence_value_is_private(
            target.temporary,
            target.object_id,
            target.state.definition_generation,
        )?;
        if let Some(positions) = self.shared_positions(target, private) {
            return self.reserve_at_position(positions, target);
        }
        if let Some((result, autonomous)) = self.mutate_persistent_value(
            target.temporary,
            private,
            "reserve sequence values",
            |catalog| {
                catalog.reserve_sequence_values(
                    &target.name,
                    target.object_id,
                    target.state.definition_generation,
                )
            },
        )? {
            return match result {
                SequenceReservationResult::Reserved(reservation) => Ok(Some(ReservedBlock {
                    reservation,
                    autonomous,
                    positioned: false,
                })),
                SequenceReservationResult::DefinitionChanged => Ok(None),
                SequenceReservationResult::Missing => {
                    Err(SequenceValueError::Undefined(target.name.clone()))
                }
                SequenceReservationResult::Exhausted => {
                    Err(exhausted(&target.relation.name, target.state))
                }
            };
        }
        let mut sequences = self.runtime.states_write();
        let sequence = sequences
            .get_mut(&target.relation)
            .ok_or_else(|| SequenceValueError::Undefined(target.name.clone()))?;
        if sequence.definition_generation != target.state.definition_generation {
            return Ok(None);
        }
        let reservation = uqa_storage::sequence_value_reservation(
            uqa_storage::SequenceValuePosition {
                current: sequence.current,
                called: sequence.called,
                log_count: sequence.log_count,
            },
            sequence.increment,
            sequence.min_value,
            sequence.max_value,
            sequence.cycle,
            sequence.cache_size,
        )
        .ok_or_else(|| exhausted(&target.relation.name, *sequence))?;
        sequence.current = reservation.last_value;
        sequence.called = true;
        sequence.log_count = reservation.log_count;
        Ok(Some(ReservedBlock {
            reservation,
            autonomous: false,
            positioned: false,
        }))
    }
    pub(super) fn install_nextval_reservation(
        &self,
        target: &NextvalTarget,
        block: ReservedBlock,
        caches: &mut BTreeMap<RelationIdentity, SessionSequenceCache>,
    ) -> Result<SequenceState, SequenceValueError> {
        let ReservedBlock {
            reservation,
            autonomous,
            positioned,
        } = block;
        let mut physical = target.state;
        physical.current = reservation.last_value;
        physical.called = true;
        physical.log_count = reservation.log_count;
        if !positioned {
            if let Some(state) = self
                .runtime
                .states_write()
                .get_mut(&target.relation)
                .filter(|state| state.definition_generation == target.state.definition_generation)
            {
                state.current = reservation.last_value;
                state.called = true;
                state.log_count = reservation.log_count;
            }
        }
        if reservation.count > 1 {
            let next_value = reservation
                .first_value
                .checked_add(target.state.increment)
                .ok_or_else(|| {
                    SequenceValueError::Internal(format!(
                        "cached sequence `{}` value overflow",
                        target.name
                    ))
                })?;
            caches.insert(
                target.relation.clone(),
                SessionSequenceCache {
                    object_id: target.object_id,
                    definition_generation: target.state.definition_generation,
                    next_value,
                    remaining: reservation.count - 1,
                    autonomous,
                },
            );
        }
        Ok(physical)
    }
    pub(super) fn complete_nextval(
        &self,
        relation: &RelationIdentity,
        object_id: [u8; 16],
        current: i64,
        physical: SequenceState,
        autonomous: bool,
    ) {
        let mut session = self.runtime.session_write();
        session
            .currvals_mut()
            .retain(|_, value| value.object_id != object_id);
        session.currvals_mut().insert(
            relation.clone(),
            SessionSequenceValue {
                object_id,
                value: current,
            },
        );
        *session.last_mut() = Some(SessionLastSequenceReference {
            relation: relation.clone(),
            object_id,
        });
        drop(session);
        self.runtime.record_nontransactional_sequence_value(
            physical.definition_generation,
            NontransactionalSequenceValue {
                object_id,
                current: physical.current,
                called: physical.called,
                log_count: physical.log_count,
                autonomous,
            },
            true,
        );
    }
}
pub(super) fn exhausted(name: &str, state: SequenceState) -> SequenceValueError {
    SequenceValueError::Exhausted {
        name: name.to_string(),
        bound: if state.increment > 0 {
            "maximum"
        } else {
            "minimum"
        },
        value: if state.increment > 0 {
            state.max_value
        } else {
            state.min_value
        },
    }
}

#[cfg(test)]
mod tests;