use std::net::SocketAddr;
use crate::control::security::identity::AuthenticatedIdentity;
use crate::control::state::SharedState;
use super::ddl_buffer;
use super::outcome::TxnDataPlane;
use super::overlay_drop::drop_txn_overlay;
use super::store::SessionStore;
pub fn run_begin(
sessions: &SessionStore,
addr: &SocketAddr,
state: &SharedState,
) -> Result<(), crate::Error> {
let snapshot_lsn = {
let next = state.wal.next_lsn();
crate::types::Lsn::new(next.as_u64().saturating_sub(1))
};
let snapshot_epoch = state
.last_applied_calvin_epoch
.load(std::sync::atomic::Ordering::Acquire);
ddl_buffer::activate();
sessions
.begin(addr, snapshot_lsn, snapshot_epoch)
.map_err(|msg| crate::Error::BadRequest {
detail: msg.to_owned(),
})
}
pub async fn run_rollback(
sessions: &SessionStore,
addr: &SocketAddr,
identity: &AuthenticatedIdentity,
state: &SharedState,
dp: &impl TxnDataPlane,
) {
ddl_buffer::discard();
let (overlay_txn_id, overlay_vshards) = sessions.txn_identity(addr);
super::reservation_release::release_session_reservations(
state,
sessions,
addr,
nodedb_cluster::calvin::types::ReleaseReason::Abort,
)
.await;
let reservations = sessions.rollback(addr).unwrap_or_default();
for handle in &reservations {
let key = &handle.sequence_key;
let registry = &state.sequence_registry;
registry.gap_free_manager().rollback(handle, || {
let map = registry.sequences_read();
if let Some(h) = map.get(key.as_str()) {
h.rollback_one();
}
});
{
let catalog = state.credentials.catalog();
crate::control::sequence::log::log_reservation(
catalog,
&crate::control::sequence::log::rolled_back(
key,
handle.value,
&identity.username,
identity.tenant_id.as_u64(),
),
);
}
}
sessions.close_non_hold_cursors(addr);
sessions.discard_pending_notifies(addr);
if let Some(txn_id) = overlay_txn_id {
for vshard_id in overlay_vshards {
if let Err(e) = drop_txn_overlay(state, dp, identity.tenant_id, vshard_id, txn_id).await
{
tracing::error!(
vshard = vshard_id.as_u32(),
error = %e,
"failed to release per-transaction staging overlay on rollback"
);
}
}
}
}