diff --git a/crates/tracedecay-global-db/src/session_temporal_handle.rs b/crates/tracedecay-global-db/src/session_temporal_handle.rs index b1e1fd0d96..a825119e83 100644 --- a/crates/tracedecay-global-db/src/session_temporal_handle.rs +++ b/crates/tracedecay-global-db/src/session_temporal_handle.rs @@ -53,6 +53,10 @@ impl SessionTemporalWriteTxn for RegisteredGlobalDbWriteTransaction<'_> { fn commit(self) -> impl Future> + Send { RegisteredGlobalDbWriteTransaction::commit(self) } + + fn rollback(self) -> impl Future> + Send { + RegisteredGlobalDbWriteTransaction::rollback(self) + } } impl SessionTemporalRegisteredDb for RegisteredGlobalDb { diff --git a/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs b/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs index 3fc874b600..f8f8070898 100644 --- a/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs +++ b/crates/tracedecay-session-runtime/src/session_temporal_refresh_scheduler/worker.rs @@ -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, }; @@ -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); @@ -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> { + 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:?}")); @@ -722,6 +751,7 @@ pub async fn process_refresh_begin_requests( } } report.saturated |= state.has_requests(); + prepared } #[tracing::instrument( @@ -729,16 +759,16 @@ pub async fn process_refresh_begin_requests( 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> { 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); @@ -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(); @@ -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( @@ -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, @@ -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>, +} + 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()), @@ -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; }; @@ -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, + } + } => {} } } @@ -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, bool)> { +) -> Option<( + Vec, + Vec>, + 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, @@ -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 { @@ -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::>(); + 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::>(); @@ -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; diff --git a/crates/tracedecay-session-runtime/tests/session_store_read_cost.rs b/crates/tracedecay-session-runtime/tests/session_store_read_cost.rs index 6d4e257963..1ffb8280bf 100644 --- a/crates/tracedecay-session-runtime/tests/session_store_read_cost.rs +++ b/crates/tracedecay-session-runtime/tests/session_store_read_cost.rs @@ -922,9 +922,9 @@ fn wal_commits(store: &Path, from: WalMark, to: WalMark) -> Option { /// receipt, and the message's Git evidence span before the host is /// acknowledged. The drain commits the external-source replay, the /// observation projection, and the Git evidence convergence. The temporal -/// refresh commits its operation, the projected batch, the pending relation -/// receipt the native graph write recovers from, and the activation that -/// settles that receipt. +/// refresh commits its projected batch: the operation's begin folds into +/// that same commit, the pending relation receipt the native graph write +/// recovers from, and the activation that settles that receipt. #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn streamed_message_commits_once_per_durability_boundary() { const MESSAGES: u64 = 8; @@ -970,7 +970,7 @@ async fn streamed_message_commits_once_per_durability_boundary() { "most messages must land within one log generation: {measured:?}" ); assert!( - measured.iter().all(|commits| *commits == (4, 3, 4)), + measured.iter().all(|commits| *commits == (4, 3, 3)), "one streamed message must commit once per durability boundary: {measured:?}" ); } diff --git a/crates/tracedecay-session-temporal-store/src/handle.rs b/crates/tracedecay-session-temporal-store/src/handle.rs index 1ffce50574..c3bdf4d346 100644 --- a/crates/tracedecay-session-temporal-store/src/handle.rs +++ b/crates/tracedecay-session-temporal-store/src/handle.rs @@ -45,6 +45,8 @@ pub trait SessionTemporalExec: SessionTemporalQuery { /// Write transaction the session-temporal store can query, mutate, and commit. pub trait SessionTemporalWriteTxn: SessionTemporalExec { fn commit(self) -> impl Future> + Send; + + fn rollback(self) -> impl Future> + Send; } impl SessionTemporalQuery for DatabaseEngineReadSnapshot { diff --git a/crates/tracedecay-session-temporal-store/src/lib.rs b/crates/tracedecay-session-temporal-store/src/lib.rs index 2bb767502b..5d8002b63c 100644 --- a/crates/tracedecay-session-temporal-store/src/lib.rs +++ b/crates/tracedecay-session-temporal-store/src/lib.rs @@ -106,7 +106,10 @@ pub use doctor_health::{ SessionTemporalHealthFindingKind, SessionTemporalHealthReport, SessionTemporalHealthStatus, }; pub use projection::{record_canonical_observation_effect, request_session_temporal_reset}; -pub use refresh::{SessionRefreshRecoveryV1, SessionRefreshRestartStateV1}; +pub use refresh::{ + SessionRefreshBeginBatchOutcomeV1, SessionRefreshBeginPlanV1, SessionRefreshRecoveryV1, + SessionRefreshRestartStateV1, +}; pub use store::SessionTemporalStore; impl SessionTemporalAccess<'_, D> { diff --git a/crates/tracedecay-session-temporal-store/src/refresh.rs b/crates/tracedecay-session-temporal-store/src/refresh.rs index 0c6e85a51a..56569d1405 100644 --- a/crates/tracedecay-session-temporal-store/src/refresh.rs +++ b/crates/tracedecay-session-temporal-store/src/refresh.rs @@ -34,10 +34,7 @@ use super::rebuild::{ validate_candidate_frontier, }; use super::relation_receipts::acknowledge_relation_receipt; -use crate::handle::{ - SessionTemporalAccess, SessionTemporalExec, SessionTemporalRegisteredDb, - SessionTemporalWriteTxn, -}; +use crate::handle::{SessionTemporalAccess, SessionTemporalRegisteredDb, SessionTemporalWriteTxn}; const BEGIN_REFRESH: &str = "begin or join session refresh"; const PERSIST_REFRESH: &str = "persist session refresh progress"; @@ -72,6 +69,7 @@ pub struct SessionRefreshRecoveryV1 { binding_digest: String, progress: Option, restart_state: SessionRefreshRestartStateV1, + accepted_at: UtcMicros, } impl SessionRefreshRecoveryV1 { @@ -131,6 +129,13 @@ impl SessionRefreshRecoveryV1 { self.restart_state } + /// When the begin accepted (or will accept) the operation: the + /// `session_refresh_operations.created_at` the committing replay must + /// reuse so rows written between plan and commit stay ordered. + pub const fn accepted_at(&self) -> UtcMicros { + self.accepted_at + } + pub fn source_coverage( &self, committed_through: u64, @@ -156,166 +161,544 @@ impl SessionRefreshRecoveryV1 { } } +/// What a refresh begin would commit, computed by replaying it inside a +/// write transaction that is then rolled back. +#[derive(Debug)] +pub enum SessionRefreshBeginPlanV1 { + /// The begin starts a fresh running operation; the prepared recovery + /// describes the rows the begin will write so the projector can build + /// the first batch before anything commits. + Prepared(Box), + /// The begin committed durably instead of folding: a pending reset + /// deletes the base rows the first batch would project from, so the + /// operation resumes through durable recovery, which reads the + /// post-reset state, in the same pass. + Begun, + /// An equivalent operation already exists; the request attaches to it. + Joined, +} + +/// Outcome of committing a refresh begin and its first projected batch in +/// one write transaction. +#[derive(Debug)] +pub enum SessionRefreshBeginBatchOutcomeV1 { + /// The begin and the batch committed together. + Persisted { + progress: Box, + receipt: Box, + }, + /// The begin committed alone because the replayed begin no longer + /// matched the projected batch, or the batch was refused; the durable + /// running operation resumes from recovery next pass. + BeganOnly, + /// An equivalent operation already exists; nothing new began. + Joined, +} + +/// The begin's durable decision inside an open write transaction: it either +/// wrote a fresh running operation or attached to an equivalent existing +/// one. The caller chooses whether the transaction commits. +enum SessionRefreshBeginTxnOutcome { + Started { + operation_id: SessionRefreshOperationIdV1, + target_frontier: SessionRefreshFrontierV1, + accepted_at: UtcMicros, + binding: RefreshBinding, + }, + Joined { + operation_id: SessionRefreshOperationIdV1, + target_frontier: SessionRefreshFrontierV1, + created_at: UtcMicros, + }, +} + +/// Runs every begin-or-join decision and write inside `transaction` without +/// committing it, so callers can commit the begin alone, roll it back as a +/// plan, or fold it into the first projected batch's commit. +async fn begin_session_refresh_in_transaction( + transaction: &impl SessionTemporalWriteTxn, + request: SessionRefreshBeginOrJoinRequestV1, + accepted_at: UtcMicros, +) -> SessionStoreResult { + apply_requested_reset(transaction, request.session_id()).await?; + let request = match read_active_generation(transaction, request.session_id()).await? { + Some((_, active_watermarks)) => { + rebase_on_committed_frontier(request, active_watermarks.projection_frontier())? + } + None => request, + }; + let reset_generation = read_reset_generation(transaction, request.session_id()).await?; + let request_digest = refresh_binding_digest(&request, reset_generation)?; + + if let Some(existing) = + read_joinable_operation_by_digest(transaction, request.session_id(), &request_digest) + .await? + { + return Ok(SessionRefreshBeginTxnOutcome::Joined { + operation_id: existing.operation_id, + target_frontier: request.target_frontier(), + created_at: existing.created_at, + }); + } + if read_running_operation(transaction, request.session_id()) + .await? + .is_some() + { + return Err(SessionStoreError::IdempotencyConflict { + context: "session refresh busy", + }); + } + require_no_open_candidate(transaction, request.session_id(), None, BEGIN_REFRESH).await?; + + let provisioned_cursor_key = + ensure_active_session_cursor_key_in_transaction(transaction).await?; + let (active_generation, active_watermarks) = + ensure_active_generation(transaction, &request).await?; + let candidate_generation = next_generation(transaction, request.session_id()).await?; + let mut frozen_watermarks = SessionFrozenWatermarksV1::new( + active_generation, + request.target_frontier().observed_through(), + request.target_frontier().observed_through(), + active_watermarks.summary_frontier(), + ); + if let Some(cursor_key) = active_watermarks.cursor_key() { + frozen_watermarks = frozen_watermarks.with_cursor_key(cursor_key.clone()); + } else { + frozen_watermarks = frozen_watermarks.with_cursor_key(provisioned_cursor_key); + } + let frozen_watermarks_json = encode_watermarks(&frozen_watermarks, BEGIN_REFRESH)?; + let attempt = + next_operation_attempt(transaction, request.session_id(), &request_digest).await?; + let operation_id = operation_id_for_digest(&request_digest, attempt)?; + + // The running-operation and attempt reads above share this writer transaction, so + // they already decided one-running ownership and the operation id; a constraint + // failure here is a storage fault, never a busy refresh. + transaction + .execute( + "INSERT INTO session_refresh_operations ( + session_id, operation_id, request_digest, target_frontier_json, + state, created_at, updated_at + ) VALUES (?1, ?2, ?3, ?4, 'running', ?5, ?5)", + params![ + request.session_id().as_str(), + operation_id.as_str(), + request_digest.as_str(), + encode_refresh_target(&request)?, + accepted_at.0, + ], + ) + .await + .map_err(|error| storage(BEGIN_REFRESH, error))?; + transaction + .execute( + "INSERT INTO session_temporal_generations ( + session_id, generation, state, frozen_watermarks_json, created_at + ) VALUES (?1, ?2, 'building', ?3, ?4)", + params![ + request.session_id().as_str(), + generation_i64(candidate_generation, BEGIN_REFRESH)?, + frozen_watermarks_json.as_str(), + accepted_at.0, + ], + ) + .await + .map_err(|error| storage(BEGIN_REFRESH, error))?; + // Summary availability is generation-bound: the candidate inherits the + // active generation's rows exactly like the summary-publication + // route's generation builder, otherwise activating this refresh would + // silently drop every published summary from generation-bound reads. + transaction + .execute( + "INSERT INTO session_summary_availability ( + session_id, generation, summary_id, availability, + source_horizon_json, reason, checked_at + ) + SELECT session_id, ?2, summary_id, availability, + source_horizon_json, reason, ?3 + FROM session_summary_availability + WHERE session_id = ?1 AND generation = ?4", + params![ + request.session_id().as_str(), + generation_i64(candidate_generation, BEGIN_REFRESH)?, + accepted_at.0, + generation_i64(active_generation, BEGIN_REFRESH)?, + ], + ) + .await + .map_err(|error| storage(BEGIN_REFRESH, error))?; + transaction + .execute( + "INSERT INTO session_refresh_bindings ( + session_id, operation_id, scope_kind, source_frontier, target_frontier, + projector_version, config_digest, generation, frozen_watermarks_json, + binding_digest, created_at + ) VALUES (?1, ?2, 'session_store', ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)", + params![ + request.session_id().as_str(), + operation_id.as_str(), + frontier_i64(request.target_frontier().committed_through(), BEGIN_REFRESH,)?, + frontier_i64(request.target_frontier().observed_through(), BEGIN_REFRESH)?, + PROJECTOR_VERSION, + config_digest(), + generation_i64(candidate_generation, BEGIN_REFRESH)?, + frozen_watermarks_json, + request_digest.clone(), + accepted_at.0, + ], + ) + .await + .map_err(|error| storage(BEGIN_REFRESH, error))?; + Ok(SessionRefreshBeginTxnOutcome::Started { + operation_id, + target_frontier: request.target_frontier(), + accepted_at, + binding: RefreshBinding { + generation: candidate_generation, + source_frontier: request.target_frontier().committed_through(), + target_frontier: request.target_frontier().observed_through(), + watermarks: frozen_watermarks, + projector_version: PROJECTOR_VERSION.to_owned(), + config_digest: config_digest(), + binding_digest: request_digest, + }, + }) +} + +/// Resolves the single-source refresh target the recovery decode derives +/// for a begin request without a refresh key: the session's lowest provider +/// name, or `all` when none is registered. +async fn default_refresh_source_targets( + conn: &impl crate::handle::SessionTemporalQuery, + session_id: &SessionId, + target_frontier: SessionRefreshFrontierV1, +) -> SessionStoreResult> { + let mut rows = conn + .query( + "SELECT COALESCE(( + SELECT MIN(source.provider) + FROM sessions AS source + WHERE source.session_id = ?1 + ), 'all')", + params![session_id.as_str()], + ) + .await + .map_err(|error| storage(BEGIN_REFRESH, error))?; + let provider: String = rows + .next() + .await + .map_err(|error| storage(BEGIN_REFRESH, error))? + .ok_or_else(|| storage_message(BEGIN_REFRESH, "refresh source provider returned no row"))? + .get(0) + .map_err(|error| storage(BEGIN_REFRESH, error))?; + Ok(vec![ + SessionRefreshSourceTargetV1::new( + SessionSourceIdV1::new(format!("{}:{provider}", session_id.as_str())) + .map_err(SessionStoreError::from)?, + SessionSourceFrontierV1::new(target_frontier.observed_through()), + SessionSourceFrontierV1::new(target_frontier.observed_through()), + ) + .map_err(SessionStoreError::from)?, + ]) +} + impl SessionTemporalAccess<'_, D> { #[tracing::instrument(name = "session_temporal.txn.begin_refresh", level = "trace", skip_all)] pub async fn begin_or_join_session_refresh_result( &self, request: SessionRefreshBeginOrJoinRequestV1, ) -> SessionStoreResult { + let session_id = request.session_id().clone(); let transaction = self .begin_write_transaction() .instrument(tracing::trace_span!("session_temporal.txn.begin")) .await .map_err(|error| storage(BEGIN_REFRESH, error))?; - apply_requested_reset(&transaction, request.session_id()).await?; - let request = match read_active_generation(&transaction, request.session_id()).await? { - Some((_, active_watermarks)) => { - rebase_on_committed_frontier(request, active_watermarks.projection_frontier())? - } - None => request, + let receipt = match begin_session_refresh_in_transaction( + &transaction, + request, + now_micros(BEGIN_REFRESH)?, + ) + .await? + { + SessionRefreshBeginTxnOutcome::Started { + operation_id, + target_frontier, + accepted_at, + .. + } => SessionRefreshBeginOrJoinReceiptV1::new( + operation_id, + session_id, + target_frontier, + SessionRefreshDispositionV1::Started, + accepted_at, + ), + SessionRefreshBeginTxnOutcome::Joined { + operation_id, + target_frontier, + created_at, + } => SessionRefreshBeginOrJoinReceiptV1::new( + operation_id, + session_id, + target_frontier, + SessionRefreshDispositionV1::Joined, + created_at, + ), }; - let reset_generation = read_reset_generation(&transaction, request.session_id()).await?; - let request_digest = refresh_binding_digest(&request, reset_generation)?; + transaction + .commit() + .instrument(tracing::trace_span!("session_temporal.txn.commit")) + .await + .map_err(|error| storage(BEGIN_REFRESH, error))?; + Ok(receipt) + } - if let Some(existing) = - read_joinable_operation_by_digest(&transaction, request.session_id(), &request_digest) - .await? - { + /// Replays a refresh begin inside a write transaction that is rolled + /// back, so the worker can project against the recovery the begin would + /// commit and fold both into one write with + /// [`Self::commit_session_refresh_begin_batch_result`]. + /// + /// Rollback keeps every begin-time decision identical to the committing + /// path: reset application, frontier rebase, join detection, cursor-key + /// provisioning, and generation allocation, while leaving nothing durable. + /// Pending refresh work is derived from `session_temporal_observation_effects`, + /// so a plan that never commits loses no work: discovery re-queues an + /// equivalent request next pass. + #[tracing::instrument(name = "session_temporal.txn.plan_refresh", level = "trace", skip_all)] + pub async fn plan_session_refresh_begin_result( + &self, + request: SessionRefreshBeginOrJoinRequestV1, + ) -> SessionStoreResult { + let session_id = request.session_id().clone(); + let coverage_request = request.coverage_request().clone(); + let refresh_key = request.refresh_key().cloned(); + let transaction = self + .begin_write_transaction() + .instrument(tracing::trace_span!("session_temporal.txn.begin")) + .await + .map_err(|error| storage(BEGIN_REFRESH, error))?; + if session_reset_is_pending(&transaction, &session_id).await? { + // The begin deletes the base rows a folded first batch would + // project from, so planning here can only produce a batch the + // committing replay refuses. Commit the begin instead; durable + // recovery picks the operation up this pass and projects the + // post-reset state. + let outcome = begin_session_refresh_in_transaction( + &transaction, + request, + now_micros(BEGIN_REFRESH)?, + ) + .await; + let outcome = match outcome { + Ok(outcome) => outcome, + Err(error) => { + transaction + .rollback() + .instrument(tracing::trace_span!("session_temporal.txn.rollback")) + .await + .map_err(|rollback| storage(BEGIN_REFRESH, rollback))?; + return Err(error); + } + }; transaction .commit() .instrument(tracing::trace_span!("session_temporal.txn.commit")) .await .map_err(|error| storage(BEGIN_REFRESH, error))?; - return Ok(SessionRefreshBeginOrJoinReceiptV1::new( - existing.operation_id, - request.session_id().clone(), - request.target_frontier(), - SessionRefreshDispositionV1::Joined, - existing.created_at, - )); - } - if read_running_operation(&transaction, request.session_id()) - .await? - .is_some() - { - return Err(SessionStoreError::IdempotencyConflict { - context: "session refresh busy", + return Ok(match outcome { + SessionRefreshBeginTxnOutcome::Started { .. } => SessionRefreshBeginPlanV1::Begun, + SessionRefreshBeginTxnOutcome::Joined { .. } => SessionRefreshBeginPlanV1::Joined, }); } - require_no_open_candidate(&transaction, request.session_id(), None, BEGIN_REFRESH).await?; - - let provisioned_cursor_key = - ensure_active_session_cursor_key_in_transaction(&transaction).await?; - let (active_generation, active_watermarks) = - ensure_active_generation(&transaction, &request).await?; - let candidate_generation = next_generation(&transaction, request.session_id()).await?; - let mut frozen_watermarks = SessionFrozenWatermarksV1::new( - active_generation, - request.target_frontier().observed_through(), - request.target_frontier().observed_through(), - active_watermarks.summary_frontier(), - ); - if let Some(cursor_key) = active_watermarks.cursor_key() { - frozen_watermarks = frozen_watermarks.with_cursor_key(cursor_key.clone()); - } else { - frozen_watermarks = frozen_watermarks.with_cursor_key(provisioned_cursor_key); - } - let frozen_watermarks_json = encode_watermarks(&frozen_watermarks, BEGIN_REFRESH)?; - let accepted_at = now_micros(BEGIN_REFRESH)?; - let attempt = - next_operation_attempt(&transaction, request.session_id(), &request_digest).await?; - let operation_id = operation_id_for_digest(&request_digest, attempt)?; - - // The running-operation and attempt reads above share this writer transaction, so - // they already decided one-running ownership and the operation id; a constraint - // failure here is a storage fault, never a busy refresh. - transaction - .execute( - "INSERT INTO session_refresh_operations ( - session_id, operation_id, request_digest, target_frontier_json, - state, created_at, updated_at - ) VALUES (?1, ?2, ?3, ?4, 'running', ?5, ?5)", - params![ - request.session_id().as_str(), - operation_id.as_str(), - request_digest.as_str(), - encode_refresh_target(&request)?, - accepted_at.0, - ], - ) - .await - .map_err(|error| storage(BEGIN_REFRESH, error))?; + // The begin mints a random cursor key when no active key exists, and + // a rolled-back mint leaves the commit replay reading a different key. + // Provisioning the shared key first makes both replays deterministic; + // with an active key already committed this transaction writes + // nothing and its commit appends no WAL frames. + ensure_active_session_cursor_key_in_transaction(&transaction).await?; transaction - .execute( - "INSERT INTO session_temporal_generations ( - session_id, generation, state, frozen_watermarks_json, created_at - ) VALUES (?1, ?2, 'building', ?3, ?4)", - params![ - request.session_id().as_str(), - generation_i64(candidate_generation, BEGIN_REFRESH)?, - frozen_watermarks_json.as_str(), - accepted_at.0, - ], - ) + .commit() + .instrument(tracing::trace_span!("session_temporal.txn.commit")) .await .map_err(|error| storage(BEGIN_REFRESH, error))?; - // Summary availability is generation-bound: the candidate inherits the - // active generation's rows exactly like the summary-publication - // route's generation builder, otherwise activating this refresh would - // silently drop every published summary from generation-bound reads. - transaction - .execute( - "INSERT INTO session_summary_availability ( - session_id, generation, summary_id, availability, - source_horizon_json, reason, checked_at - ) - SELECT session_id, ?2, summary_id, availability, - source_horizon_json, reason, ?3 - FROM session_summary_availability - WHERE session_id = ?1 AND generation = ?4", - params![ - request.session_id().as_str(), - generation_i64(candidate_generation, BEGIN_REFRESH)?, - accepted_at.0, - generation_i64(active_generation, BEGIN_REFRESH)?, - ], - ) + let transaction = self + .begin_write_transaction() + .instrument(tracing::trace_span!("session_temporal.txn.begin")) .await .map_err(|error| storage(BEGIN_REFRESH, error))?; + let outcome = + begin_session_refresh_in_transaction(&transaction, request, now_micros(BEGIN_REFRESH)?) + .await; + let plan = match outcome? { + SessionRefreshBeginTxnOutcome::Joined { .. } => SessionRefreshBeginPlanV1::Joined, + SessionRefreshBeginTxnOutcome::Started { + operation_id, + target_frontier, + accepted_at, + binding, + } => { + let source_targets = match &refresh_key { + Some(refresh_key) => refresh_key.sources().to_vec(), + None => { + default_refresh_source_targets(&transaction, &session_id, target_frontier) + .await? + } + }; + SessionRefreshBeginPlanV1::Prepared(Box::new(SessionRefreshRecoveryV1 { + operation_id, + session_id, + source_targets, + coverage_request, + source_frontier: binding.source_frontier, + target_frontier, + candidate_generation: binding.generation, + frozen_watermarks: binding.watermarks, + projector_version: binding.projector_version, + config_digest: binding.config_digest, + binding_digest: binding.binding_digest, + progress: None, + restart_state: SessionRefreshRestartStateV1::BeginProjection, + accepted_at, + })) + } + }; transaction - .execute( - "INSERT INTO session_refresh_bindings ( - session_id, operation_id, scope_kind, source_frontier, target_frontier, - projector_version, config_digest, generation, frozen_watermarks_json, - binding_digest, created_at - ) VALUES (?1, ?2, 'session_store', ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)", - params![ - request.session_id().as_str(), - operation_id.as_str(), - frontier_i64(request.target_frontier().committed_through(), BEGIN_REFRESH,)?, - frontier_i64(request.target_frontier().observed_through(), BEGIN_REFRESH)?, - PROJECTOR_VERSION, - config_digest(), - generation_i64(candidate_generation, BEGIN_REFRESH)?, - frozen_watermarks_json, - request_digest, - accepted_at.0, - ], - ) + .rollback() + .instrument(tracing::trace_span!("session_temporal.txn.rollback")) .await .map_err(|error| storage(BEGIN_REFRESH, error))?; - transaction - .commit() - .instrument(tracing::trace_span!("session_temporal.txn.commit")) + Ok(plan) + } + + /// Commits a refresh begin and its first projected batch in one write + /// transaction. The batch was projected against a + /// [`SessionRefreshBeginPlanV1::Prepared`] recovery, so the begin is + /// replayed here and the batch persists only when the replayed binding + /// matches it; when the durable state moved between plan and commit the + /// begin commits alone and the running operation resumes from durable + /// recovery next pass, the same state a crash between the two commits + /// left behind. + #[tracing::instrument( + name = "session_temporal.txn.begin_refresh_batch", + level = "trace", + skip_all + )] + pub async fn commit_session_refresh_begin_batch_result( + &self, + request: SessionRefreshBeginOrJoinRequestV1, + accepted_at: UtcMicros, + progress: SessionRefreshProgressV1, + batch: SessionTemporalProjectionBatchV1, + execution_control: ExecutionControl, + ) -> SessionStoreResult { + validate_progress_batch_identity(&progress, &batch)?; + let authoritative_validation_time = now_micros(PERSIST_REFRESH)?; + let transaction = self + .begin_write_transaction() + .instrument(tracing::trace_span!("session_temporal.txn.begin")) .await - .map_err(|error| storage(BEGIN_REFRESH, error))?; - Ok(SessionRefreshBeginOrJoinReceiptV1::new( + .map_err(|error| storage(PERSIST_REFRESH, error))?; + let SessionRefreshBeginTxnOutcome::Started { operation_id, - request.session_id().clone(), - request.target_frontier(), - SessionRefreshDispositionV1::Started, - accepted_at, - )) + binding, + .. + } = begin_session_refresh_in_transaction(&transaction, request.clone(), accepted_at) + .await? + else { + transaction + .commit() + .instrument(tracing::trace_span!("session_temporal.txn.commit")) + .await + .map_err(|error| storage(BEGIN_REFRESH, error))?; + return Ok(SessionRefreshBeginBatchOutcomeV1::Joined); + }; + // The projected batch was built against the planned recovery, so it + // persists only when the begin this transaction replayed agrees with + // it. Batch writes that a later check refuses roll back with the + // transaction; the begin is then replayed alone so a refused batch + // still leaves a durable running operation the next pass retires, + // the same state a crash between the two commits left behind. + let diverged = operation_id != *progress.operation_id() + || validate_batch_binding(&binding, &batch).is_err() + || validate_progress_binding(&binding, &progress).is_err(); + if diverged { + transaction + .commit() + .instrument(tracing::trace_span!("session_temporal.txn.commit")) + .await + .map_err(|error| storage(PERSIST_REFRESH, error))?; + return Ok(SessionRefreshBeginBatchOutcomeV1::BeganOnly); + } + let persisted = async { + require_progress_timestamp(&progress, authoritative_validation_time)?; + let receipt = persist_session_temporal_projection_batch_in_transaction( + &transaction, + &batch, + &execution_control, + ) + .await?; + validate_next_progress( + &transaction, + &progress, + batch.generation(), + batch.batch_ordinal(), + batch.item_count(), + ) + .await?; + validate_progress_binding(&binding, &progress)?; + insert_progress_and_binding(&transaction, &progress, &batch).await?; + touch_running_operation( + &transaction, + progress.session_id(), + progress.operation_id(), + progress.updated_at(), + ) + .await?; + checkpoint_relation_rebuild_control(&execution_control)?; + Ok::<_, SessionStoreError>(receipt) + } + .await; + match persisted { + Ok(receipt) => { + transaction + .commit() + .instrument(tracing::trace_span!("session_temporal.txn.commit")) + .await + .map_err(|error| storage(PERSIST_REFRESH, error))?; + Ok(SessionRefreshBeginBatchOutcomeV1::Persisted { + progress: Box::new(progress), + receipt: Box::new(receipt), + }) + } + Err(_) => { + transaction + .rollback() + .instrument(tracing::trace_span!("session_temporal.txn.rollback")) + .await + .map_err(|error| storage(PERSIST_REFRESH, error))?; + let transaction = self + .begin_write_transaction() + .instrument(tracing::trace_span!("session_temporal.txn.begin")) + .await + .map_err(|error| storage(BEGIN_REFRESH, error))?; + let outcome = + begin_session_refresh_in_transaction(&transaction, request, accepted_at) + .await?; + transaction + .commit() + .instrument(tracing::trace_span!("session_temporal.txn.commit")) + .await + .map_err(|error| storage(BEGIN_REFRESH, error))?; + Ok(match outcome { + SessionRefreshBeginTxnOutcome::Started { .. } => { + SessionRefreshBeginBatchOutcomeV1::BeganOnly + } + SessionRefreshBeginTxnOutcome::Joined { .. } => { + SessionRefreshBeginBatchOutcomeV1::Joined + } + }) + } + } } pub async fn persist_session_refresh_projection_batch_result( @@ -1303,10 +1686,13 @@ async fn ensure_active_generation( /// activates in its place, so the refresh begun next projects every live /// effect against an empty base. Receipts at or below the reset generation /// describe the deleted rows, so base frontiers and coverage skip them. -async fn apply_requested_reset( - conn: &impl crate::handle::SessionTemporalExec, +/// Whether the session has a requested reset waiting for a begin to apply +/// it. [`apply_requested_reset`] performs the same check; the plan path reads +/// it first because a reset changes what the first batch must project from. +async fn session_reset_is_pending( + conn: &impl crate::handle::SessionTemporalQuery, session_id: &SessionId, -) -> SessionStoreResult<()> { +) -> SessionStoreResult { let mut rows = conn .query( "SELECT 1 FROM session_temporal_resets @@ -1321,6 +1707,14 @@ async fn apply_requested_reset( .map_err(|error| storage(BEGIN_REFRESH, error))? .is_some(); drop(rows); + Ok(requested) +} + +async fn apply_requested_reset( + conn: &impl crate::handle::SessionTemporalExec, + session_id: &SessionId, +) -> SessionStoreResult<()> { + let requested = session_reset_is_pending(conn, session_id).await?; if !requested || read_running_operation(conn, session_id).await?.is_some() { return Ok(()); } @@ -2314,7 +2708,8 @@ async fn read_running_recoveries( operation.target_frontier_json, binding.generation, binding.source_frontier, binding.target_frontier, binding.frozen_watermarks_json, binding.projector_version, - binding.config_digest, binding.binding_digest + binding.config_digest, binding.binding_digest, + operation.created_at FROM session_refresh_operations AS operation JOIN session_refresh_bindings AS binding ON binding.session_id = operation.session_id @@ -2388,6 +2783,7 @@ async fn read_running_recoveries( config_digest: row.get(9).map_err(|error| storage(READ_REFRESH, error))?, binding_digest: row.get(10).map_err(|error| storage(READ_REFRESH, error))?, }; + let accepted_at = UtcMicros(row.get(11).map_err(|error| storage(READ_REFRESH, error))?); pending.push(( operation_id, session_id, @@ -2395,13 +2791,21 @@ async fn read_running_recoveries( coverage_request, target_frontier, binding, + accepted_at, )); } drop(rows); let mut recoveries = Vec::with_capacity(pending.len()); - for (operation_id, session_id, source_targets, coverage_request, target_frontier, binding) in - pending + for ( + operation_id, + session_id, + source_targets, + coverage_request, + target_frontier, + binding, + accepted_at, + ) in pending { let progress = read_progress(conn, &session_id, &operation_id).await?; let restart_state = match progress.as_ref() { @@ -2437,6 +2841,7 @@ async fn read_running_recoveries( binding_digest: binding.binding_digest, progress, restart_state, + accepted_at, }); } Ok(recoveries) diff --git a/crates/tracedecay-session-temporal-store/src/store.rs b/crates/tracedecay-session-temporal-store/src/store.rs index 929535b589..edf058b404 100644 --- a/crates/tracedecay-session-temporal-store/src/store.rs +++ b/crates/tracedecay-session-temporal-store/src/store.rs @@ -96,6 +96,39 @@ impl<'a, D: SessionTemporalRegisteredDb + Sync> SessionTemporalStore<'a, D> { .await } + /// Plans a refresh begin without committing it so the worker can fold + /// the begin into the first projected batch's commit. + pub async fn plan_session_refresh_begin( + &self, + request: SessionRefreshBeginOrJoinRequestV1, + ) -> SessionStoreResult { + self.access() + .plan_session_refresh_begin_result(request) + .await + } + + /// Commits a refresh begin and its first projected batch in one write + /// transaction; see + /// [`SessionTemporalAccess::commit_session_refresh_begin_batch_result`]. + pub async fn commit_session_refresh_begin_batch( + &self, + request: SessionRefreshBeginOrJoinRequestV1, + accepted_at: tracedecay_domain::UtcMicros, + progress: SessionRefreshProgressV1, + batch: SessionTemporalProjectionBatchV1, + execution_control: ExecutionControl, + ) -> SessionStoreResult { + self.access() + .commit_session_refresh_begin_batch_result( + request, + accepted_at, + progress, + batch, + execution_control, + ) + .await + } + pub async fn session_refresh_recovery( &self, session_id: &tracedecay_domain::SessionId, diff --git a/crates/tracedecay-session-temporal-store/src/test_registered_impls.rs b/crates/tracedecay-session-temporal-store/src/test_registered_impls.rs index ce6e9b6d62..caf537f245 100644 --- a/crates/tracedecay-session-temporal-store/src/test_registered_impls.rs +++ b/crates/tracedecay-session-temporal-store/src/test_registered_impls.rs @@ -118,6 +118,10 @@ impl SessionTemporalWriteTxn for RegisteredGlobalDbWriteTransaction<'_> { fn commit(self) -> impl Future> + Send { RegisteredGlobalDbWriteTransaction::commit(self) } + + fn rollback(self) -> impl Future> + Send { + RegisteredGlobalDbWriteTransaction::rollback(self) + } } impl SessionTemporalQuery for RegisteredGlobalDbWriterConnection<'_> { diff --git a/crates/tracedecay/tests/session_suite/session_runtime/temporal_refresh.rs b/crates/tracedecay/tests/session_suite/session_runtime/temporal_refresh.rs index aed22105ee..d7672a0c5d 100644 --- a/crates/tracedecay/tests/session_suite/session_runtime/temporal_refresh.rs +++ b/crates/tracedecay/tests/session_suite/session_runtime/temporal_refresh.rs @@ -509,8 +509,14 @@ async fn retained_begin_retry_prevents_discovery_queue_growth_and_cursor_advance .unwrap(); transaction.commit().await.unwrap(); - let mut recovered = SessionTemporalRefreshPassReport::default(); - process_refresh_begin_requests(&store, &state, 1, &mut recovered).await; + let state = Arc::new(state); + let recovered = authority + .run_pass( + &state, + &CanonicalSessionTemporalProjector, + SessionTemporalRefreshPolicy::default(), + ) + .await; assert_eq!( recovered.begun, 1, "the exact retained request must admit once"