use std::time::Duration;
use tracing::warn;
use nodedb_types::sync::wire::{EngineKind, SyncProvenance, stream_id_for};
use crate::control::state::SharedState;
use crate::types::TenantId;
use super::super::wire::{
CompensationHint, DeltaPushMsg, DeltaRejectMsg, SyncFrame, SyncMessageType,
};
pub(crate) async fn apply_delta_and_finalize(
shared: &SharedState,
delta_msg: &DeltaPushMsg,
ack_frame: SyncFrame,
session_tenant: TenantId,
session_producer_id: u64,
session_epoch: u64,
) -> Option<SyncFrame> {
use crate::bridge::envelope::PhysicalPlan;
use nodedb_physical::physical_plan::CrdtOp;
let tenant_id = session_tenant;
if let Err(e) = shared.check_tenant_quota(tenant_id) {
warn!(error = %e, "sync: delta validation rejected by quota");
let reject = DeltaRejectMsg {
mutation_id: delta_msg.mutation_id,
reason: e.to_string(),
compensation: Some(CompensationHint::Custom {
constraint: "quota".into(),
detail: e.to_string(),
}),
};
return SyncFrame::try_encode(SyncMessageType::DeltaReject, &reject);
}
let surrogate = match shared.surrogate_assigner.assign(
crate::types::DatabaseId::DEFAULT,
tenant_id,
&delta_msg.collection,
delta_msg.document_id.as_bytes(),
) {
Ok(s) => s,
Err(e) => {
warn!(error = %e, "sync: surrogate assignment failed");
let reject = DeltaRejectMsg {
mutation_id: delta_msg.mutation_id,
reason: e.to_string(),
compensation: Some(CompensationHint::Custom {
constraint: "surrogate".into(),
detail: e.to_string(),
}),
};
return SyncFrame::try_encode(SyncMessageType::DeltaReject, &reject);
}
};
let prov = SyncProvenance {
producer_id: session_producer_id,
epoch: session_epoch,
stream_id: stream_id_for(EngineKind::Crdt, &delta_msg.collection),
seq: delta_msg.seq,
};
let constraint_version_required = shared
.credentials
.catalog()
.get_collection(
crate::types::DatabaseId::DEFAULT,
tenant_id.as_u64(),
&delta_msg.collection,
)
.ok()
.flatten()
.map(|col| col.constraint_version)
.unwrap_or(0);
let plan = PhysicalPlan::Crdt(CrdtOp::Apply {
collection: delta_msg.collection.clone(),
document_id: delta_msg.document_id.clone(),
delta: delta_msg.delta.clone(),
peer_id: delta_msg.peer_id,
mutation_id: delta_msg.mutation_id,
surrogate,
provenance: Some(prov),
constraint_version_required,
});
shared.tenant_request_start(tenant_id);
let dispatch_result = super::super::raft_dispatch::dispatch_sync_bytes(
shared,
tenant_id,
&delta_msg.collection,
plan,
Duration::from_secs(10),
crate::event::EventSource::CrdtSync,
)
.await;
shared.tenant_request_end(tenant_id);
match dispatch_result {
Ok(payload) => {
let gate_result = match zerompk::from_msgpack::<nodedb_types::sync::wire::SyncAckResult>(
&payload,
) {
Ok(r) => r,
Err(err) => {
warn!(
collection = %delta_msg.collection,
error = %err,
"sync: failed to decode SyncAckResult from Data Plane; using default ack"
);
return Some(ack_frame);
}
};
if let Some(violation) = gate_result.reject {
let hint = violation.to_compensation_hint();
warn!(
collection = %delta_msg.collection,
doc = %delta_msg.document_id,
violation = %violation,
"sync: delta applied but rejected by CRDT validator"
);
let reject = DeltaRejectMsg {
mutation_id: delta_msg.mutation_id,
reason: violation.to_string(),
compensation: Some(hint),
};
return SyncFrame::try_encode(SyncMessageType::DeltaReject, &reject);
}
let (mutation_id, clock_skew_warning_ms) = if let Some(existing_ack) =
ack_frame.decode_body::<super::super::wire::DeltaAckMsg>()
{
(existing_ack.mutation_id, existing_ack.clock_skew_warning_ms)
} else {
(delta_msg.mutation_id, None)
};
let ack = super::super::wire::DeltaAckMsg {
mutation_id,
lsn: 0, clock_skew_warning_ms,
applied_seq: gate_result.applied_seq,
status: gate_result.status,
};
SyncFrame::try_encode(SyncMessageType::DeltaAck, &ack)
}
Err(e) => {
let hint = compensation_hint_for_dispatch_error(&e);
warn!(
collection = %delta_msg.collection,
doc = %delta_msg.document_id,
hint = hint.code(),
error = %e,
"sync: delta rejected by Data Plane"
);
let reject = DeltaRejectMsg {
mutation_id: delta_msg.mutation_id,
reason: e.to_string(),
compensation: Some(hint),
};
SyncFrame::try_encode(SyncMessageType::DeltaReject, &reject)
}
}
}
fn compensation_hint_for_dispatch_error(e: &crate::Error) -> CompensationHint {
use crate::bridge::envelope::ErrorCode;
match e {
crate::Error::DataPlane(code) => match code {
ErrorCode::RejectedConstraint { constraint, detail } => CompensationHint::Custom {
constraint: constraint.clone(),
detail: detail.clone(),
},
ErrorCode::RejectedPrevalidation { reason } => CompensationHint::Custom {
constraint: "prevalidation".into(),
detail: reason.clone(),
},
ErrorCode::RejectedAuthz => CompensationHint::PermissionDenied,
ErrorCode::RateExceeded { retry_after_ms, .. } => CompensationHint::RateLimited {
retry_after_ms: *retry_after_ms,
},
other => CompensationHint::Custom {
constraint: "apply_failed".into(),
detail: format!("{other:?}"),
},
},
crate::Error::RejectedConstraint {
constraint, detail, ..
} => CompensationHint::Custom {
constraint: constraint.clone(),
detail: detail.clone(),
},
crate::Error::RejectedPrevalidation { constraint, reason } => CompensationHint::Custom {
constraint: constraint.clone(),
detail: reason.clone(),
},
crate::Error::RejectedAuthz { .. } => CompensationHint::PermissionDenied,
crate::Error::RateExceeded { retry_after_ms, .. } => CompensationHint::RateLimited {
retry_after_ms: *retry_after_ms,
},
other => CompensationHint::Custom {
constraint: "apply_failed".into(),
detail: other.to_string(),
},
}
}
#[cfg(test)]
mod tests {
use super::compensation_hint_for_dispatch_error;
use crate::bridge::envelope::ErrorCode;
use crate::types::TenantId;
use nodedb_types::sync::compensation::CompensationHint;
#[test]
fn preserved_data_plane_constraint_maps_to_custom_with_real_name() {
let e = crate::Error::DataPlane(ErrorCode::RejectedConstraint {
constraint: "users_email_unique".into(),
detail: "value 'a@b.com' already exists".into(),
});
match compensation_hint_for_dispatch_error(&e) {
CompensationHint::Custom { constraint, detail } => {
assert_eq!(constraint, "users_email_unique");
assert!(detail.contains("a@b.com"));
}
other => panic!("expected Custom, got {other:?}"),
}
}
#[test]
fn data_plane_authz_maps_to_permission_denied() {
let e = crate::Error::DataPlane(ErrorCode::RejectedAuthz);
assert_eq!(
compensation_hint_for_dispatch_error(&e),
CompensationHint::PermissionDenied
);
}
#[test]
fn data_plane_rate_exceeded_preserves_retry_after() {
let e = crate::Error::DataPlane(ErrorCode::RateExceeded {
gate: "writes".into(),
retry_after_ms: 1500,
});
assert_eq!(
compensation_hint_for_dispatch_error(&e),
CompensationHint::RateLimited {
retry_after_ms: 1500
}
);
}
#[test]
fn import_failure_maps_to_apply_failed_not_fabricated_constraint() {
let e = crate::Error::DataPlane(ErrorCode::Internal {
detail: "loro import failed".into(),
});
match compensation_hint_for_dispatch_error(&e) {
CompensationHint::Custom { constraint, .. } => assert_eq!(constraint, "apply_failed"),
other => panic!("expected Custom apply_failed, got {other:?}"),
}
}
#[test]
fn typed_authz_error_also_maps_to_permission_denied() {
let e = crate::Error::RejectedAuthz {
tenant_id: TenantId::new(0),
resource: "users".into(),
};
assert_eq!(
compensation_hint_for_dispatch_error(&e),
CompensationHint::PermissionDenied
);
}
}