use serde::{Deserialize, Deserializer, Serialize, Serializer};
use uuid::Uuid;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum IndexDdlReceipt {
Accepted {
#[serde(with = "uuid_string")]
operation_id: Uuid,
#[serde(with = "positive_u64_string")]
index_id: u64,
#[serde(with = "positive_u64_string")]
generation: u64,
},
ExistingOperation {
#[serde(with = "uuid_string")]
operation_id: Uuid,
},
AlreadyActive {
#[serde(with = "positive_u64_string")]
index_id: u64,
#[serde(with = "positive_u64_string")]
generation: u64,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum IndexOperationKind {
Build,
Drop,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum IndexFamily {
Secondary,
Vector,
Text,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum IndexOperationStage {
Scan,
ScanPartitions,
CatchUp,
Validate,
ValidateDescriptor,
ValidateLegacyPhysical,
Compact,
PrepareManifests,
ValidateManifests,
Activate,
DeleteEntries,
RetireCache,
DeletePhysical,
DeleteDeltas,
DeleteMetadata,
Finalize,
AbortingDeleteEntries,
AbortingRetireCache,
AbortingDeletePhysical,
AbortingDeleteDeltas,
AbortingDeleteMetadata,
AbortingFinalize,
}
impl IndexOperationStage {
pub const fn is_aborting(self) -> bool {
matches!(
self,
Self::AbortingDeleteEntries
| Self::AbortingRetireCache
| Self::AbortingDeletePhysical
| Self::AbortingDeleteDeltas
| Self::AbortingDeleteMetadata
| Self::AbortingFinalize
)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum IndexOperationBlockerCode {
InvalidSourceData,
UniquenessViolation,
OversizedEntity,
ManifestLimit,
ObjectStoreConfigurationUnavailable,
InvariantViolation,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum IndexErrorCode {
IndexLifecycleUnavailable,
IndexAlreadyExists,
IndexDefinitionConflict,
IndexBusy,
IndexNotFound,
IndexOperationNotFound,
IndexOperationNotAbortable,
IndexIdExhausted,
VectorPhysicalIdExhausted,
IndexGenerationExhausted,
IndexRevisionExhausted,
IndexOperationRevisionExhausted,
StaleIndexGeneration,
WriterFencedCommitOutcomeUnknown,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct IndexOperationProgress {
#[serde(with = "u64_string")]
pub entities: u64,
#[serde(with = "u64_string")]
pub input_bytes: u64,
#[serde(with = "u64_string")]
pub output_operations: u64,
#[serde(with = "u64_string")]
pub output_bytes: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct IndexOperationStatusCommon {
#[serde(with = "uuid_string")]
pub operation_id: Uuid,
#[serde(with = "positive_u64_string")]
pub index_id: u64,
#[serde(with = "positive_u64_string")]
pub generation: u64,
pub operation_kind: IndexOperationKind,
pub family: IndexFamily,
pub stage: IndexOperationStage,
pub attempt: u32,
pub progress: IndexOperationProgress,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(tag = "status", rename_all = "snake_case")]
pub enum IndexOperationStatus {
Queued {
#[serde(flatten)]
common: IndexOperationStatusCommon,
},
Running {
#[serde(flatten)]
common: IndexOperationStatusCommon,
},
Blocked {
#[serde(flatten)]
common: IndexOperationStatusCommon,
blocker_code: IndexOperationBlockerCode,
#[serde(default)]
message: Option<String>,
},
Succeeded {
#[serde(flatten)]
common: IndexOperationStatusCommon,
},
Aborted {
#[serde(flatten)]
common: IndexOperationStatusCommon,
},
}
#[derive(Deserialize)]
#[serde(tag = "status", rename_all = "snake_case")]
enum IndexOperationStatusWire {
Queued {
#[serde(flatten)]
common: IndexOperationStatusCommon,
},
Running {
#[serde(flatten)]
common: IndexOperationStatusCommon,
},
Blocked {
#[serde(flatten)]
common: IndexOperationStatusCommon,
blocker_code: IndexOperationBlockerCode,
#[serde(default)]
message: Option<String>,
},
Succeeded {
#[serde(flatten)]
common: IndexOperationStatusCommon,
},
Aborted {
#[serde(flatten)]
common: IndexOperationStatusCommon,
},
}
impl<'de> Deserialize<'de> for IndexOperationStatus {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: Deserializer<'de>,
{
let status = IndexOperationStatusWire::deserialize(deserializer)?;
Ok(match status {
IndexOperationStatusWire::Queued { common } => Self::Queued { common },
IndexOperationStatusWire::Running { common } => Self::Running { common },
IndexOperationStatusWire::Blocked {
common,
blocker_code,
message,
} => Self::Blocked {
common,
blocker_code,
message,
},
IndexOperationStatusWire::Succeeded { common } => Self::Succeeded { common },
IndexOperationStatusWire::Aborted { common } => {
if common.operation_kind != IndexOperationKind::Build || !common.stage.is_aborting()
{
return Err(serde::de::Error::custom(
"aborted status must describe build cleanup",
));
}
Self::Aborted { common }
}
})
}
}
mod u64_string {
use super::*;
pub(super) fn serialize<S>(value: &u64, serializer: S) -> Result<S::Ok, S::Error>
where
S: Serializer,
{
serializer.serialize_str(&value.to_string())
}
pub(super) fn deserialize<'de, D>(deserializer: D) -> Result<u64, D::Error>
where
D: Deserializer<'de>,
{
let value = String::deserialize(deserializer)?;
let parsed = value.parse::<u64>().map_err(serde::de::Error::custom)?;
if parsed.to_string() != value {
return Err(serde::de::Error::custom(
"expected a canonical unsigned decimal string",
));
}
Ok(parsed)
}
}
mod positive_u64_string {
use super::*;
pub(super) fn serialize<S>(value: &u64, serializer: S) -> Result<S::Ok, S::Error>
where
S: Serializer,
{
serializer.serialize_str(&value.to_string())
}
pub(super) fn deserialize<'de, D>(deserializer: D) -> Result<u64, D::Error>
where
D: Deserializer<'de>,
{
let value = u64_string::deserialize(deserializer)?;
if value == 0 {
return Err(serde::de::Error::custom("identifier must be non-zero"));
}
Ok(value)
}
}
mod uuid_string {
use super::*;
pub(super) fn serialize<S>(value: &Uuid, serializer: S) -> Result<S::Ok, S::Error>
where
S: Serializer,
{
serializer.serialize_str(&value.to_string())
}
pub(super) fn deserialize<'de, D>(deserializer: D) -> Result<Uuid, D::Error>
where
D: Deserializer<'de>,
{
let value = String::deserialize(deserializer)?;
let parsed = Uuid::parse_str(&value).map_err(serde::de::Error::custom)?;
if parsed.is_nil() || parsed.to_string() != value {
return Err(serde::de::Error::custom(
"operation ID must be a canonical lowercase non-nil UUID",
));
}
Ok(parsed)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn response_decoders_accept_additive_fields_and_reject_invalid_required_fields() {
let receipt: IndexDdlReceipt = sonic_rs::from_str(
r#"{"kind":"accepted","operation_id":"018f0c58-6bc7-7c56-8d3d-9c5f18a0f001","index_id":"42","generation":"3","future":true}"#,
)
.unwrap();
assert!(matches!(
receipt,
IndexDdlReceipt::Accepted { index_id: 42, .. }
));
let status: IndexOperationStatus = sonic_rs::from_str(
r#"{"status":"blocked","operation_id":"018f0c58-6bc7-7c56-8d3d-9c5f18a0f001","index_id":"42","generation":"3","operation_kind":"build","family":"secondary","stage":"scan","attempt":2,"progress":{"entities":"9","input_bytes":"10","output_operations":"11","output_bytes":"12","future":true},"blocker_code":"uniqueness_violation","future":true}"#,
)
.unwrap();
assert!(matches!(status, IndexOperationStatus::Blocked { .. }));
for (stage, expected) in [
(
"validate_legacy_physical",
IndexOperationStage::ValidateLegacyPhysical,
),
("validate_manifests", IndexOperationStage::ValidateManifests),
] {
let status: IndexOperationStatus = sonic_rs::from_str(&format!(
r#"{{"status":"queued","operation_id":"018f0c58-6bc7-7c56-8d3d-9c5f18a0f001","index_id":"42","generation":"3","operation_kind":"build","family":"text","stage":"{stage}","attempt":0,"progress":{{"entities":"0","input_bytes":"0","output_operations":"0","output_bytes":"0"}}}}"#,
))
.unwrap();
let IndexOperationStatus::Queued { common } = status else {
panic!("valid text build stage must decode as queued");
};
assert_eq!(common.stage, expected);
}
let aborted: IndexOperationStatus = sonic_rs::from_str(
r#"{"status":"aborted","operation_id":"018f0c58-6bc7-7c56-8d3d-9c5f18a0f001","index_id":"42","generation":"3","operation_kind":"build","family":"secondary","stage":"aborting_finalize","attempt":2,"progress":{"entities":"9","input_bytes":"10","output_operations":"11","output_bytes":"12"}}"#,
)
.unwrap();
assert!(matches!(aborted, IndexOperationStatus::Aborted { .. }));
assert!(sonic_rs::from_str::<IndexDdlReceipt>(
r#"{"kind":"accepted","operation_id":"018F0C58-6BC7-7C56-8D3D-9C5F18A0F001","index_id":"42","generation":"3"}"#,
)
.is_err());
assert!(sonic_rs::from_str::<IndexDdlReceipt>(
r#"{"kind":"already_active","index_id":"0","generation":"03"}"#,
)
.is_err());
assert!(sonic_rs::from_str::<IndexOperationStatus>(r#"{"status":"future"}"#).is_err());
assert!(sonic_rs::from_str::<IndexOperationStatus>(
r#"{"status":"queued","operation_id":"018f0c58-6bc7-7c56-8d3d-9c5f18a0f001","index_id":"42","generation":"3","operation_kind":"build","family":"secondary","stage":"future","attempt":0,"progress":{"entities":"0","input_bytes":"0","output_operations":"0","output_bytes":"0"}}"#,
)
.is_err());
assert!(sonic_rs::from_str::<IndexOperationStatus>(
r#"{"status":"aborted","operation_id":"018f0c58-6bc7-7c56-8d3d-9c5f18a0f001","index_id":"42","generation":"3","operation_kind":"drop","family":"secondary","stage":"finalize","attempt":0,"progress":{"entities":"0","input_bytes":"0","output_operations":"0","output_bytes":"0"}}"#,
)
.is_err());
}
}