mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-15 17:43:13 +00:00
Compare commits
6 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 4609e78fcd | |||
| 98a96776c6 | |||
| 2a9034ed14 | |||
| d3c2b7d67e | |||
| 71e83aeec4 | |||
| 9138c24571 |
@@ -94,6 +94,7 @@ jobs:
|
||||
short_sha: ${{ steps.check.outputs.short_sha }}
|
||||
is_prerelease: ${{ steps.check.outputs.is_prerelease }}
|
||||
create_latest: ${{ steps.check.outputs.create_latest }}
|
||||
source_ref: ${{ steps.check.outputs.source_ref }}
|
||||
steps:
|
||||
- name: Checkout repository
|
||||
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
|
||||
@@ -118,6 +119,7 @@ jobs:
|
||||
short_sha=""
|
||||
is_prerelease=false
|
||||
create_latest=false
|
||||
source_ref="$GITHUB_SHA"
|
||||
|
||||
if [[ "${{ github.event_name }}" == "workflow_run" ]]; then
|
||||
# Triggered by build workflow completion
|
||||
@@ -137,6 +139,7 @@ jobs:
|
||||
# Extract version info from commit message or use commit SHA
|
||||
# Use Git to generate consistent short SHA (ensures uniqueness like build.yml)
|
||||
short_sha=$(git rev-parse --short "$HEAD_SHA")
|
||||
source_ref="$HEAD_SHA"
|
||||
|
||||
# Determine build type based on triggering workflow event and ref
|
||||
triggering_event="$TRIGGERING_EVENT"
|
||||
@@ -261,6 +264,23 @@ jobs:
|
||||
echo "⚠️ Only release versions (latest, v1.0.0, 1.0.0) and prereleases (v1.0.0-alpha1, 1.0.0-beta2) are supported"
|
||||
;;
|
||||
esac
|
||||
|
||||
if [[ "$should_build" == true && "$input_version" != "latest" ]]; then
|
||||
tag_ref="refs/tags/$input_version"
|
||||
if ! git ls-remote --exit-code origin "$tag_ref" >/dev/null 2>&1; then
|
||||
if [[ "$input_version" == v* ]]; then
|
||||
tag_ref="refs/tags/${input_version#v}"
|
||||
else
|
||||
tag_ref="refs/tags/v$input_version"
|
||||
fi
|
||||
fi
|
||||
|
||||
if ! git ls-remote --exit-code origin "$tag_ref" >/dev/null 2>&1; then
|
||||
echo "❌ Release tag not found for Docker build: $input_version"
|
||||
exit 1
|
||||
fi
|
||||
source_ref="$tag_ref"
|
||||
fi
|
||||
fi
|
||||
|
||||
{
|
||||
@@ -271,6 +291,7 @@ jobs:
|
||||
echo "short_sha=$short_sha"
|
||||
echo "is_prerelease=$is_prerelease"
|
||||
echo "create_latest=$create_latest"
|
||||
echo "source_ref=$source_ref"
|
||||
} >> "$GITHUB_OUTPUT"
|
||||
|
||||
echo "🐳 Docker Build Summary:"
|
||||
@@ -281,6 +302,7 @@ jobs:
|
||||
echo " - Short SHA: $short_sha"
|
||||
echo " - Is prerelease: $is_prerelease"
|
||||
echo " - Create latest: $create_latest"
|
||||
echo " - Source ref: $source_ref"
|
||||
|
||||
# Build multi-arch Docker images
|
||||
# Strategy: Build images using pre-built binaries from dl.rustfs.com
|
||||
@@ -308,6 +330,7 @@ jobs:
|
||||
uses: actions/checkout@9c091bb21b7c1c1d1991bb908d89e4e9dddfe3e0 # v7
|
||||
with:
|
||||
persist-credentials: false
|
||||
ref: ${{ needs.build-check.outputs.source_ref }}
|
||||
|
||||
- name: Login to Docker Hub
|
||||
uses: docker/login-action@c94ce9fb468520275223c153574b00df6fe4bcc9 # v3
|
||||
@@ -397,7 +420,8 @@ jobs:
|
||||
LABELS="org.opencontainers.image.title=RustFS"
|
||||
LABELS="$LABELS,org.opencontainers.image.description=RustFS distributed object storage system"
|
||||
LABELS="$LABELS,org.opencontainers.image.version=$VERSION"
|
||||
LABELS="$LABELS,org.opencontainers.image.revision=${{ github.sha }}"
|
||||
SOURCE_REVISION="$(git rev-parse HEAD)"
|
||||
LABELS="$LABELS,org.opencontainers.image.revision=$SOURCE_REVISION"
|
||||
LABELS="$LABELS,org.opencontainers.image.source=${{ github.server_url }}/${{ github.repository }}"
|
||||
LABELS="$LABELS,org.opencontainers.image.created=$(date -u +'%Y-%m-%dT%H:%M:%SZ')"
|
||||
LABELS="$LABELS,org.opencontainers.image.build-type=$BUILD_TYPE"
|
||||
|
||||
@@ -1153,7 +1153,7 @@ struct MrfResponse {
|
||||
fn build_mrf_response(
|
||||
bucket: String,
|
||||
bucket_stats: &BucketStats,
|
||||
durable: crate::admin::storage_api::replication::DurableMrfBacklog,
|
||||
durable: &crate::admin::storage_api::replication::DurableMrfBacklog,
|
||||
) -> MrfResponse {
|
||||
let observation_scope = if bucket_stats.replication_stats.cluster_complete {
|
||||
"cluster_aggregated"
|
||||
@@ -1223,7 +1223,10 @@ fn build_mrf_response(
|
||||
total_failed_size,
|
||||
queued_count: queued.count,
|
||||
queued_size: queued.bytes,
|
||||
per_object_entries_available: false,
|
||||
// The default (non-aggregate) response mode streams the durable
|
||||
// backlog per object, so the enumerable API exists whenever the
|
||||
// backlog is readable.
|
||||
per_object_entries_available: durable.available,
|
||||
runtime_stats_available: bucket_stats.replication_stats.provider_available,
|
||||
cluster_complete: bucket_stats.replication_stats.cluster_complete,
|
||||
observed_node_count: bucket_stats.replication_stats.observed_node_count,
|
||||
@@ -1235,23 +1238,155 @@ fn build_mrf_response(
|
||||
}
|
||||
}
|
||||
|
||||
/// One durable MRF backlog entry rendered for the default (madmin-compatible)
|
||||
/// stream. Field names are the exact json tags of madmin-go `ReplicationMRF`
|
||||
/// (replication-api.go), which `mc replicate backlog` decodes one JSON
|
||||
/// document at a time. `Size` and `TargetARNs` are RustFS extension keys with
|
||||
/// no madmin counterpart; Go decoders ignore unknown keys.
|
||||
#[derive(Debug, Serialize)]
|
||||
struct MrfEntryDocument {
|
||||
/// The durable backlog is a cluster-shared ledger with no per-node
|
||||
/// attribution, so the madmin `nodeName` tag is always empty.
|
||||
#[serde(rename = "nodeName")]
|
||||
node_name: String,
|
||||
#[serde(rename = "bucket")]
|
||||
bucket: String,
|
||||
#[serde(rename = "object")]
|
||||
object: String,
|
||||
#[serde(rename = "versionId")]
|
||||
version_id: String,
|
||||
#[serde(rename = "retryCount")]
|
||||
retry_count: i32,
|
||||
#[serde(rename = "Size")]
|
||||
size: i64,
|
||||
#[serde(rename = "TargetARNs", skip_serializing_if = "Vec::is_empty")]
|
||||
target_arns: Vec<String>,
|
||||
}
|
||||
|
||||
/// Upper bound on the number of documents one stream response emits. The
|
||||
/// durable ledger is not bounded by the in-memory pending cap (recovery can
|
||||
/// persist far larger generations), and the body is buffered before send, so
|
||||
/// an unbounded read could stage hundreds of MB per request. A truncated
|
||||
/// stream is signalled via `x-rustfs-replication-mrf-truncated`.
|
||||
const REPLICATION_MRF_MAX_STREAM_ENTRIES: usize = 10_000;
|
||||
|
||||
/// Project the durable backlog into madmin `ReplicationMRF` documents,
|
||||
/// scoped to `bucket` when it is non-empty (madmin allows an empty bucket to
|
||||
/// mean "across all buckets"), bounded by
|
||||
/// [`REPLICATION_MRF_MAX_STREAM_ENTRIES`]. Returns the documents and whether
|
||||
/// the backlog was truncated.
|
||||
fn mrf_entry_documents(
|
||||
bucket: &str,
|
||||
durable: &crate::admin::storage_api::replication::DurableMrfBacklog,
|
||||
) -> (Vec<MrfEntryDocument>, bool) {
|
||||
let mut documents = Vec::new();
|
||||
let mut truncated = false;
|
||||
for entry in durable
|
||||
.entries
|
||||
.iter()
|
||||
.filter(|entry| bucket.is_empty() || entry.bucket == bucket)
|
||||
{
|
||||
if documents.len() >= REPLICATION_MRF_MAX_STREAM_ENTRIES {
|
||||
truncated = true;
|
||||
break;
|
||||
}
|
||||
documents.push(MrfEntryDocument {
|
||||
node_name: String::new(),
|
||||
bucket: entry.bucket.clone(),
|
||||
object: entry.object.clone(),
|
||||
// Delete-marker purge entries track the marker version separately;
|
||||
// fall back to it so those rows still carry a version identity.
|
||||
// The nil UUID is RustFS's in-memory null-version sentinel and
|
||||
// must leave as the S3 wire token, not a zero UUID.
|
||||
version_id: entry
|
||||
.version_id
|
||||
.or(entry.delete_marker_version_id)
|
||||
.map(|v| {
|
||||
if v.is_nil() {
|
||||
rustfs_filemeta::NULL_VERSION_ID.to_string()
|
||||
} else {
|
||||
v.to_string()
|
||||
}
|
||||
})
|
||||
.unwrap_or_default(),
|
||||
retry_count: entry.retry_count,
|
||||
size: entry.size,
|
||||
target_arns: entry.target_arns.clone(),
|
||||
});
|
||||
}
|
||||
(documents, truncated)
|
||||
}
|
||||
|
||||
/// Render the MRF backlog as a response body.
|
||||
///
|
||||
/// Default (madmin-compatible) mode emits one `ReplicationMRF` JSON document
|
||||
/// per line with no envelope — madmin's `BucketReplicationMRF` reads the body
|
||||
/// with a `json.Decoder` loop, so an envelope object would decode as a single
|
||||
/// entry whose `"Bucket"` key case-insensitively matches
|
||||
/// `ReplicationMRF.Bucket` (a phantom row in `mc replicate backlog`), and an
|
||||
/// empty backlog must render an empty body so the loop ends on io.EOF with
|
||||
/// zero rows.
|
||||
///
|
||||
/// `aggregate=true` (RustFS extension) keeps the enveloped counter shape;
|
||||
/// backlog-source health (`RuntimeStatsAvailable`/`DurableBacklogAvailable`)
|
||||
/// is only representable there — an unreadable ledger fails the stream
|
||||
/// request outright in the handler (madmin only decodes the body of a 200,
|
||||
/// so an empty stream would read as a healthy zero-row backlog).
|
||||
fn render_mrf_backlog(
|
||||
response: &MrfResponse,
|
||||
durable: &crate::admin::storage_api::replication::DurableMrfBacklog,
|
||||
aggregate: bool,
|
||||
) -> Result<(Vec<u8>, bool), serde_json::Error> {
|
||||
if aggregate {
|
||||
return Ok((serde_json::to_vec(response)?, false));
|
||||
}
|
||||
|
||||
let (documents, truncated) = mrf_entry_documents(&response.bucket, durable);
|
||||
let mut data = Vec::new();
|
||||
for entry in documents {
|
||||
serde_json::to_writer(&mut data, &entry)?;
|
||||
data.push(b'\n');
|
||||
}
|
||||
Ok((data, truncated))
|
||||
}
|
||||
|
||||
/// `GET /v3/replication/mrf`
|
||||
///
|
||||
/// Reports the failed-replication backlog (MinIO's MRF concept) for a bucket.
|
||||
///
|
||||
/// Compatibility note: MinIO returns a stream of individual MRF entries. RustFS
|
||||
/// deliberately returns aggregate runtime and durable counters instead.
|
||||
/// `PerObjectEntriesAvailable` remains false until an enumerable API exists.
|
||||
/// `PerTargetDurableEntriesAvailable` is false when the durable backlog includes
|
||||
/// older entries that cannot be attributed to a target.
|
||||
/// The default response is a madmin-compatible stream of `ReplicationMRF`
|
||||
/// documents built from the durable backlog ledger (in-memory failures that
|
||||
/// have not been flushed yet — the persister runs every few seconds — are not
|
||||
/// visible). `?aggregate=true` (RustFS extension) returns the enveloped
|
||||
/// runtime + durable counter shape instead; `PerTargetDurableEntriesAvailable`
|
||||
/// is false there when the durable backlog includes older entries that cannot
|
||||
/// be attributed to a target.
|
||||
///
|
||||
/// The madmin `node` parameter is accepted but has no filtering effect: the
|
||||
/// durable ledger is cluster-shared with no per-node attribution, so every
|
||||
/// node serves the same (complete) backlog.
|
||||
///
|
||||
/// Authorization: the stream requires `admin:ReplicationDiff` (it enumerates
|
||||
/// object names and version ids, MinIO parity); `?aggregate=true` carries no
|
||||
/// object identities and requires only `admin:GetReplicationMetrics`.
|
||||
pub struct ReplicationMrfHandler {}
|
||||
|
||||
#[async_trait::async_trait]
|
||||
impl Operation for ReplicationMrfHandler {
|
||||
async fn call(&self, req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
|
||||
validate_replication_admin_request(&req, AdminAction::GetReplicationMetricsAction).await?;
|
||||
|
||||
let queries = extract_query_params(&req.uri);
|
||||
let aggregate = queries.get("aggregate").map(String::as_str) == Some("true");
|
||||
// The default stream enumerates object names and version ids, which
|
||||
// a metrics-only principal must not see; gate it on the same action
|
||||
// MinIO uses for this endpoint. The aggregate counters carry no
|
||||
// object identities and keep the metrics action.
|
||||
let action = if aggregate {
|
||||
AdminAction::GetReplicationMetricsAction
|
||||
} else {
|
||||
AdminAction::ReplicationDiff
|
||||
};
|
||||
validate_replication_admin_request(&req, action).await?;
|
||||
|
||||
let Some(bucket) = queries.get("bucket").filter(|b| !b.is_empty()).cloned() else {
|
||||
return Err(s3_error!(InvalidRequest, "bucket is required"));
|
||||
};
|
||||
@@ -1275,14 +1410,47 @@ impl Operation for ReplicationMrfHandler {
|
||||
return Err(ApiError::from(err).into());
|
||||
}
|
||||
|
||||
if let Some(node) = queries.get("node").filter(|node| !node.is_empty() && node.as_str() != "all") {
|
||||
// The durable backlog ledger is cluster-shared with no per-node
|
||||
// attribution, so a node-scoped request still sees the complete
|
||||
// (superset) backlog.
|
||||
debug!(node = %node, "replication mrf node filter has no effect on the cluster-shared backlog");
|
||||
}
|
||||
|
||||
let durable = crate::admin::storage_api::replication::read_durable_mrf_backlog(store).await;
|
||||
let bucket_stats = cluster_replication_stats(&bucket, app_context_from_req(&req)).await;
|
||||
let response = build_mrf_response(bucket, &bucket_stats, durable);
|
||||
let response = build_mrf_response(bucket, &bucket_stats, &durable);
|
||||
|
||||
let data = serde_json::to_vec(&response)
|
||||
if !durable.available && !aggregate {
|
||||
// The madmin stream has no envelope to carry source health, and
|
||||
// madmin only decodes the body of a 200 — an empty stream would
|
||||
// read as a clean, healthy zero-row backlog. Fail loudly instead;
|
||||
// aggregate mode still reports the availability fields.
|
||||
tracing::warn!(
|
||||
bucket = %response.bucket,
|
||||
"durable MRF backlog is unreadable; failing the stream request — use aggregate=true to see source health"
|
||||
);
|
||||
return Err(S3Error::with_message(
|
||||
S3ErrorCode::ServiceUnavailable,
|
||||
"durable MRF backlog is unreadable; retry, or use aggregate=true for source health".to_string(),
|
||||
));
|
||||
}
|
||||
|
||||
let (data, truncated) = render_mrf_backlog(&response, &durable, aggregate)
|
||||
.map_err(|e| S3Error::with_message(S3ErrorCode::InternalError, format!("serialize failed: {e}")))?;
|
||||
let mut headers = HeaderMap::new();
|
||||
headers.insert(CONTENT_TYPE, HeaderValue::from_static("application/json"));
|
||||
if truncated {
|
||||
// The madmin stream has no envelope to carry truncation; signal
|
||||
// it out-of-band (madmin/mc ignore unknown headers), mirroring
|
||||
// x-rustfs-replication-diff-truncated.
|
||||
tracing::warn!(
|
||||
bucket = %response.bucket,
|
||||
max_entries = REPLICATION_MRF_MAX_STREAM_ENTRIES,
|
||||
"replication mrf stream truncated; narrow with ?bucket= or drain the backlog"
|
||||
);
|
||||
headers.insert("x-rustfs-replication-mrf-truncated", HeaderValue::from_static("true"));
|
||||
}
|
||||
Ok(S3Response::with_headers((StatusCode::OK, Body::from(data)), headers))
|
||||
}
|
||||
}
|
||||
@@ -1292,7 +1460,8 @@ mod tests {
|
||||
use super::{
|
||||
REMOTE_TARGET_UNSUPPORTED_FIELDS, REMOTE_TARGET_WRITABLE_FIELDS, RemoteTargetCredentialsRequest, RemoteTargetRequest,
|
||||
ReplicationDiffEntry, SUPPORTED_REMOTE_TARGET_API, TargetUpdateOp, build_mrf_response, extract_query_params,
|
||||
parse_remote_target_update_ops, render_replication_diff, unique_replication_peers, validate_remote_target_tls_settings,
|
||||
parse_remote_target_update_ops, render_mrf_backlog, render_replication_diff, unique_replication_peers,
|
||||
validate_remote_target_tls_settings,
|
||||
};
|
||||
use crate::admin::storage_api::bucket::target::{BucketTarget, LatencyStat};
|
||||
use crate::admin::storage_api::replication::{BucketStats, DurableMrfBacklog, MrfOpKind, MrfReplicateEntry};
|
||||
@@ -1510,7 +1679,7 @@ mod tests {
|
||||
],
|
||||
};
|
||||
|
||||
let response = build_mrf_response("bucket-a".to_string(), &stats, durable);
|
||||
let response = build_mrf_response("bucket-a".to_string(), &stats, &durable);
|
||||
let json = serde_json::to_value(response).expect("MRF response should serialize");
|
||||
|
||||
assert_eq!(json["TotalFailedCount"], 3);
|
||||
@@ -1522,7 +1691,9 @@ mod tests {
|
||||
assert_eq!(json["RuntimeStatsAvailable"], true);
|
||||
assert_eq!(json["ClusterComplete"], false);
|
||||
assert_eq!(json["Targets"][0]["ObservationScope"], "partial_cluster");
|
||||
assert_eq!(json["PerObjectEntriesAvailable"], false);
|
||||
// The bare stream enumerates the durable backlog per object, so a
|
||||
// readable backlog advertises the enumerable API.
|
||||
assert_eq!(json["PerObjectEntriesAvailable"], true);
|
||||
assert_eq!(json["PerTargetDurableEntriesAvailable"], true);
|
||||
|
||||
let targets = json["Targets"].as_array().expect("targets should serialize as an array");
|
||||
@@ -1569,7 +1740,7 @@ mod tests {
|
||||
}],
|
||||
};
|
||||
|
||||
let response = build_mrf_response("bucket-a".to_string(), &stats, durable);
|
||||
let response = build_mrf_response("bucket-a".to_string(), &stats, &durable);
|
||||
let json = serde_json::to_value(response).expect("MRF response should serialize");
|
||||
|
||||
assert_eq!(json["DurableBacklogAvailable"], true);
|
||||
@@ -1587,7 +1758,7 @@ mod tests {
|
||||
|
||||
#[test]
|
||||
fn mrf_response_distinguishes_unavailable_sources_from_valid_zero() {
|
||||
let unavailable = build_mrf_response("bucket-a".to_string(), &BucketStats::default(), DurableMrfBacklog::default());
|
||||
let unavailable = build_mrf_response("bucket-a".to_string(), &BucketStats::default(), &DurableMrfBacklog::default());
|
||||
let unavailable_json = serde_json::to_value(unavailable).expect("unavailable response should serialize");
|
||||
assert_eq!(unavailable_json["RuntimeStatsAvailable"], false);
|
||||
assert_eq!(unavailable_json["DurableBacklogAvailable"], false);
|
||||
@@ -1600,7 +1771,7 @@ mod tests {
|
||||
let valid_empty = build_mrf_response(
|
||||
"bucket-a".to_string(),
|
||||
&valid_empty_stats,
|
||||
DurableMrfBacklog {
|
||||
&DurableMrfBacklog {
|
||||
available: true,
|
||||
entries: Vec::new(),
|
||||
},
|
||||
@@ -1613,6 +1784,151 @@ mod tests {
|
||||
assert_eq!(valid_empty_json["PerTargetDurableEntriesAvailable"], true);
|
||||
}
|
||||
|
||||
fn sample_durable_backlog() -> DurableMrfBacklog {
|
||||
DurableMrfBacklog {
|
||||
available: true,
|
||||
entries: vec![
|
||||
MrfReplicateEntry {
|
||||
bucket: "bucket-a".to_string(),
|
||||
object: "object-a".to_string(),
|
||||
version_id: Some(uuid::Uuid::from_u128(7)),
|
||||
retry_count: 2,
|
||||
size: 250,
|
||||
op: MrfOpKind::Object,
|
||||
target_arns: vec!["arn-a".to_string()],
|
||||
..Default::default()
|
||||
},
|
||||
MrfReplicateEntry {
|
||||
bucket: "other-bucket".to_string(),
|
||||
object: "object-b".to_string(),
|
||||
version_id: None,
|
||||
retry_count: 0,
|
||||
size: 999,
|
||||
op: MrfOpKind::Object,
|
||||
target_arns: Vec::new(),
|
||||
..Default::default()
|
||||
},
|
||||
],
|
||||
}
|
||||
}
|
||||
|
||||
/// madmin's `BucketReplicationMRF` decodes the body one `ReplicationMRF`
|
||||
/// JSON document at a time; the default response must therefore be a bare
|
||||
/// document stream with madmin's exact json tags, not an envelope.
|
||||
#[test]
|
||||
fn mrf_stream_renders_bare_madmin_documents() {
|
||||
let durable = sample_durable_backlog();
|
||||
let response = build_mrf_response("bucket-a".to_string(), &BucketStats::default(), &durable);
|
||||
|
||||
let (body, _) = render_mrf_backlog(&response, &durable, false).expect("stream body should serialize");
|
||||
let text = String::from_utf8(body).expect("body should be utf-8");
|
||||
let lines: Vec<&str> = text.lines().filter(|line| !line.trim().is_empty()).collect();
|
||||
|
||||
// Only the entry matching the requested bucket is streamed.
|
||||
assert_eq!(lines.len(), 1, "expected one MRF document, got: {text}");
|
||||
let doc: serde_json::Value = serde_json::from_str(lines[0]).expect("each line should be a JSON document");
|
||||
assert_eq!(doc["bucket"], "bucket-a");
|
||||
assert_eq!(doc["object"], "object-a");
|
||||
assert_eq!(doc["versionId"], uuid::Uuid::from_u128(7).to_string());
|
||||
assert_eq!(doc["retryCount"], 2);
|
||||
// madmin `ReplicationMRF` has a `nodeName` tag; the durable backlog is
|
||||
// cluster-shared, so RustFS reports an empty node name.
|
||||
assert_eq!(doc["nodeName"], "");
|
||||
// The envelope keys must not leak into the stream: a `"Bucket"` key
|
||||
// would case-insensitively populate `ReplicationMRF.Bucket` and render
|
||||
// a phantom row in `mc replicate backlog`.
|
||||
assert!(doc.get("Bucket").is_none());
|
||||
assert!(doc.get("Targets").is_none());
|
||||
}
|
||||
|
||||
/// An empty backlog must produce an empty body: madmin's decoder loop then
|
||||
/// terminates on io.EOF with zero rows instead of one phantom row.
|
||||
#[test]
|
||||
fn mrf_stream_renders_empty_body_for_no_entries() {
|
||||
let durable = DurableMrfBacklog {
|
||||
available: true,
|
||||
entries: Vec::new(),
|
||||
};
|
||||
let response = build_mrf_response("bucket-a".to_string(), &BucketStats::default(), &durable);
|
||||
|
||||
let (body, _) = render_mrf_backlog(&response, &durable, false).expect("stream body should serialize");
|
||||
assert!(
|
||||
body.is_empty(),
|
||||
"empty backlog must serialize to an empty body, got: {}",
|
||||
String::from_utf8_lossy(&body)
|
||||
);
|
||||
}
|
||||
|
||||
/// `?aggregate=true` (RustFS extension) keeps the enveloped counter shape.
|
||||
#[test]
|
||||
fn mrf_aggregate_envelope_retains_counters() {
|
||||
let durable = sample_durable_backlog();
|
||||
let response = build_mrf_response("bucket-a".to_string(), &BucketStats::default(), &durable);
|
||||
|
||||
let (body, _) = render_mrf_backlog(&response, &durable, true).expect("aggregate body should serialize");
|
||||
let json: serde_json::Value = serde_json::from_slice(&body).expect("aggregate body should be one JSON object");
|
||||
assert_eq!(json["Bucket"], "bucket-a");
|
||||
assert_eq!(json["DurableCount"], 1);
|
||||
assert_eq!(json["DurableBacklogAvailable"], true);
|
||||
// The bare stream is an enumerable per-object API, so the aggregate
|
||||
// shell now truthfully advertises it whenever the backlog is readable.
|
||||
assert_eq!(json["PerObjectEntriesAvailable"], true);
|
||||
}
|
||||
|
||||
/// The nil UUID is RustFS's in-memory null-version sentinel; the wire
|
||||
/// token is `null`, never the zero UUID (second review round).
|
||||
#[test]
|
||||
fn mrf_stream_maps_nil_version_to_null_token() {
|
||||
let durable = DurableMrfBacklog {
|
||||
available: true,
|
||||
entries: vec![MrfReplicateEntry {
|
||||
bucket: "bucket-a".to_string(),
|
||||
object: "null-version-object".to_string(),
|
||||
version_id: Some(uuid::Uuid::nil()),
|
||||
retry_count: 1,
|
||||
size: 10,
|
||||
op: MrfOpKind::Object,
|
||||
..Default::default()
|
||||
}],
|
||||
};
|
||||
let response = build_mrf_response("bucket-a".to_string(), &BucketStats::default(), &durable);
|
||||
|
||||
let (body, truncated) = render_mrf_backlog(&response, &durable, false).expect("stream body should serialize");
|
||||
assert!(!truncated);
|
||||
let doc: serde_json::Value =
|
||||
serde_json::from_str(String::from_utf8(body).expect("utf-8").lines().next().expect("one line"))
|
||||
.expect("line should be a JSON document");
|
||||
assert_eq!(doc["versionId"], "null");
|
||||
}
|
||||
|
||||
/// The durable ledger is not bounded by the in-memory pending cap; the
|
||||
/// stream must stop at the documented bound and signal truncation
|
||||
/// (second review round).
|
||||
#[test]
|
||||
fn mrf_stream_truncates_at_the_documented_bound() {
|
||||
let entries = (0..super::REPLICATION_MRF_MAX_STREAM_ENTRIES + 1)
|
||||
.map(|index| MrfReplicateEntry {
|
||||
bucket: "bucket-a".to_string(),
|
||||
object: format!("object-{index}"),
|
||||
retry_count: 1,
|
||||
op: MrfOpKind::Object,
|
||||
..Default::default()
|
||||
})
|
||||
.collect();
|
||||
let durable = DurableMrfBacklog {
|
||||
available: true,
|
||||
entries,
|
||||
};
|
||||
let response = build_mrf_response("bucket-a".to_string(), &BucketStats::default(), &durable);
|
||||
|
||||
let (body, truncated) = render_mrf_backlog(&response, &durable, false).expect("stream body should serialize");
|
||||
assert!(truncated, "one entry past the bound must signal truncation");
|
||||
assert_eq!(
|
||||
String::from_utf8(body).expect("utf-8").lines().count(),
|
||||
super::REPLICATION_MRF_MAX_STREAM_ENTRIES
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_extract_query_params_decodes_percent_encoded_values() {
|
||||
let uri: Uri = "/rustfs/admin/v3/list-remote-targets?bucket=foo%2Fbar&flag=a+b"
|
||||
|
||||
@@ -1067,9 +1067,13 @@ fn parse_site_replication_state(data: &[u8]) -> S3Result<SiteReplicationState> {
|
||||
state.peers = normalize_peer_map_by_identity(state.peers);
|
||||
// A peer-edit high-water mark only fences a CURRENT peer. A site that
|
||||
// leaves drops below two peers, which clears its own state object and
|
||||
// restarts its generation counter at zero — a mark left over from the
|
||||
// previous membership would then reject every edit it sends after it
|
||||
// rejoins. Dropping departed origins on load also keeps the map bounded.
|
||||
// restarts its generation counter — a mark left over from the previous
|
||||
// membership must not reject the edits it sends after it rejoins. This
|
||||
// pruning covers departures THIS site observed; an origin removed
|
||||
// unilaterally elsewhere stays in this peer map with its mark, and the
|
||||
// wall-clock floor in `next_peer_edit_generation` is what lifts its
|
||||
// restarted counter over that mark. Dropping departed origins on load
|
||||
// also keeps the map bounded.
|
||||
state
|
||||
.applied_edit_generations
|
||||
.retain(|origin, _| state.peers.contains_key(origin));
|
||||
@@ -5935,11 +5939,51 @@ fn summarize_peer_error_detail(detail: &str) -> String {
|
||||
summary
|
||||
}
|
||||
|
||||
/// Allocate the next peer-edit generation. Called inside the state
|
||||
/// transaction, so the counter is handed out under the distributed
|
||||
/// state-object lock and two nodes of this site can never take the same one.
|
||||
/// The wall clock in unix nanoseconds, clamped into u64. A pre-1970 (or
|
||||
/// post-2554) clock yields 0, which makes the hybrid allocation below
|
||||
/// degrade to the plain `previous + 1` counter — monotone, never panicking.
|
||||
fn edit_generation_wall_clock() -> u64 {
|
||||
u64::try_from(OffsetDateTime::now_utc().unix_timestamp_nanos()).unwrap_or(0)
|
||||
}
|
||||
|
||||
/// Allocate the next peer-edit generation as a hybrid logical clock:
|
||||
/// `max(wall clock in unix nanoseconds, previous + 1)`. Called inside the
|
||||
/// state transaction, so the value is handed out under the distributed
|
||||
/// state-object lock and two nodes of this site can never take the same one
|
||||
/// (`previous + 1` keeps the sequence strictly increasing even when two
|
||||
/// allocations land in one clock tick, and keeps it monotone on a node
|
||||
/// whose clock stepped backwards mid-lifetime).
|
||||
///
|
||||
/// The wall-clock floor is what survives the counter's death. A site
|
||||
/// removed while unreachable — the receiver never dropped it from its peer
|
||||
/// map, so the load-time mark pruning in `parse_site_replication_state`
|
||||
/// never fired — that later rejoins recreates its state object with the
|
||||
/// counter back at zero. A plain counter would then hand out generations
|
||||
/// below the receiver's stale high-water mark and every delivery would be
|
||||
/// silently fenced until the counter caught up. Jumping to wall time clears
|
||||
/// that mark: every value the deleted lifetime handed out was capped by the
|
||||
/// wall clock at its own allocation (or by a prior lifetime's cap, applied
|
||||
/// inductively), so the recreated lifetime's first allocation exceeds them
|
||||
/// all — while a pre-removal delivery still in flight stays below the new
|
||||
/// floor and remains correctly fenced. Marks recorded by pre-hybrid
|
||||
/// receivers (small plain-counter values) sit far below any wall-clock
|
||||
/// value, so a restarted origin passes those too — the fix needs only the
|
||||
/// sender upgraded, nothing on the wire or in the receiver changed.
|
||||
///
|
||||
/// A wall clock that regresses across a delete/recreate (the recreating
|
||||
/// node's clock behind the clock that fed the previous lifetime) mints
|
||||
/// below the stale mark and the origin stays fenced — but only until real
|
||||
/// time passes the previous lifetime's last allocation, because every later
|
||||
/// allocation takes the wall-clock floor again. Bounded by the skew,
|
||||
/// self-healing, and no rollback window beyond the plain counter's: a
|
||||
/// delivery applies only at or above the receiver's mark, so the one
|
||||
/// cross-lifetime interleaving that can apply stale content — a
|
||||
/// pre-removal delivery whose generation lands above everything the
|
||||
/// regressed new lifetime has minted — required the same straggler landing
|
||||
/// above the mark under the plain counter, where the recreated counter's
|
||||
/// low restart made it strictly easier to hit.
|
||||
fn next_peer_edit_generation(state: &mut SiteReplicationState) -> u64 {
|
||||
state.edit_generation = state.edit_generation.saturating_add(1);
|
||||
state.edit_generation = edit_generation_wall_clock().max(state.edit_generation.saturating_add(1));
|
||||
state.edit_generation
|
||||
}
|
||||
|
||||
@@ -13244,6 +13288,104 @@ mod tests {
|
||||
assert!(!peer_edit_delivery_is_stale(&reloaded, "origin-site", 1));
|
||||
}
|
||||
|
||||
/// The unilateral-removal rejoin gap the hybrid clock closes. The origin
|
||||
/// was removed while unreachable, but THIS site never dropped it from
|
||||
/// its peer map, so the load-time mark pruning never fired and the mark
|
||||
/// from the previous membership survives. The origin's recreated state
|
||||
/// object restarts its counter, and with a plain `previous + 1` counter
|
||||
/// every delivery it sent — generations 1, 2, … below the stale mark —
|
||||
/// would be silently acked-and-dropped until the counter caught up. The
|
||||
/// wall-clock floor in `next_peer_edit_generation` lifts the restarted
|
||||
/// counter over every value the deleted lifetime handed out. Reverting
|
||||
/// the allocation to the plain counter (dropping the wall-clock max)
|
||||
/// turns the not-stale assertion red.
|
||||
#[test]
|
||||
fn hybrid_generation_unfences_a_rejoined_origin_whose_counter_restarted() {
|
||||
// First lifetime of the origin's state object: two allocations, both
|
||||
// capped by the wall clock at their own allocation.
|
||||
let mut first_life = SiteReplicationState::default();
|
||||
let straggler = next_peer_edit_generation(&mut first_life);
|
||||
let last_applied = next_peer_edit_generation(&mut first_life);
|
||||
assert!(last_applied > straggler, "allocations must be strictly increasing");
|
||||
|
||||
// The receiver applied up to `last_applied` and keeps the origin in
|
||||
// its peer map across the unilateral removal — reloading must keep
|
||||
// the mark, which is exactly why pruning cannot cover this case.
|
||||
let mut receiver = SiteReplicationState::default();
|
||||
receiver.peers.insert(
|
||||
"origin-site".to_string(),
|
||||
PeerInfo {
|
||||
deployment_id: "origin-site".to_string(),
|
||||
..peer("origin", "https://origin.example:9000")
|
||||
},
|
||||
);
|
||||
record_applied_peer_edit_generation(&mut receiver, "origin-site", last_applied);
|
||||
let mut receiver = parse_site_replication_state(&serde_json::to_vec(&receiver).expect("serialize")).expect("reload");
|
||||
assert_eq!(receiver.applied_edit_generations.get("origin-site"), Some(&last_applied));
|
||||
|
||||
// The origin rejoins with a RECREATED state object: counter back at
|
||||
// zero. The wall-clock floor must lift its first allocation over the
|
||||
// previous lifetime's mark…
|
||||
let mut second_life = SiteReplicationState::default();
|
||||
let restarted = next_peer_edit_generation(&mut second_life);
|
||||
assert!(
|
||||
!peer_edit_delivery_is_stale(&receiver, "origin-site", restarted),
|
||||
"the recreated lifetime's first allocation ({restarted}) must not be fenced by the previous lifetime's mark ({last_applied})"
|
||||
);
|
||||
record_applied_peer_edit_generation(&mut receiver, "origin-site", restarted);
|
||||
|
||||
// …while a pre-removal delivery still in flight stays below the new
|
||||
// floor and remains correctly fenced — the rollback the fence exists
|
||||
// to reject.
|
||||
assert!(
|
||||
peer_edit_delivery_is_stale(&receiver, "origin-site", straggler),
|
||||
"a pre-removal in-flight delivery ({straggler}) must stay fenced after the rejoin"
|
||||
);
|
||||
}
|
||||
|
||||
/// Marks recorded before the hybrid clock existed are small plain-counter
|
||||
/// values, far below any wall-clock allocation: a restarted origin passes
|
||||
/// them as soon as the SENDER runs the hybrid clock — nothing changes on
|
||||
/// the wire or in the receiver, so pre-hybrid receivers get the fix too.
|
||||
/// The other direction is unchanged: among plain-counter values the
|
||||
/// generation order still fences the delivery that lost the race.
|
||||
#[test]
|
||||
fn hybrid_generation_passes_marks_recorded_by_plain_counter_receivers() {
|
||||
let mut receiver = SiteReplicationState::default();
|
||||
record_applied_peer_edit_generation(&mut receiver, "origin-site", 57);
|
||||
assert!(peer_edit_delivery_is_stale(&receiver, "origin-site", 56));
|
||||
assert!(!peer_edit_delivery_is_stale(&receiver, "origin-site", 57));
|
||||
|
||||
let mut rejoined = SiteReplicationState::default();
|
||||
let restarted = next_peer_edit_generation(&mut rejoined);
|
||||
assert!(
|
||||
!peer_edit_delivery_is_stale(&receiver, "origin-site", restarted),
|
||||
"a wall-clock allocation ({restarted}) must clear a plain-counter mark (57)"
|
||||
);
|
||||
}
|
||||
|
||||
/// The `previous + 1` half of the hybrid clock: allocations stay strictly
|
||||
/// increasing even when the wall clock cannot move them forward — two
|
||||
/// allocations inside one clock tick, or a clock that stepped backwards
|
||||
/// mid-lifetime (a counter already ahead of the wall clock advances by
|
||||
/// exactly one per allocation instead of jumping back). Dropping the
|
||||
/// `previous + 1` half (allocating bare wall time) turns this red.
|
||||
#[test]
|
||||
fn hybrid_generation_is_strictly_increasing_when_the_clock_stalls() {
|
||||
let mut state = SiteReplicationState {
|
||||
// A counter far ahead of any wall clock this test will see.
|
||||
edit_generation: u64::MAX / 2,
|
||||
..Default::default()
|
||||
};
|
||||
assert_eq!(next_peer_edit_generation(&mut state), u64::MAX / 2 + 1);
|
||||
assert_eq!(next_peer_edit_generation(&mut state), u64::MAX / 2 + 2);
|
||||
// Saturation pins at the ceiling instead of wrapping; the equal-value
|
||||
// escape (`applied > generation` is false for equal) keeps deliveries
|
||||
// applying rather than fencing the origin out.
|
||||
state.edit_generation = u64::MAX;
|
||||
assert_eq!(next_peer_edit_generation(&mut state), u64::MAX);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_retry_stats_for_state_counts_pending_and_failed() {
|
||||
let state = SiteReplicationState {
|
||||
@@ -16044,10 +16186,77 @@ mod tests {
|
||||
generations.len(),
|
||||
"two nodes took the same edit generation, so their deliveries cannot be ordered: {generations:?}"
|
||||
);
|
||||
// The hybrid clock allocates `max(wall nanos, previous + 1)` — the
|
||||
// persisted counter is the largest allocation, and the `+ 1` half
|
||||
// keeps allocations distinct even inside one clock tick.
|
||||
assert_eq!(
|
||||
Some(&load_site_replication_state().await.expect("reload").edit_generation),
|
||||
unique.last(),
|
||||
"the persisted counter must be the largest allocation handed out"
|
||||
);
|
||||
}
|
||||
|
||||
/// The unilateral-removal rejoin, end to end across the state object's
|
||||
/// real lifecycle: dropping below two peers clears the object (the
|
||||
/// counter dies with it), and the recreated object's first allocation —
|
||||
/// raced by two nodes — must clear the previous lifetime's values via
|
||||
/// the wall-clock floor, so a receiver still holding the old mark
|
||||
/// accepts the restarted counter instead of fencing it.
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
|
||||
#[serial]
|
||||
async fn test_recreated_state_object_allocates_over_the_previous_lifetimes_mark() {
|
||||
publish_ready_iam_context().await;
|
||||
let seed = || SiteReplicationState {
|
||||
peers: ["site-a", "site-b"]
|
||||
.into_iter()
|
||||
.map(|name| (name.to_string(), peer(name, &format!("https://{name}.example:9000"))))
|
||||
.collect(),
|
||||
..Default::default()
|
||||
};
|
||||
|
||||
save_site_replication_state(&seed()).await.expect("seed state");
|
||||
let straggler = update_site_replication_state(|state| Ok(next_peer_edit_generation(state)))
|
||||
.await
|
||||
.expect("first-life allocation");
|
||||
let last_applied = update_site_replication_state(|state| Ok(next_peer_edit_generation(state)))
|
||||
.await
|
||||
.expect("first-life allocation");
|
||||
// A receiver that never dropped this site from its peer map holds
|
||||
// this mark across the removal.
|
||||
let mut receiver = SiteReplicationState::default();
|
||||
record_applied_peer_edit_generation(&mut receiver, "origin-site", last_applied);
|
||||
|
||||
// Unilateral removal: the site drops below two peers, which clears
|
||||
// its state object and the counter with it.
|
||||
let mut departed = seed();
|
||||
departed.peers.remove("site-b");
|
||||
save_site_replication_state(&departed).await.expect("clear state");
|
||||
assert_eq!(
|
||||
load_site_replication_state().await.expect("reload").edit_generation,
|
||||
generations.len() as u64,
|
||||
"the persisted counter must account for every allocation"
|
||||
0,
|
||||
"clearing the state object must take the counter with it"
|
||||
);
|
||||
|
||||
// Rejoin recreates the state object; two nodes race the first
|
||||
// allocation of the new life.
|
||||
save_site_replication_state(&seed()).await.expect("recreate state");
|
||||
let node_a = tokio::spawn(update_site_replication_state(|state| Ok(next_peer_edit_generation(state))));
|
||||
let node_b = tokio::spawn(update_site_replication_state(|state| Ok(next_peer_edit_generation(state))));
|
||||
let generation_a = node_a.await.expect("node a task").expect("node a allocation");
|
||||
let generation_b = node_b.await.expect("node b task").expect("node b allocation");
|
||||
assert_ne!(generation_a, generation_b, "racing allocations must stay distinct");
|
||||
|
||||
// The receiver's stale mark must not fence the restarted counter…
|
||||
let restarted = generation_a.min(generation_b);
|
||||
assert!(
|
||||
!peer_edit_delivery_is_stale(&receiver, "origin-site", restarted),
|
||||
"the recreated life's first allocation ({restarted}) must clear the previous life's mark ({last_applied})"
|
||||
);
|
||||
record_applied_peer_edit_generation(&mut receiver, "origin-site", restarted);
|
||||
// …while the cleared life's in-flight leftovers stay fenced.
|
||||
assert!(
|
||||
peer_edit_delivery_is_stale(&receiver, "origin-site", straggler),
|
||||
"a pre-removal in-flight delivery ({straggler}) must stay fenced after the rejoin"
|
||||
);
|
||||
}
|
||||
|
||||
|
||||
@@ -1459,10 +1459,13 @@ pub const ADMIN_ROUTE_POLICY_SPECS: &[AdminRouteSpec] = &[
|
||||
REPLICATION_DIFF,
|
||||
RouteRiskLevel::Sensitive,
|
||||
),
|
||||
// The default stream enumerates object names/version ids and requires
|
||||
// ReplicationDiff (MinIO parity); only ?aggregate=true relaxes to
|
||||
// GetReplicationMetrics in the handler.
|
||||
admin(
|
||||
HttpMethod::Get,
|
||||
"/rustfs/admin/v3/replication/mrf",
|
||||
GET_REPLICATION_METRICS,
|
||||
REPLICATION_DIFF,
|
||||
RouteRiskLevel::Sensitive,
|
||||
),
|
||||
];
|
||||
|
||||
@@ -195,6 +195,13 @@ IFS= read -r -d '' expected_docker_automatic_guard <<'EOF' || true
|
||||
EOF
|
||||
expected_docker_automatic_guard=${expected_docker_automatic_guard%$'\n'}
|
||||
require_job_if "$docker_workflow" "build-check" "$expected_docker_automatic_guard"
|
||||
require_line "$docker_workflow" ' source_ref: ${{ steps.check.outputs.source_ref }}' "Docker source ref output"
|
||||
require_line "$docker_workflow" ' source_ref="$HEAD_SHA"' "automatic Docker source ref"
|
||||
require_line "$docker_workflow" ' source_ref="$tag_ref"' "manual Docker source ref"
|
||||
require_line "$docker_workflow" ' ref: ${{ needs.build-check.outputs.source_ref }}' "Docker release source checkout"
|
||||
require_line "$docker_workflow" ' SOURCE_REVISION="$(git rev-parse HEAD)"' "Docker source revision resolution"
|
||||
require_line "$docker_workflow" ' LABELS="$LABELS,org.opencontainers.image.revision=$SOURCE_REVISION"' "Docker revision label"
|
||||
require_absent "$docker_workflow" 'org.opencontainers.image.revision=${{ github.sha }}' "Docker revision must not use the workflow branch SHA"
|
||||
|
||||
docker_manual_guard=$(awk '
|
||||
$0 == " *-preview*)" { in_preview = 1 }
|
||||
|
||||
Reference in New Issue
Block a user