mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-07 23:25:51 +00:00
Compare commits
8
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
597d383ca4 | ||
|
|
b5cdd71600 | ||
|
|
2d4ea8c665 | ||
|
|
b3e50bb12f | ||
|
|
2a6f27eb08 | ||
|
|
08f48e62c9 | ||
|
|
e29b685c20 | ||
|
|
4287b7b12a |
@@ -39,8 +39,16 @@ jobs:
|
||||
- name: Install cross-compilation tools
|
||||
if: matrix.cross
|
||||
run: |
|
||||
sudo apt-get install -y gcc-aarch64-linux-gnu
|
||||
sudo dpkg --add-architecture arm64
|
||||
sudo sed -i 's/^deb /deb [arch=amd64] /' /etc/apt/sources.list
|
||||
echo "deb [arch=arm64] http://ports.ubuntu.com/ jammy main restricted universe multiverse" | sudo tee /etc/apt/sources.list.d/arm64.list
|
||||
echo "deb [arch=arm64] http://ports.ubuntu.com/ jammy-updates main restricted universe multiverse" | sudo tee -a /etc/apt/sources.list.d/arm64.list
|
||||
sudo apt-get update
|
||||
sudo apt-get install -y gcc-aarch64-linux-gnu libssl-dev:arm64
|
||||
echo "CARGO_TARGET_AARCH64_UNKNOWN_LINUX_GNU_LINKER=aarch64-linux-gnu-gcc" >> "$GITHUB_ENV"
|
||||
echo "OPENSSL_DIR=/usr" >> "$GITHUB_ENV"
|
||||
echo "OPENSSL_INCLUDE_DIR=/usr/include" >> "$GITHUB_ENV"
|
||||
echo "OPENSSL_LIB_DIR=/usr/lib/aarch64-linux-gnu" >> "$GITHUB_ENV"
|
||||
|
||||
- name: Cache cargo registry and target
|
||||
uses: actions/cache@v5
|
||||
@@ -80,6 +88,7 @@ jobs:
|
||||
rm weed-volume-normal
|
||||
|
||||
- name: Upload release assets
|
||||
if: startsWith(github.ref, 'refs/tags/')
|
||||
uses: softprops/action-gh-release@v2
|
||||
with:
|
||||
files: |
|
||||
@@ -88,6 +97,15 @@ jobs:
|
||||
env:
|
||||
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
|
||||
|
||||
- name: Upload artifacts
|
||||
if: ${{ !startsWith(github.ref, 'refs/tags/') }}
|
||||
uses: actions/upload-artifact@v4
|
||||
with:
|
||||
name: rust-volume-${{ matrix.asset_suffix }}
|
||||
path: |
|
||||
weed-volume_large_disk_${{ matrix.asset_suffix }}.tar.gz
|
||||
weed-volume_${{ matrix.asset_suffix }}.tar.gz
|
||||
|
||||
build-rust-volume-darwin:
|
||||
permissions:
|
||||
contents: write
|
||||
@@ -147,6 +165,7 @@ jobs:
|
||||
rm weed-volume-normal
|
||||
|
||||
- name: Upload release assets
|
||||
if: startsWith(github.ref, 'refs/tags/')
|
||||
uses: softprops/action-gh-release@v2
|
||||
with:
|
||||
files: |
|
||||
@@ -155,6 +174,15 @@ jobs:
|
||||
env:
|
||||
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
|
||||
|
||||
- name: Upload artifacts
|
||||
if: ${{ !startsWith(github.ref, 'refs/tags/') }}
|
||||
uses: actions/upload-artifact@v4
|
||||
with:
|
||||
name: rust-volume-${{ matrix.asset_suffix }}
|
||||
path: |
|
||||
weed-volume_large_disk_${{ matrix.asset_suffix }}.tar.gz
|
||||
weed-volume_${{ matrix.asset_suffix }}.tar.gz
|
||||
|
||||
build-rust-volume-windows:
|
||||
permissions:
|
||||
contents: write
|
||||
@@ -206,6 +234,7 @@ jobs:
|
||||
rm weed-volume-normal.exe
|
||||
|
||||
- name: Upload release assets
|
||||
if: startsWith(github.ref, 'refs/tags/')
|
||||
uses: softprops/action-gh-release@v2
|
||||
with:
|
||||
files: |
|
||||
@@ -213,3 +242,12 @@ jobs:
|
||||
weed-volume_windows_amd64.zip
|
||||
env:
|
||||
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
|
||||
|
||||
- name: Upload artifacts
|
||||
if: ${{ !startsWith(github.ref, 'refs/tags/') }}
|
||||
uses: actions/upload-artifact@v4
|
||||
with:
|
||||
name: rust-volume-windows_amd64
|
||||
path: |
|
||||
weed-volume_large_disk_windows_amd64.zip
|
||||
weed-volume_windows_amd64.zip
|
||||
|
||||
@@ -65,8 +65,6 @@ reed-solomon-erasure = "6"
|
||||
# Logging
|
||||
tracing = "0.1"
|
||||
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
|
||||
pprof = { version = "0.15", features = ["prost-codec"] }
|
||||
|
||||
# Config
|
||||
toml = "0.8"
|
||||
serde = { version = "1", features = ["derive"] }
|
||||
@@ -127,6 +125,10 @@ aws-sdk-s3 = { version = "1.125.0", default-features = false, features = ["sigv4
|
||||
aws-credential-types = "1"
|
||||
aws-types = "1"
|
||||
|
||||
# pprof is Unix-only (requires libc/nix APIs not available on Windows)
|
||||
[target.'cfg(unix)'.dependencies]
|
||||
pprof = { version = "0.15", features = ["prost-codec"] }
|
||||
|
||||
[dev-dependencies]
|
||||
tempfile = "3"
|
||||
|
||||
|
||||
@@ -10,9 +10,11 @@ use seaweed_volume::security::tls::{
|
||||
GrpcClientAuthPolicy, TlsPolicy,
|
||||
};
|
||||
use seaweed_volume::security::{Guard, SigningKey};
|
||||
#[cfg(unix)]
|
||||
use seaweed_volume::server::debug::build_debug_router;
|
||||
use seaweed_volume::server::grpc_client::load_outgoing_grpc_tls;
|
||||
use seaweed_volume::server::grpc_server::VolumeGrpcService;
|
||||
#[cfg(unix)]
|
||||
use seaweed_volume::server::profiling::CpuProfileSession;
|
||||
use seaweed_volume::server::request_id::GrpcRequestIdLayer;
|
||||
use seaweed_volume::server::volume_server::{
|
||||
@@ -24,6 +26,11 @@ use seaweed_volume::storage::types::DiskType;
|
||||
|
||||
use tokio_rustls::TlsAcceptor;
|
||||
|
||||
#[cfg(unix)]
|
||||
type CpuProfileParam = Option<CpuProfileSession>;
|
||||
#[cfg(not(unix))]
|
||||
type CpuProfileParam = Option<()>;
|
||||
|
||||
const GRPC_MAX_MESSAGE_SIZE: usize = 1 << 30;
|
||||
const GRPC_KEEPALIVE_INTERVAL: std::time::Duration = std::time::Duration::from_secs(60);
|
||||
const GRPC_KEEPALIVE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(20);
|
||||
@@ -42,6 +49,7 @@ fn main() {
|
||||
|
||||
let config = config::parse_cli();
|
||||
seaweed_volume::server::server_stats::init_process_start();
|
||||
#[cfg(unix)]
|
||||
let cpu_profile = match CpuProfileSession::start(&config) {
|
||||
Ok(session) => session,
|
||||
Err(e) => {
|
||||
@@ -49,6 +57,8 @@ fn main() {
|
||||
std::process::exit(1);
|
||||
}
|
||||
};
|
||||
#[cfg(not(unix))]
|
||||
let cpu_profile: Option<()> = None;
|
||||
info!(
|
||||
"SeaweedFS Volume Server (Rust) v{}",
|
||||
seaweed_volume::version::full_version()
|
||||
@@ -257,7 +267,7 @@ where
|
||||
|
||||
async fn run(
|
||||
config: VolumeServerConfig,
|
||||
cpu_profile: Option<CpuProfileSession>,
|
||||
#[allow(unused_variables)] cpu_profile: CpuProfileParam,
|
||||
) -> Result<(), Box<dyn std::error::Error>> {
|
||||
// Initialize the store
|
||||
let mut store = Store::new(config.index_type);
|
||||
@@ -431,10 +441,12 @@ async fn run(
|
||||
}
|
||||
|
||||
// Build HTTP routers
|
||||
#[allow(unused_mut)]
|
||||
let mut admin_router = seaweed_volume::server::volume_server::build_admin_router_with_ui(
|
||||
state.clone(),
|
||||
config.ui_enabled,
|
||||
);
|
||||
#[cfg(unix)]
|
||||
if config.pprof {
|
||||
admin_router = admin_router.merge(build_debug_router());
|
||||
}
|
||||
@@ -721,6 +733,7 @@ async fn run(
|
||||
None
|
||||
};
|
||||
|
||||
#[cfg(unix)]
|
||||
let debug_handle = if config.debug {
|
||||
let debug_addr = format!("0.0.0.0:{}", config.debug_port);
|
||||
info!("Debug pprof server listening on {}", debug_addr);
|
||||
@@ -742,6 +755,8 @@ async fn run(
|
||||
} else {
|
||||
None
|
||||
};
|
||||
#[cfg(not(unix))]
|
||||
let debug_handle: Option<tokio::task::JoinHandle<()>> = None;
|
||||
|
||||
let metrics_push_handle = {
|
||||
let push_state = state.clone();
|
||||
@@ -774,6 +789,7 @@ async fn run(
|
||||
// Close all volumes (flush and release file handles) matching Go's Shutdown()
|
||||
state.store.write().unwrap().close();
|
||||
|
||||
#[cfg(unix)]
|
||||
if let Some(cpu_profile) = cpu_profile {
|
||||
cpu_profile.finish().map_err(std::io::Error::other)?;
|
||||
}
|
||||
|
||||
@@ -1,9 +1,11 @@
|
||||
#[cfg(unix)]
|
||||
pub mod debug;
|
||||
pub mod grpc_client;
|
||||
pub mod grpc_server;
|
||||
pub mod handlers;
|
||||
pub mod heartbeat;
|
||||
pub mod memory_status;
|
||||
#[cfg(unix)]
|
||||
pub mod profiling;
|
||||
pub mod request_id;
|
||||
pub mod server_stats;
|
||||
|
||||
@@ -6,7 +6,7 @@
|
||||
use std::fs::File;
|
||||
use std::io;
|
||||
#[cfg(not(unix))]
|
||||
use std::io::{Seek, SeekFrom};
|
||||
use std::io::{Read, Seek, SeekFrom};
|
||||
|
||||
use reed_solomon_erasure::galois_8::ReedSolomon;
|
||||
|
||||
|
||||
@@ -359,6 +359,7 @@ func doSubscribeFilerMetaChanges(clientId int32, clientEpoch int32, sourceGrpcDi
|
||||
}
|
||||
|
||||
var lastLogTsNs = time.Now().UnixNano()
|
||||
var lastProgressedTsNs int64
|
||||
var clientName = fmt.Sprintf("syncFrom_%s_To_%s", string(sourceFiler), string(targetFiler))
|
||||
processEventFnWithOffset := pb.AddOffsetFunc(func(resp *filer_pb.SubscribeMetadataResponse) error {
|
||||
processor.AddSyncJob(resp)
|
||||
@@ -372,6 +373,18 @@ func doSubscribeFilerMetaChanges(clientId int32, clientEpoch int32, sourceGrpcDi
|
||||
now := time.Now().UnixNano()
|
||||
glog.V(0).Infof("sync %s to %s progressed to %v %0.2f/sec", sourceFiler, targetFiler, time.Unix(0, offsetTsNs), float64(counter)/(float64(now-lastLogTsNs)/1e9))
|
||||
lastLogTsNs = now
|
||||
if offsetTsNs == lastProgressedTsNs {
|
||||
for _, t := range filerSink.ActiveTransfers() {
|
||||
if t.LastErr != "" {
|
||||
glog.V(0).Infof(" %s %s: %d bytes received, %s, last error: %s",
|
||||
t.ChunkFileId, t.Path, t.BytesReceived, t.Status, t.LastErr)
|
||||
} else {
|
||||
glog.V(0).Infof(" %s %s: %d bytes received, %s",
|
||||
t.ChunkFileId, t.Path, t.BytesReceived, t.Status)
|
||||
}
|
||||
}
|
||||
}
|
||||
lastProgressedTsNs = offsetTsNs
|
||||
// collect synchronous offset
|
||||
statsCollect.FilerSyncOffsetGauge.WithLabelValues(sourceFiler.String(), targetFiler.String(), clientName, sourcePath).Set(float64(offsetTsNs))
|
||||
return setOffset(targetGrpcDialOption, targetFiler, getSignaturePrefixByPath(sourcePath), sourceFilerSignature, offsetTsNs)
|
||||
|
||||
@@ -241,6 +241,14 @@ func (fs *FilerSink) fetchAndWrite(sourceChunk *filer_pb.FileChunk, path string,
|
||||
return "", fmt.Errorf("upload data: %w", err)
|
||||
}
|
||||
|
||||
transferStatus := &ChunkTransferStatus{
|
||||
ChunkFileId: sourceChunk.GetFileIdString(),
|
||||
Path: path,
|
||||
Status: "downloading",
|
||||
}
|
||||
fs.activeTransfers.Store(sourceChunk.GetFileIdString(), transferStatus)
|
||||
defer fs.activeTransfers.Delete(sourceChunk.GetFileIdString())
|
||||
|
||||
eofBackoff := time.Duration(0)
|
||||
var partialData []byte
|
||||
var savedFilename string
|
||||
@@ -282,6 +290,11 @@ func (fs *FilerSink) fetchAndWrite(sourceChunk *filer_pb.FileChunk, path string,
|
||||
fullData = data
|
||||
}
|
||||
|
||||
transferStatus.mu.Lock()
|
||||
transferStatus.BytesReceived = int64(len(fullData))
|
||||
transferStatus.Status = "uploading"
|
||||
transferStatus.mu.Unlock()
|
||||
|
||||
currentFileId, uploadResult, uploadErr, _ := uploader.UploadWithRetry(
|
||||
fs,
|
||||
&filer_pb.AssignVolumeRequest{
|
||||
@@ -324,11 +337,21 @@ func (fs *FilerSink) fetchAndWrite(sourceChunk *filer_pb.FileChunk, path string,
|
||||
glog.V(1).Infof("skip retrying stale source %s for %s: %v", sourceChunk.GetFileIdString(), path, retryErr)
|
||||
return false
|
||||
}
|
||||
transferStatus.mu.Lock()
|
||||
transferStatus.LastErr = retryErr.Error()
|
||||
transferStatus.mu.Unlock()
|
||||
if isEofError(retryErr) {
|
||||
eofBackoff = nextEofBackoff(eofBackoff)
|
||||
transferStatus.mu.Lock()
|
||||
transferStatus.BytesReceived = int64(len(partialData))
|
||||
transferStatus.Status = fmt.Sprintf("waiting %v", eofBackoff)
|
||||
transferStatus.mu.Unlock()
|
||||
glog.V(0).Infof("source connection interrupted while replicating %s for %s (%d bytes received so far), backing off %v: %v",
|
||||
sourceChunk.GetFileIdString(), path, len(partialData), eofBackoff, retryErr)
|
||||
time.Sleep(eofBackoff)
|
||||
transferStatus.mu.Lock()
|
||||
transferStatus.Status = "downloading"
|
||||
transferStatus.mu.Unlock()
|
||||
} else {
|
||||
glog.V(0).Infof("replicate %s for %s: %v", sourceChunk.GetFileIdString(), path, retryErr)
|
||||
}
|
||||
|
||||
@@ -4,6 +4,7 @@ import (
|
||||
"context"
|
||||
"fmt"
|
||||
"math"
|
||||
"sync"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/wdclient"
|
||||
@@ -20,6 +21,19 @@ import (
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
)
|
||||
|
||||
// ChunkTransferStatus tracks the progress of a single chunk being replicated.
|
||||
// Fields are guarded by mu: ChunkFileId and Path are immutable after creation,
|
||||
// while BytesReceived, Status, and LastErr are updated by fetchAndWrite and
|
||||
// read by ActiveTransfers.
|
||||
type ChunkTransferStatus struct {
|
||||
mu sync.RWMutex
|
||||
ChunkFileId string
|
||||
Path string
|
||||
BytesReceived int64
|
||||
Status string // "downloading", "uploading", or "waiting 10s" etc.
|
||||
LastErr string
|
||||
}
|
||||
|
||||
type FilerSink struct {
|
||||
filerSource *source.FilerSource
|
||||
grpcAddress string
|
||||
@@ -35,6 +49,7 @@ type FilerSink struct {
|
||||
isIncremental bool
|
||||
executor *util.LimitedConcurrentExecutor
|
||||
signature int32
|
||||
activeTransfers sync.Map // chunkFileId -> *ChunkTransferStatus
|
||||
}
|
||||
|
||||
func init() {
|
||||
@@ -101,6 +116,25 @@ func (fs *FilerSink) SetChunkConcurrency(concurrency int) {
|
||||
}
|
||||
}
|
||||
|
||||
// ActiveTransfers returns an immutable snapshot of all in-progress chunk transfers.
|
||||
func (fs *FilerSink) ActiveTransfers() []ChunkTransferStatus {
|
||||
var transfers []ChunkTransferStatus
|
||||
fs.activeTransfers.Range(func(key, value any) bool {
|
||||
t := value.(*ChunkTransferStatus)
|
||||
t.mu.RLock()
|
||||
transfers = append(transfers, ChunkTransferStatus{
|
||||
ChunkFileId: t.ChunkFileId,
|
||||
Path: t.Path,
|
||||
BytesReceived: t.BytesReceived,
|
||||
Status: t.Status,
|
||||
LastErr: t.LastErr,
|
||||
})
|
||||
t.mu.RUnlock()
|
||||
return true
|
||||
})
|
||||
return transfers
|
||||
}
|
||||
|
||||
func (fs *FilerSink) DeleteEntry(key string, isDirectory, deleteIncludeChunks bool, signatures []int32) error {
|
||||
|
||||
dir, name := util.FullPath(key).DirAndName()
|
||||
|
||||
@@ -58,9 +58,9 @@ var (
|
||||
|
||||
// SSECustomerKey represents a customer-provided encryption key for SSE-C
|
||||
type SSECustomerKey struct {
|
||||
Algorithm string
|
||||
Key []byte
|
||||
KeyMD5 string
|
||||
Algorithm string
|
||||
Key []byte
|
||||
KeyMD5 string
|
||||
}
|
||||
|
||||
// IsSSECRequest checks if the request contains SSE-C headers
|
||||
@@ -119,8 +119,8 @@ func validateAndParseSSECHeaders(algorithm, key, keyMD5 string) (*SSECustomerKey
|
||||
sum := md5.Sum(keyBytes)
|
||||
expectedMD5 := base64.StdEncoding.EncodeToString(sum[:])
|
||||
|
||||
// Debug logging for MD5 validation
|
||||
glog.V(4).Infof("SSE-C MD5 validation: provided='%s', expected='%s', keyBytes=%x", keyMD5, expectedMD5, keyBytes)
|
||||
// Debug logging for MD5 validation (never log key material)
|
||||
glog.V(4).Infof("SSE-C MD5 validation: provided='%s', expected='%s'", keyMD5, expectedMD5)
|
||||
|
||||
if keyMD5 != expectedMD5 {
|
||||
glog.Errorf("SSE-C MD5 mismatch: provided='%s', expected='%s'", keyMD5, expectedMD5)
|
||||
|
||||
Reference in New Issue
Block a user