drop common/error

This commit is contained in:
weisd
2025-06-06 16:34:59 +08:00
parent c589972fa7
commit f199718e3b
12 changed files with 81 additions and 72 deletions
+13 -9
View File
@@ -1,10 +1,10 @@
use common::error::Result;
use rsa::Pkcs1v15Encrypt; use rsa::Pkcs1v15Encrypt;
use rsa::{ use rsa::{
pkcs8::{DecodePrivateKey, DecodePublicKey},
RsaPrivateKey, RsaPublicKey, RsaPrivateKey, RsaPublicKey,
pkcs8::{DecodePrivateKey, DecodePublicKey},
}; };
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
use std::io::{Error, Result};
#[derive(Serialize, Deserialize, Debug, Default, Clone)] #[derive(Serialize, Deserialize, Debug, Default, Clone)]
pub struct Token { pub struct Token {
@@ -18,8 +18,10 @@ pub struct Token {
// 返回 base64 处理的加密字符串 // 返回 base64 处理的加密字符串
pub fn gencode(token: &Token, key: &str) -> Result<String> { pub fn gencode(token: &Token, key: &str) -> Result<String> {
let data = serde_json::to_vec(token)?; let data = serde_json::to_vec(token)?;
let public_key = RsaPublicKey::from_public_key_pem(key)?; let public_key = RsaPublicKey::from_public_key_pem(key).map_err(Error::other)?;
let encrypted_data = public_key.encrypt(&mut rand::thread_rng(), Pkcs1v15Encrypt, &data)?; let encrypted_data = public_key
.encrypt(&mut rand::thread_rng(), Pkcs1v15Encrypt, &data)
.map_err(Error::other)?;
Ok(base64_simd::URL_SAFE_NO_PAD.encode_to_string(&encrypted_data)) Ok(base64_simd::URL_SAFE_NO_PAD.encode_to_string(&encrypted_data))
} }
@@ -28,9 +30,11 @@ pub fn gencode(token: &Token, key: &str) -> Result<String> {
// [key] 私钥字符串 // [key] 私钥字符串
// 返回 Token 对象 // 返回 Token 对象
pub fn parse(token: &str, key: &str) -> Result<Token> { pub fn parse(token: &str, key: &str) -> Result<Token> {
let encrypted_data = base64_simd::URL_SAFE_NO_PAD.decode_to_vec(token.as_bytes())?; let encrypted_data = base64_simd::URL_SAFE_NO_PAD
let private_key = RsaPrivateKey::from_pkcs8_pem(key)?; .decode_to_vec(token.as_bytes())
let decrypted_data = private_key.decrypt(Pkcs1v15Encrypt, &encrypted_data)?; .map_err(Error::other)?;
let private_key = RsaPrivateKey::from_pkcs8_pem(key).map_err(Error::other)?;
let decrypted_data = private_key.decrypt(Pkcs1v15Encrypt, &encrypted_data).map_err(Error::other)?;
let res: Token = serde_json::from_slice(&decrypted_data)?; let res: Token = serde_json::from_slice(&decrypted_data)?;
Ok(res) Ok(res)
} }
@@ -49,14 +53,14 @@ pub fn parse_license(license: &str) -> Result<Token> {
// } // }
} }
static TEST_PRIVATE_KEY:&str ="-----BEGIN PRIVATE KEY-----\nMIIEvAIBADANBgkqhkiG9w0BAQEFAASCBKYwggSiAgEAAoIBAQCj86SrJIuxSxR6\nBJ/dlJEUIj6NeBRnhLQlCDdovuz61+7kJXVcxaR66w4m8W7SLEUP+IlPtnn6vmiG\n7XMhGNHIr7r1JsEVVLhZmL3tKI66DEZl786ZhG81BWqUlmcooIPS8UEPZNqJXLuz\nVGhxNyVGbj/tV7QC2pSISnKaixc+nrhxvo7w56p5qrm9tik0PjTgfZsUePkoBsSN\npoRkAauS14MAzK6HGB75CzG3dZqXUNWSWVocoWtQbZUwFGXyzU01ammsHQDvc2xu\nK1RQpd1qYH5bOWZ0N0aPFwT0r59HztFXg9sbjsnuhO1A7OiUOkc6iGVuJ0wm/9nA\nwZIBqzgjAgMBAAECggEAPMpeSEbotPhNw2BrllE76ec4omPfzPJbiU+em+wPGoNu\nRJHPDnMKJbl6Kd5jZPKdOOrCnxfd6qcnQsBQa/kz7+GYxMV12l7ra+1Cnujm4v0i\nLTHZvPpp8ZLsjeOmpF3AAzsJEJgon74OqtOlVjVIUPEYKvzV9ijt4gsYq0zfdYv0\nhrTMzyrGM4/UvKLsFIBROAfCeWfA7sXLGH8JhrRAyDrtCPzGtyyAmzoHKHtHafcB\nuyPFw/IP8otAgpDk5iiQPNkH0WwzAQIm12oHuNUa66NwUK4WEjXTnDg8KeWLHHNv\nIfN8vdbZchMUpMIvvkr7is315d8f2cHCB5gEO+GWAQKBgQDR/0xNll+FYaiUKCPZ\nvkOCAd3l5mRhsqnjPQ/6Ul1lAyYWpoJSFMrGGn/WKTa/FVFJRTGbBjwP+Mx10bfb\ngUg2GILDTISUh54fp4zngvTi9w4MWGKXrb7I1jPkM3vbJfC/v2fraQ/r7qHPpO2L\nf6ZbGxasIlSvr37KeGoelwcAQQKBgQDH3hmOTS2Hl6D4EXdq5meHKrfeoicGN7m8\noQK7u8iwn1R9zK5nh6IXxBhKYNXNwdCQtBZVRvFjjZ56SZJb7lKqa1BcTsgJfZCy\nnI3Uu4UykrECAH8AVCVqBXUDJmeA2yE+gDAtYEjvhSDHpUfWxoGHr0B/Oqk2Lxc/\npRy1qV5fYwKBgBWSL/hYVf+RhIuTg/s9/BlCr9SJ0g3nGGRrRVTlWQqjRCpXeFOO\nJzYqSq9pFGKUggEQxoOyJEFPwVDo9gXqRcyov+Xn2kaXl7qQr3yoixc1YZALFDWY\nd1ySBEqQr0xXnV9U/gvEgwotPRnjSzNlLWV2ZuHPtPtG/7M0o1H5GZMBAoGAKr3N\nW0gX53o+my4pCnxRQW+aOIsWq1a5aqRIEFudFGBOUkS2Oz+fI1P1GdrRfhnnfzpz\n2DK+plp/vIkFOpGhrf4bBlJ2psjqa7fdANRFLMaAAfyXLDvScHTQTCcnVUAHQPVq\n2BlSH56pnugyj7SNuLV6pnql+wdhAmRN2m9o1h8CgYAbX2juSr4ioXwnYjOUdrIY\n4+ERvHcXdjoJmmPcAm4y5NbSqLXyU0FQmplNMt2A5LlniWVJ9KNdjAQUt60FZw/+\nr76LdxXaHNZghyx0BOs7mtq5unSQXamZ8KixasfhE9uz3ij1jXjG6hafWkS8/68I\nuWbaZqgvy7a9oPHYlKH7Jg==\n-----END PRIVATE KEY-----\n"; static TEST_PRIVATE_KEY: &str = "-----BEGIN PRIVATE KEY-----\nMIIEvAIBADANBgkqhkiG9w0BAQEFAASCBKYwggSiAgEAAoIBAQCj86SrJIuxSxR6\nBJ/dlJEUIj6NeBRnhLQlCDdovuz61+7kJXVcxaR66w4m8W7SLEUP+IlPtnn6vmiG\n7XMhGNHIr7r1JsEVVLhZmL3tKI66DEZl786ZhG81BWqUlmcooIPS8UEPZNqJXLuz\nVGhxNyVGbj/tV7QC2pSISnKaixc+nrhxvo7w56p5qrm9tik0PjTgfZsUePkoBsSN\npoRkAauS14MAzK6HGB75CzG3dZqXUNWSWVocoWtQbZUwFGXyzU01ammsHQDvc2xu\nK1RQpd1qYH5bOWZ0N0aPFwT0r59HztFXg9sbjsnuhO1A7OiUOkc6iGVuJ0wm/9nA\nwZIBqzgjAgMBAAECggEAPMpeSEbotPhNw2BrllE76ec4omPfzPJbiU+em+wPGoNu\nRJHPDnMKJbl6Kd5jZPKdOOrCnxfd6qcnQsBQa/kz7+GYxMV12l7ra+1Cnujm4v0i\nLTHZvPpp8ZLsjeOmpF3AAzsJEJgon74OqtOlVjVIUPEYKvzV9ijt4gsYq0zfdYv0\nhrTMzyrGM4/UvKLsFIBROAfCeWfA7sXLGH8JhrRAyDrtCPzGtyyAmzoHKHtHafcB\nuyPFw/IP8otAgpDk5iiQPNkH0WwzAQIm12oHuNUa66NwUK4WEjXTnDg8KeWLHHNv\nIfN8vdbZchMUpMIvvkr7is315d8f2cHCB5gEO+GWAQKBgQDR/0xNll+FYaiUKCPZ\nvkOCAd3l5mRhsqnjPQ/6Ul1lAyYWpoJSFMrGGn/WKTa/FVFJRTGbBjwP+Mx10bfb\ngUg2GILDTISUh54fp4zngvTi9w4MWGKXrb7I1jPkM3vbJfC/v2fraQ/r7qHPpO2L\nf6ZbGxasIlSvr37KeGoelwcAQQKBgQDH3hmOTS2Hl6D4EXdq5meHKrfeoicGN7m8\noQK7u8iwn1R9zK5nh6IXxBhKYNXNwdCQtBZVRvFjjZ56SZJb7lKqa1BcTsgJfZCy\nnI3Uu4UykrECAH8AVCVqBXUDJmeA2yE+gDAtYEjvhSDHpUfWxoGHr0B/Oqk2Lxc/\npRy1qV5fYwKBgBWSL/hYVf+RhIuTg/s9/BlCr9SJ0g3nGGRrRVTlWQqjRCpXeFOO\nJzYqSq9pFGKUggEQxoOyJEFPwVDo9gXqRcyov+Xn2kaXl7qQr3yoixc1YZALFDWY\nd1ySBEqQr0xXnV9U/gvEgwotPRnjSzNlLWV2ZuHPtPtG/7M0o1H5GZMBAoGAKr3N\nW0gX53o+my4pCnxRQW+aOIsWq1a5aqRIEFudFGBOUkS2Oz+fI1P1GdrRfhnnfzpz\n2DK+plp/vIkFOpGhrf4bBlJ2psjqa7fdANRFLMaAAfyXLDvScHTQTCcnVUAHQPVq\n2BlSH56pnugyj7SNuLV6pnql+wdhAmRN2m9o1h8CgYAbX2juSr4ioXwnYjOUdrIY\n4+ERvHcXdjoJmmPcAm4y5NbSqLXyU0FQmplNMt2A5LlniWVJ9KNdjAQUt60FZw/+\nr76LdxXaHNZghyx0BOs7mtq5unSQXamZ8KixasfhE9uz3ij1jXjG6hafWkS8/68I\nuWbaZqgvy7a9oPHYlKH7Jg==\n-----END PRIVATE KEY-----\n";
#[cfg(test)] #[cfg(test)]
mod tests { mod tests {
use super::*; use super::*;
use rsa::{ use rsa::{
pkcs8::{EncodePrivateKey, EncodePublicKey, LineEnding},
RsaPrivateKey, RsaPrivateKey,
pkcs8::{EncodePrivateKey, EncodePublicKey, LineEnding},
}; };
use std::time::{SystemTime, UNIX_EPOCH}; use std::time::{SystemTime, UNIX_EPOCH};
#[test] #[test]
+1 -1
View File
@@ -1,4 +1,4 @@
pub mod error; // pub mod error;
pub mod globals; pub mod globals;
pub mod last_minute; pub mod last_minute;
+17 -14
View File
@@ -3,7 +3,7 @@ use std::time::{Duration, Instant};
use tokio::{sync::mpsc::Sender, time::sleep}; use tokio::{sync::mpsc::Sender, time::sleep};
use tracing::{info, warn}; use tracing::{info, warn};
use crate::{lock_args::LockArgs, LockApi, Locker}; use crate::{LockApi, Locker, lock_args::LockArgs};
const DRW_MUTEX_REFRESH_INTERVAL: Duration = Duration::from_secs(10); const DRW_MUTEX_REFRESH_INTERVAL: Duration = Duration::from_secs(10);
const LOCK_RETRY_MIN_INTERVAL: Duration = Duration::from_millis(250); const LOCK_RETRY_MIN_INTERVAL: Duration = Duration::from_millis(250);
@@ -117,7 +117,10 @@ impl DRWMutex {
quorum += 1; quorum += 1;
} }
} }
info!("lockBlocking {}/{} for {:?}: lockType readLock({}), additional opts: {:?}, quorum: {}, tolerance: {}, lockClients: {}\n", id, source, self.names, is_read_lock, opts, quorum, tolerance, locker_len); info!(
"lockBlocking {}/{} for {:?}: lockType readLock({}), additional opts: {:?}, quorum: {}, tolerance: {}, lockClients: {}\n",
id, source, self.names, is_read_lock, opts, quorum, tolerance, locker_len
);
// Recalculate tolerance after potential quorum adjustment // Recalculate tolerance after potential quorum adjustment
// Use saturating_sub to prevent underflow // Use saturating_sub to prevent underflow
@@ -376,8 +379,8 @@ mod tests {
use super::*; use super::*;
use crate::local_locker::LocalLocker; use crate::local_locker::LocalLocker;
use async_trait::async_trait; use async_trait::async_trait;
use common::error::{Error, Result};
use std::collections::HashMap; use std::collections::HashMap;
use std::io::{Error, Result};
use std::sync::{Arc, Mutex}; use std::sync::{Arc, Mutex};
// Mock locker for testing // Mock locker for testing
@@ -436,10 +439,10 @@ mod tests {
async fn lock(&mut self, args: &LockArgs) -> Result<bool> { async fn lock(&mut self, args: &LockArgs) -> Result<bool> {
let mut state = self.state.lock().unwrap(); let mut state = self.state.lock().unwrap();
if state.should_fail { if state.should_fail {
return Err(Error::from_string("Mock lock failure")); return Err(Error::other("Mock lock failure"));
} }
if !state.is_online { if !state.is_online {
return Err(Error::from_string("Mock locker offline")); return Err(Error::other("Mock locker offline"));
} }
// Check if already locked // Check if already locked
@@ -454,7 +457,7 @@ mod tests {
async fn unlock(&mut self, args: &LockArgs) -> Result<bool> { async fn unlock(&mut self, args: &LockArgs) -> Result<bool> {
let mut state = self.state.lock().unwrap(); let mut state = self.state.lock().unwrap();
if state.should_fail { if state.should_fail {
return Err(Error::from_string("Mock unlock failure")); return Err(Error::other("Mock unlock failure"));
} }
Ok(state.locks.remove(&args.uid).is_some()) Ok(state.locks.remove(&args.uid).is_some())
@@ -463,10 +466,10 @@ mod tests {
async fn rlock(&mut self, args: &LockArgs) -> Result<bool> { async fn rlock(&mut self, args: &LockArgs) -> Result<bool> {
let mut state = self.state.lock().unwrap(); let mut state = self.state.lock().unwrap();
if state.should_fail { if state.should_fail {
return Err(Error::from_string("Mock rlock failure")); return Err(Error::other("Mock rlock failure"));
} }
if !state.is_online { if !state.is_online {
return Err(Error::from_string("Mock locker offline")); return Err(Error::other("Mock locker offline"));
} }
// Check if write lock exists // Check if write lock exists
@@ -481,7 +484,7 @@ mod tests {
async fn runlock(&mut self, args: &LockArgs) -> Result<bool> { async fn runlock(&mut self, args: &LockArgs) -> Result<bool> {
let mut state = self.state.lock().unwrap(); let mut state = self.state.lock().unwrap();
if state.should_fail { if state.should_fail {
return Err(Error::from_string("Mock runlock failure")); return Err(Error::other("Mock runlock failure"));
} }
Ok(state.read_locks.remove(&args.uid).is_some()) Ok(state.read_locks.remove(&args.uid).is_some())
@@ -490,7 +493,7 @@ mod tests {
async fn refresh(&mut self, _args: &LockArgs) -> Result<bool> { async fn refresh(&mut self, _args: &LockArgs) -> Result<bool> {
let state = self.state.lock().unwrap(); let state = self.state.lock().unwrap();
if state.should_fail { if state.should_fail {
return Err(Error::from_string("Mock refresh failure")); return Err(Error::other("Mock refresh failure"));
} }
Ok(true) Ok(true)
} }
@@ -880,8 +883,8 @@ mod tests {
// Case 1: Even number of lockers // Case 1: Even number of lockers
let locks = vec!["uid1".to_string(), "uid2".to_string(), "uid3".to_string(), "uid4".to_string()]; let locks = vec!["uid1".to_string(), "uid2".to_string(), "uid3".to_string(), "uid4".to_string()];
let tolerance = 2; // locks.len() / 2 = 4 / 2 = 2 let tolerance = 2; // locks.len() / 2 = 4 / 2 = 2
// locks.len() - tolerance = 4 - 2 = 2, which equals tolerance // locks.len() - tolerance = 4 - 2 = 2, which equals tolerance
// So the special case applies: un_locks_failed >= tolerance // So the special case applies: un_locks_failed >= tolerance
// All 4 failed unlocks // All 4 failed unlocks
assert!(check_failed_unlocks(&locks, tolerance)); // 4 >= 2 = true assert!(check_failed_unlocks(&locks, tolerance)); // 4 >= 2 = true
@@ -897,8 +900,8 @@ mod tests {
// Case 2: Odd number of lockers // Case 2: Odd number of lockers
let locks = vec!["uid1".to_string(), "uid2".to_string(), "uid3".to_string()]; let locks = vec!["uid1".to_string(), "uid2".to_string(), "uid3".to_string()];
let tolerance = 1; // locks.len() / 2 = 3 / 2 = 1 let tolerance = 1; // locks.len() / 2 = 3 / 2 = 1
// locks.len() - tolerance = 3 - 1 = 2, which does NOT equal tolerance (1) // locks.len() - tolerance = 3 - 1 = 2, which does NOT equal tolerance (1)
// So the normal case applies: un_locks_failed > tolerance // So the normal case applies: un_locks_failed > tolerance
// 3 failed unlocks // 3 failed unlocks
assert!(check_failed_unlocks(&locks, tolerance)); // 3 > 1 = true assert!(check_failed_unlocks(&locks, tolerance)); // 3 > 1 = true
+1 -1
View File
@@ -3,11 +3,11 @@
use std::sync::Arc; use std::sync::Arc;
use async_trait::async_trait; use async_trait::async_trait;
use common::error::Result;
use lazy_static::lazy_static; use lazy_static::lazy_static;
use local_locker::LocalLocker; use local_locker::LocalLocker;
use lock_args::LockArgs; use lock_args::LockArgs;
use remote_client::RemoteClient; use remote_client::RemoteClient;
use std::io::Result;
use tokio::sync::RwLock; use tokio::sync::RwLock;
pub mod drwmutex; pub mod drwmutex;
+9 -9
View File
@@ -1,11 +1,11 @@
use async_trait::async_trait; use async_trait::async_trait;
use common::error::{Error, Result}; use std::io::{Error, Result};
use std::{ use std::{
collections::HashMap, collections::HashMap,
time::{Duration, Instant}, time::{Duration, Instant},
}; };
use crate::{lock_args::LockArgs, Locker}; use crate::{Locker, lock_args::LockArgs};
const MAX_DELETE_LIST: usize = 1000; const MAX_DELETE_LIST: usize = 1000;
@@ -116,7 +116,7 @@ impl LocalLocker {
impl Locker for LocalLocker { impl Locker for LocalLocker {
async fn lock(&mut self, args: &LockArgs) -> Result<bool> { async fn lock(&mut self, args: &LockArgs) -> Result<bool> {
if args.resources.len() > MAX_DELETE_LIST { if args.resources.len() > MAX_DELETE_LIST {
return Err(Error::from_string(format!( return Err(Error::other(format!(
"internal error: LocalLocker.lock called with more than {} resources", "internal error: LocalLocker.lock called with more than {} resources",
MAX_DELETE_LIST MAX_DELETE_LIST
))); )));
@@ -152,7 +152,7 @@ impl Locker for LocalLocker {
async fn unlock(&mut self, args: &LockArgs) -> Result<bool> { async fn unlock(&mut self, args: &LockArgs) -> Result<bool> {
if args.resources.len() > MAX_DELETE_LIST { if args.resources.len() > MAX_DELETE_LIST {
return Err(Error::from_string(format!( return Err(Error::other(format!(
"internal error: LocalLocker.unlock called with more than {} resources", "internal error: LocalLocker.unlock called with more than {} resources",
MAX_DELETE_LIST MAX_DELETE_LIST
))); )));
@@ -197,7 +197,7 @@ impl Locker for LocalLocker {
async fn rlock(&mut self, args: &LockArgs) -> Result<bool> { async fn rlock(&mut self, args: &LockArgs) -> Result<bool> {
if args.resources.len() != 1 { if args.resources.len() != 1 {
return Err(Error::from_string("internal error: localLocker.RLock called with more than one resource")); return Err(Error::other("internal error: localLocker.RLock called with more than one resource"));
} }
let resource = &args.resources[0]; let resource = &args.resources[0];
@@ -241,7 +241,7 @@ impl Locker for LocalLocker {
async fn runlock(&mut self, args: &LockArgs) -> Result<bool> { async fn runlock(&mut self, args: &LockArgs) -> Result<bool> {
if args.resources.len() != 1 { if args.resources.len() != 1 {
return Err(Error::from_string("internal error: localLocker.RLock called with more than one resource")); return Err(Error::other("internal error: localLocker.RLock called with more than one resource"));
} }
let mut reply = false; let mut reply = false;
@@ -249,7 +249,7 @@ impl Locker for LocalLocker {
match self.lock_map.get_mut(resource) { match self.lock_map.get_mut(resource) {
Some(lris) => { Some(lris) => {
if is_write_lock(lris) { if is_write_lock(lris) {
return Err(Error::from_string(format!("runlock attempted on a write locked entity: {}", resource))); return Err(Error::other(format!("runlock attempted on a write locked entity: {}", resource)));
} else { } else {
lris.retain(|lri| { lris.retain(|lri| {
if lri.uid == args.uid && (args.owner.is_empty() || lri.owner == args.owner) { if lri.uid == args.uid && (args.owner.is_empty() || lri.owner == args.owner) {
@@ -389,8 +389,8 @@ fn format_uuid(s: &mut String, idx: &usize) {
#[cfg(test)] #[cfg(test)]
mod test { mod test {
use super::LocalLocker; use super::LocalLocker;
use crate::{lock_args::LockArgs, Locker}; use crate::{Locker, lock_args::LockArgs};
use common::error::Result; use std::io::Result;
use tokio; use tokio;
#[tokio::test] #[tokio::test]
+1 -1
View File
@@ -125,7 +125,7 @@ impl LRWMutex {
mod test { mod test {
use std::{sync::Arc, time::Duration}; use std::{sync::Arc, time::Duration};
use common::error::Result; use std::io::Result;
use tokio::time::sleep; use tokio::time::sleep;
use crate::lrwmutex::LRWMutex; use crate::lrwmutex::LRWMutex;
+4 -4
View File
@@ -5,11 +5,11 @@ use tokio::sync::RwLock;
use uuid::Uuid; use uuid::Uuid;
use crate::{ use crate::{
LockApi,
drwmutex::{DRWMutex, Options}, drwmutex::{DRWMutex, Options},
lrwmutex::LRWMutex, lrwmutex::LRWMutex,
LockApi,
}; };
use common::error::Result; use std::io::Result;
pub type RWLockerImpl = Box<dyn RWLocker + Send + Sync>; pub type RWLockerImpl = Box<dyn RWLocker + Send + Sync>;
@@ -258,12 +258,12 @@ impl RWLocker for LocalLockInstance {
mod test { mod test {
use std::{sync::Arc, time::Duration}; use std::{sync::Arc, time::Duration};
use common::error::Result; use std::io::Result;
use tokio::sync::RwLock; use tokio::sync::RwLock;
use crate::{ use crate::{
drwmutex::Options, drwmutex::Options,
namespace_lock::{new_nslock, NsLockMap}, namespace_lock::{NsLockMap, new_nslock},
}; };
#[tokio::test] #[tokio::test]
+20 -20
View File
@@ -1,10 +1,10 @@
use async_trait::async_trait; use async_trait::async_trait;
use common::error::{Error, Result};
use protos::{node_service_time_out_client, proto_gen::node_service::GenerallyLockRequest}; use protos::{node_service_time_out_client, proto_gen::node_service::GenerallyLockRequest};
use std::io::{Error, Result};
use tonic::Request; use tonic::Request;
use tracing::info; use tracing::info;
use crate::{lock_args::LockArgs, Locker}; use crate::{Locker, lock_args::LockArgs};
#[derive(Debug, Clone)] #[derive(Debug, Clone)]
pub struct RemoteClient { pub struct RemoteClient {
@@ -25,13 +25,13 @@ impl Locker for RemoteClient {
let args = serde_json::to_string(args)?; let args = serde_json::to_string(args)?;
let mut client = node_service_time_out_client(&self.addr) let mut client = node_service_time_out_client(&self.addr)
.await .await
.map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; .map_err(|err| Error::other(format!("can not get client, err: {}", err)))?;
let request = Request::new(GenerallyLockRequest { args }); let request = Request::new(GenerallyLockRequest { args });
let response = client.lock(request).await?.into_inner(); let response = client.lock(request).await.map_err(Error::other)?.into_inner();
if let Some(error_info) = response.error_info { if let Some(error_info) = response.error_info {
return Err(Error::from_string(error_info)); return Err(Error::other(error_info));
} }
Ok(response.success) Ok(response.success)
@@ -42,13 +42,13 @@ impl Locker for RemoteClient {
let args = serde_json::to_string(args)?; let args = serde_json::to_string(args)?;
let mut client = node_service_time_out_client(&self.addr) let mut client = node_service_time_out_client(&self.addr)
.await .await
.map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; .map_err(|err| Error::other(format!("can not get client, err: {}", err)))?;
let request = Request::new(GenerallyLockRequest { args }); let request = Request::new(GenerallyLockRequest { args });
let response = client.un_lock(request).await?.into_inner(); let response = client.un_lock(request).await.map_err(Error::other)?.into_inner();
if let Some(error_info) = response.error_info { if let Some(error_info) = response.error_info {
return Err(Error::from_string(error_info)); return Err(Error::other(error_info));
} }
Ok(response.success) Ok(response.success)
@@ -59,13 +59,13 @@ impl Locker for RemoteClient {
let args = serde_json::to_string(args)?; let args = serde_json::to_string(args)?;
let mut client = node_service_time_out_client(&self.addr) let mut client = node_service_time_out_client(&self.addr)
.await .await
.map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; .map_err(|err| Error::other(format!("can not get client, err: {}", err)))?;
let request = Request::new(GenerallyLockRequest { args }); let request = Request::new(GenerallyLockRequest { args });
let response = client.r_lock(request).await?.into_inner(); let response = client.r_lock(request).await.map_err(Error::other)?.into_inner();
if let Some(error_info) = response.error_info { if let Some(error_info) = response.error_info {
return Err(Error::from_string(error_info)); return Err(Error::other(error_info));
} }
Ok(response.success) Ok(response.success)
@@ -76,13 +76,13 @@ impl Locker for RemoteClient {
let args = serde_json::to_string(args)?; let args = serde_json::to_string(args)?;
let mut client = node_service_time_out_client(&self.addr) let mut client = node_service_time_out_client(&self.addr)
.await .await
.map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; .map_err(|err| Error::other(format!("can not get client, err: {}", err)))?;
let request = Request::new(GenerallyLockRequest { args }); let request = Request::new(GenerallyLockRequest { args });
let response = client.r_un_lock(request).await?.into_inner(); let response = client.r_un_lock(request).await.map_err(Error::other)?.into_inner();
if let Some(error_info) = response.error_info { if let Some(error_info) = response.error_info {
return Err(Error::from_string(error_info)); return Err(Error::other(error_info));
} }
Ok(response.success) Ok(response.success)
@@ -93,13 +93,13 @@ impl Locker for RemoteClient {
let args = serde_json::to_string(args)?; let args = serde_json::to_string(args)?;
let mut client = node_service_time_out_client(&self.addr) let mut client = node_service_time_out_client(&self.addr)
.await .await
.map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; .map_err(|err| Error::other(format!("can not get client, err: {}", err)))?;
let request = Request::new(GenerallyLockRequest { args }); let request = Request::new(GenerallyLockRequest { args });
let response = client.refresh(request).await?.into_inner(); let response = client.refresh(request).await.map_err(Error::other)?.into_inner();
if let Some(error_info) = response.error_info { if let Some(error_info) = response.error_info {
return Err(Error::from_string(error_info)); return Err(Error::other(error_info));
} }
Ok(response.success) Ok(response.success)
@@ -110,13 +110,13 @@ impl Locker for RemoteClient {
let args = serde_json::to_string(args)?; let args = serde_json::to_string(args)?;
let mut client = node_service_time_out_client(&self.addr) let mut client = node_service_time_out_client(&self.addr)
.await .await
.map_err(|err| Error::from_string(format!("can not get client, err: {}", err)))?; .map_err(|err| Error::other(format!("can not get client, err: {}", err)))?;
let request = Request::new(GenerallyLockRequest { args }); let request = Request::new(GenerallyLockRequest { args });
let response = client.force_un_lock(request).await?.into_inner(); let response = client.force_un_lock(request).await.map_err(Error::other)?.into_inner();
if let Some(error_info) = response.error_info { if let Some(error_info) = response.error_info {
return Err(Error::from_string(error_info)); return Err(Error::other(error_info));
} }
Ok(response.success) Ok(response.success)
+7 -6
View File
@@ -1,8 +1,9 @@
use crate::error::{Error, Result};
use crate::{ use crate::{
disk::endpoint::Endpoint, disk::endpoint::Endpoint,
global::{GLOBAL_Endpoints, GLOBAL_BOOT_TIME}, global::{GLOBAL_BOOT_TIME, GLOBAL_Endpoints},
heal::{ heal::{
data_usage::{load_data_usage_from_backend, DATA_USAGE_CACHE_NAME, DATA_USAGE_ROOT}, data_usage::{DATA_USAGE_CACHE_NAME, DATA_USAGE_ROOT, load_data_usage_from_backend},
data_usage_cache::DataUsageCache, data_usage_cache::DataUsageCache,
heal_commands::{DRIVE_STATE_OK, DRIVE_STATE_UNFORMATTED}, heal_commands::{DRIVE_STATE_OK, DRIVE_STATE_UNFORMATTED},
}, },
@@ -11,10 +12,10 @@ use crate::{
store_api::StorageAPI, store_api::StorageAPI,
}; };
use common::{ use common::{
error::{Error, Result}, // error::{Error, Result},
globals::GLOBAL_Local_Node_Name, globals::GLOBAL_Local_Node_Name,
}; };
use madmin::{BackendDisks, Disk, ErasureSetInfo, InfoMessage, ServerProperties, ITEM_INITIALIZING, ITEM_OFFLINE, ITEM_ONLINE}; use madmin::{BackendDisks, Disk, ErasureSetInfo, ITEM_INITIALIZING, ITEM_OFFLINE, ITEM_ONLINE, InfoMessage, ServerProperties};
use protos::{ use protos::{
models::{PingBody, PingBodyBuilder}, models::{PingBody, PingBodyBuilder},
node_service_time_out_client, node_service_time_out_client,
@@ -87,7 +88,7 @@ async fn is_server_resolvable(endpoint: &Endpoint) -> Result<()> {
// 创建客户端 // 创建客户端
let mut client = node_service_time_out_client(&addr) let mut client = node_service_time_out_client(&addr)
.await .await
.map_err(|err| Error::msg(err.to_string()))?; .map_err(|err| Error::other(err.to_string()))?;
// 构造 PingRequest // 构造 PingRequest
let request = Request::new(PingRequest { let request = Request::new(PingRequest {
@@ -332,7 +333,7 @@ fn get_online_offline_disks_stats(disks_info: &[Disk]) -> (BackendDisks, Backend
async fn get_pools_info(all_disks: &[Disk]) -> Result<HashMap<i32, HashMap<i32, ErasureSetInfo>>> { async fn get_pools_info(all_disks: &[Disk]) -> Result<HashMap<i32, HashMap<i32, ErasureSetInfo>>> {
let Some(store) = new_object_layer_fn() else { let Some(store) = new_object_layer_fn() else {
return Err(Error::msg("ServerNotInitialized")); return Err(Error::other("ServerNotInitialized"));
}; };
let mut pools_info: HashMap<i32, HashMap<i32, ErasureSetInfo>> = HashMap::new(); let mut pools_info: HashMap<i32, HashMap<i32, ErasureSetInfo>> = HashMap::new();
+2 -2
View File
@@ -6,15 +6,15 @@ use std::{
pin::Pin, pin::Pin,
ptr, ptr,
sync::{ sync::{
atomic::{AtomicPtr, AtomicU64, Ordering},
Arc, Arc,
atomic::{AtomicPtr, AtomicU64, Ordering},
}, },
time::{Duration, SystemTime, UNIX_EPOCH}, time::{Duration, SystemTime, UNIX_EPOCH},
}; };
use tokio::{spawn, sync::Mutex}; use tokio::{spawn, sync::Mutex};
use common::error::Result; use std::io::Result;
pub type UpdateFn<T> = Box<dyn Fn() -> Pin<Box<dyn Future<Output = Result<T>> + Send>> + Send + Sync + 'static>; pub type UpdateFn<T> = Box<dyn Fn() -> Pin<Box<dyn Future<Output = Result<T>> + Send>> + Send + Sync + 'static>;
+1 -1
View File
@@ -1,2 +1,2 @@
pub mod cache; // pub mod cache;
pub mod metacache_set; pub mod metacache_set;
+5 -4
View File
@@ -19,7 +19,7 @@ use bytes::Bytes;
use chrono::Datelike; use chrono::Datelike;
use clap::Parser; use clap::Parser;
use common::{ use common::{
error::{Error, Result}, // error::{Error, Result},
globals::set_global_addr, globals::set_global_addr,
}; };
use ecstore::StorageAPI; use ecstore::StorageAPI;
@@ -54,6 +54,7 @@ use rustls::ServerConfig;
use s3s::{host::MultiDomain, service::S3ServiceBuilder}; use s3s::{host::MultiDomain, service::S3ServiceBuilder};
use service::hybrid; use service::hybrid;
use socket2::SockRef; use socket2::SockRef;
use std::io::{Error, Result};
use std::net::SocketAddr; use std::net::SocketAddr;
use std::sync::Arc; use std::sync::Arc;
use std::time::Duration; use std::time::Duration;
@@ -107,7 +108,7 @@ async fn main() -> Result<()> {
let (_logger, guard) = init_obs(Some(opt.clone().obs_endpoint)).await; let (_logger, guard) = init_obs(Some(opt.clone().obs_endpoint)).await;
// Store in global storage // Store in global storage
set_global_guard(guard)?; set_global_guard(guard).map_err(Error::other)?;
// Run parameters // Run parameters
run(opt).await run(opt).await
@@ -211,7 +212,7 @@ async fn run(opt: config::Opt) -> Result<()> {
if !opt.server_domains.is_empty() { if !opt.server_domains.is_empty() {
info!("virtual-hosted-style requests are enabled use domain_name {:?}", &opt.server_domains); info!("virtual-hosted-style requests are enabled use domain_name {:?}", &opt.server_domains);
b.set_host(MultiDomain::new(&opt.server_domains)?); b.set_host(MultiDomain::new(&opt.server_domains).map_err(Error::other)?);
} }
// // Enable parsing virtual-hosted-style requests // // Enable parsing virtual-hosted-style requests
@@ -506,7 +507,7 @@ async fn run(opt: config::Opt) -> Result<()> {
.await .await
.map_err(|err| { .map_err(|err| {
error!("ECStore::new {:?}", &err); error!("ECStore::new {:?}", &err);
Error::from_string(err.to_string()) err
})?; })?;
ecconfig::init(); ecconfig::init();