use std::collections::{BTreeMap, VecDeque};
use std::fmt;
use std::sync::{Arc, Mutex};
use std::time::Duration;
use mongreldb_types::errors::ErrorCategory;
use mongreldb_types::hlc::{ClockSkewError, HlcClock, HlcTimestamp};
use mongreldb_types::ids::{DatabaseId, MetadataVersion, SchemaVersion, TableId, TabletId};
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use crate::meta::{MetaRejectionReason, SchemaJobState};
use crate::split::{RecordStream, SnapshotPin, TabletDataError};
use crate::tablet::Key;
pub const DDL_STORE_FORMAT_VERSION: u32 = 1;
pub const MIN_SUPPORTED_DDL_STORE_FORMAT_VERSION: u32 = 1;
pub const DDL_REJECTION_LIMIT: usize = 256;
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, thiserror::Error)]
pub enum DdlRejection {
#[error(transparent)]
Meta(#[from] MetaRejectionReason),
#[error("illegal DDL phase transition {from} -> {to} for job {job_id}")]
IllegalPhaseTransition {
job_id: u64,
from: DdlPhase,
to: DdlPhase,
},
#[error("DDL job {job_id} is {state:?}, not Running")]
JobNotRunning {
job_id: u64,
state: SchemaJobState,
},
#[error(
"schema version mismatch on table {table_id}: job pinned {expected}, \
table is now {found}; refresh schema metadata and resubmit"
)]
SchemaVersionMismatch {
table_id: TableId,
expected: SchemaVersion,
found: SchemaVersion,
},
#[error("tablet {tablet} progress regressed for job {job_id}: {reason}")]
TabletProgressRegression {
job_id: u64,
tablet: TabletId,
reason: String,
},
#[error("DDL job {job_id} cannot publish: tablets pending validation: {pending:?}")]
ValidationIncomplete {
job_id: u64,
pending: Vec<TabletId>,
},
#[error(
"index `{index_name}` on table {table_id} cannot be reclaimed: oldest reader \
pins metadata version {oldest_reader:?}, retirement requires {required}"
)]
ReclaimBlocked {
table_id: TableId,
index_name: String,
oldest_reader: Option<MetadataVersion>,
required: MetadataVersion,
},
}
#[derive(Debug, thiserror::Error)]
pub enum DdlError {
#[error(transparent)]
Rejection(#[from] DdlRejection),
#[error(transparent)]
TabletData(#[from] TabletDataError),
#[error(transparent)]
Clock(#[from] ClockSkewError),
#[error("DDL job {job_id} failed validation: {reason}")]
ValidationFailed {
job_id: u64,
reason: String,
},
#[error("DDL job {job_id} cancelled")]
Cancelled {
job_id: u64,
},
}
impl DdlError {
pub fn category(&self) -> ErrorCategory {
match self {
Self::Rejection(DdlRejection::SchemaVersionMismatch { .. }) => {
ErrorCategory::SchemaVersionMismatch
}
Self::Rejection(DdlRejection::Meta(MetaRejectionReason::StaleWrite { .. })) => {
ErrorCategory::StaleMetadata
}
Self::Cancelled { .. } => ErrorCategory::Cancelled,
Self::Rejection(_)
| Self::TabletData(_)
| Self::Clock(_)
| Self::ValidationFailed { .. } => ErrorCategory::ResourceExhausted,
}
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
pub enum DdlPhase {
Pending,
WriteOnly,
Backfilling,
Validating,
Public,
Dropping,
}
impl DdlPhase {
pub const ALL: [DdlPhase; 6] = [
DdlPhase::Pending,
DdlPhase::WriteOnly,
DdlPhase::Backfilling,
DdlPhase::Validating,
DdlPhase::Public,
DdlPhase::Dropping,
];
pub fn can_transition(self, next: Self) -> bool {
use DdlPhase::{Backfilling, Dropping, Pending, Public, Validating, WriteOnly};
matches!(
(self, next),
(Pending, WriteOnly)
| (WriteOnly, Backfilling)
| (Backfilling, Validating)
| (Validating, Public)
| (Public, Dropping)
)
}
}
impl fmt::Display for DdlPhase {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
let name = match self {
Self::Pending => "Pending",
Self::WriteOnly => "WriteOnly",
Self::Backfilling => "Backfilling",
Self::Validating => "Validating",
Self::Public => "Public",
Self::Dropping => "Dropping",
};
f.write_str(name)
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub enum DdlJobKind {
AddIndex,
DropIndex,
AlterSchema,
}
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
pub enum DdlDefinition {
AddIndex {
index_name: String,
spec: serde_json::Value,
},
DropIndex {
index_name: String,
},
AlterSchema {
target: serde_json::Value,
},
}
impl DdlDefinition {
pub fn index_name(&self) -> Option<&str> {
match self {
Self::AddIndex { index_name, .. } | Self::DropIndex { index_name } => Some(index_name),
Self::AlterSchema { .. } => None,
}
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
pub enum TabletDdlStage {
Pending,
Backfilling,
CaughtUp,
Validated,
}
impl fmt::Display for TabletDdlStage {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
let name = match self {
Self::Pending => "Pending",
Self::Backfilling => "Backfilling",
Self::CaughtUp => "CaughtUp",
Self::Validated => "Validated",
};
f.write_str(name)
}
}
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
pub struct TabletDdlProgress {
pub tablet_id: TabletId,
pub stage: TabletDdlStage,
pub rows_scanned: u64,
pub caught_up_through: Option<HlcTimestamp>,
pub validation: Option<TabletValidationReport>,
}
impl TabletDdlProgress {
pub fn pending(tablet_id: TabletId) -> Self {
Self {
tablet_id,
stage: TabletDdlStage::Pending,
rows_scanned: 0,
caught_up_through: None,
validation: None,
}
}
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct TabletValidationReport {
pub tablet_id: TabletId,
pub watermark: HlcTimestamp,
pub expected_rows: u64,
pub actual_rows: u64,
pub expected_checksum: [u8; 32],
pub actual_checksum: [u8; 32],
}
impl TabletValidationReport {
pub fn passed(&self) -> bool {
self.expected_rows == self.actual_rows && self.expected_checksum == self.actual_checksum
}
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct JobValidationReport {
pub job_id: u64,
pub tablets: Vec<TabletValidationReport>,
pub total_expected: u64,
pub total_actual: u64,
pub passed: bool,
}
pub fn generation_checksum(entries: &BTreeMap<Key, Vec<u8>>) -> [u8; 32] {
let mut hasher = Sha256::new();
for (key, entry) in entries {
hasher.update(key.as_bytes());
hasher.update(entry);
}
hasher.finalize().into()
}
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
pub struct DdlJobRecord {
pub job_id: u64,
pub database_id: DatabaseId,
pub table_id: TableId,
pub kind: DdlJobKind,
pub state: SchemaJobState,
pub phase: DdlPhase,
pub definition: DdlDefinition,
pub source_schema_version: SchemaVersion,
pub created_at: HlcTimestamp,
pub updated_at: HlcTimestamp,
pub pinned_snapshot: Option<HlcTimestamp>,
#[serde(with = "tablet_progress_map")]
pub tablet_progress: BTreeMap<TabletId, TabletDdlProgress>,
pub error: Option<String>,
#[serde(default = "zero_metadata_version")]
pub metadata_version: MetadataVersion,
}
fn zero_metadata_version() -> MetadataVersion {
MetadataVersion::ZERO
}
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
pub struct DdlIndexRecord {
pub table_id: TableId,
pub index_name: String,
pub definition: serde_json::Value,
pub phase: DdlPhase,
pub job_id: u64,
pub created_at: HlcTimestamp,
pub publication_version: Option<MetadataVersion>,
pub published_at: Option<HlcTimestamp>,
pub dropping_since: Option<MetadataVersion>,
#[serde(default = "zero_metadata_version")]
pub metadata_version: MetadataVersion,
}
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
pub struct TableAnchor {
pub table_id: TableId,
pub database_id: DatabaseId,
pub schema_version: SchemaVersion,
pub schema: serde_json::Value,
#[serde(default = "zero_metadata_version")]
pub metadata_version: MetadataVersion,
}
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
pub enum DdlCommand {
RegisterTable {
anchor: TableAnchor,
},
SubmitJob {
job: DdlJobRecord,
},
AdvancePhase {
job_id: u64,
to: DdlPhase,
pinned_snapshot: Option<HlcTimestamp>,
expected_version: Option<MetadataVersion>,
},
UpdateTabletProgress {
job_id: u64,
progress: TabletDdlProgress,
},
ForgetTabletProgress {
job_id: u64,
tablet_id: TabletId,
},
ReportTabletValidation {
job_id: u64,
report: TabletValidationReport,
},
PublishJob {
job_id: u64,
published_at: HlcTimestamp,
},
SetJobState {
job_id: u64,
state: SchemaJobState,
updated_at: HlcTimestamp,
error: Option<String>,
expected_version: Option<MetadataVersion>,
},
RemoveIndexRecord {
job_id: u64,
},
ReclaimIndex {
table_id: TableId,
index_name: String,
},
PinReader {
reader_id: u64,
version: MetadataVersion,
},
ReleaseReader {
reader_id: u64,
},
}
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
pub struct DdlStoreRejection {
pub command_id: Option<[u8; 16]>,
pub reason: DdlRejection,
}
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
#[serde(default)]
pub struct DdlJobStore {
pub format_version: u32,
pub metadata_version: MetadataVersion,
pub next_job_id: u64,
pub tables: BTreeMap<TableId, TableAnchor>,
pub jobs: BTreeMap<u64, DdlJobRecord>,
#[serde(with = "index_record_map")]
pub indexes: BTreeMap<(TableId, String), DdlIndexRecord>,
pub reader_pins: BTreeMap<u64, MetadataVersion>,
pub rejections: VecDeque<DdlStoreRejection>,
}
impl Default for DdlJobStore {
fn default() -> Self {
Self {
format_version: DDL_STORE_FORMAT_VERSION,
metadata_version: MetadataVersion::ZERO,
next_job_id: 1,
tables: BTreeMap::new(),
jobs: BTreeMap::new(),
indexes: BTreeMap::new(),
reader_pins: BTreeMap::new(),
rejections: VecDeque::new(),
}
}
}
mod index_record_map {
use super::{DdlIndexRecord, TableId};
use serde::{Deserialize, Deserializer, Serialize, Serializer};
use std::collections::BTreeMap;
pub fn serialize<S>(
map: &BTreeMap<(TableId, String), DdlIndexRecord>,
serializer: S,
) -> Result<S::Ok, S::Error>
where
S: Serializer,
{
let triples: Vec<(TableId, String, DdlIndexRecord)> = map
.iter()
.map(|((table, name), record)| (*table, name.clone(), record.clone()))
.collect();
triples.serialize(serializer)
}
pub fn deserialize<'de, D>(
deserializer: D,
) -> Result<BTreeMap<(TableId, String), DdlIndexRecord>, D::Error>
where
D: Deserializer<'de>,
{
let triples: Vec<(TableId, String, DdlIndexRecord)> = Vec::deserialize(deserializer)?;
Ok(triples
.into_iter()
.map(|(table, name, record)| ((table, name), record))
.collect())
}
}
mod tablet_progress_map {
use super::{TabletDdlProgress, TabletId};
use serde::{Deserialize, Deserializer, Serialize, Serializer};
use std::collections::BTreeMap;
pub fn serialize<S>(
map: &BTreeMap<TabletId, TabletDdlProgress>,
serializer: S,
) -> Result<S::Ok, S::Error>
where
S: Serializer,
{
let pairs: Vec<(TabletId, TabletDdlProgress)> = map
.iter()
.map(|(tablet, progress)| (*tablet, progress.clone()))
.collect();
pairs.serialize(serializer)
}
pub fn deserialize<'de, D>(
deserializer: D,
) -> Result<BTreeMap<TabletId, TabletDdlProgress>, D::Error>
where
D: Deserializer<'de>,
{
let pairs: Vec<(TabletId, TabletDdlProgress)> = Vec::deserialize(deserializer)?;
Ok(pairs.into_iter().collect())
}
}
impl DdlJobStore {
pub fn apply(
&mut self,
command: &DdlCommand,
command_id: Option<[u8; 16]>,
commit_ts: HlcTimestamp,
) -> Result<(), DdlRejection> {
self.metadata_version = MetadataVersion(self.metadata_version.get() + 1);
let version = self.metadata_version;
let result = self.dispatch(command, commit_ts, version);
if let Err(reason) = &result {
self.rejections.push_back(DdlStoreRejection {
command_id,
reason: reason.clone(),
});
while self.rejections.len() > DDL_REJECTION_LIMIT {
self.rejections.pop_front();
}
}
result
}
fn dispatch(
&mut self,
command: &DdlCommand,
commit_ts: HlcTimestamp,
version: MetadataVersion,
) -> Result<(), DdlRejection> {
match command {
DdlCommand::RegisterTable { anchor } => self.apply_register_table(anchor, version),
DdlCommand::SubmitJob { job } => self.apply_submit_job(job, commit_ts, version),
DdlCommand::AdvancePhase {
job_id,
to,
pinned_snapshot,
expected_version,
} => self.apply_advance_phase(
*job_id,
*to,
*pinned_snapshot,
*expected_version,
commit_ts,
version,
),
DdlCommand::UpdateTabletProgress { job_id, progress } => {
self.apply_update_tablet_progress(*job_id, progress, commit_ts, version)
}
DdlCommand::ForgetTabletProgress { job_id, tablet_id } => {
self.apply_forget_tablet_progress(*job_id, *tablet_id, commit_ts, version)
}
DdlCommand::ReportTabletValidation { job_id, report } => {
self.apply_report_tablet_validation(*job_id, report, commit_ts, version)
}
DdlCommand::PublishJob {
job_id,
published_at,
} => self.apply_publish_job(*job_id, *published_at, commit_ts, version),
DdlCommand::SetJobState {
job_id,
state,
updated_at,
error,
expected_version,
} => self.apply_set_job_state(
*job_id,
*state,
*updated_at,
error,
*expected_version,
version,
),
DdlCommand::RemoveIndexRecord { job_id } => {
self.apply_remove_index_record(*job_id, version)
}
DdlCommand::ReclaimIndex {
table_id,
index_name,
} => self.apply_reclaim_index(*table_id, index_name, version),
DdlCommand::PinReader { reader_id, version } => {
self.apply_pin_reader(*reader_id, *version)
}
DdlCommand::ReleaseReader { reader_id } => {
self.reader_pins.remove(reader_id);
Ok(())
}
}
}
fn apply_register_table(
&mut self,
anchor: &TableAnchor,
version: MetadataVersion,
) -> Result<(), DdlRejection> {
if anchor.table_id == TableId::ZERO {
return Err(MetaRejectionReason::Invalid {
reason: "reserved zero table id".to_owned(),
}
.into());
}
if anchor.schema_version == SchemaVersion::ZERO {
return Err(MetaRejectionReason::Invalid {
reason: "reserved zero schema version".to_owned(),
}
.into());
}
match self.tables.get(&anchor.table_id) {
Some(existing) => {
if anchor.schema_version > existing.schema_version {
let mut anchor = anchor.clone();
anchor.metadata_version = version;
self.tables.insert(anchor.table_id, anchor);
Ok(())
} else if anchor.schema_version == existing.schema_version {
if existing.database_id == anchor.database_id
&& existing.schema == anchor.schema
{
Ok(())
} else {
Err(MetaRejectionReason::Conflict {
resource: format!("table {}", anchor.table_id),
reason: "schema version already used for different content".to_owned(),
}
.into())
}
} else {
Err(MetaRejectionReason::StaleWrite {
resource: format!("table {}", anchor.table_id),
current: MetadataVersion(existing.schema_version.get()),
attempted: MetadataVersion(anchor.schema_version.get()),
}
.into())
}
}
None => {
let mut anchor = anchor.clone();
anchor.metadata_version = version;
self.tables.insert(anchor.table_id, anchor);
Ok(())
}
}
}
fn apply_submit_job(
&mut self,
job: &DdlJobRecord,
commit_ts: HlcTimestamp,
version: MetadataVersion,
) -> Result<(), DdlRejection> {
if job.job_id == 0 {
return Err(MetaRejectionReason::Invalid {
reason: "reserved zero job id".to_owned(),
}
.into());
}
if let Some(existing) = self.jobs.get(&job.job_id) {
let mut comparable = existing.clone();
comparable.metadata_version = job.metadata_version;
return if comparable == *job {
Ok(())
} else {
Err(MetaRejectionReason::Conflict {
resource: format!("DDL job {}", job.job_id),
reason: "job id already exists with different content".to_owned(),
}
.into())
};
}
let anchor = self
.tables
.get(&job.table_id)
.ok_or(MetaRejectionReason::NotFound {
resource: format!("table {}", job.table_id),
})?;
if anchor.database_id != job.database_id {
return Err(MetaRejectionReason::Conflict {
resource: format!("table {}", job.table_id),
reason: "table belongs to a different database".to_owned(),
}
.into());
}
if anchor.schema_version != job.source_schema_version {
return Err(DdlRejection::SchemaVersionMismatch {
table_id: job.table_id,
expected: job.source_schema_version,
found: anchor.schema_version,
});
}
if job.state != SchemaJobState::Pending {
return Err(MetaRejectionReason::Invalid {
reason: "submitted jobs start Pending".to_owned(),
}
.into());
}
let expected_phase = match job.kind {
DdlJobKind::AddIndex | DdlJobKind::AlterSchema => DdlPhase::Pending,
DdlJobKind::DropIndex => DdlPhase::Public,
};
if job.phase != expected_phase {
return Err(MetaRejectionReason::Invalid {
reason: format!(
"{:?} jobs start in phase {expected_phase}, not {}",
job.kind, job.phase
),
}
.into());
}
if let Some(active) = self
.jobs
.values()
.find(|existing| existing.table_id == job.table_id && !existing.state.is_terminal())
{
return Err(MetaRejectionReason::Conflict {
resource: format!("table {}", job.table_id),
reason: format!("table already has active DDL job {}", active.job_id),
}
.into());
}
match &job.definition {
DdlDefinition::AddIndex { index_name, spec } => {
if index_name.trim().is_empty() {
return Err(MetaRejectionReason::Invalid {
reason: "index name is empty".to_owned(),
}
.into());
}
let key = (job.table_id, index_name.clone());
if self.indexes.contains_key(&key) {
return Err(MetaRejectionReason::Conflict {
resource: format!("index `{index_name}` on table {}", job.table_id),
reason: "index name already exists".to_owned(),
}
.into());
}
self.indexes.insert(
key,
DdlIndexRecord {
table_id: job.table_id,
index_name: index_name.clone(),
definition: spec.clone(),
phase: DdlPhase::Pending,
job_id: job.job_id,
created_at: commit_ts,
publication_version: None,
published_at: None,
dropping_since: None,
metadata_version: version,
},
);
}
DdlDefinition::DropIndex { index_name } => {
let key = (job.table_id, index_name.clone());
match self.indexes.get(&key) {
None => {
return Err(MetaRejectionReason::NotFound {
resource: format!("index `{index_name}` on table {}", job.table_id),
}
.into());
}
Some(record) if record.phase != DdlPhase::Public => {
return Err(MetaRejectionReason::Conflict {
resource: format!("index `{index_name}` on table {}", job.table_id),
reason: format!("index is {}, not Public", record.phase),
}
.into());
}
Some(_) => {}
}
}
DdlDefinition::AlterSchema { .. } => {}
}
let mut job = job.clone();
job.metadata_version = version;
job.updated_at = commit_ts;
self.next_job_id = self.next_job_id.max(job.job_id + 1);
self.jobs.insert(job.job_id, job);
Ok(())
}
fn apply_advance_phase(
&mut self,
job_id: u64,
to: DdlPhase,
pinned_snapshot: Option<HlcTimestamp>,
expected_version: Option<MetadataVersion>,
commit_ts: HlcTimestamp,
version: MetadataVersion,
) -> Result<(), DdlRejection> {
let job = self
.jobs
.get(&job_id)
.ok_or(MetaRejectionReason::NotFound {
resource: format!("DDL job {job_id}"),
})?;
if let Some(expected) = expected_version {
if expected != job.metadata_version {
return Err(MetaRejectionReason::StaleWrite {
resource: format!("DDL job {job_id}"),
current: job.metadata_version,
attempted: expected,
}
.into());
}
}
if job.state.is_terminal() {
return Err(MetaRejectionReason::Conflict {
resource: format!("DDL job {job_id}"),
reason: format!("terminal state {:?}", job.state),
}
.into());
}
if job.state != SchemaJobState::Running {
return Err(DdlRejection::JobNotRunning {
job_id,
state: job.state,
});
}
if job.phase == to {
if to == DdlPhase::Backfilling && job.pinned_snapshot != pinned_snapshot {
return Err(MetaRejectionReason::Conflict {
resource: format!("DDL job {job_id}"),
reason: "phase already reached with a different pinned snapshot".to_owned(),
}
.into());
}
return Ok(());
}
if !job.phase.can_transition(to) {
return Err(DdlRejection::IllegalPhaseTransition {
job_id,
from: job.phase,
to,
});
}
if to == DdlPhase::Public && job.kind != DdlJobKind::DropIndex {
return Err(MetaRejectionReason::Invalid {
reason: "Public is reached through PublishJob, never AdvancePhase".to_owned(),
}
.into());
}
if to == DdlPhase::Backfilling && pinned_snapshot.is_none() {
return Err(MetaRejectionReason::Invalid {
reason: "Backfilling requires the job-wide pinned snapshot".to_owned(),
}
.into());
}
let job = self.jobs.get_mut(&job_id).expect("job existence checked");
job.phase = to;
if to == DdlPhase::Backfilling {
job.pinned_snapshot = pinned_snapshot;
}
job.updated_at = commit_ts;
job.metadata_version = version;
if let Some(index_name) = job.definition.index_name() {
let key = (job.table_id, index_name.to_owned());
if let Some(record) = self.indexes.get_mut(&key) {
record.phase = to;
if to == DdlPhase::Dropping {
record.dropping_since = Some(version);
}
record.metadata_version = version;
}
}
Ok(())
}
fn apply_update_tablet_progress(
&mut self,
job_id: u64,
progress: &TabletDdlProgress,
commit_ts: HlcTimestamp,
version: MetadataVersion,
) -> Result<(), DdlRejection> {
let job = self
.jobs
.get_mut(&job_id)
.ok_or(MetaRejectionReason::NotFound {
resource: format!("DDL job {job_id}"),
})?;
if job.phase != DdlPhase::Backfilling {
let dominated = job
.tablet_progress
.get(&progress.tablet_id)
.is_some_and(|existing| {
existing.stage >= progress.stage
&& existing.rows_scanned >= progress.rows_scanned
});
if dominated {
return Ok(());
}
return Err(MetaRejectionReason::Conflict {
resource: format!("DDL job {job_id}"),
reason: format!("job is {}, not Backfilling", job.phase),
}
.into());
}
if let Some(existing) = job.tablet_progress.get(&progress.tablet_id) {
if progress.stage < existing.stage {
return Err(DdlRejection::TabletProgressRegression {
job_id,
tablet: progress.tablet_id,
reason: format!("stage {} -> {}", existing.stage, progress.stage),
});
}
if progress.stage == existing.stage && progress.rows_scanned < existing.rows_scanned {
return Err(DdlRejection::TabletProgressRegression {
job_id,
tablet: progress.tablet_id,
reason: format!(
"rows_scanned {} -> {}",
existing.rows_scanned, progress.rows_scanned
),
});
}
}
job.tablet_progress
.insert(progress.tablet_id, progress.clone());
job.updated_at = commit_ts;
job.metadata_version = version;
Ok(())
}
fn apply_forget_tablet_progress(
&mut self,
job_id: u64,
tablet_id: TabletId,
commit_ts: HlcTimestamp,
version: MetadataVersion,
) -> Result<(), DdlRejection> {
let job = self
.jobs
.get_mut(&job_id)
.ok_or(MetaRejectionReason::NotFound {
resource: format!("DDL job {job_id}"),
})?;
if !matches!(job.phase, DdlPhase::Backfilling | DdlPhase::Validating) {
return Err(MetaRejectionReason::Conflict {
resource: format!("DDL job {job_id}"),
reason: format!(
"job is {}; tablet cursors are only mutable while driving",
job.phase
),
}
.into());
}
if job.tablet_progress.remove(&tablet_id).is_some() {
job.updated_at = commit_ts;
job.metadata_version = version;
}
Ok(())
}
fn apply_report_tablet_validation(
&mut self,
job_id: u64,
report: &TabletValidationReport,
commit_ts: HlcTimestamp,
version: MetadataVersion,
) -> Result<(), DdlRejection> {
let job = self
.jobs
.get_mut(&job_id)
.ok_or(MetaRejectionReason::NotFound {
resource: format!("DDL job {job_id}"),
})?;
if job.phase != DdlPhase::Validating {
return Err(MetaRejectionReason::Conflict {
resource: format!("DDL job {job_id}"),
reason: format!("job is {}, not Validating", job.phase),
}
.into());
}
let progress = job.tablet_progress.get_mut(&report.tablet_id).ok_or(
MetaRejectionReason::NotFound {
resource: format!("tablet {} progress of job {job_id}", report.tablet_id),
},
)?;
if progress.stage < TabletDdlStage::CaughtUp {
return Err(DdlRejection::TabletProgressRegression {
job_id,
tablet: report.tablet_id,
reason: "validation reported before catch-up completed".to_owned(),
});
}
if progress.stage == TabletDdlStage::Validated
&& progress.validation.as_ref() == Some(report)
{
return Ok(());
}
progress.validation = Some(report.clone());
progress.caught_up_through = Some(report.watermark);
if report.passed() {
progress.stage = TabletDdlStage::Validated;
}
job.updated_at = commit_ts;
job.metadata_version = version;
Ok(())
}
fn apply_publish_job(
&mut self,
job_id: u64,
published_at: HlcTimestamp,
commit_ts: HlcTimestamp,
version: MetadataVersion,
) -> Result<(), DdlRejection> {
let job = self
.jobs
.get(&job_id)
.ok_or(MetaRejectionReason::NotFound {
resource: format!("DDL job {job_id}"),
})?;
if job.state == SchemaJobState::Succeeded
&& job.phase == DdlPhase::Public
&& job.kind != DdlJobKind::DropIndex
{
return Ok(()); }
if job.kind == DdlJobKind::DropIndex {
return Err(MetaRejectionReason::Invalid {
reason: "drop jobs ride AdvancePhase(Public -> Dropping), never PublishJob"
.to_owned(),
}
.into());
}
if job.state != SchemaJobState::Running {
return Err(DdlRejection::JobNotRunning {
job_id,
state: job.state,
});
}
if job.phase != DdlPhase::Validating {
return Err(DdlRejection::IllegalPhaseTransition {
job_id,
from: job.phase,
to: DdlPhase::Public,
});
}
let anchor = self
.tables
.get(&job.table_id)
.ok_or(MetaRejectionReason::NotFound {
resource: format!("table {}", job.table_id),
})?;
if anchor.schema_version != job.source_schema_version {
return Err(DdlRejection::SchemaVersionMismatch {
table_id: job.table_id,
expected: job.source_schema_version,
found: anchor.schema_version,
});
}
let pending: Vec<TabletId> = job
.tablet_progress
.values()
.filter(|progress| {
progress.stage != TabletDdlStage::Validated
|| progress
.validation
.as_ref()
.is_none_or(|report| !report.passed())
})
.map(|progress| progress.tablet_id)
.collect();
if job.tablet_progress.is_empty() || !pending.is_empty() {
return Err(DdlRejection::ValidationIncomplete { job_id, pending });
}
match &job.definition {
DdlDefinition::AddIndex { index_name, .. } => {
let key = (job.table_id, index_name.clone());
let record = self
.indexes
.get_mut(&key)
.ok_or(MetaRejectionReason::NotFound {
resource: format!("index `{index_name}` on table {}", job.table_id),
})?;
record.phase = DdlPhase::Public;
record.publication_version = Some(version);
record.published_at = Some(published_at);
record.metadata_version = version;
}
DdlDefinition::AlterSchema { target } => {
let anchor = self
.tables
.get_mut(&job.table_id)
.expect("anchor existence checked");
anchor.schema_version = SchemaVersion(job.source_schema_version.get() + 1);
anchor.schema = target.clone();
anchor.metadata_version = version;
}
DdlDefinition::DropIndex { .. } => unreachable!("drop kind refused above"),
}
let job = self.jobs.get_mut(&job_id).expect("job existence checked");
job.phase = DdlPhase::Public;
job.state = SchemaJobState::Succeeded;
job.updated_at = commit_ts;
job.metadata_version = version;
Ok(())
}
fn apply_set_job_state(
&mut self,
job_id: u64,
state: SchemaJobState,
updated_at: HlcTimestamp,
error: &Option<String>,
expected_version: Option<MetadataVersion>,
version: MetadataVersion,
) -> Result<(), DdlRejection> {
let record = self
.jobs
.get_mut(&job_id)
.ok_or(MetaRejectionReason::NotFound {
resource: format!("DDL job {job_id}"),
})?;
if let Some(expected) = expected_version {
if expected != record.metadata_version {
return Err(MetaRejectionReason::StaleWrite {
resource: format!("DDL job {job_id}"),
current: record.metadata_version,
attempted: expected,
}
.into());
}
}
if record.state.is_terminal() {
return Err(MetaRejectionReason::Conflict {
resource: format!("DDL job {job_id}"),
reason: format!("terminal state {:?}", record.state),
}
.into());
}
if record.state != state && !record.state.can_transition(state) {
return Err(MetaRejectionReason::Conflict {
resource: format!("DDL job {job_id}"),
reason: format!("illegal transition {:?} -> {:?}", record.state, state),
}
.into());
}
if record.state == state && record.error == *error {
return Ok(());
}
record.state = state;
record.updated_at = updated_at;
record.error = error.clone();
record.metadata_version = version;
Ok(())
}
fn apply_remove_index_record(
&mut self,
job_id: u64,
version: MetadataVersion,
) -> Result<(), DdlRejection> {
let job = self
.jobs
.get(&job_id)
.ok_or(MetaRejectionReason::NotFound {
resource: format!("DDL job {job_id}"),
})?;
if !matches!(
job.state,
SchemaJobState::RollingBack | SchemaJobState::Failed
) {
return Err(MetaRejectionReason::Conflict {
resource: format!("DDL job {job_id}"),
reason: "index records are removed only during rollback".to_owned(),
}
.into());
}
let Some(index_name) = job.definition.index_name() else {
return Ok(());
};
let key = (job.table_id, index_name.to_owned());
match self.indexes.get(&key) {
None => Ok(()),
Some(record) if record.phase == DdlPhase::Public => {
Err(MetaRejectionReason::Conflict {
resource: format!("index `{index_name}` on table {}", job.table_id),
reason: "refusing to remove a published index".to_owned(),
}
.into())
}
Some(_) => {
self.indexes.remove(&key);
let job = self.jobs.get_mut(&job_id).expect("job existence checked");
job.metadata_version = version;
Ok(())
}
}
}
fn apply_reclaim_index(
&mut self,
table_id: TableId,
index_name: &str,
_version: MetadataVersion,
) -> Result<(), DdlRejection> {
let key = (table_id, index_name.to_owned());
let Some(record) = self.indexes.get(&key) else {
return Ok(()); };
if record.phase != DdlPhase::Dropping {
return Err(MetaRejectionReason::Conflict {
resource: format!("index `{index_name}` on table {table_id}"),
reason: format!("index is {}, not Dropping", record.phase),
}
.into());
}
let required = record
.dropping_since
.expect("Dropping records carry dropping_since");
let oldest_reader = self.oldest_reader_version();
if oldest_reader.is_some_and(|oldest| oldest < required) {
return Err(DdlRejection::ReclaimBlocked {
table_id,
index_name: index_name.to_owned(),
oldest_reader,
required,
});
}
self.indexes.remove(&key);
Ok(())
}
fn apply_pin_reader(
&mut self,
reader_id: u64,
version: MetadataVersion,
) -> Result<(), DdlRejection> {
if reader_id == 0 {
return Err(MetaRejectionReason::Invalid {
reason: "reserved zero reader id".to_owned(),
}
.into());
}
if let Some(existing) = self.reader_pins.get(&reader_id) {
if version < *existing {
return Err(MetaRejectionReason::Conflict {
resource: format!("reader {reader_id}"),
reason: format!("reader pins only move forward ({} -> {version})", *existing),
}
.into());
}
}
self.reader_pins.insert(reader_id, version);
Ok(())
}
pub fn job(&self, job_id: u64) -> Option<&DdlJobRecord> {
self.jobs.get(&job_id)
}
pub fn index(&self, table_id: TableId, index_name: &str) -> Option<&DdlIndexRecord> {
self.indexes.get(&(table_id, index_name.to_owned()))
}
pub fn table_anchor(&self, table_id: TableId) -> Option<&TableAnchor> {
self.tables.get(&table_id)
}
pub fn oldest_reader_version(&self) -> Option<MetadataVersion> {
self.reader_pins.values().copied().min()
}
pub fn planner_visible_at(
&self,
table_id: TableId,
version: MetadataVersion,
) -> Vec<DdlIndexRecord> {
self.indexes
.values()
.filter(|record| {
record.table_id == table_id
&& record
.publication_version
.is_some_and(|publication| publication <= version)
&& record
.dropping_since
.is_none_or(|dropping| version < dropping)
})
.cloned()
.collect()
}
pub fn planner_visible(&self, table_id: TableId) -> Vec<DdlIndexRecord> {
self.planner_visible_at(table_id, self.metadata_version)
}
pub fn write_maintained(&self, table_id: TableId) -> Vec<DdlIndexRecord> {
self.indexes
.values()
.filter(|record| {
record.table_id == table_id
&& matches!(
record.phase,
DdlPhase::WriteOnly
| DdlPhase::Backfilling
| DdlPhase::Validating
| DdlPhase::Public
)
})
.cloned()
.collect()
}
}
pub type DeltaStream<'a> = Box<dyn Iterator<Item = (HlcTimestamp, Key, Option<Vec<u8>>)> + 'a>;
pub trait BackfillKeyspace {
fn pin_snapshot(&self, ts: HlcTimestamp) -> Result<Box<dyn SnapshotPin>, TabletDataError>;
fn snapshot_at(&self, ts: HlcTimestamp) -> Result<RecordStream<'_>, TabletDataError>;
fn deltas_after(&self, ts: HlcTimestamp) -> Result<DeltaStream<'_>, TabletDataError>;
}
pub trait HiddenIndexSink {
fn begin_build(&mut self, index: &str) -> Result<(), TabletDataError>;
fn stage_entry(&mut self, index: &str, key: &Key, entry: &[u8]) -> Result<(), TabletDataError>;
fn install_staged(&mut self, index: &str) -> Result<(), TabletDataError>;
fn apply_delta(
&mut self,
index: &str,
key: &Key,
entry: Option<&[u8]>,
) -> Result<(), TabletDataError>;
fn drop_generation(&mut self, index: &str) -> Result<(), TabletDataError>;
fn generation_entries(&self, index: &str) -> Result<BTreeMap<Key, Vec<u8>>, TabletDataError>;
fn entry_count(&self, index: &str) -> Result<u64, TabletDataError> {
Ok(self.generation_entries(index)?.len() as u64)
}
}
pub trait ApplySideIndexMaintainer {
fn sync_definitions(&mut self, maintained: Vec<String>) -> Result<(), TabletDataError>;
fn apply_committed_write(
&mut self,
key: &Key,
row: Option<&[u8]>,
) -> Result<(), TabletDataError>;
}
pub trait DdlTabletProvider {
fn tablets_of(&self, table: TableId) -> Vec<TabletId>;
fn keyspace(&self, tablet: TabletId) -> Result<Box<dyn BackfillKeyspace>, TabletDataError>;
fn sink(&self, tablet: TabletId) -> Result<Box<dyn HiddenIndexSink>, TabletDataError>;
fn maintainer(
&self,
tablet: TabletId,
) -> Result<Box<dyn ApplySideIndexMaintainer>, TabletDataError>;
}
pub type IndexProjection = Arc<dyn Fn(&DdlIndexRecord, &Key, &[u8]) -> Vec<u8> + Send + Sync>;
pub fn identity_projection() -> IndexProjection {
Arc::new(|_index, _key, value| value.to_vec())
}
type VersionedRows = BTreeMap<Key, Vec<(HlcTimestamp, Option<Vec<u8>>)>>;
#[derive(Clone, Default)]
pub struct InMemoryDdlKeyspace {
state: Arc<Mutex<VersionedRows>>,
}
impl InMemoryDdlKeyspace {
pub fn new() -> Self {
Self::default()
}
pub fn insert(&self, key: Key, ts: HlcTimestamp, value: Vec<u8>) {
let mut rows = self.state.lock().expect("keyspace lock poisoned");
let chain = rows.entry(key).or_default();
chain.push((ts, Some(value)));
chain.sort_by_key(|(version, _)| *version);
}
pub fn delete(&self, key: Key, ts: HlcTimestamp) {
let mut rows = self.state.lock().expect("keyspace lock poisoned");
let chain = rows.entry(key).or_default();
chain.push((ts, None));
chain.sort_by_key(|(version, _)| *version);
}
pub fn rows_at(&self, ts: HlcTimestamp) -> BTreeMap<Key, Vec<u8>> {
let rows = self.state.lock().expect("keyspace lock poisoned");
rows.iter()
.filter_map(|(key, chain)| {
let (_, value) = chain.iter().rfind(|(version, _)| *version <= ts)?;
value.clone().map(|value| (key.clone(), value))
})
.collect()
}
}
impl BackfillKeyspace for InMemoryDdlKeyspace {
fn pin_snapshot(&self, ts: HlcTimestamp) -> Result<Box<dyn SnapshotPin>, TabletDataError> {
Ok(Box::new(DdlSnapshotPin { ts }))
}
fn snapshot_at(&self, ts: HlcTimestamp) -> Result<RecordStream<'_>, TabletDataError> {
Ok(Box::new(self.rows_at(ts).into_iter()))
}
fn deltas_after(&self, ts: HlcTimestamp) -> Result<DeltaStream<'_>, TabletDataError> {
let rows = self.state.lock().expect("keyspace lock poisoned");
let mut deltas: Vec<(HlcTimestamp, Key, Option<Vec<u8>>)> = Vec::new();
for (key, chain) in rows.iter() {
for (version, value) in chain {
if *version > ts {
deltas.push((*version, key.clone(), value.clone()));
}
}
}
deltas.sort_by(|(left_ts, left_key, _), (right_ts, right_key, _)| {
(left_ts, left_key).cmp(&(right_ts, right_key))
});
Ok(Box::new(deltas.into_iter()))
}
}
struct DdlSnapshotPin {
ts: HlcTimestamp,
}
impl SnapshotPin for DdlSnapshotPin {
fn pinned_at(&self) -> HlcTimestamp {
self.ts
}
}
#[derive(Clone, Default)]
pub struct InMemoryTabletIndexes {
state: Arc<Mutex<IndexesState>>,
}
#[derive(Default)]
struct IndexesState {
maintained: Vec<String>,
staged: BTreeMap<String, BTreeMap<Key, Vec<u8>>>,
generations: BTreeMap<String, BTreeMap<Key, Vec<u8>>>,
begin_builds: usize,
}
impl InMemoryTabletIndexes {
pub fn new() -> Self {
Self::default()
}
pub fn entries(&self, index: &str) -> BTreeMap<Key, Vec<u8>> {
self.state
.lock()
.expect("indexes lock poisoned")
.generations
.get(index)
.cloned()
.unwrap_or_default()
}
pub fn begin_build_count(&self) -> usize {
self.state
.lock()
.expect("indexes lock poisoned")
.begin_builds
}
pub fn is_installed(&self, index: &str) -> bool {
self.state
.lock()
.expect("indexes lock poisoned")
.generations
.contains_key(index)
}
pub fn maintained(&self) -> Vec<String> {
self.state
.lock()
.expect("indexes lock poisoned")
.maintained
.clone()
}
#[cfg(test)]
pub fn seed_staged(&self, index: &str, key: Key, entry: Vec<u8>) {
self.state
.lock()
.expect("indexes lock poisoned")
.staged
.entry(index.to_owned())
.or_default()
.insert(key, entry);
}
}
impl HiddenIndexSink for InMemoryTabletIndexes {
fn begin_build(&mut self, index: &str) -> Result<(), TabletDataError> {
let mut state = self.state.lock().expect("indexes lock poisoned");
state.staged.insert(index.to_owned(), BTreeMap::new());
state.begin_builds += 1;
Ok(())
}
fn stage_entry(&mut self, index: &str, key: &Key, entry: &[u8]) -> Result<(), TabletDataError> {
let mut state = self.state.lock().expect("indexes lock poisoned");
let Some(staged) = state.staged.get_mut(index) else {
return Err(TabletDataError::NoStagedBuild);
};
staged.insert(key.clone(), entry.to_vec());
Ok(())
}
fn install_staged(&mut self, index: &str) -> Result<(), TabletDataError> {
let mut state = self.state.lock().expect("indexes lock poisoned");
let Some(staged) = state.staged.remove(index) else {
return Err(TabletDataError::NoStagedBuild);
};
state.generations.insert(index.to_owned(), staged);
Ok(())
}
fn apply_delta(
&mut self,
index: &str,
key: &Key,
entry: Option<&[u8]>,
) -> Result<(), TabletDataError> {
let mut state = self.state.lock().expect("indexes lock poisoned");
let Some(generation) = state.generations.get_mut(index) else {
return Err(TabletDataError::Sink(format!(
"generation `{index}` is not installed"
)));
};
match entry {
Some(entry) => {
generation.insert(key.clone(), entry.to_vec());
}
None => {
generation.remove(key);
}
}
Ok(())
}
fn drop_generation(&mut self, index: &str) -> Result<(), TabletDataError> {
let mut state = self.state.lock().expect("indexes lock poisoned");
state.generations.remove(index);
state.staged.remove(index);
Ok(())
}
fn generation_entries(&self, index: &str) -> Result<BTreeMap<Key, Vec<u8>>, TabletDataError> {
Ok(self.entries(index))
}
}
impl ApplySideIndexMaintainer for InMemoryTabletIndexes {
fn sync_definitions(&mut self, maintained: Vec<String>) -> Result<(), TabletDataError> {
self.state.lock().expect("indexes lock poisoned").maintained = maintained;
Ok(())
}
fn apply_committed_write(
&mut self,
key: &Key,
row: Option<&[u8]>,
) -> Result<(), TabletDataError> {
let mut state = self.state.lock().expect("indexes lock poisoned");
for index in state.maintained.clone() {
if let Some(generation) = state.generations.get_mut(&index) {
match row {
Some(row) => {
generation.insert(key.clone(), row.to_vec());
}
None => {
generation.remove(key);
}
}
}
}
Ok(())
}
}
struct InMemoryTablet {
table_id: TableId,
keyspace: InMemoryDdlKeyspace,
indexes: InMemoryTabletIndexes,
}
#[derive(Clone, Default)]
pub struct InMemoryDdlTablets {
tablets: Arc<Mutex<BTreeMap<TabletId, InMemoryTablet>>>,
}
impl InMemoryDdlTablets {
pub fn new() -> Self {
Self::default()
}
pub fn add_tablet(&self, tablet: TabletId, table: TableId) {
self.tablets
.lock()
.expect("tablets lock poisoned")
.entry(tablet)
.or_insert_with(|| InMemoryTablet {
table_id: table,
keyspace: InMemoryDdlKeyspace::new(),
indexes: InMemoryTabletIndexes::new(),
});
}
#[cfg(test)]
pub fn remove_tablet(&self, tablet: TabletId) {
self.tablets
.lock()
.expect("tablets lock poisoned")
.remove(&tablet);
}
pub fn keyspace_handle(&self, tablet: TabletId) -> InMemoryDdlKeyspace {
self.tablets
.lock()
.expect("tablets lock poisoned")
.get(&tablet)
.expect("unknown tablet")
.keyspace
.clone()
}
pub fn indexes_handle(&self, tablet: TabletId) -> InMemoryTabletIndexes {
self.tablets
.lock()
.expect("tablets lock poisoned")
.get(&tablet)
.expect("unknown tablet")
.indexes
.clone()
}
pub fn commit_write(
&self,
tablet: TabletId,
key: Key,
ts: HlcTimestamp,
row: Option<Vec<u8>>,
) -> Result<(), TabletDataError> {
let (keyspace, mut indexes) = {
let tablets = self.tablets.lock().expect("tablets lock poisoned");
let tablet = tablets.get(&tablet).expect("unknown tablet");
(tablet.keyspace.clone(), tablet.indexes.clone())
};
match row {
Some(row) => {
keyspace.insert(key.clone(), ts, row.clone());
indexes.apply_committed_write(&key, Some(&row))
}
None => {
keyspace.delete(key.clone(), ts);
indexes.apply_committed_write(&key, None)
}
}
}
}
impl DdlTabletProvider for InMemoryDdlTablets {
fn tablets_of(&self, table: TableId) -> Vec<TabletId> {
self.tablets
.lock()
.expect("tablets lock poisoned")
.iter()
.filter(|(_, tablet)| tablet.table_id == table)
.map(|(tablet_id, _)| *tablet_id)
.collect()
}
fn keyspace(&self, tablet: TabletId) -> Result<Box<dyn BackfillKeyspace>, TabletDataError> {
let tablets = self.tablets.lock().expect("tablets lock poisoned");
let Some(tablet) = tablets.get(&tablet) else {
return Err(TabletDataError::Keyspace(format!(
"unknown tablet {tablet}"
)));
};
Ok(Box::new(tablet.keyspace.clone()))
}
fn sink(&self, tablet: TabletId) -> Result<Box<dyn HiddenIndexSink>, TabletDataError> {
let tablets = self.tablets.lock().expect("tablets lock poisoned");
let Some(tablet) = tablets.get(&tablet) else {
return Err(TabletDataError::Sink(format!("unknown tablet {tablet}")));
};
Ok(Box::new(tablet.indexes.clone()))
}
fn maintainer(
&self,
tablet: TabletId,
) -> Result<Box<dyn ApplySideIndexMaintainer>, TabletDataError> {
let tablets = self.tablets.lock().expect("tablets lock poisoned");
let Some(tablet) = tablets.get(&tablet) else {
return Err(TabletDataError::Sink(format!("unknown tablet {tablet}")));
};
Ok(Box::new(tablet.indexes.clone()))
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum DriveOutcome {
Completed,
Parked,
RolledBack,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum DriveStep {
Admitted,
PhaseAdvanced {
from: DdlPhase,
to: DdlPhase,
},
TabletBackfilled {
tablet: TabletId,
},
TabletValidated {
tablet: TabletId,
},
Published {
version: MetadataVersion,
},
Parked,
RolledBack,
Terminal,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum ReclaimOutcome {
Reclaimed,
Blocked {
oldest_reader: Option<MetadataVersion>,
required: MetadataVersion,
},
}
#[derive(Clone, Debug, PartialEq)]
pub struct DdlJobStatus {
pub record: DdlJobRecord,
pub tablets_registered: usize,
pub tablets_caught_up: usize,
pub tablets_validated: usize,
pub rows_backfilled: u64,
}
pub struct DdlDriver<P: DdlTabletProvider> {
store: Arc<Mutex<DdlJobStore>>,
provider: P,
clock: HlcClock,
projection: IndexProjection,
}
impl<P: DdlTabletProvider> DdlDriver<P> {
pub fn new(store: Arc<Mutex<DdlJobStore>>, provider: P) -> Self {
Self::with_projection(store, provider, identity_projection())
}
pub fn with_projection(
store: Arc<Mutex<DdlJobStore>>,
provider: P,
projection: IndexProjection,
) -> Self {
Self {
store,
provider,
clock: HlcClock::new(0, Duration::from_secs(30)),
projection,
}
}
pub fn store_handle(&self) -> Arc<Mutex<DdlJobStore>> {
self.store.clone()
}
pub fn submit_add_index(
&self,
database_id: DatabaseId,
table_id: TableId,
index_name: impl Into<String>,
spec: serde_json::Value,
source_schema_version: SchemaVersion,
) -> Result<u64, DdlError> {
let created_at = self.clock.now()?;
let job = DdlJobRecord {
job_id: self.peek_next_job_id(),
database_id,
table_id,
kind: DdlJobKind::AddIndex,
state: SchemaJobState::Pending,
phase: DdlPhase::Pending,
definition: DdlDefinition::AddIndex {
index_name: index_name.into(),
spec,
},
source_schema_version,
created_at,
updated_at: created_at,
pinned_snapshot: None,
tablet_progress: BTreeMap::new(),
error: None,
metadata_version: MetadataVersion::ZERO,
};
self.apply(DdlCommand::SubmitJob { job: job.clone() })?;
Ok(job.job_id)
}
pub fn submit_drop_index(
&self,
database_id: DatabaseId,
table_id: TableId,
index_name: impl Into<String>,
source_schema_version: SchemaVersion,
) -> Result<u64, DdlError> {
let created_at = self.clock.now()?;
let job = DdlJobRecord {
job_id: self.peek_next_job_id(),
database_id,
table_id,
kind: DdlJobKind::DropIndex,
state: SchemaJobState::Pending,
phase: DdlPhase::Public,
definition: DdlDefinition::DropIndex {
index_name: index_name.into(),
},
source_schema_version,
created_at,
updated_at: created_at,
pinned_snapshot: None,
tablet_progress: BTreeMap::new(),
error: None,
metadata_version: MetadataVersion::ZERO,
};
self.apply(DdlCommand::SubmitJob { job: job.clone() })?;
Ok(job.job_id)
}
pub fn submit_alter_schema(
&self,
database_id: DatabaseId,
table_id: TableId,
target: serde_json::Value,
source_schema_version: SchemaVersion,
) -> Result<u64, DdlError> {
let created_at = self.clock.now()?;
let job = DdlJobRecord {
job_id: self.peek_next_job_id(),
database_id,
table_id,
kind: DdlJobKind::AlterSchema,
state: SchemaJobState::Pending,
phase: DdlPhase::Pending,
definition: DdlDefinition::AlterSchema { target },
source_schema_version,
created_at,
updated_at: created_at,
pinned_snapshot: None,
tablet_progress: BTreeMap::new(),
error: None,
metadata_version: MetadataVersion::ZERO,
};
self.apply(DdlCommand::SubmitJob { job: job.clone() })?;
Ok(job.job_id)
}
pub fn pause_job(&self, job_id: u64) -> Result<(), DdlError> {
let record = self.job_record(job_id)?;
if record.state != SchemaJobState::Running {
return Err(DdlError::Rejection(DdlRejection::JobNotRunning {
job_id,
state: record.state,
}));
}
self.apply(DdlCommand::SetJobState {
job_id,
state: SchemaJobState::Paused,
updated_at: self.clock.now()?,
error: None,
expected_version: Some(record.metadata_version),
})
}
pub fn resume_job(&self, job_id: u64) -> Result<(), DdlError> {
let record = self.job_record(job_id)?;
if record.state != SchemaJobState::Paused {
return Err(DdlError::Rejection(DdlRejection::Meta(
MetaRejectionReason::Conflict {
resource: format!("DDL job {job_id}"),
reason: format!("cannot resume from {:?}", record.state),
},
)));
}
self.apply(DdlCommand::SetJobState {
job_id,
state: SchemaJobState::Pending,
updated_at: self.clock.now()?,
error: None,
expected_version: Some(record.metadata_version),
})
}
pub fn cancel_job(&self, job_id: u64) -> Result<(), DdlError> {
let record = self.job_record(job_id)?;
if record.state.is_terminal() {
return Err(DdlError::Rejection(DdlRejection::Meta(
MetaRejectionReason::Conflict {
resource: format!("DDL job {job_id}"),
reason: format!("terminal state {:?}", record.state),
},
)));
}
self.apply(DdlCommand::SetJobState {
job_id,
state: SchemaJobState::Cancelling,
updated_at: self.clock.now()?,
error: None,
expected_version: Some(record.metadata_version),
})
}
pub fn job_status(&self, job_id: u64) -> Option<DdlJobStatus> {
let record = self
.store
.lock()
.expect("store lock poisoned")
.job(job_id)?
.clone();
let tablets_registered = record.tablet_progress.len();
let tablets_caught_up = record
.tablet_progress
.values()
.filter(|progress| progress.stage >= TabletDdlStage::CaughtUp)
.count();
let tablets_validated = record
.tablet_progress
.values()
.filter(|progress| progress.stage == TabletDdlStage::Validated)
.count();
let rows_backfilled = record
.tablet_progress
.values()
.map(|progress| progress.rows_scanned)
.sum();
Some(DdlJobStatus {
record,
tablets_registered,
tablets_caught_up,
tablets_validated,
rows_backfilled,
})
}
pub fn validation_report(&self, job_id: u64) -> Option<JobValidationReport> {
let record = self
.store
.lock()
.expect("store lock poisoned")
.job(job_id)?
.clone();
let tablets: Vec<TabletValidationReport> = record
.tablet_progress
.values()
.filter_map(|progress| progress.validation.clone())
.collect();
if tablets.is_empty() {
return None;
}
let total_expected = tablets.iter().map(|report| report.expected_rows).sum();
let total_actual = tablets.iter().map(|report| report.actual_rows).sum();
let passed = tablets.iter().all(TabletValidationReport::passed);
Some(JobValidationReport {
job_id,
tablets,
total_expected,
total_actual,
passed,
})
}
pub fn pin_reader(&self, reader_id: u64, version: MetadataVersion) -> Result<(), DdlError> {
self.apply(DdlCommand::PinReader { reader_id, version })
}
pub fn release_reader(&self, reader_id: u64) -> Result<(), DdlError> {
self.apply(DdlCommand::ReleaseReader { reader_id })
}
pub fn reclaim_index(
&self,
table_id: TableId,
index_name: &str,
) -> Result<ReclaimOutcome, DdlError> {
match self.apply(DdlCommand::ReclaimIndex {
table_id,
index_name: index_name.to_owned(),
}) {
Ok(()) => {
for tablet in self.provider.tablets_of(table_id) {
self.provider.sink(tablet)?.drop_generation(index_name)?;
}
Ok(ReclaimOutcome::Reclaimed)
}
Err(DdlError::Rejection(DdlRejection::ReclaimBlocked {
oldest_reader,
required,
..
})) => Ok(ReclaimOutcome::Blocked {
oldest_reader,
required,
}),
Err(error) => Err(error),
}
}
pub fn drive_job(&self, job_id: u64) -> Result<DriveOutcome, DdlError> {
loop {
match self.step_job(job_id)? {
DriveStep::Parked => return Ok(DriveOutcome::Parked),
DriveStep::RolledBack => return Ok(DriveOutcome::RolledBack),
DriveStep::Terminal => {
let state = self.job_record(job_id)?.state;
return Ok(match state {
SchemaJobState::Succeeded => DriveOutcome::Completed,
_ => DriveOutcome::RolledBack,
});
}
DriveStep::Admitted
| DriveStep::PhaseAdvanced { .. }
| DriveStep::TabletBackfilled { .. }
| DriveStep::TabletValidated { .. }
| DriveStep::Published { .. } => {}
}
}
}
fn step_job(&self, job_id: u64) -> Result<DriveStep, DdlError> {
let job = self.job_record(job_id)?;
match job.state {
SchemaJobState::Pending => {
self.set_state(job_id, SchemaJobState::Running, None)?;
return Ok(DriveStep::Admitted);
}
SchemaJobState::Paused => return Ok(DriveStep::Parked),
SchemaJobState::Cancelling | SchemaJobState::RollingBack => {
self.rollback(&job, "cancelled by operator")?;
return Ok(DriveStep::RolledBack);
}
SchemaJobState::Failed | SchemaJobState::Succeeded => {
return Ok(DriveStep::Terminal);
}
SchemaJobState::Running => {}
}
match job.phase {
DdlPhase::Pending => {
self.advance_phase(job_id, DdlPhase::WriteOnly, None)?;
self.sync_maintainers(job.table_id)?;
Ok(DriveStep::PhaseAdvanced {
from: DdlPhase::Pending,
to: DdlPhase::WriteOnly,
})
}
DdlPhase::WriteOnly => {
let pin = self.clock.now()?;
self.advance_phase(job_id, DdlPhase::Backfilling, Some(pin))?;
Ok(DriveStep::PhaseAdvanced {
from: DdlPhase::WriteOnly,
to: DdlPhase::Backfilling,
})
}
DdlPhase::Backfilling => {
let tablets = self.provider.tablets_of(job.table_id);
for tablet in &tablets {
if !job.tablet_progress.contains_key(tablet) {
self.update_progress(job_id, TabletDdlProgress::pending(*tablet))?;
}
}
let mut job = self.job_record(job_id)?;
let departed: Vec<TabletId> = job
.tablet_progress
.keys()
.filter(|tablet| !tablets.contains(tablet))
.copied()
.collect();
for tablet in departed {
self.apply(DdlCommand::ForgetTabletProgress {
job_id,
tablet_id: tablet,
})?;
}
job = self.job_record(job_id)?;
let next = tablets.iter().copied().find(|tablet| {
job.tablet_progress
.get(tablet)
.is_none_or(|progress| progress.stage < TabletDdlStage::CaughtUp)
});
match next {
Some(tablet) => {
let progress = job.tablet_progress.get(&tablet);
if progress.is_none_or(|p| p.stage < TabletDdlStage::Backfilling) {
let mut marker = TabletDdlProgress::pending(tablet);
marker.stage = TabletDdlStage::Backfilling;
self.update_progress(job_id, marker)?;
job = self.job_record(job_id)?;
}
let progress = self.backfill_tablet(&job, tablet)?;
self.update_progress(job_id, progress)?;
Ok(DriveStep::TabletBackfilled { tablet })
}
None => {
self.advance_phase(job_id, DdlPhase::Validating, None)?;
Ok(DriveStep::PhaseAdvanced {
from: DdlPhase::Backfilling,
to: DdlPhase::Validating,
})
}
}
}
DdlPhase::Validating => {
let next = job
.tablet_progress
.values()
.find(|progress| progress.stage < TabletDdlStage::Validated)
.map(|progress| progress.tablet_id);
match next {
Some(tablet) => {
let report = self.validate_tablet(&job, tablet)?;
let passed = report.passed();
self.apply(DdlCommand::ReportTabletValidation {
job_id,
report: report.clone(),
})?;
if !passed {
let reason = format!(
"tablet {tablet}: expected {} rows, built {} rows",
report.expected_rows, report.actual_rows
);
self.rollback(&job, &reason)?;
return Err(DdlError::ValidationFailed { job_id, reason });
}
Ok(DriveStep::TabletValidated { tablet })
}
None => {
let published_at = self.clock.now()?;
match self.apply(DdlCommand::PublishJob {
job_id,
published_at,
}) {
Ok(()) => Ok(DriveStep::Published {
version: self
.store
.lock()
.expect("store lock poisoned")
.metadata_version,
}),
Err(DdlError::Rejection(
reason @ DdlRejection::SchemaVersionMismatch { .. },
)) => {
self.rollback(&job, &reason.to_string())?;
Err(DdlError::Rejection(reason))
}
Err(error) => Err(error),
}
}
}
}
DdlPhase::Public => match job.kind {
DdlJobKind::DropIndex => {
self.advance_phase(job_id, DdlPhase::Dropping, None)?;
self.sync_maintainers(job.table_id)?;
Ok(DriveStep::PhaseAdvanced {
from: DdlPhase::Public,
to: DdlPhase::Dropping,
})
}
DdlJobKind::AddIndex | DdlJobKind::AlterSchema => Ok(DriveStep::Terminal),
},
DdlPhase::Dropping => {
self.set_state(job_id, SchemaJobState::Succeeded, None)?;
Ok(DriveStep::Terminal)
}
}
}
fn backfill_tablet(
&self,
job: &DdlJobRecord,
tablet: TabletId,
) -> Result<TabletDdlProgress, DdlError> {
let pin = job
.pinned_snapshot
.ok_or(DdlRejection::Meta(MetaRejectionReason::Invalid {
reason: format!("job {} entered Backfilling without a pin", job.job_id),
}))?;
let keyspace = self.provider.keyspace(tablet)?;
let _pin = keyspace.pin_snapshot(pin)?;
let mut rows_scanned = 0_u64;
match &job.definition {
DdlDefinition::AddIndex { index_name, .. } => {
let mut sink = self.provider.sink(tablet)?;
let record = self.index_record(job, index_name)?;
sink.begin_build(index_name)?;
for (key, value) in keyspace.snapshot_at(pin)? {
let entry = (self.projection)(&record, &key, &value);
sink.stage_entry(index_name, &key, &entry)?;
rows_scanned += 1;
}
sink.install_staged(index_name)?;
}
DdlDefinition::AlterSchema { .. } => {
for (_key, _value) in keyspace.snapshot_at(pin)? {
rows_scanned += 1;
}
}
DdlDefinition::DropIndex { .. } => {
return Err(DdlRejection::Meta(MetaRejectionReason::Invalid {
reason: "drop jobs never backfill".to_owned(),
})
.into());
}
}
let watermark = self.catch_up_tablet(job, tablet, pin)?;
Ok(TabletDdlProgress {
tablet_id: tablet,
stage: TabletDdlStage::CaughtUp,
rows_scanned,
caught_up_through: Some(watermark),
validation: None,
})
}
fn catch_up_tablet(
&self,
job: &DdlJobRecord,
tablet: TabletId,
from: HlcTimestamp,
) -> Result<HlcTimestamp, DdlError> {
let keyspace = self.provider.keyspace(tablet)?;
let mut sink = match &job.definition {
DdlDefinition::AddIndex { index_name, .. } => {
Some((self.provider.sink(tablet)?, index_name.clone()))
}
_ => None,
};
let mut watermark = from;
loop {
let mut saw_any = false;
let mut max_ts = watermark;
for (ts, key, value) in keyspace.deltas_after(watermark)? {
saw_any = true;
max_ts = max_ts.max(ts);
if let Some((sink, index_name)) = &mut sink {
let entry = match &value {
Some(value) => {
let record = self.index_record(job, index_name)?;
Some((self.projection)(&record, &key, value))
}
None => None,
};
sink.apply_delta(index_name, &key, entry.as_deref())?;
}
}
if !saw_any {
return Ok(watermark);
}
watermark = max_ts;
}
}
fn validate_tablet(
&self,
job: &DdlJobRecord,
tablet: TabletId,
) -> Result<TabletValidationReport, DdlError> {
let pin = job.pinned_snapshot.expect("Validating implies a pin");
let from = job
.tablet_progress
.get(&tablet)
.and_then(|progress| progress.caught_up_through)
.unwrap_or(pin);
let watermark = self.catch_up_tablet(job, tablet, from)?;
let keyspace = self.provider.keyspace(tablet)?;
let expected: BTreeMap<Key, Vec<u8>> = match &job.definition {
DdlDefinition::AddIndex { index_name, .. } => {
let record = self.index_record(job, index_name)?;
keyspace
.snapshot_at(watermark)?
.map(|(key, value)| {
let entry = (self.projection)(&record, &key, &value);
(key, entry)
})
.collect()
}
_ => keyspace.snapshot_at(watermark)?.collect(),
};
let actual = match &job.definition {
DdlDefinition::AddIndex { index_name, .. } => {
self.provider.sink(tablet)?.generation_entries(index_name)?
}
_ => expected.clone(),
};
Ok(TabletValidationReport {
tablet_id: tablet,
watermark,
expected_rows: expected.len() as u64,
actual_rows: actual.len() as u64,
expected_checksum: generation_checksum(&expected),
actual_checksum: generation_checksum(&actual),
})
}
fn rollback(&self, job: &DdlJobRecord, reason: &str) -> Result<(), DdlError> {
let current = self.job_record(job.job_id)?;
match current.state {
SchemaJobState::Failed => return Ok(()),
SchemaJobState::Succeeded => {
return Err(DdlRejection::Meta(MetaRejectionReason::Conflict {
resource: format!("DDL job {}", job.job_id),
reason: "cannot roll back a succeeded job".to_owned(),
})
.into());
}
SchemaJobState::Running | SchemaJobState::Cancelling => {
self.set_state(
job.job_id,
SchemaJobState::RollingBack,
Some(reason.to_owned()),
)?;
}
SchemaJobState::Pending | SchemaJobState::Paused | SchemaJobState::RollingBack => {}
}
if let DdlDefinition::AddIndex { index_name, .. } = &job.definition {
for tablet in self.provider.tablets_of(job.table_id) {
self.provider.sink(tablet)?.drop_generation(index_name)?;
}
self.apply(DdlCommand::RemoveIndexRecord { job_id: job.job_id })?;
self.sync_maintainers(job.table_id)?;
}
self.set_state(job.job_id, SchemaJobState::Failed, Some(reason.to_owned()))
}
fn apply(&self, command: DdlCommand) -> Result<(), DdlError> {
let commit_ts = self.clock.now()?;
self.store
.lock()
.expect("store lock poisoned")
.apply(&command, None, commit_ts)
.map_err(DdlError::from)
}
fn job_record(&self, job_id: u64) -> Result<DdlJobRecord, DdlError> {
self.store
.lock()
.expect("store lock poisoned")
.job(job_id)
.cloned()
.ok_or_else(|| {
DdlRejection::Meta(MetaRejectionReason::NotFound {
resource: format!("DDL job {job_id}"),
})
.into()
})
}
fn index_record(
&self,
job: &DdlJobRecord,
index_name: &str,
) -> Result<DdlIndexRecord, DdlError> {
self.store
.lock()
.expect("store lock poisoned")
.index(job.table_id, index_name)
.cloned()
.ok_or_else(|| {
DdlRejection::Meta(MetaRejectionReason::NotFound {
resource: format!("index `{index_name}` on table {}", job.table_id),
})
.into()
})
}
fn peek_next_job_id(&self) -> u64 {
self.store.lock().expect("store lock poisoned").next_job_id
}
fn set_state(
&self,
job_id: u64,
state: SchemaJobState,
error: Option<String>,
) -> Result<(), DdlError> {
let updated_at = self.clock.now()?;
self.apply(DdlCommand::SetJobState {
job_id,
state,
updated_at,
error,
expected_version: None,
})
}
fn advance_phase(
&self,
job_id: u64,
to: DdlPhase,
pinned_snapshot: Option<HlcTimestamp>,
) -> Result<(), DdlError> {
self.apply(DdlCommand::AdvancePhase {
job_id,
to,
pinned_snapshot,
expected_version: None,
})
}
fn update_progress(&self, job_id: u64, progress: TabletDdlProgress) -> Result<(), DdlError> {
self.apply(DdlCommand::UpdateTabletProgress { job_id, progress })
}
fn sync_maintainers(&self, table_id: TableId) -> Result<(), DdlError> {
let names: Vec<String> = self
.store
.lock()
.expect("store lock poisoned")
.write_maintained(table_id)
.iter()
.map(|record| record.index_name.clone())
.collect();
for tablet in self.provider.tablets_of(table_id) {
self.provider
.maintainer(tablet)?
.sync_definitions(names.clone())?;
}
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
const TABLE: u64 = 1;
fn ts(micros: u64) -> HlcTimestamp {
HlcTimestamp {
physical_micros: micros,
logical: 0,
node_tiebreaker: 0,
}
}
fn after(pin: HlcTimestamp, n: u64) -> HlcTimestamp {
HlcTimestamp {
physical_micros: pin.physical_micros + n,
logical: 0,
node_tiebreaker: 0,
}
}
fn tablet(n: u8) -> TabletId {
TabletId::from_bytes([n; 16])
}
fn database() -> DatabaseId {
DatabaseId::from_bytes([7; 16])
}
fn table() -> TableId {
TableId(TABLE)
}
fn key(n: u64) -> Key {
Key::from_bytes(n.to_be_bytes().to_vec())
}
fn value(n: u64) -> Vec<u8> {
format!("value-{n}").into_bytes()
}
fn spec() -> serde_json::Value {
serde_json::json!({"kind": "Bitmap", "columns": [1]})
}
struct Fixture {
store: Arc<Mutex<DdlJobStore>>,
tablets: InMemoryDdlTablets,
driver: DdlDriver<InMemoryDdlTablets>,
}
impl Fixture {
fn new(tablet_count: u8, rows_per_tablet: u64) -> Self {
let store = Arc::new(Mutex::new(DdlJobStore::default()));
let tablets = InMemoryDdlTablets::new();
for n in 1..=tablet_count {
tablets.add_tablet(tablet(n), table());
let keyspace = tablets.keyspace_handle(tablet(n));
for row in 1..=rows_per_tablet {
keyspace.insert(key(row), ts(row), value(row));
}
}
let driver = DdlDriver::new(store.clone(), tablets.clone());
driver
.apply(DdlCommand::RegisterTable {
anchor: TableAnchor {
table_id: table(),
database_id: database(),
schema_version: SchemaVersion(1),
schema: serde_json::json!({"columns": ["id", "v"]}),
metadata_version: MetadataVersion::ZERO,
},
})
.expect("register table");
Self {
store,
tablets,
driver,
}
}
fn store(&self) -> std::sync::MutexGuard<'_, DdlJobStore> {
self.store.lock().expect("store lock poisoned")
}
fn submit_add(&self, index_name: &str) -> Result<u64, DdlError> {
self.driver
.submit_add_index(database(), table(), index_name, spec(), SchemaVersion(1))
}
fn pin_of(&self, job_id: u64) -> HlcTimestamp {
self.store()
.job(job_id)
.and_then(|job| job.pinned_snapshot)
.expect("job is pinned")
}
fn drive_until(
&self,
job_id: u64,
mut stop: impl FnMut(DriveStep) -> bool,
) -> Vec<DriveStep> {
let mut steps = Vec::new();
for _ in 0..64 {
let step = self.driver.step_job(job_id).expect("step");
steps.push(step);
if stop(step) || matches!(step, DriveStep::Terminal | DriveStep::Parked) {
return steps;
}
}
panic!("job {job_id} did not reach the expected step");
}
fn drive_to_backfilled(&self, job_id: u64, count: usize) {
let mut backfilled = 0_usize;
self.drive_until(job_id, |step| {
if matches!(step, DriveStep::TabletBackfilled { .. }) {
backfilled += 1;
}
backfilled == count
});
}
}
#[test]
fn add_index_full_lifecycle_over_three_tablets() {
let fixture = Fixture::new(3, 5);
let job_id = fixture.submit_add("idx_a").expect("submit");
let steps = fixture.drive_until(job_id, |step| {
matches!(
step,
DriveStep::PhaseAdvanced {
from: DdlPhase::Pending,
to: DdlPhase::WriteOnly,
}
)
});
assert_eq!(steps.first(), Some(&DriveStep::Admitted));
assert_eq!(fixture.store().write_maintained(table()).len(), 1);
assert!(fixture.store().planner_visible(table()).is_empty());
for n in 1..=3 {
assert_eq!(
fixture.tablets.indexes_handle(tablet(n)).maintained(),
["idx_a"]
);
}
fixture.drive_until(job_id, |step| {
matches!(
step,
DriveStep::PhaseAdvanced {
from: DdlPhase::WriteOnly,
to: DdlPhase::Backfilling,
}
)
});
let pin = fixture.pin_of(job_id);
fixture
.tablets
.commit_write(tablet(1), key(100), after(pin, 1), Some(value(100)))
.expect("write 100");
fixture.drive_to_backfilled(job_id, 1);
assert!(fixture
.tablets
.indexes_handle(tablet(1))
.is_installed("idx_a"));
fixture
.tablets
.commit_write(tablet(1), key(102), after(pin, 3), Some(value(102)))
.expect("write 102");
assert_eq!(
fixture
.tablets
.indexes_handle(tablet(1))
.entries("idx_a")
.len(),
7,
);
fixture
.tablets
.commit_write(tablet(2), key(101), after(pin, 2), Some(value(101)))
.expect("write 101");
fixture
.tablets
.commit_write(tablet(1), key(1), after(pin, 4), None)
.expect("delete 1");
fixture.drive_to_backfilled(job_id, 2);
let steps = fixture.drive_until(job_id, |step| matches!(step, DriveStep::Published { .. }));
let publication_version = steps.iter().find_map(|step| match step {
DriveStep::Published { version } => Some(*version),
_ => None,
});
let publication_version = publication_version.expect("published");
let store = fixture.store();
let job = store.job(job_id).expect("job");
assert_eq!(job.state, SchemaJobState::Succeeded);
assert_eq!(job.phase, DdlPhase::Public);
let record = store.index(table(), "idx_a").expect("index record");
assert_eq!(record.phase, DdlPhase::Public);
assert_eq!(record.publication_version, Some(publication_version));
assert!(store
.planner_visible_at(table(), MetadataVersion(publication_version.get() - 1))
.is_empty());
assert_eq!(
store.planner_visible_at(table(), publication_version).len(),
1
);
drop(store);
let report = fixture
.driver
.validation_report(job_id)
.expect("validation report");
assert!(report.passed);
assert_eq!(report.tablets.len(), 3);
assert_eq!(report.total_expected, report.total_actual);
let per_tablet: BTreeMap<TabletId, &TabletValidationReport> = report
.tablets
.iter()
.map(|tablet_report| (tablet_report.tablet_id, tablet_report))
.collect();
assert_eq!(per_tablet[&tablet(1)].expected_rows, 6); assert_eq!(per_tablet[&tablet(2)].expected_rows, 6);
assert_eq!(per_tablet[&tablet(3)].expected_rows, 5);
let watermark = per_tablet[&tablet(1)].watermark;
let expected = fixture
.tablets
.keyspace_handle(tablet(1))
.rows_at(watermark);
let actual = fixture.tablets.indexes_handle(tablet(1)).entries("idx_a");
assert_eq!(actual, expected);
let status = fixture.driver.job_status(job_id).expect("status");
assert_eq!(status.tablets_registered, 3);
assert_eq!(status.tablets_caught_up, 3);
assert_eq!(status.tablets_validated, 3);
assert_eq!(status.rows_backfilled, 15);
}
#[test]
fn catch_up_covers_pre_install_writes_and_tombstones() {
let fixture = Fixture::new(1, 3);
let job_id = fixture.submit_add("idx_t").expect("submit");
fixture.drive_until(job_id, |step| {
matches!(
step,
DriveStep::PhaseAdvanced {
from: DdlPhase::WriteOnly,
to: DdlPhase::Backfilling,
}
)
});
let pin = fixture.pin_of(job_id);
fixture
.tablets
.commit_write(tablet(1), key(50), after(pin, 1), Some(value(50)))
.expect("write 50");
fixture
.tablets
.commit_write(tablet(1), key(50), after(pin, 2), None)
.expect("delete 50");
fixture
.tablets
.commit_write(tablet(1), key(2), after(pin, 3), None)
.expect("delete 2");
assert_eq!(
fixture.driver.drive_job(job_id).expect("drive"),
DriveOutcome::Completed
);
let report = fixture.driver.validation_report(job_id).expect("report");
assert!(report.passed);
assert_eq!(report.total_expected, 2); let actual = fixture.tablets.indexes_handle(tablet(1)).entries("idx_t");
assert_eq!(actual.len(), 2);
assert!(!actual.contains_key(&key(2)));
assert!(!actual.contains_key(&key(50)));
}
#[test]
fn atomic_publish_flips_planner_visibility_at_one_metadata_version() {
let fixture = Fixture::new(1, 2);
let job_id = fixture.submit_add("idx_v").expect("submit");
let before = fixture.store().metadata_version;
assert_eq!(
fixture.driver.drive_job(job_id).expect("drive"),
DriveOutcome::Completed
);
let store = fixture.store();
let publication = store
.index(table(), "idx_v")
.and_then(|record| record.publication_version)
.expect("publication version");
assert!(publication > before);
assert!(
store
.planner_visible_at(table(), MetadataVersion(publication.get() - 1))
.is_empty(),
"hidden before the publication version"
);
let visible = store.planner_visible_at(table(), publication);
assert_eq!(visible.len(), 1);
assert_eq!(visible[0].index_name, "idx_v");
}
#[test]
fn stale_schema_returns_structured_retry() {
let fixture = Fixture::new(1, 2);
let job_id = fixture.submit_add("idx_s").expect("submit");
fixture.drive_until(job_id, |step| {
matches!(
step,
DriveStep::PhaseAdvanced {
from: DdlPhase::Backfilling,
to: DdlPhase::Validating,
}
)
});
fixture
.driver
.apply(DdlCommand::RegisterTable {
anchor: TableAnchor {
table_id: table(),
database_id: database(),
schema_version: SchemaVersion(2),
schema: serde_json::json!({"columns": ["id", "v", "w"]}),
metadata_version: MetadataVersion::ZERO,
},
})
.expect("schema bump");
let error = fixture
.driver
.drive_job(job_id)
.expect_err("publish must refuse the stale schema");
assert_eq!(error.category(), ErrorCategory::SchemaVersionMismatch);
match &error {
DdlError::Rejection(DdlRejection::SchemaVersionMismatch {
table_id,
expected,
found,
}) => {
assert_eq!(*table_id, table());
assert_eq!(*expected, SchemaVersion(1));
assert_eq!(*found, SchemaVersion(2));
}
other => panic!("expected SchemaVersionMismatch, got {other:?}"),
}
let store = fixture.store();
let job = store.job(job_id).expect("job");
assert_eq!(job.state, SchemaJobState::Failed);
assert!(store.index(table(), "idx_s").is_none());
drop(store);
assert!(!fixture
.tablets
.indexes_handle(tablet(1))
.is_installed("idx_s"));
let stale = fixture.driver.submit_add_index(
database(),
table(),
"idx_s2",
spec(),
SchemaVersion(1),
);
assert!(matches!(
stale,
Err(DdlError::Rejection(
DdlRejection::SchemaVersionMismatch { .. }
))
));
let job_id = fixture
.driver
.submit_add_index(database(), table(), "idx_s2", spec(), SchemaVersion(2))
.expect("resubmit at the current schema version");
assert_eq!(
fixture.driver.drive_job(job_id).expect("drive"),
DriveOutcome::Completed
);
}
#[test]
fn pause_resume_mid_backfill_resumes_exactly() {
let fixture = Fixture::new(3, 4);
let job_id = fixture.submit_add("idx_p").expect("submit");
fixture.drive_to_backfilled(job_id, 1);
fixture.driver.pause_job(job_id).expect("pause");
assert!(
fixture.driver.pause_job(job_id).is_err(),
"double pause is refused"
);
assert_eq!(
fixture.driver.drive_job(job_id).expect("drive"),
DriveOutcome::Parked
);
assert_eq!(
fixture
.tablets
.indexes_handle(tablet(1))
.begin_build_count(),
1
);
assert_eq!(
fixture
.tablets
.indexes_handle(tablet(2))
.begin_build_count(),
0
);
fixture.driver.resume_job(job_id).expect("resume");
assert_eq!(
fixture.driver.drive_job(job_id).expect("drive"),
DriveOutcome::Completed
);
assert_eq!(
fixture
.tablets
.indexes_handle(tablet(1))
.begin_build_count(),
1
);
assert_eq!(
fixture
.tablets
.indexes_handle(tablet(2))
.begin_build_count(),
1
);
assert_eq!(
fixture
.tablets
.indexes_handle(tablet(3))
.begin_build_count(),
1
);
assert!(
fixture
.driver
.validation_report(job_id)
.expect("report")
.passed
);
assert!(
fixture.driver.resume_job(job_id).is_err(),
"resume of a terminal job is refused"
);
}
#[test]
fn cancel_unwinds_the_hidden_generation_and_leaves_the_table_unaffected() {
let fixture = Fixture::new(3, 4);
let job_id = fixture.submit_add("idx_c").expect("submit");
fixture.drive_to_backfilled(job_id, 1);
assert!(fixture
.tablets
.indexes_handle(tablet(1))
.is_installed("idx_c"));
fixture.driver.cancel_job(job_id).expect("cancel");
assert_eq!(
fixture.driver.drive_job(job_id).expect("drive"),
DriveOutcome::RolledBack
);
let store = fixture.store();
let job = store.job(job_id).expect("job");
assert_eq!(job.state, SchemaJobState::Failed);
assert_eq!(job.error.as_deref(), Some("cancelled by operator"));
assert!(store.index(table(), "idx_c").is_none());
assert!(store.planner_visible(table()).is_empty());
assert!(store.write_maintained(table()).is_empty());
assert_eq!(
store.table_anchor(table()).expect("anchor").schema_version,
SchemaVersion(1)
);
drop(store);
for n in 1..=3 {
let indexes = fixture.tablets.indexes_handle(tablet(n));
assert!(!indexes.is_installed("idx_c"));
assert_eq!(indexes.entries("idx_c").len(), 0);
assert_eq!(indexes.maintained(), Vec::<String>::new());
}
assert_eq!(
fixture
.tablets
.keyspace_handle(tablet(1))
.rows_at(ts(u64::MAX))
.len(),
4
);
let second = fixture.submit_add("idx_c2").expect("submit after unwind");
assert_eq!(
fixture.driver.drive_job(second).expect("drive"),
DriveOutcome::Completed
);
assert!(fixture.driver.cancel_job(second).is_err());
}
#[test]
fn drop_index_lifecycle_reclaims_after_readers_drain() {
let fixture = Fixture::new(3, 3);
let add = fixture.submit_add("idx_d").expect("submit add");
assert_eq!(
fixture.driver.drive_job(add).expect("drive"),
DriveOutcome::Completed
);
let pin = fixture.pin_of(add);
fixture
.tablets
.commit_write(tablet(1), key(40), after(pin, 1), Some(value(40)))
.expect("maintained write");
assert_eq!(
fixture
.tablets
.indexes_handle(tablet(1))
.entries("idx_d")
.len(),
4
);
let pre_drop_version = fixture.store().metadata_version;
let drop = fixture
.driver
.submit_drop_index(database(), table(), "idx_d", SchemaVersion(1))
.expect("submit drop");
assert_eq!(
fixture.driver.drive_job(drop).expect("drive"),
DriveOutcome::Completed
);
let dropping_since = {
let store = fixture.store();
let record = store.index(table(), "idx_d").expect("record");
assert_eq!(record.phase, DdlPhase::Dropping);
assert_eq!(
store.job(drop).expect("job").state,
SchemaJobState::Succeeded
);
assert!(store.planner_visible(table()).is_empty());
assert_eq!(store.planner_visible_at(table(), pre_drop_version).len(), 1);
record.dropping_since.expect("dropping version")
};
assert!(dropping_since > pre_drop_version);
assert_eq!(
fixture.tablets.indexes_handle(tablet(1)).maintained(),
Vec::<String>::new()
);
fixture
.tablets
.commit_write(tablet(1), key(41), after(pin, 2), Some(value(41)))
.expect("write after drop");
assert_eq!(
fixture
.tablets
.indexes_handle(tablet(1))
.entries("idx_d")
.len(),
4
);
fixture
.driver
.pin_reader(1, pre_drop_version)
.expect("pin reader");
let blocked = fixture
.driver
.reclaim_index(table(), "idx_d")
.expect("reclaim attempt");
assert_eq!(
blocked,
ReclaimOutcome::Blocked {
oldest_reader: Some(pre_drop_version),
required: dropping_since,
}
);
let forward_version = fixture.store().metadata_version;
fixture
.driver
.pin_reader(1, forward_version)
.expect("reader moves forward");
assert_eq!(
fixture
.driver
.reclaim_index(table(), "idx_d")
.expect("reclaim"),
ReclaimOutcome::Reclaimed
);
assert!(fixture.store().index(table(), "idx_d").is_none());
for n in 1..=3 {
assert!(!fixture
.tablets
.indexes_handle(tablet(n))
.is_installed("idx_d"));
}
fixture
.driver
.pin_reader(2, dropping_since)
.expect("pin reader 2");
assert!(fixture.driver.pin_reader(2, pre_drop_version).is_err());
}
#[test]
fn per_tablet_progress_survives_a_driver_crash() {
let fixture = Fixture::new(3, 4);
let job_id = fixture.submit_add("idx_x").expect("submit");
fixture.drive_to_backfilled(job_id, 1);
let pin = fixture.pin_of(job_id);
let crashed = fixture.driver;
fixture.tablets.indexes_handle(tablet(2)).seed_staged(
"idx_x",
key(999),
b"garbage".to_vec(),
);
fixture
.tablets
.commit_write(tablet(2), key(60), after(pin, 1), Some(value(60)))
.expect("write during outage");
drop(crashed);
let resumed = DdlDriver::new(fixture.store.clone(), fixture.tablets.clone());
assert_eq!(
resumed.drive_job(job_id).expect("drive"),
DriveOutcome::Completed
);
assert_eq!(
fixture
.tablets
.indexes_handle(tablet(1))
.begin_build_count(),
1
);
assert_eq!(
fixture
.tablets
.indexes_handle(tablet(2))
.begin_build_count(),
1
);
assert!(!fixture
.tablets
.indexes_handle(tablet(2))
.entries("idx_x")
.contains_key(&key(999)));
assert_eq!(
fixture
.tablets
.indexes_handle(tablet(2))
.entries("idx_x")
.len(),
5
);
assert!(resumed.validation_report(job_id).expect("report").passed);
}
#[test]
fn phase_graph_is_enforced() {
let edges: [(DdlPhase, DdlPhase); 5] = [
(DdlPhase::Pending, DdlPhase::WriteOnly),
(DdlPhase::WriteOnly, DdlPhase::Backfilling),
(DdlPhase::Backfilling, DdlPhase::Validating),
(DdlPhase::Validating, DdlPhase::Public),
(DdlPhase::Public, DdlPhase::Dropping),
];
for from in DdlPhase::ALL {
for to in DdlPhase::ALL {
assert_eq!(
from.can_transition(to),
edges.contains(&(from, to)),
"{from} -> {to}",
);
}
}
let fixture = Fixture::new(1, 1);
let job_id = fixture.submit_add("idx_g").expect("submit");
fixture.drive_until(job_id, |step| matches!(step, DriveStep::Admitted));
let skipped = fixture
.driver
.advance_phase(job_id, DdlPhase::Validating, None);
assert!(matches!(
skipped,
Err(DdlError::Rejection(DdlRejection::IllegalPhaseTransition {
from: DdlPhase::Pending,
to: DdlPhase::Validating,
..
}))
));
fixture.drive_to_backfilled(job_id, 1);
fixture.drive_until(job_id, |step| {
matches!(
step,
DriveStep::PhaseAdvanced {
from: DdlPhase::Backfilling,
to: DdlPhase::Validating,
}
)
});
let bypass = fixture.driver.advance_phase(job_id, DdlPhase::Public, None);
assert!(
matches!(
bypass,
Err(DdlError::Rejection(DdlRejection::Meta(
MetaRejectionReason::Invalid { .. }
)))
),
"Public rides PublishJob: {bypass:?}",
);
}
#[test]
fn admin_state_graph_is_enforced() {
let fixture = Fixture::new(1, 1);
let job_id = fixture.submit_add("idx_a").expect("submit");
assert!(
fixture.driver.pause_job(job_id).is_err(),
"Pending cannot pause"
);
assert!(
fixture.driver.resume_job(job_id).is_err(),
"Pending cannot resume"
);
fixture
.driver
.cancel_job(job_id)
.expect("cancel from Pending");
assert_eq!(
fixture.driver.drive_job(job_id).expect("drive"),
DriveOutcome::RolledBack
);
let job = fixture.store().job(job_id).expect("job").clone();
assert_eq!(job.state, SchemaJobState::Failed);
assert!(
fixture.driver.cancel_job(job_id).is_err(),
"terminal jobs refuse cancel"
);
}
#[test]
fn alter_schema_publishes_the_next_schema_version_atomically() {
let fixture = Fixture::new(2, 3);
let target = serde_json::json!({"columns": ["id", "v", "w"]});
let job_id = fixture
.driver
.submit_alter_schema(database(), table(), target.clone(), SchemaVersion(1))
.expect("submit alter");
assert_eq!(
fixture.driver.drive_job(job_id).expect("drive"),
DriveOutcome::Completed
);
let store = fixture.store();
let anchor = store.table_anchor(table()).expect("anchor");
assert_eq!(anchor.schema_version, SchemaVersion(2));
assert_eq!(anchor.schema, target);
let job = store.job(job_id).expect("job");
assert_eq!(job.state, SchemaJobState::Succeeded);
assert_eq!(job.phase, DdlPhase::Public);
assert_eq!(job.tablet_progress.len(), 2);
drop(store);
let stale = fixture.submit_add("idx_old");
assert!(matches!(
stale,
Err(DdlError::Rejection(
DdlRejection::SchemaVersionMismatch { .. }
))
));
let fresh = fixture
.driver
.submit_add_index(database(), table(), "idx_new", spec(), SchemaVersion(2))
.expect("submit at v2");
assert_eq!(
fixture.driver.drive_job(fresh).expect("drive"),
DriveOutcome::Completed
);
}
#[test]
fn topology_change_mid_backfill_is_reconciled() {
let fixture = Fixture::new(3, 3);
let job_id = fixture.submit_add("idx_t").expect("submit");
fixture.drive_to_backfilled(job_id, 1);
fixture.tablets.add_tablet(tablet(4), table());
let keyspace = fixture.tablets.keyspace_handle(tablet(4));
for row in 1..=3 {
keyspace.insert(key(row + 100), ts(row), value(row + 100));
}
fixture.tablets.remove_tablet(tablet(3));
assert_eq!(
fixture.driver.drive_job(job_id).expect("drive"),
DriveOutcome::Completed
);
let status = fixture.driver.job_status(job_id).expect("status");
assert_eq!(status.tablets_registered, 3);
assert!(status.record.tablet_progress.contains_key(&tablet(4)));
assert!(!status.record.tablet_progress.contains_key(&tablet(3)));
assert!(
fixture
.driver
.validation_report(job_id)
.expect("report")
.passed
);
}
#[test]
fn conflicting_and_duplicate_submissions_are_refused() {
let fixture = Fixture::new(1, 2);
let job_id = fixture.submit_add("idx_1").expect("submit");
let concurrent = fixture.submit_add("idx_2");
assert!(matches!(
concurrent,
Err(DdlError::Rejection(DdlRejection::Meta(
MetaRejectionReason::Conflict { .. }
)))
));
assert_eq!(
fixture.driver.drive_job(job_id).expect("drive"),
DriveOutcome::Completed
);
let duplicate = fixture.submit_add("idx_1");
assert!(matches!(
duplicate,
Err(DdlError::Rejection(DdlRejection::Meta(
MetaRejectionReason::Conflict { .. }
)))
));
let missing =
fixture
.driver
.submit_drop_index(database(), table(), "idx_missing", SchemaVersion(1));
assert!(matches!(
missing,
Err(DdlError::Rejection(DdlRejection::Meta(
MetaRejectionReason::NotFound { .. }
)))
));
let job = fixture.store().job(job_id).expect("job").clone();
fixture
.driver
.apply(DdlCommand::SubmitJob { job })
.expect("identical replay is a no-op");
assert_eq!(fixture.store().jobs.len(), 1);
}
#[test]
fn store_roundtrips_through_json() {
let fixture = Fixture::new(2, 3);
let job_id = fixture.submit_add("idx_j").expect("submit");
fixture.drive_to_backfilled(job_id, 1);
let snapshot = {
let store = fixture.store();
serde_json::to_string(&*store).expect("encode")
};
let decoded: DdlJobStore = serde_json::from_str(&snapshot).expect("decode");
assert_eq!(decoded, *fixture.store());
}
#[test]
fn command_roundtrips_through_json() {
let job = DdlJobRecord {
job_id: 9,
database_id: database(),
table_id: table(),
kind: DdlJobKind::AddIndex,
state: SchemaJobState::Pending,
phase: DdlPhase::Pending,
definition: DdlDefinition::AddIndex {
index_name: "idx_serde".to_owned(),
spec: spec(),
},
source_schema_version: SchemaVersion(1),
created_at: ts(10),
updated_at: ts(10),
pinned_snapshot: None,
tablet_progress: BTreeMap::new(),
error: None,
metadata_version: MetadataVersion::ZERO,
};
let command = DdlCommand::SubmitJob { job };
let encoded = serde_json::to_string(&command).expect("encode");
let decoded: DdlCommand = serde_json::from_str(&encoded).expect("decode");
assert_eq!(decoded, command);
}
}