Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 13 additions & 0 deletions libsql-replication/proto/metadata.proto
Original file line number Diff line number Diff line change
Expand Up @@ -27,4 +27,17 @@ message DatabaseConfig {
optional bool shared_schema = 11;
optional string shared_schema_name = 12;
optional DurabilityMode durability_mode = 13;
// The namespace fence as seen by the primary when it answered. Only ever filled by the
// primary's replication `Hello`, and only while a fence is active; it is never part of a
// stored configuration. Absent from older primaries; older replicas ignore it and still
// see the legacy `block_*` fields above.
optional ReplicatedFence fence = 14;
}

// The part of a namespace fence a replica needs to apply the primary's read admission.
message ReplicatedFence {
// The fence state name, e.g. "SOURCE_READ_FENCED".
string state = 1;
// The revision of the fence record the state belongs to.
uint64 revision = 2;
}
5 changes: 5 additions & 0 deletions libsql-replication/proto/proxy.proto
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,11 @@ message Error {
ErrorCode code = 1;
string message = 2;
int32 extended_code = 3;
// Stable machine-readable outcome of the error, e.g. "MIGRATION_WRITE_FENCED" for a
// request refused by a namespace fence. Absent when the error has no typed outcome, and
// always absent from older servers: a receiver treats an absent field as "no typed
// outcome" and falls back to `code`.
optional string stable_code = 4;
}

message ResultRows {
Expand Down
17 changes: 17 additions & 0 deletions libsql-replication/src/generated/metadata.rs
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,23 @@ pub struct DatabaseConfig {
pub shared_schema_name: ::core::option::Option<::prost::alloc::string::String>,
#[prost(enumeration = "DurabilityMode", optional, tag = "13")]
pub durability_mode: ::core::option::Option<i32>,
/// The namespace fence as seen by the primary when it answered. Only ever filled by the
/// primary's replication `Hello`, and only while a fence is active; it is never part of a
/// stored configuration. Absent from older primaries; older replicas ignore it and still
/// see the legacy `block_*` fields above.
#[prost(message, optional, tag = "14")]
pub fence: ::core::option::Option<ReplicatedFence>,
}
/// The part of a namespace fence a replica needs to apply the primary's read admission.
#[allow(clippy::derive_partial_eq_without_eq)]
#[derive(Clone, PartialEq, ::prost::Message)]
pub struct ReplicatedFence {
/// The fence state name, e.g. "SOURCE_READ_FENCED".
#[prost(string, tag = "1")]
pub state: ::prost::alloc::string::String,
/// The revision of the fence record the state belongs to.
#[prost(uint64, tag = "2")]
pub revision: u64,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord, ::prost::Enumeration)]
#[repr(i32)]
Expand Down
6 changes: 6 additions & 0 deletions libsql-replication/src/generated/proxy.rs
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,12 @@ pub struct Error {
pub message: ::prost::alloc::string::String,
#[prost(int32, tag = "3")]
pub extended_code: i32,
/// Stable machine-readable outcome of the error, e.g. "MIGRATION_WRITE_FENCED" for a
/// request refused by a namespace fence. Absent when the error has no typed outcome, and
/// always absent from older servers: a receiver treats an absent field as "no typed
/// outcome" and falls back to `code`.
#[prost(string, optional, tag = "4")]
pub stable_code: ::core::option::Option<::prost::alloc::string::String>,
}
/// Nested message and enum types in `Error`.
pub mod error {
Expand Down
115 changes: 115 additions & 0 deletions libsql-replication/src/rpc.rs
Original file line number Diff line number Diff line change
Expand Up @@ -105,3 +105,118 @@ pub mod metadata {
#![allow(clippy::all)]
include!("generated/metadata.rs");
}

#[cfg(test)]
mod test {
use prost::Message;

use super::metadata::{DatabaseConfig, ReplicatedFence};
use super::proxy::{error::ErrorCode, Error};

/// `proxy.Error` as a peer built before `stable_code` existed knows it.
#[derive(Clone, PartialEq, ::prost::Message)]
struct ErrorWithoutStableCode {
#[prost(enumeration = "ErrorCode", tag = "1")]
code: i32,
#[prost(string, tag = "2")]
message: String,
#[prost(int32, tag = "3")]
extended_code: i32,
}

/// The legacy part of `metadata.DatabaseConfig`, as a peer built before `fence` existed
/// knows it (the fields in between are skipped the same way as `fence` is).
#[derive(Clone, PartialEq, ::prost::Message)]
struct DatabaseConfigWithoutFence {
#[prost(bool, tag = "1")]
block_reads: bool,
#[prost(bool, tag = "2")]
block_writes: bool,
#[prost(string, optional, tag = "3")]
block_reason: Option<String>,
#[prost(uint64, tag = "4")]
max_db_pages: u64,
}

fn error(stable_code: Option<&str>) -> Error {
Error {
code: ErrorCode::SqlError as i32,
message: "writes are fenced".into(),
extended_code: 23,
stable_code: stable_code.map(Into::into),
}
}

#[test]
fn proxy_error_stable_code_is_additive() {
// A newer server's error, read by an older replica: the known fields are intact and
// the stable code is skipped.
let new = error(Some("MIGRATION_WRITE_FENCED"));
let old = ErrorWithoutStableCode::decode(&new.encode_to_vec()[..]).unwrap();
assert_eq!(
old,
ErrorWithoutStableCode {
code: ErrorCode::SqlError as i32,
message: "writes are fenced".into(),
extended_code: 23,
}
);

// An older server's error, read by a newer replica: no typed outcome.
let decoded = Error::decode(&old.encode_to_vec()[..]).unwrap();
assert_eq!(decoded, error(None));

// Without a stable code the encoding is exactly the older one, so an error that has no
// typed outcome is unchanged on the wire.
assert_eq!(error(None).encode_to_vec(), old.encode_to_vec());

// And the field round-trips between newer peers.
assert_eq!(Error::decode(&new.encode_to_vec()[..]).unwrap(), new);
}

fn config(fence: Option<ReplicatedFence>) -> DatabaseConfig {
DatabaseConfig {
block_reads: true,
block_writes: true,
block_reason: Some("namespace fence".into()),
max_db_pages: 1024,
fence,
..Default::default()
}
}

#[test]
fn replicated_fence_is_additive() {
let fence = ReplicatedFence {
state: "SOURCE_READ_FENCED".into(),
revision: 3,
};

// A newer primary's config, read by an older replica: the legacy block fields are
// intact, so the older replica still applies the legacy mirror of the fence.
let new = config(Some(fence.clone()));
let old = DatabaseConfigWithoutFence::decode(&new.encode_to_vec()[..]).unwrap();
assert_eq!(
old,
DatabaseConfigWithoutFence {
block_reads: true,
block_writes: true,
block_reason: Some("namespace fence".into()),
max_db_pages: 1024,
}
);

// An older primary's config, read by a newer replica: no fence.
let decoded = DatabaseConfig::decode(&old.encode_to_vec()[..]).unwrap();
assert_eq!(decoded.fence, None);
assert_eq!(decoded, config(None));

// Without a fence the encoding is exactly the older one: a stored configuration, which
// never carries a fence, is unchanged.
assert_eq!(config(None).encode_to_vec(), old.encode_to_vec());

// And the fence round-trips between newer peers.
let round_trip = DatabaseConfig::decode(&new.encode_to_vec()[..]).unwrap();
assert_eq!(round_trip.fence, Some(fence));
}
}
3 changes: 3 additions & 0 deletions libsql-server/src/connection/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,9 @@ impl From<&DatabaseConfig> for metadata::DatabaseConfig {
shared_schema: Some(value.is_shared_schema),
shared_schema_name: value.shared_schema_name.as_ref().map(|s| s.to_string()),
durability_mode: Some(metadata::DurabilityMode::from(value.durability_mode).into()),
// Never part of a configuration: the replication `hello` fills it from the live
// fence gate, and it is not stored.
fence: None,
}
}
}
Expand Down
8 changes: 5 additions & 3 deletions libsql-server/src/connection/write_proxy.rs
Original file line number Diff line number Diff line change
Expand Up @@ -266,13 +266,15 @@ impl RemoteConnection {
let response_stream = match client.stream_exec(req).await {
Ok(i) => i.into_inner(),
Err(e) => {
// Only `UNAVAILABLE` is retried. A fence denial is `FAILED_PRECONDITION`
// with its stable code, answered to the client as the primary's denial.
if e.code() == Code::Unavailable {
tracing::error!("retrying proxy connection: {}", e);
tokio::time::sleep(Duration::from_millis(500) * 2u32.pow(retries)).await;
retries += 1;
continue;
} else {
return Err(e.into());
return Err(Error::from_proxy_status(e));
}
}
};
Expand Down Expand Up @@ -387,7 +389,7 @@ where
)
}
exec_resp::Response::DescribeResp(_) => Err(Error::PrimaryStreamMisuse),
exec_resp::Response::Error(e) => Err(Error::RpcQueryError(e)),
exec_resp::Response::Error(e) => Err(Error::from_proxy_error(e)),
}
};

Expand Down Expand Up @@ -430,7 +432,7 @@ where

Ok(false)
}
exec_resp::Response::Error(e) => Err(Error::RpcQueryError(e)),
exec_resp::Response::Error(e) => Err(Error::from_proxy_error(e)),
exec_resp::Response::ProgramResp(_) => Err(Error::PrimaryStreamMisuse),
};

Expand Down
111 changes: 110 additions & 1 deletion libsql-server/src/error.rs
Original file line number Diff line number Diff line change
Expand Up @@ -151,6 +151,55 @@ pub trait ResponseError: std::error::Error {

impl ResponseError for Error {}

impl Error {
/// The fence denial this error carries, looking through the wrappers it can arrive in.
pub(crate) fn fence_error(&self) -> Option<&crate::namespace::fence::outcome::FenceError> {
match self {
Error::NamespaceFence(e) => Some(e),
Error::Migration(crate::schema::Error::NamespaceFence(e)) => Some(e),
Error::Ref(this) => this.fence_error(),
Error::Anyhow(e) => e.downcast_ref::<Error>().and_then(Error::fence_error),
_ => None,
}
}

/// A step or program error the primary returned through the write proxy. A fence denial
/// from a primary that fills the additive `stable_code` field
/// (`docs/NAMESPACE_FENCE.md` section 6.1) becomes the same [`Error::NamespaceFence`] a
/// local denial is, so the replica answers its client exactly as the primary would. Any
/// other error, and every error from a primary that does not fill the field, stays
/// [`Error::RpcQueryError`].
pub(crate) fn from_proxy_error(e: crate::rpc::proxy::rpc::Error) -> Self {
let fence = e.stable_code.as_deref().and_then(|code| {
crate::namespace::fence::outcome::FenceError::from_proxy_stable_code(code, &e.message)
});
match fence {
Some(fence) => Error::NamespaceFence(fence),
None => Error::RpcQueryError(e),
}
}

/// A gRPC status from the primary's proxy service: its typed fence denial
/// (`FAILED_PRECONDITION` with the stable code, section 6.1) as [`Error::NamespaceFence`],
/// anything else unchanged.
pub(crate) fn from_proxy_status(status: tonic::Status) -> Self {
match crate::namespace::fence::outcome::FenceError::from_grpc_status(&status) {
Some(fence) => Error::NamespaceFence(fence),
None => Error::RpcQueryExecutionError(status),
}
}
}

/// The HTTP response for a fence denial (`docs/NAMESPACE_FENCE.md` section 6): the fence status
/// and the JSON error body with the additive `code` (and `detail`) fields.
pub(crate) fn fence_error_response(
e: &crate::namespace::fence::outcome::FenceError,
) -> axum::response::Response {
let status = e.http_status();
tracing::debug!("HTTP API: {status}, {e}");
(status, axum::Json(e.http_error_body())).into_response()
}

impl IntoResponse for Error {
fn into_response(self) -> axum::response::Response {
(&self).into_response()
Expand Down Expand Up @@ -226,7 +275,7 @@ impl IntoResponse for &Error {
AttachInMigration => self.format_err(StatusCode::BAD_REQUEST),
RuntimeTaskJoinError(_) => self.format_err(StatusCode::INTERNAL_SERVER_ERROR),
NotAPrimary => self.format_err(StatusCode::BAD_REQUEST),
NamespaceFence(e) => self.format_err(e.outcome().admin_http_status()),
NamespaceFence(e) => fence_error_response(e),
}
}
}
Expand Down Expand Up @@ -338,3 +387,63 @@ impl IntoResponse for &ForkError {
}
}
}

#[cfg(test)]
mod fence_tests {
use super::*;
use crate::namespace::fence::outcome::{FenceDetail, FenceError, FenceOutcome};

async fn response(e: &Error) -> (StatusCode, serde_json::Value) {
let response = e.into_response();
let status = response.status();
let body = hyper::body::to_bytes(response.into_body()).await.unwrap();
(status, serde_json::from_slice(&body).unwrap())
}

/// Section 6: a fence denial is `423` with the additive `code` (and `detail`) field, found
/// through every wrapper the error can arrive in; other errors keep their shape.
#[tokio::test]
async fn fence_errors_carry_code() {
for outcome in [
FenceOutcome::MigrationWriteFenced,
FenceOutcome::MigrationReadFenced,
FenceOutcome::MigrationTargetQuarantined,
FenceOutcome::FenceStateUnavailable,
] {
let fence = FenceError::new(outcome, "denied");
let wrapped = [
Error::NamespaceFence(fence.clone()),
Error::Ref(std::sync::Arc::new(Error::NamespaceFence(fence.clone()))),
Error::Anyhow(anyhow::anyhow!(Error::NamespaceFence(fence.clone()))),
Error::Migration(crate::schema::Error::NamespaceFence(fence.clone())),
];
for e in &wrapped {
assert_eq!(e.fence_error(), Some(&fence), "{e:?}");
let (status, body) = response(e).await;
assert_eq!(status, StatusCode::LOCKED, "{e:?}");
assert_eq!(body["code"], outcome.as_str(), "{e:?}");
assert_eq!(body["error"], fence.to_string(), "{e:?}");
assert!(body.get("detail").is_none(), "{body}");
}
}

let unavailable = FenceError::new(FenceOutcome::FenceStateUnavailable, "corrupt")
.with_detail(FenceDetail::CorruptRecord);
let (_, body) = response(&Error::NamespaceFence(unavailable)).await;
assert_eq!(body["detail"], "corrupt_record");

let precondition = FenceError::new(FenceOutcome::FencePreconditionFailed, "no")
.with_detail(FenceDetail::NotPrimary);
let (status, body) = response(&Error::NamespaceFence(precondition)).await;
assert_eq!(status, StatusCode::PRECONDITION_FAILED);
assert_eq!(body["code"], "FENCE_PRECONDITION_FAILED");
assert_eq!(body["detail"], "not_primary");

// The legacy `block_*` refusal keeps its mapping and has no code.
let blocked = Error::Blocked(Some("maintenance".into()));
assert_eq!(blocked.fence_error(), None);
let (status, body) = response(&blocked).await;
assert_eq!(status, StatusCode::INTERNAL_SERVER_ERROR);
assert!(body.get("code").is_none(), "{body}");
}
}
6 changes: 6 additions & 0 deletions libsql-server/src/hrana/batch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,10 @@ pub enum BatchError {
ResponseTooLarge,
#[error("Schema migration error: {message}")]
SchemaError { message: String },
/// The whole batch was refused by the namespace fence (`docs/NAMESPACE_FENCE.md`
/// section 6), for example a read under a read fence.
#[error(transparent)]
Fence(crate::namespace::fence::outcome::FenceError),
}

fn proto_cond_to_cond(
Expand Down Expand Up @@ -183,6 +187,7 @@ pub fn batch_error_from_sqld_error(sqld_error: SqldError) -> Result<BatchError,
SqldError::BuilderError(QueryResultBuilderError::ResponseTooLarge(_)) => {
BatchError::ResponseTooLarge
}
SqldError::NamespaceFence(e) => BatchError::Fence(e),
sqld_error => return Err(sqld_error),
})
}
Expand All @@ -201,6 +206,7 @@ impl BatchError {
Self::TransactionBusy => "TRANSACTION_BUSY",
Self::ResponseTooLarge => "RESPONSE_TOO_LARGE",
Self::SchemaError { message: _ } => "SCHEMA_MIGRATION_ERROR",
Self::Fence(e) => e.outcome().as_str(),
}
}
}
Loading
Loading