Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions crates/tracedecay-global-db/src/session_temporal_handle.rs
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,10 @@ impl SessionTemporalWriteTxn for RegisteredGlobalDbWriteTransaction<'_> {
fn commit(self) -> impl Future<Output = Result<(), EngineError>> + Send {
RegisteredGlobalDbWriteTransaction::commit(self)
}

fn rollback(self) -> impl Future<Output = Result<(), EngineError>> + Send {
RegisteredGlobalDbWriteTransaction::rollback(self)
}
}

impl SessionTemporalRegisteredDb for RegisteredGlobalDb {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,8 @@ use std::time::Duration;

use tracedecay_lcm::LcmError;
use tracedecay_store::{
SessionRefreshCompletionRequestV1, SessionRefreshFailureRequestV1, SessionRefreshFrontierV1,
SessionRefreshBeginOrJoinRequestV1, SessionRefreshCompletionRequestV1,
SessionRefreshDispositionV1, SessionRefreshFailureRequestV1, SessionRefreshFrontierV1,
SessionRefreshProgressV1, SessionRefreshStore, SessionStoreError,
};

Expand All @@ -27,8 +28,8 @@ use super::wake::{
use tracedecay_global_db::{RegisteredGlobalDb, RegisteredGlobalDbLeaseV1};
use tracedecay_runtime_core::db::engine::Error as EngineError;
use tracedecay_session_temporal_store::{
SessionRefreshRecoveryV1, SessionRefreshRestartStateV1, SessionTemporalAccess,
SessionTemporalStore,
SessionRefreshBeginBatchOutcomeV1, SessionRefreshBeginPlanV1, SessionRefreshRecoveryV1,
SessionRefreshRestartStateV1, SessionTemporalAccess, SessionTemporalStore,
};

const HISTORY_IDLE_RECHECK_INTERVAL: Duration = Duration::from_mins(1);
Expand Down Expand Up @@ -684,30 +685,58 @@ fn is_deterministic_refusal(error: &SessionStoreError) -> bool {
}
}

pub async fn process_refresh_begin_requests(
/// A queued begin request whose begin has been replayed and rolled back so
/// the projector can build the first batch before anything commits. The
/// operation becomes durable when the first batch's commit folds the begin's
/// rows in; until then the guard retains the request so a dropped pass
/// re-queues it for the durable-state discovery the request came from.
pub struct PreparedSessionRefresh<'a> {
request: SessionRefreshBeginOrJoinRequestV1,
recovery: SessionRefreshRecoveryV1,
pending: PendingBeginRequestGuard<'a>,
}

impl PreparedSessionRefresh<'_> {
fn disarm(&mut self) {
self.pending.disarm();
}
}

pub async fn process_refresh_begin_requests<'a>(
store: &SessionTemporalStore<'_, tracedecay_global_db::RegisteredGlobalDb>,
state: &SessionTemporalRefreshWakeState,
state: &'a SessionTemporalRefreshWakeState,
limit: usize,
report: &mut SessionTemporalRefreshPassReport,
) {
) -> Vec<PreparedSessionRefresh<'a>> {
let mut prepared = Vec::new();
for _ in 0..limit {
let Some(request) = state.take_requests(1).pop() else {
break;
};
let mut pending = PendingBeginRequestGuard::new(state, request);
if state.cancelled.load(Ordering::Acquire) {
return;
return prepared;
}
match store
.begin_or_join_session_refresh(pending.request().clone())
.plan_session_refresh_begin(pending.request().clone())
.await
{
Ok(receipt) => {
Ok(SessionRefreshBeginPlanV1::Prepared(recovery)) => {
prepared.push(PreparedSessionRefresh {
request: pending.request().clone(),
recovery: *recovery,
pending,
});
}
// A pending reset committed the begin instead of folding; the
// running operation joins this pass's durable recoveries.
Ok(SessionRefreshBeginPlanV1::Begun) => {
pending.disarm();
match receipt.disposition() {
tracedecay_store::SessionRefreshDispositionV1::Started => report.begun += 1,
tracedecay_store::SessionRefreshDispositionV1::Joined => report.joined += 1,
}
report.begun += 1;
}
Ok(SessionRefreshBeginPlanV1::Joined) => {
pending.disarm();
report.joined += 1;
}
Err(error) if is_retryable_storage(&error) => {
report.last_error = Some(format!("{error:?}"));
Expand All @@ -722,23 +751,24 @@ pub async fn process_refresh_begin_requests(
}
}
report.saturated |= state.has_requests();
prepared
}

#[tracing::instrument(
name = "daemon.scheduler.session_temporal.begin_admitted",
level = "trace",
skip_all
)]
pub async fn begin_admitted_session_refreshes(
pub async fn begin_admitted_session_refreshes<'a>(
database: &RegisteredGlobalDb,
store: &SessionTemporalStore<'_, tracedecay_global_db::RegisteredGlobalDb>,
state: &SessionTemporalRefreshWakeState,
state: &'a SessionTemporalRefreshWakeState,
limit: usize,
report: &mut SessionTemporalRefreshPassReport,
) {
) -> Vec<PreparedSessionRefresh<'a>> {
if state.has_requests() {
report.saturated = true;
return;
return Vec::new();
}
let cursor = state.projection_discovery_cursor();
let active_scan_slots = state.projection_discovery_active_slots(limit);
Expand All @@ -755,7 +785,7 @@ pub async fn begin_admitted_session_refreshes(
} else {
report.terminal_errors += 1;
}
return;
return Vec::new();
}
};
let (requests, next_cursor, has_more) = page.into_parts();
Expand All @@ -766,7 +796,7 @@ pub async fn begin_admitted_session_refreshes(
}
state.update_projection_discovery_cursor(next_cursor);
report.saturated |= has_more;
process_refresh_begin_requests(store, state, limit, report).await;
process_refresh_begin_requests(store, state, limit, report).await
}

async fn complete_ready_refresh(
Expand Down Expand Up @@ -951,6 +981,90 @@ pub async fn apply_refresh_effect(
}
}

/// Applies the projected effect for a begin that has not committed yet.
/// The first projected batch folds the begin's rows into one commit;
/// deferral, refusal, and durable failure keep the request queued or retire
/// it through the same paths a durable operation uses.
async fn apply_prepared_refresh_effect(
store: &SessionTemporalStore<'_, tracedecay_global_db::RegisteredGlobalDb>,
state: &SessionTemporalRefreshWakeState,
mut prepared: PreparedSessionRefresh<'_>,
effect: SessionTemporalRefreshEffect,
report: &mut SessionTemporalRefreshPassReport,
) {
match effect {
SessionTemporalRefreshEffect::Projection { progress, batch } => {
match store
.commit_session_refresh_begin_batch(
prepared.request.clone(),
prepared.recovery.accepted_at(),
progress,
batch,
state.completion_control(),
)
.await
{
Ok(SessionRefreshBeginBatchOutcomeV1::Persisted { .. }) => {
prepared.disarm();
report.begun += 1;
report.projected_batches += 1;
}
// The begin committed alone; durable recovery resumes the
// operation next pass.
Ok(SessionRefreshBeginBatchOutcomeV1::BeganOnly) => {
prepared.disarm();
report.begun += 1;
}
Ok(SessionRefreshBeginBatchOutcomeV1::Joined) => {
prepared.disarm();
report.joined += 1;
}
Err(error) if is_retryable_storage(&error) => {
report.last_error = Some(format!("{error:?}"));
report.retryable_errors += 1;
report.observe_retry(SessionTemporalRefreshRetryClass::Storage);
}
Err(error) => {
report.last_error = Some(format!("{error:?}"));
prepared.disarm();
report.terminal_errors += 1;
}
}
}
SessionTemporalRefreshEffect::Fail(request) => {
// The durable failure marker needs an operation row: begin through
// the committing path, then retire it like any durable refresh.
match store
.begin_or_join_session_refresh(prepared.request.clone())
.await
{
Ok(receipt) => {
prepared.disarm();
match receipt.disposition() {
SessionRefreshDispositionV1::Started => report.begun += 1,
SessionRefreshDispositionV1::Joined => report.joined += 1,
}
if receipt.operation_id() == prepared.recovery.operation_id() {
apply_fail_effect(store, state, &prepared.recovery, request, report).await;
} else {
report.terminal_errors += 1;
}
}
Err(error) if is_retryable_storage(&error) => {
report.last_error = Some(format!("{error:?}"));
report.retryable_errors += 1;
report.observe_retry(SessionTemporalRefreshRetryClass::Storage);
}
Err(_) => {
prepared.disarm();
report.terminal_errors += 1;
}
}
}
SessionTemporalRefreshEffect::Deferred => report.deferred += 1,
}
}

async fn apply_fail_effect(
store: &SessionTemporalStore<'_, tracedecay_global_db::RegisteredGlobalDb>,
state: &SessionTemporalRefreshWakeState,
Expand Down Expand Up @@ -987,15 +1101,24 @@ async fn apply_fail_effect(
}
}

/// One projection unit for a pass: the recovery the projector reads and,
/// when its begin has not committed yet, the prepared begin whose rows fold
/// into the first projected batch's commit.
struct ProjectionWorkItem<'a> {
recovery: SessionRefreshRecoveryV1,
prepared: Option<PreparedSessionRefresh<'a>>,
}

async fn project_running_refresh(
database: &RegisteredGlobalDbLeaseV1,
store: &SessionTemporalStore<'_, tracedecay_global_db::RegisteredGlobalDb>,
state: &SessionTemporalRefreshWakeState,
projector: &dyn SessionTemporalRefreshProjector,
policy: SessionTemporalRefreshPolicy,
recovery: &SessionRefreshRecoveryV1,
work: ProjectionWorkItem<'_>,
report: &mut SessionTemporalRefreshPassReport,
) {
let ProjectionWorkItem { recovery, prepared } = work;
let deadline_at = tokio::time::Instant::now() + policy.operation_deadline;
let projection = tracing::Instrument::instrument(
projector.project(database, recovery.clone()),
Expand Down Expand Up @@ -1027,7 +1150,7 @@ async fn project_running_refresh(
Err(error) => {
let failure_code = durable_projector_failure_code(&error.code);
report.last_error = Some(failure_code.clone());
let Some(request) = durable_failure_request(recovery, failure_code) else {
let Some(request) = durable_failure_request(&recovery, failure_code) else {
report.terminal_errors += 1;
return;
};
Expand All @@ -1043,7 +1166,14 @@ async fn project_running_refresh(
tokio::select! {
biased;
() = tracing::Instrument::instrument(state.wait_for_cancellation(), tracing::trace_span!("daemon.scheduler.session_temporal.effect_apply_cancel")) => {}
() = apply_refresh_effect(store, state, recovery, effect, report) => {}
() = async {
match prepared {
Some(prepared) => {
apply_prepared_refresh_effect(store, state, prepared, effect, report).await;
}
None => apply_refresh_effect(store, state, &recovery, effect, report).await,
}
} => {}
}
}

Expand All @@ -1066,20 +1196,24 @@ async fn running_refreshes(
}
}

async fn recoveries_for_pass(
async fn recoveries_for_pass<'a>(
database: &RegisteredGlobalDbLeaseV1,
store: &SessionTemporalStore<'_, tracedecay_global_db::RegisteredGlobalDb>,
state: &SessionTemporalRefreshWakeState,
state: &'a SessionTemporalRefreshWakeState,
policy: SessionTemporalRefreshPolicy,
report: &mut SessionTemporalRefreshPassReport,
) -> Option<(Vec<SessionRefreshRecoveryV1>, bool)> {
) -> Option<(
Vec<SessionRefreshRecoveryV1>,
Vec<PreparedSessionRefresh<'a>>,
bool,
)> {
let mut recoveries = running_refreshes(store, report).await?;
if !recoveries.is_empty() {
// Existing durable work owns this pass; discovery waits until these
// recoveries drain.
return Some((recoveries, true));
return Some((recoveries, Vec::new(), true));
}
begin_admitted_session_refreshes(
let prepared = begin_admitted_session_refreshes(
database,
store,
state,
Expand All @@ -1088,7 +1222,7 @@ async fn recoveries_for_pass(
)
.await;
recoveries = running_refreshes(store, report).await?;
Some((recoveries, false))
Some((recoveries, prepared, false))
}

fn recovery_key(recovery: &SessionRefreshRecoveryV1) -> String {
Expand All @@ -1110,18 +1244,38 @@ pub async fn run_session_temporal_refresh_pass(
if state.cancelled.load(Ordering::Acquire) {
return report;
}
process_refresh_begin_requests(
let mut prepared = process_refresh_begin_requests(
&store,
state,
policy.max_begin_requests_per_pass,
&mut report,
)
.await;
let Some((mut recoveries, discovery_deferred)) =
let Some((mut recoveries, admitted, discovery_deferred)) =
recoveries_for_pass(database, &store, state, policy, &mut report).await
else {
return report;
};
prepared.extend(admitted);
// Planned begins are not durable: a durable operation holding the same
// key owns the request already, and a second plan of an equivalent
// request attaches to the first plan.
let durable_keys = recoveries.iter().map(recovery_key).collect::<HashSet<_>>();
let mut prepared_by_key = HashMap::new();
for mut prepared_begin in prepared {
let key = recovery_key(&prepared_begin.recovery);
if durable_keys.contains(&key) || prepared_by_key.contains_key(&key) {
prepared_begin.disarm();
report.joined += 1;
continue;
}
prepared_by_key.insert(key, prepared_begin);
}
recoveries.extend(
prepared_by_key
.values()
.map(|prepared_begin| prepared_begin.recovery.clone()),
);
recoveries.sort_by_cached_key(recovery_key);
state.observe_durable_backlog(recoveries.len());
let ordered_keys = recoveries.iter().map(recovery_key).collect::<Vec<_>>();
Expand Down Expand Up @@ -1184,7 +1338,10 @@ pub async fn run_session_temporal_refresh_pass(
state,
projector,
policy,
&recovery,
ProjectionWorkItem {
recovery,
prepared: prepared_by_key.remove(&operation),
},
&mut report,
)
.await;
Expand Down
Loading
Loading