fix(admin): report refreshed lock lease TTLs (#6230)

This commit is contained in:
cxymds
2026-08-19 11:02:07 +08:00
committed by GitHub
parent aa6b9001f1
commit 13a9d19505
7 changed files with 491 additions and 69 deletions
+330 -43
View File
@@ -28,17 +28,19 @@ use crate::admin::router::{AdminOperation, Operation, S3Router};
use crate::admin::storage_api::access::spawn_traced;
use crate::auth::{check_key_valid, get_session_token};
use crate::server::{ADMIN_PREFIX, RemoteAddr};
use crate::storage::storage_api::get_global_lock_clients;
use bytes::Bytes;
use futures::{Stream, StreamExt};
use futures::{Stream, StreamExt, future::join_all};
use http::{HeaderMap, HeaderValue, Uri, header::CONTENT_LENGTH};
use hyper::{Method, StatusCode};
use matchit::Params;
use rustfs_lock::{LockMode, ObjectKey, get_global_lock_manager};
use rustfs_lock::{LockLeaseInfo, LockMode, LockType, ObjectKey, get_global_lock_manager};
use rustfs_policy::policy::action::{Action, AdminAction};
use s3s::header::CONTENT_TYPE;
use s3s::stream::{ByteStream, DynByteStream};
use s3s::{Body, S3Error, S3ErrorCode, S3Request, S3Response, S3Result, StdError, s3_error};
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use std::pin::Pin;
use std::task::{Context, Poll};
use std::time::{Duration, SystemTime};
@@ -212,7 +214,142 @@ fn system_time_to_rfc3339(t: SystemTime) -> Option<String> {
dt.format(&time::format_description::well_known::Rfc3339).ok()
}
fn collect_top_locks(limit: usize) -> TopLocksResponse {
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
enum TopLockMode {
Read,
Write,
}
impl TopLockMode {
fn label(self) -> &'static str {
match self {
Self::Read => "READ",
Self::Write => "WRITE",
}
}
}
#[derive(Debug, PartialEq, Eq, Hash)]
struct LockHolderKey {
resource: ObjectKey,
mode: TopLockMode,
owner: String,
}
#[derive(Debug)]
struct TopLockState {
acquired_at: SystemTime,
ttl_secs: u64,
priority: &'static str,
}
#[derive(Debug)]
struct LeaseHolderState {
acquired_at: SystemTime,
ttl_secs: u64,
holder_count: u32,
}
fn build_top_locks_response(
limit: usize,
now: SystemTime,
lease_infos: Vec<LockLeaseInfo>,
fast_infos: Vec<(rustfs_lock::ObjectLockInfo, u32)>,
) -> TopLocksResponse {
let mut lease_holders = HashMap::with_capacity(lease_infos.len());
for info in lease_infos {
let mode = match info.lock_type {
LockType::Shared => TopLockMode::Read,
LockType::Exclusive => TopLockMode::Write,
};
let key = LockHolderKey {
resource: info.resource,
mode,
owner: info.owner,
};
let ttl_secs = info.remaining_ttl.as_secs();
lease_holders
.entry(key)
.and_modify(|state: &mut LeaseHolderState| {
if info.acquired_at < state.acquired_at {
state.acquired_at = info.acquired_at;
}
state.ttl_secs = state.ttl_secs.max(ttl_secs);
state.holder_count = state.holder_count.saturating_add(1);
})
.or_insert(LeaseHolderState {
acquired_at: info.acquired_at,
ttl_secs,
holder_count: 1,
});
}
let mut infos: Vec<_> = fast_infos
.into_iter()
.map(|(info, holder_count)| {
let mode = match info.mode {
LockMode::Shared => TopLockMode::Read,
LockMode::Exclusive => TopLockMode::Write,
};
let key = LockHolderKey {
resource: info.key,
mode,
owner: info.owner.to_string(),
};
let priority = lock_priority_label(info.priority);
// Shared-owner timestamps do not roll back when a newer sibling releases, so only their count is stable.
let state = match lease_holders.remove(&key) {
Some(lease)
if lease.holder_count == holder_count
&& (mode == TopLockMode::Read || info.acquired_at <= lease.acquired_at) =>
{
TopLockState {
acquired_at: lease.acquired_at,
ttl_secs: lease.ttl_secs,
priority,
}
}
_ => TopLockState {
acquired_at: info.acquired_at,
ttl_secs: info.expires_at.duration_since(now).unwrap_or(Duration::ZERO).as_secs(),
priority,
},
};
(key, state)
})
.collect();
// Longest-held first, matching MinIO's `top locks` ordering intent.
infos.sort_by_key(|(_, state)| state.acquired_at);
let total = infos.len();
let truncated = total > limit;
let locks = infos
.into_iter()
.take(limit)
.map(|(holder, state)| LockEntry {
resource: format!("{}/{}", holder.resource.bucket, holder.resource.object),
bucket: holder.resource.bucket.to_string(),
object: holder.resource.object.to_string(),
version: holder.resource.version.as_ref().map(|version| version.to_string()),
lock_type: holder.mode.label(),
owner: holder.owner,
priority: state.priority,
since: system_time_to_rfc3339(state.acquired_at),
elapsed_secs: now.duration_since(state.acquired_at).unwrap_or(Duration::ZERO).as_secs(),
ttl_secs: state.ttl_secs,
})
.collect();
TopLocksResponse {
total,
truncated,
locks,
capability_note: None,
}
}
async fn collect_top_locks(limit: usize) -> TopLocksResponse {
let manager = get_global_lock_manager();
let Some(fast) = manager.as_fast_lock_manager() else {
return TopLocksResponse {
@@ -225,43 +362,19 @@ fn collect_top_locks(limit: usize) -> TopLocksResponse {
};
};
let now = SystemTime::now();
let mut infos = fast.list_locks();
// Longest-held first, matching MinIO's `top locks` ordering intent.
infos.sort_by_key(|i| i.acquired_at);
let total = infos.len();
let truncated = total > limit;
let lease_infos = if let Some(clients) = get_global_lock_clients() {
join_all(clients.values().map(|client| client.list_lock_leases()))
.await
.into_iter()
.flatten()
.collect()
} else {
Vec::new()
};
// Capture holders last so released or replaced lease guards fail the merge checks.
let fast_infos = fast.list_locks_with_holder_counts();
let locks = infos
.into_iter()
.take(limit)
.map(|info| {
let elapsed_secs = now.duration_since(info.acquired_at).unwrap_or(Duration::ZERO).as_secs();
let ttl_secs = info.expires_at.duration_since(now).unwrap_or(Duration::ZERO).as_secs();
LockEntry {
resource: format!("{}/{}", info.key.bucket, info.key.object),
bucket: info.key.bucket.to_string(),
object: info.key.object.to_string(),
version: info.key.version.as_ref().map(|v| v.to_string()),
lock_type: match info.mode {
LockMode::Exclusive => "WRITE",
LockMode::Shared => "READ",
},
owner: info.owner.to_string(),
priority: lock_priority_label(info.priority),
since: system_time_to_rfc3339(info.acquired_at),
elapsed_secs,
ttl_secs,
}
})
.collect();
TopLocksResponse {
total,
truncated,
locks,
capability_note: None,
}
build_top_locks_response(limit, SystemTime::now(), lease_infos, fast_infos)
}
fn parse_top_locks_limit(uri: &Uri) -> usize {
@@ -279,7 +392,7 @@ impl Operation for TopLocksHandler {
async fn call(&self, req: S3Request<Body>, _params: Params<'_, '_>) -> S3Result<S3Response<(StatusCode, Body)>> {
authorize(&req, AdminAction::TopLocksAdminAction).await?;
let limit = parse_top_locks_limit(&req.uri);
let response = collect_top_locks(limit);
let response = collect_top_locks(limit).await;
json_response(StatusCode::OK, &response)
}
}
@@ -1107,6 +1220,180 @@ mod tests {
assert_eq!(parse_top_locks_limit(&Uri::from_static("/x?count=999999")), TOP_LOCKS_MAX_LIMIT);
}
#[test]
fn top_locks_prefers_renewable_lease_deadlines() {
let now = SystemTime::UNIX_EPOCH + Duration::from_secs(100);
let leased_resource = ObjectKey::new("bucket", "shared-object");
let exclusive_resource = ObjectKey::new("bucket", "write-object");
let direct_resource = ObjectKey::new("bucket", "direct-object");
let mixed_resource = ObjectKey::new("bucket", "mixed-object");
let replaced_resource = ObjectKey::new("bucket", "replaced-object");
let remaining_shared_resource = ObjectKey::new("bucket", "remaining-shared-object");
let response = build_top_locks_response(
TOP_LOCKS_DEFAULT_LIMIT,
now,
vec![
LockLeaseInfo {
resource: leased_resource.clone(),
lock_type: LockType::Shared,
owner: "owner-a".to_string(),
acquired_at: now - Duration::from_secs(50),
remaining_ttl: Duration::from_secs(5),
},
LockLeaseInfo {
resource: leased_resource.clone(),
lock_type: LockType::Shared,
owner: "owner-a".to_string(),
acquired_at: now - Duration::from_secs(40),
remaining_ttl: Duration::from_secs(20),
},
LockLeaseInfo {
resource: mixed_resource.clone(),
lock_type: LockType::Shared,
owner: "owner-c".to_string(),
acquired_at: now - Duration::from_secs(30),
remaining_ttl: Duration::from_secs(25),
},
LockLeaseInfo {
resource: exclusive_resource.clone(),
lock_type: LockType::Exclusive,
owner: "owner-d".to_string(),
acquired_at: now - Duration::from_secs(15),
remaining_ttl: Duration::from_secs(18),
},
LockLeaseInfo {
resource: replaced_resource.clone(),
lock_type: LockType::Exclusive,
owner: "owner-e".to_string(),
acquired_at: now - Duration::from_secs(30),
remaining_ttl: Duration::from_secs(25),
},
LockLeaseInfo {
resource: remaining_shared_resource.clone(),
lock_type: LockType::Shared,
owner: "owner-f".to_string(),
acquired_at: now - Duration::from_secs(30),
remaining_ttl: Duration::from_secs(22),
},
],
vec![
(
rustfs_lock::ObjectLockInfo {
key: replaced_resource,
mode: LockMode::Exclusive,
owner: "owner-e".into(),
acquired_at: now - Duration::from_secs(5),
expires_at: now + Duration::from_secs(4),
priority: rustfs_lock::fast_lock::LockPriority::Normal,
},
1,
),
(
rustfs_lock::ObjectLockInfo {
key: remaining_shared_resource,
mode: LockMode::Shared,
owner: "owner-f".into(),
acquired_at: now - Duration::from_secs(5),
expires_at: now + Duration::from_secs(3),
priority: rustfs_lock::fast_lock::LockPriority::Normal,
},
1,
),
(
rustfs_lock::ObjectLockInfo {
key: exclusive_resource,
mode: LockMode::Exclusive,
owner: "owner-d".into(),
acquired_at: now - Duration::from_secs(15),
expires_at: now + Duration::from_secs(2),
priority: rustfs_lock::fast_lock::LockPriority::Normal,
},
1,
),
(
rustfs_lock::ObjectLockInfo {
key: leased_resource,
mode: LockMode::Shared,
owner: "owner-a".into(),
acquired_at: now - Duration::from_secs(50),
expires_at: now + Duration::from_secs(1),
priority: rustfs_lock::fast_lock::LockPriority::Normal,
},
2,
),
(
rustfs_lock::ObjectLockInfo {
key: direct_resource,
mode: LockMode::Exclusive,
owner: "owner-b".into(),
acquired_at: now - Duration::from_secs(10),
expires_at: now + Duration::from_secs(7),
priority: rustfs_lock::fast_lock::LockPriority::Normal,
},
1,
),
(
rustfs_lock::ObjectLockInfo {
key: mixed_resource,
mode: LockMode::Shared,
owner: "owner-c".into(),
acquired_at: now - Duration::from_secs(30),
expires_at: now + Duration::from_secs(9),
priority: rustfs_lock::fast_lock::LockPriority::Normal,
},
2,
),
],
);
assert_eq!(response.total, 6);
let leased = response
.locks
.iter()
.find(|entry| entry.object == "shared-object")
.expect("lease-backed shared owner should be listed once");
assert_eq!(leased.lock_type, "READ");
assert_eq!(leased.elapsed_secs, 50);
assert_eq!(leased.ttl_secs, 20);
let exclusive = response
.locks
.iter()
.find(|entry| entry.object == "write-object")
.expect("lease-backed exclusive holder should be listed");
assert_eq!(exclusive.lock_type, "WRITE");
assert_eq!(exclusive.ttl_secs, 18);
let direct = response
.locks
.iter()
.find(|entry| entry.object == "direct-object")
.expect("direct fast lock should remain visible");
assert_eq!(direct.ttl_secs, 7);
let mixed = response
.locks
.iter()
.find(|entry| entry.object == "mixed-object")
.expect("mixed direct and leased shared holders should remain visible");
assert_eq!(mixed.ttl_secs, 9);
let replaced = response
.locks
.iter()
.find(|entry| entry.object == "replaced-object")
.expect("a replaced lease holder should remain visible");
assert_eq!(replaced.ttl_secs, 4);
let remaining_shared = response
.locks
.iter()
.find(|entry| entry.object == "remaining-shared-object")
.expect("an older surviving shared lease should remain lease-backed");
assert_eq!(remaining_shared.ttl_secs, 22);
}
#[tokio::test]
async fn collect_top_locks_reports_live_lock() {
// Acquire a real lock through the global manager and confirm it surfaces.
@@ -1115,20 +1402,20 @@ mod tests {
// The fast-lock manager exposes the acquire API; if the lock subsystem is
// disabled in this environment, the response must carry a capability note.
let Some(fast) = manager.as_fast_lock_manager() else {
let response = collect_top_locks(TOP_LOCKS_DEFAULT_LIMIT);
let response = collect_top_locks(TOP_LOCKS_DEFAULT_LIMIT).await;
assert!(response.capability_note.is_some() || response.locks.is_empty());
return;
};
let guard = match fast.acquire_write_lock(key.clone(), "diag-owner").await {
Ok(g) => g,
Err(_) => {
let response = collect_top_locks(TOP_LOCKS_DEFAULT_LIMIT);
let response = collect_top_locks(TOP_LOCKS_DEFAULT_LIMIT).await;
assert!(response.capability_note.is_some() || response.locks.is_empty());
return;
}
};
let response = collect_top_locks(TOP_LOCKS_DEFAULT_LIMIT);
let response = collect_top_locks(TOP_LOCKS_DEFAULT_LIMIT).await;
let found = response
.locks
.iter()