use mkit_core::write_auth::MAX_CLOCK_LEAD_MS;
use super::{Authenticated, Pipeline, meta_error, ms};
use crate::error::{AbortCause, Code, ServerError};
use crate::pipeline::HookSet;
use crate::store::codec::{self, AbortReason, PendingOp, ReservationV1};
use crate::store::outbox::{OutboxBuilder, Terminal};
use crate::store::{
Batch, BatchOutcome, Key, MultipartBlobStore, NamespaceStore, Partition, Precondition, Value,
keys,
};
#[derive(Debug, Clone)]
pub(crate) struct PendingGuard {
pub(crate) rid: String,
pub(crate) key: Key,
pub(crate) value: Value,
pub(crate) apply_deadline_ms: u64,
pub(crate) repository: String,
}
const RECONCILE_MARGIN_MS: u64 = 1_000;
#[cfg(any(test, feature = "http-objects"))]
pub(crate) fn read_pending(
repository: String,
created_at_ms: u64,
deadline_ms: u64,
read_reconcile_grace: core::time::Duration,
) -> ReservationV1 {
ReservationV1::Pending {
repository,
created_at_ms,
reconcile_at_ms: deadline_ms
.saturating_add(u64::try_from(read_reconcile_grace.as_millis()).unwrap_or(u64::MAX)),
op: PendingOp::Read,
}
}
pub(crate) fn abort_reason(err: &ServerError) -> (AbortReason, String) {
let message = err.public_message();
match err.abort_cause() {
Some(AbortCause::EpochMismatch) => return (AbortReason::EpochMismatch, String::new()),
Some(AbortCause::ReplayRace) => return (AbortReason::ReplayRace, String::new()),
Some(AbortCause::QuotaWindow) => return (AbortReason::Unspecified, message.to_owned()),
Some(AbortCause::Contention) => return (AbortReason::Internal, String::new()),
None => {}
}
if err.code() == Code::Aborted {
(AbortReason::ReplayRace, String::new())
} else if matches!(
err.code(),
Code::PermissionDenied | Code::ResourceExhausted | Code::FailedPrecondition
) {
let detail = if message.len() <= 512 && !message.chars().any(char::is_control) {
message.to_owned()
} else {
"write denied".into()
};
(AbortReason::Unspecified, detail)
} else {
(AbortReason::Internal, String::new())
}
}
const ABORT_ATTEMPTS: usize = 5;
impl<B: MultipartBlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H> {
pub(super) async fn abort_unsupported_stream(
&self,
a: &Authenticated,
partition: &Partition,
rid: &str,
) -> Result<(), ServerError> {
let row = keys::reservation(rid).map_err(meta_error)?;
for _ in 0..ABORT_ATTEMPTS {
let read = [
keys::outbox_sequence(),
keys::outcome_backlog(),
row.clone(),
];
let values = self
.meta
.get_many(partition, &read)
.await
.map_err(meta_error)?;
if values.get(2).is_some_and(Option::is_some) {
return Ok(());
}
let mut builder = OutboxBuilder::new(
values.first().and_then(Option::as_ref),
values.get(1).and_then(Option::as_ref),
)
.map_err(meta_error)?;
let terminal = Terminal::new(ReservationV1::Aborted {
repository: a.repo().identity.clone(),
occurred_at_ms: ms(self.clock.now_ms()),
reason: AbortReason::Unspecified,
detail: "reservations unsupported on this transport".into(),
})
.map_err(meta_error)?;
builder.abort_direct(rid, terminal);
let mut batch = Batch::new();
builder
.try_finish(&mut batch.preconditions, &mut batch.writes)
.map_err(meta_error)?;
match self.apply_meta(partition, batch).await? {
BatchOutcome::Committed => return Ok(()),
BatchOutcome::PreconditionFailed { .. } => {}
BatchOutcome::DeadlinePassed { .. } => {
return Err(ServerError::unavailable("outcome commit deadline passed"));
}
}
}
Err(ServerError::unavailable("admission unavailable"))
}
pub(super) async fn record_pending(
&self,
a: &Authenticated,
partition: &Partition,
rid: &str,
) -> Result<PendingGuard, ServerError> {
let now = ms(self.clock.now_ms());
let apply_deadline_ms = now
.saturating_add(
u64::try_from(self.cfg.max_apply_window.as_millis()).unwrap_or(u64::MAX),
)
.min(a.auth.as_ref().map_or(u64::MAX, |auth| {
ms(auth.expires_at_ms).saturating_add(MAX_CLOCK_LEAD_MS.unsigned_abs())
}));
let reconcile_at_ms = apply_deadline_ms
.saturating_add(MAX_CLOCK_LEAD_MS.unsigned_abs())
.saturating_add(RECONCILE_MARGIN_MS);
let pending = ReservationV1::Pending {
repository: a.repo().identity.clone(),
created_at_ms: now,
reconcile_at_ms,
op: PendingOp::Write,
};
let value = codec::encode_reservation(&pending);
let mut builder = OutboxBuilder::new(None, None).map_err(meta_error)?;
builder.pending(rid, None, &pending);
let mut batch = Batch::new();
batch
.preconditions
.push(Precondition::NotAfter(apply_deadline_ms));
builder
.try_finish(&mut batch.preconditions, &mut batch.writes)
.map_err(|_| ServerError::unavailable("admission unavailable"))?;
match self.apply_meta(partition, batch).await? {
BatchOutcome::Committed => Ok(PendingGuard {
rid: rid.to_owned(),
key: keys::reservation(rid).map_err(meta_error)?,
value,
apply_deadline_ms,
repository: a.repo().identity.clone(),
}),
BatchOutcome::PreconditionFailed { .. } => {
Err(ServerError::unavailable("admission unavailable"))
}
BatchOutcome::DeadlinePassed { .. } => {
Err(ServerError::unavailable("commit deadline passed; retry"))
}
}
}
pub(super) async fn resolve_pending(
&self,
partition: &Partition,
pending: &PendingGuard,
reason: AbortReason,
detail: String,
) {
if let Err(err) = self.abort_pending(partition, pending, reason, detail).await {
tracing::warn!(error = %err, "pending reservation abort failed; reconcile will retry");
}
}
async fn abort_pending(
&self,
partition: &Partition,
pending: &PendingGuard,
reason: AbortReason,
detail: String,
) -> Result<(), ServerError> {
for _ in 0..ABORT_ATTEMPTS {
let read = [
keys::outbox_sequence(),
keys::outcome_backlog(),
pending.key.clone(),
];
let values = self
.meta
.get_many(partition, &read)
.await
.map_err(meta_error)?;
if values.get(2).and_then(Option::as_ref) != Some(&pending.value) {
return Ok(());
}
let mut builder = OutboxBuilder::new(
values.first().and_then(Option::as_ref),
values.get(1).and_then(Option::as_ref),
)
.map_err(meta_error)?;
let record = ReservationV1::Aborted {
repository: pending.repository.clone(),
occurred_at_ms: ms(self.clock.now_ms()),
reason,
detail: detail.clone(),
};
builder.outcome(
&pending.rid,
&pending.value,
Terminal::new(record).map_err(meta_error)?,
);
let mut batch = Batch::new();
builder
.try_finish(&mut batch.preconditions, &mut batch.writes)
.map_err(meta_error)?;
match self.apply_meta(partition, batch).await? {
BatchOutcome::Committed => return Ok(()),
BatchOutcome::PreconditionFailed { .. } => {}
BatchOutcome::DeadlinePassed { .. } => {
return Err(ServerError::unavailable(
"reservation abort deadline passed",
));
}
}
}
Err(ServerError::unavailable(
"reservation abort contended; reconcile will retry",
))
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn read_planner_uses_sixty_second_reconcile_grace() {
let pending = read_pending("repo".into(), 1, 10_000, core::time::Duration::from_mins(1));
assert_eq!(
pending,
ReservationV1::Pending {
repository: "repo".into(),
created_at_ms: 1,
reconcile_at_ms: 70_000,
op: PendingOp::Read
}
);
let value = codec::encode_reservation(&pending);
let mut builder = OutboxBuilder::new(None, None).unwrap();
builder.outcome(
"rid",
&value,
Terminal::new(ReservationV1::ReadServed {
repository: "repo".into(),
occurred_at_ms: 10_500,
object: [1; 32],
bytes_served: 7,
})
.unwrap(),
);
let mut batch = Batch::new();
builder
.try_finish(&mut batch.preconditions, &mut batch.writes)
.unwrap();
assert!(batch.preconditions.contains(&Precondition::Equals(
keys::reservation("rid").unwrap(),
value
)));
}
}