mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-08 15:45:50 +00:00
Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
77b4395261 | ||
|
|
0ffb50887a |
@@ -39,16 +39,8 @@ jobs:
|
||||
- name: Install cross-compilation tools
|
||||
if: matrix.cross
|
||||
run: |
|
||||
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
|
||||
sudo apt-get install -y gcc-aarch64-linux-gnu
|
||||
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
|
||||
@@ -88,7 +80,6 @@ jobs:
|
||||
rm weed-volume-normal
|
||||
|
||||
- name: Upload release assets
|
||||
if: startsWith(github.ref, 'refs/tags/')
|
||||
uses: softprops/action-gh-release@v2
|
||||
with:
|
||||
files: |
|
||||
@@ -97,15 +88,6 @@ 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
|
||||
@@ -165,7 +147,6 @@ jobs:
|
||||
rm weed-volume-normal
|
||||
|
||||
- name: Upload release assets
|
||||
if: startsWith(github.ref, 'refs/tags/')
|
||||
uses: softprops/action-gh-release@v2
|
||||
with:
|
||||
files: |
|
||||
@@ -174,15 +155,6 @@ 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
|
||||
@@ -234,7 +206,6 @@ jobs:
|
||||
rm weed-volume-normal.exe
|
||||
|
||||
- name: Upload release assets
|
||||
if: startsWith(github.ref, 'refs/tags/')
|
||||
uses: softprops/action-gh-release@v2
|
||||
with:
|
||||
files: |
|
||||
@@ -242,12 +213,3 @@ 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
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
apiVersion: v1
|
||||
description: SeaweedFS
|
||||
name: seaweedfs
|
||||
appVersion: "4.18"
|
||||
appVersion: "4.17"
|
||||
# Dev note: Trigger a helm chart release by `git tag -a helm-<version>`
|
||||
version: 4.18.0
|
||||
version: 4.17.0
|
||||
|
||||
@@ -65,6 +65,8 @@ 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"] }
|
||||
@@ -125,10 +127,6 @@ 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,11 +10,9 @@ 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::{
|
||||
@@ -26,11 +24,6 @@ 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);
|
||||
@@ -49,7 +42,6 @@ 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) => {
|
||||
@@ -57,8 +49,6 @@ 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()
|
||||
@@ -267,7 +257,7 @@ where
|
||||
|
||||
async fn run(
|
||||
config: VolumeServerConfig,
|
||||
#[allow(unused_variables)] cpu_profile: CpuProfileParam,
|
||||
cpu_profile: Option<CpuProfileSession>,
|
||||
) -> Result<(), Box<dyn std::error::Error>> {
|
||||
// Initialize the store
|
||||
let mut store = Store::new(config.index_type);
|
||||
@@ -441,12 +431,10 @@ 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());
|
||||
}
|
||||
@@ -733,7 +721,6 @@ 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);
|
||||
@@ -755,8 +742,6 @@ 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();
|
||||
@@ -789,7 +774,6 @@ 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,11 +1,9 @@
|
||||
#[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::{Read, Seek, SeekFrom};
|
||||
use std::io::{Seek, SeekFrom};
|
||||
|
||||
use reed_solomon_erasure::galois_8::ReedSolomon;
|
||||
|
||||
|
||||
@@ -1,294 +0,0 @@
|
||||
package volume_server_grpc_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net/http"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/test/volume_server/framework"
|
||||
"github.com/seaweedfs/seaweedfs/test/volume_server/matrix"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/volume_server_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/needle"
|
||||
)
|
||||
|
||||
// TestEcDecodePreservesDeletedNeedles verifies that needles deleted via
|
||||
// VolumeEcBlobDelete (recorded in .ecj) are correctly excluded from the
|
||||
// decoded volume produced by VolumeEcShardsToVolume.
|
||||
func TestEcDecodePreservesDeletedNeedles(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("skipping integration test in short mode")
|
||||
}
|
||||
|
||||
cluster := framework.StartVolumeCluster(t, matrix.P1())
|
||||
conn, client := framework.DialVolumeServer(t, cluster.VolumeGRPCAddress())
|
||||
defer conn.Close()
|
||||
|
||||
const (
|
||||
volumeID = uint32(140)
|
||||
keyA = uint64(990020)
|
||||
cookieA = uint32(0xDA001122)
|
||||
keyB = uint64(990021)
|
||||
cookieB = uint32(0xDA003344)
|
||||
)
|
||||
|
||||
framework.AllocateVolume(t, client, volumeID, "")
|
||||
|
||||
httpClient := framework.NewHTTPClient()
|
||||
fidA := framework.NewFileID(volumeID, keyA, cookieA)
|
||||
fidB := framework.NewFileID(volumeID, keyB, cookieB)
|
||||
payloadA := []byte("needle-A-should-be-deleted-after-decode")
|
||||
payloadB := []byte("needle-B-should-survive-decode")
|
||||
|
||||
// Upload two needles.
|
||||
for _, tc := range []struct {
|
||||
fid string
|
||||
payload []byte
|
||||
}{
|
||||
{fidA, payloadA},
|
||||
{fidB, payloadB},
|
||||
} {
|
||||
resp := framework.UploadBytes(t, httpClient, cluster.VolumeAdminURL(), tc.fid, tc.payload)
|
||||
_ = framework.ReadAllAndClose(t, resp)
|
||||
if resp.StatusCode != http.StatusCreated {
|
||||
t.Fatalf("upload %s: expected 201, got %d", tc.fid, resp.StatusCode)
|
||||
}
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second)
|
||||
defer cancel()
|
||||
|
||||
// EC encode.
|
||||
_, err := client.VolumeEcShardsGenerate(ctx, &volume_server_pb.VolumeEcShardsGenerateRequest{
|
||||
VolumeId: volumeID, Collection: "",
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("VolumeEcShardsGenerate: %v", err)
|
||||
}
|
||||
|
||||
// Mount all data shards so the EC volume is usable.
|
||||
_, err = client.VolumeEcShardsMount(ctx, &volume_server_pb.VolumeEcShardsMountRequest{
|
||||
VolumeId: volumeID, Collection: "",
|
||||
ShardIds: []uint32{0, 1, 2, 3, 4, 5, 6, 7, 8, 9},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("VolumeEcShardsMount: %v", err)
|
||||
}
|
||||
|
||||
// Delete needle A via EC path (writes to .ecj).
|
||||
_, err = client.VolumeEcBlobDelete(ctx, &volume_server_pb.VolumeEcBlobDeleteRequest{
|
||||
VolumeId: volumeID, Collection: "",
|
||||
FileKey: keyA, Version: uint32(needle.GetCurrentVersion()),
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("VolumeEcBlobDelete needle A: %v", err)
|
||||
}
|
||||
|
||||
// Unmount the normal volume so decode writes fresh files.
|
||||
_, err = client.VolumeUnmount(ctx, &volume_server_pb.VolumeUnmountRequest{VolumeId: volumeID})
|
||||
if err != nil {
|
||||
t.Fatalf("VolumeUnmount: %v", err)
|
||||
}
|
||||
|
||||
// Decode EC shards back to a normal volume.
|
||||
_, err = client.VolumeEcShardsToVolume(ctx, &volume_server_pb.VolumeEcShardsToVolumeRequest{
|
||||
VolumeId: volumeID, Collection: "",
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("VolumeEcShardsToVolume: %v", err)
|
||||
}
|
||||
|
||||
// Re-mount the decoded volume.
|
||||
_, err = client.VolumeMount(ctx, &volume_server_pb.VolumeMountRequest{VolumeId: volumeID})
|
||||
if err != nil {
|
||||
t.Fatalf("VolumeMount: %v", err)
|
||||
}
|
||||
|
||||
// Needle A should be gone (deleted via .ecj before decode).
|
||||
respA := framework.ReadBytes(t, httpClient, cluster.VolumeAdminURL(), fidA)
|
||||
bodyA := framework.ReadAllAndClose(t, respA)
|
||||
if respA.StatusCode >= 500 {
|
||||
t.Fatalf("needle A read: server error %d: %s", respA.StatusCode, bodyA)
|
||||
}
|
||||
if respA.StatusCode != http.StatusNotFound {
|
||||
t.Fatalf("needle A should be 404 after decode, got %d", respA.StatusCode)
|
||||
}
|
||||
|
||||
// Needle B should still be readable.
|
||||
respB := framework.ReadBytes(t, httpClient, cluster.VolumeAdminURL(), fidB)
|
||||
bodyB := framework.ReadAllAndClose(t, respB)
|
||||
if respB.StatusCode != http.StatusOK {
|
||||
t.Fatalf("needle B read: expected 200, got %d", respB.StatusCode)
|
||||
}
|
||||
if string(bodyB) != string(payloadB) {
|
||||
t.Fatalf("needle B payload mismatch: got %q, want %q", bodyB, payloadB)
|
||||
}
|
||||
}
|
||||
|
||||
// TestEcDecodeCollectsEcjFromPeer verifies that .ecj deletion entries from a
|
||||
// peer server that contributes no new data shards are still collected during
|
||||
// decode. This is the regression test for the fix in collectEcShards that
|
||||
// always copies .ecj from every shard location.
|
||||
//
|
||||
// Scenario:
|
||||
// - Server 0 holds all 10 data shards (decode target).
|
||||
// - Server 1 holds a copy of shard 0 (no new shards for server 0).
|
||||
// - A needle is deleted ONLY on server 1 (server 0's .ecj is empty).
|
||||
// - During decode on server 0, server 1's .ecj must be collected and applied.
|
||||
func TestEcDecodeCollectsEcjFromPeer(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("skipping integration test in short mode")
|
||||
}
|
||||
|
||||
cluster := framework.StartMultiVolumeClusterAuto(t, matrix.P1(), 2)
|
||||
conn0, client0 := framework.DialVolumeServer(t, cluster.VolumeGRPCAddress(0))
|
||||
defer conn0.Close()
|
||||
conn1, client1 := framework.DialVolumeServer(t, cluster.VolumeGRPCAddress(1))
|
||||
defer conn1.Close()
|
||||
|
||||
const (
|
||||
volumeID = uint32(141)
|
||||
keyA = uint64(990030)
|
||||
cookieA = uint32(0xDB001122)
|
||||
keyB = uint64(990031)
|
||||
cookieB = uint32(0xDB003344)
|
||||
)
|
||||
|
||||
// Allocate and upload on server 0.
|
||||
framework.AllocateVolume(t, client0, volumeID, "")
|
||||
|
||||
httpClient := framework.NewHTTPClient()
|
||||
fidA := framework.NewFileID(volumeID, keyA, cookieA)
|
||||
fidB := framework.NewFileID(volumeID, keyB, cookieB)
|
||||
payloadB := []byte("needle-B-should-survive-peer-ecj-decode")
|
||||
|
||||
resp := framework.UploadBytes(t, httpClient, cluster.VolumeAdminURL(0), fidA, []byte("needle-A-deleted-on-peer"))
|
||||
_ = framework.ReadAllAndClose(t, resp)
|
||||
if resp.StatusCode != http.StatusCreated {
|
||||
t.Fatalf("upload A: expected 201, got %d", resp.StatusCode)
|
||||
}
|
||||
resp = framework.UploadBytes(t, httpClient, cluster.VolumeAdminURL(0), fidB, payloadB)
|
||||
_ = framework.ReadAllAndClose(t, resp)
|
||||
if resp.StatusCode != http.StatusCreated {
|
||||
t.Fatalf("upload B: expected 201, got %d", resp.StatusCode)
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second)
|
||||
defer cancel()
|
||||
|
||||
// EC encode on server 0.
|
||||
_, err := client0.VolumeEcShardsGenerate(ctx, &volume_server_pb.VolumeEcShardsGenerateRequest{
|
||||
VolumeId: volumeID, Collection: "",
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("VolumeEcShardsGenerate on server 0: %v", err)
|
||||
}
|
||||
|
||||
// Build the SourceDataNode address for server 0 (format: host:adminPort.grpcPort).
|
||||
sourceDataNode := cluster.VolumeAdminAddress(0) + "." +
|
||||
strings.Split(cluster.VolumeGRPCAddress(0), ":")[1]
|
||||
|
||||
// Copy shard 0 + ecx + ecj from server 0 → server 1.
|
||||
_, err = client1.VolumeEcShardsCopy(ctx, &volume_server_pb.VolumeEcShardsCopyRequest{
|
||||
VolumeId: volumeID,
|
||||
Collection: "",
|
||||
SourceDataNode: sourceDataNode,
|
||||
ShardIds: []uint32{0},
|
||||
CopyEcxFile: true,
|
||||
CopyEcjFile: true,
|
||||
CopyVifFile: true,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("VolumeEcShardsCopy 0→1: %v", err)
|
||||
}
|
||||
|
||||
// Mount shard 0 on server 1 so the EC volume can accept deletions.
|
||||
_, err = client1.VolumeEcShardsMount(ctx, &volume_server_pb.VolumeEcShardsMountRequest{
|
||||
VolumeId: volumeID, Collection: "",
|
||||
ShardIds: []uint32{0},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("VolumeEcShardsMount on server 1: %v", err)
|
||||
}
|
||||
|
||||
// Delete needle A on server 1 only (creates .ecj entry on server 1).
|
||||
_, err = client1.VolumeEcBlobDelete(ctx, &volume_server_pb.VolumeEcBlobDeleteRequest{
|
||||
VolumeId: volumeID, Collection: "",
|
||||
FileKey: keyA, Version: uint32(needle.GetCurrentVersion()),
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("VolumeEcBlobDelete needle A on server 1: %v", err)
|
||||
}
|
||||
|
||||
// Mount all data shards on server 0 (the decode target).
|
||||
_, err = client0.VolumeEcShardsMount(ctx, &volume_server_pb.VolumeEcShardsMountRequest{
|
||||
VolumeId: volumeID, Collection: "",
|
||||
ShardIds: []uint32{0, 1, 2, 3, 4, 5, 6, 7, 8, 9},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("VolumeEcShardsMount on server 0: %v", err)
|
||||
}
|
||||
|
||||
// Collect .ecj from server 1 → server 0 with NO new shard IDs.
|
||||
// This is the critical path: server 1 has shard 0 which server 0 already
|
||||
// has, so needToCopyShardsInfo would be empty. Before the fix in
|
||||
// collectEcShards, this copy would be skipped entirely, losing server 1's
|
||||
// deletion entries.
|
||||
server1DataNode := cluster.VolumeAdminAddress(1) + "." +
|
||||
strings.Split(cluster.VolumeGRPCAddress(1), ":")[1]
|
||||
|
||||
_, err = client0.VolumeEcShardsCopy(ctx, &volume_server_pb.VolumeEcShardsCopyRequest{
|
||||
VolumeId: volumeID,
|
||||
Collection: "",
|
||||
SourceDataNode: server1DataNode,
|
||||
ShardIds: []uint32{}, // No new shards — just .ecj.
|
||||
CopyEcxFile: false,
|
||||
CopyEcjFile: true,
|
||||
CopyVifFile: false,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("VolumeEcShardsCopy .ecj from server 1→0: %v", err)
|
||||
}
|
||||
|
||||
// Unmount the normal volume before decode.
|
||||
_, err = client0.VolumeUnmount(ctx, &volume_server_pb.VolumeUnmountRequest{VolumeId: volumeID})
|
||||
if err != nil {
|
||||
t.Fatalf("VolumeUnmount on server 0: %v", err)
|
||||
}
|
||||
|
||||
// Decode on server 0. RebuildEcxFile should see needle A's deletion from
|
||||
// the .ecj that was collected from server 1.
|
||||
_, err = client0.VolumeEcShardsToVolume(ctx, &volume_server_pb.VolumeEcShardsToVolumeRequest{
|
||||
VolumeId: volumeID, Collection: "",
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("VolumeEcShardsToVolume on server 0: %v", err)
|
||||
}
|
||||
|
||||
// Re-mount the decoded normal volume.
|
||||
_, err = client0.VolumeMount(ctx, &volume_server_pb.VolumeMountRequest{VolumeId: volumeID})
|
||||
if err != nil {
|
||||
t.Fatalf("VolumeMount on server 0: %v", err)
|
||||
}
|
||||
|
||||
// Needle A should be gone — its deletion was in server 1's .ecj.
|
||||
respA := framework.ReadBytes(t, httpClient, cluster.VolumeAdminURL(0), fidA)
|
||||
bodyA := framework.ReadAllAndClose(t, respA)
|
||||
if respA.StatusCode >= 500 {
|
||||
t.Fatalf("needle A read: server error %d: %s", respA.StatusCode, bodyA)
|
||||
}
|
||||
if respA.StatusCode != http.StatusNotFound {
|
||||
t.Fatalf("needle A should be 404 (ecj from peer), got %d", respA.StatusCode)
|
||||
}
|
||||
|
||||
// Needle B should still be readable.
|
||||
respB := framework.ReadBytes(t, httpClient, cluster.VolumeAdminURL(0), fidB)
|
||||
bodyB := framework.ReadAllAndClose(t, respB)
|
||||
if respB.StatusCode != http.StatusOK {
|
||||
t.Fatalf("needle B read: expected 200, got %d", respB.StatusCode)
|
||||
}
|
||||
if string(bodyB) != string(payloadB) {
|
||||
t.Fatalf("needle B payload mismatch: got %q, want %q", bodyB, payloadB)
|
||||
}
|
||||
}
|
||||
@@ -943,7 +943,7 @@ func (s *AdminServer) GetClusterMasters() (*ClusterMastersData, error) {
|
||||
leaderCount++
|
||||
}
|
||||
|
||||
masterMap[masterInfo.Address] = masterInfo
|
||||
masterMap[master.Address] = masterInfo
|
||||
}
|
||||
|
||||
// Then, get additional master information from Raft cluster
|
||||
@@ -955,11 +955,11 @@ func (s *AdminServer) GetClusterMasters() (*ClusterMastersData, error) {
|
||||
|
||||
// Process each raft server
|
||||
for _, server := range resp.ClusterServers {
|
||||
// Raft stores gRPC addresses, convert to HTTP address
|
||||
httpAddress := pb.GrpcAddressToServerAddress(server.Address)
|
||||
address := server.Address
|
||||
httpAddress := pb.ServerAddress(address).ToHttpAddress()
|
||||
|
||||
// Update existing master info or create new one
|
||||
if masterInfo, exists := masterMap[httpAddress]; exists {
|
||||
if masterInfo, exists := masterMap[address]; exists {
|
||||
// Update existing master with raft data
|
||||
masterInfo.IsLeader = server.IsLeader
|
||||
masterInfo.Suffrage = server.Suffrage
|
||||
@@ -970,7 +970,7 @@ func (s *AdminServer) GetClusterMasters() (*ClusterMastersData, error) {
|
||||
IsLeader: server.IsLeader,
|
||||
Suffrage: server.Suffrage,
|
||||
}
|
||||
masterMap[httpAddress] = masterInfo
|
||||
masterMap[address] = masterInfo
|
||||
}
|
||||
|
||||
if server.IsLeader {
|
||||
@@ -1646,6 +1646,159 @@ func (as *AdminServer) GetConfigPersistence() *ConfigPersistence {
|
||||
return as.configPersistence
|
||||
}
|
||||
|
||||
// convertJSONToMaintenanceConfig converts JSON map to protobuf MaintenanceConfig
|
||||
func convertJSONToMaintenanceConfig(jsonConfig map[string]interface{}) (*maintenance.MaintenanceConfig, error) {
|
||||
config := &maintenance.MaintenanceConfig{}
|
||||
|
||||
// Helper function to get int32 from interface{}
|
||||
getInt32 := func(key string) (int32, error) {
|
||||
if val, ok := jsonConfig[key]; ok {
|
||||
switch v := val.(type) {
|
||||
case int:
|
||||
return int32(v), nil
|
||||
case int32:
|
||||
return v, nil
|
||||
case int64:
|
||||
return int32(v), nil
|
||||
case float64:
|
||||
return int32(v), nil
|
||||
default:
|
||||
return 0, fmt.Errorf("invalid type for %s: expected number, got %T", key, v)
|
||||
}
|
||||
}
|
||||
return 0, nil
|
||||
}
|
||||
|
||||
// Helper function to get bool from interface{}
|
||||
getBool := func(key string) bool {
|
||||
if val, ok := jsonConfig[key]; ok {
|
||||
if b, ok := val.(bool); ok {
|
||||
return b
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
var err error
|
||||
|
||||
// Convert basic fields
|
||||
config.Enabled = getBool("enabled")
|
||||
|
||||
if config.ScanIntervalSeconds, err = getInt32("scan_interval_seconds"); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if config.WorkerTimeoutSeconds, err = getInt32("worker_timeout_seconds"); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if config.TaskTimeoutSeconds, err = getInt32("task_timeout_seconds"); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if config.RetryDelaySeconds, err = getInt32("retry_delay_seconds"); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if config.MaxRetries, err = getInt32("max_retries"); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if config.CleanupIntervalSeconds, err = getInt32("cleanup_interval_seconds"); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if config.TaskRetentionSeconds, err = getInt32("task_retention_seconds"); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// Convert policy if present
|
||||
if policyData, ok := jsonConfig["policy"]; ok {
|
||||
if policyMap, ok := policyData.(map[string]interface{}); ok {
|
||||
policy := &maintenance.MaintenancePolicy{}
|
||||
|
||||
if globalMaxConcurrent, err := getInt32FromMap(policyMap, "global_max_concurrent"); err != nil {
|
||||
return nil, err
|
||||
} else {
|
||||
policy.GlobalMaxConcurrent = globalMaxConcurrent
|
||||
}
|
||||
|
||||
if defaultRepeatIntervalSeconds, err := getInt32FromMap(policyMap, "default_repeat_interval_seconds"); err != nil {
|
||||
return nil, err
|
||||
} else {
|
||||
policy.DefaultRepeatIntervalSeconds = defaultRepeatIntervalSeconds
|
||||
}
|
||||
|
||||
if defaultCheckIntervalSeconds, err := getInt32FromMap(policyMap, "default_check_interval_seconds"); err != nil {
|
||||
return nil, err
|
||||
} else {
|
||||
policy.DefaultCheckIntervalSeconds = defaultCheckIntervalSeconds
|
||||
}
|
||||
|
||||
// Convert task policies if present
|
||||
if taskPoliciesData, ok := policyMap["task_policies"]; ok {
|
||||
if taskPoliciesMap, ok := taskPoliciesData.(map[string]interface{}); ok {
|
||||
policy.TaskPolicies = make(map[string]*maintenance.TaskPolicy)
|
||||
|
||||
for taskType, taskPolicyData := range taskPoliciesMap {
|
||||
if taskPolicyMap, ok := taskPolicyData.(map[string]interface{}); ok {
|
||||
taskPolicy := &maintenance.TaskPolicy{}
|
||||
|
||||
taskPolicy.Enabled = getBoolFromMap(taskPolicyMap, "enabled")
|
||||
|
||||
if maxConcurrent, err := getInt32FromMap(taskPolicyMap, "max_concurrent"); err != nil {
|
||||
return nil, err
|
||||
} else {
|
||||
taskPolicy.MaxConcurrent = maxConcurrent
|
||||
}
|
||||
|
||||
if repeatIntervalSeconds, err := getInt32FromMap(taskPolicyMap, "repeat_interval_seconds"); err != nil {
|
||||
return nil, err
|
||||
} else {
|
||||
taskPolicy.RepeatIntervalSeconds = repeatIntervalSeconds
|
||||
}
|
||||
|
||||
if checkIntervalSeconds, err := getInt32FromMap(taskPolicyMap, "check_interval_seconds"); err != nil {
|
||||
return nil, err
|
||||
} else {
|
||||
taskPolicy.CheckIntervalSeconds = checkIntervalSeconds
|
||||
}
|
||||
|
||||
policy.TaskPolicies[taskType] = taskPolicy
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
config.Policy = policy
|
||||
}
|
||||
}
|
||||
|
||||
return config, nil
|
||||
}
|
||||
|
||||
// Helper functions for map conversion
|
||||
func getInt32FromMap(m map[string]interface{}, key string) (int32, error) {
|
||||
if val, ok := m[key]; ok {
|
||||
switch v := val.(type) {
|
||||
case int:
|
||||
return int32(v), nil
|
||||
case int32:
|
||||
return v, nil
|
||||
case int64:
|
||||
return int32(v), nil
|
||||
case float64:
|
||||
return int32(v), nil
|
||||
default:
|
||||
return 0, fmt.Errorf("invalid type for %s: expected number, got %T", key, v)
|
||||
}
|
||||
}
|
||||
return 0, nil
|
||||
}
|
||||
|
||||
func getBoolFromMap(m map[string]interface{}, key string) bool {
|
||||
if val, ok := m[key]; ok {
|
||||
if b, ok := val.(bool); ok {
|
||||
return b
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
type collectionStats struct {
|
||||
PhysicalSize int64
|
||||
LogicalSize int64
|
||||
|
||||
@@ -361,6 +361,26 @@ func normalizeQuotaUnit(unit string) (string, error) {
|
||||
}
|
||||
}
|
||||
|
||||
// Helper function to convert bytes to appropriate unit and size
|
||||
func convertBytesToQuota(bytes int64) (int64, string) {
|
||||
if bytes == 0 {
|
||||
return 0, "MB"
|
||||
}
|
||||
|
||||
// Convert to TB if >= 1TB
|
||||
if bytes >= 1024*1024*1024*1024 && bytes%(1024*1024*1024*1024) == 0 {
|
||||
return bytes / (1024 * 1024 * 1024 * 1024), "TB"
|
||||
}
|
||||
|
||||
// Convert to GB if >= 1GB
|
||||
if bytes >= 1024*1024*1024 && bytes%(1024*1024*1024) == 0 {
|
||||
return bytes / (1024 * 1024 * 1024), "GB"
|
||||
}
|
||||
|
||||
// Convert to MB (default)
|
||||
return bytes / (1024 * 1024), "MB"
|
||||
}
|
||||
|
||||
// SetBucketQuota sets the quota for a bucket
|
||||
func (s *AdminServer) SetBucketQuota(bucketName string, quotaBytes int64, quotaEnabled bool) error {
|
||||
return s.WithFilerClient(func(client filer_pb.SeaweedFilerClient) error {
|
||||
|
||||
@@ -506,6 +506,18 @@ func getShardCount(ecIndexBits uint32) int {
|
||||
return count
|
||||
}
|
||||
|
||||
// getMissingShards returns a slice of missing shard IDs for a volume
|
||||
// Assumes default 10+4 EC configuration (14 total shards)
|
||||
func getMissingShards(ecIndexBits uint32) []int {
|
||||
var missing []int
|
||||
for i := 0; i < erasure_coding.TotalShardsCount; i++ {
|
||||
if (ecIndexBits & (1 << uint(i))) == 0 {
|
||||
missing = append(missing, i)
|
||||
}
|
||||
}
|
||||
return missing
|
||||
}
|
||||
|
||||
// sortEcShards sorts EC shards based on the specified field and order
|
||||
func sortEcShards(shards []EcShardWithInfo, sortBy string, sortOrder string) {
|
||||
sort.Slice(shards, func(i, j int) bool {
|
||||
|
||||
@@ -430,6 +430,67 @@ func (s *AdminServer) GetConsumerGroupOffsets(namespace, topicName string) ([]Co
|
||||
return offsets, nil
|
||||
}
|
||||
|
||||
// convertRecordTypeToSchemaFields converts a protobuf RecordType to SchemaFieldInfo slice
|
||||
func convertRecordTypeToSchemaFields(recordType *schema_pb.RecordType) []SchemaFieldInfo {
|
||||
var schemaFields []SchemaFieldInfo
|
||||
|
||||
if recordType == nil || recordType.Fields == nil {
|
||||
return schemaFields
|
||||
}
|
||||
|
||||
for _, field := range recordType.Fields {
|
||||
schemaField := SchemaFieldInfo{
|
||||
Name: field.Name,
|
||||
Type: getFieldTypeString(field.Type),
|
||||
Required: field.IsRequired,
|
||||
}
|
||||
schemaFields = append(schemaFields, schemaField)
|
||||
}
|
||||
|
||||
return schemaFields
|
||||
}
|
||||
|
||||
// getFieldTypeString converts a protobuf Type to a human-readable string
|
||||
func getFieldTypeString(fieldType *schema_pb.Type) string {
|
||||
if fieldType == nil {
|
||||
return "unknown"
|
||||
}
|
||||
|
||||
switch kind := fieldType.Kind.(type) {
|
||||
case *schema_pb.Type_ScalarType:
|
||||
return getScalarTypeString(kind.ScalarType)
|
||||
case *schema_pb.Type_RecordType:
|
||||
return "record"
|
||||
case *schema_pb.Type_ListType:
|
||||
elementType := getFieldTypeString(kind.ListType.ElementType)
|
||||
return fmt.Sprintf("list<%s>", elementType)
|
||||
default:
|
||||
return "unknown"
|
||||
}
|
||||
}
|
||||
|
||||
// getScalarTypeString converts a protobuf ScalarType to a string
|
||||
func getScalarTypeString(scalarType schema_pb.ScalarType) string {
|
||||
switch scalarType {
|
||||
case schema_pb.ScalarType_BOOL:
|
||||
return "bool"
|
||||
case schema_pb.ScalarType_INT32:
|
||||
return "int32"
|
||||
case schema_pb.ScalarType_INT64:
|
||||
return "int64"
|
||||
case schema_pb.ScalarType_FLOAT:
|
||||
return "float"
|
||||
case schema_pb.ScalarType_DOUBLE:
|
||||
return "double"
|
||||
case schema_pb.ScalarType_BYTES:
|
||||
return "bytes"
|
||||
case schema_pb.ScalarType_STRING:
|
||||
return "string"
|
||||
default:
|
||||
return "unknown"
|
||||
}
|
||||
}
|
||||
|
||||
// convertTopicPublishers converts protobuf TopicPublisher slice to PublisherInfo slice
|
||||
func convertTopicPublishers(publishers []*mq_pb.TopicPublisher) []PublisherInfo {
|
||||
publisherInfos := make([]PublisherInfo, 0, len(publishers))
|
||||
|
||||
@@ -2,6 +2,8 @@ package dash
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/rand"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
@@ -849,6 +851,43 @@ func normalizeTimeout(timeoutSeconds int, defaultTimeout, maxTimeout time.Durati
|
||||
return timeout
|
||||
}
|
||||
|
||||
func buildJobSpecFromProposal(jobType string, proposal *plugin_pb.JobProposal, index int) *plugin_pb.JobSpec {
|
||||
now := timestamppb.Now()
|
||||
suffix := make([]byte, 4)
|
||||
if _, err := rand.Read(suffix); err != nil {
|
||||
// Fallback to simpler ID if rand fails
|
||||
suffix = []byte(fmt.Sprintf("%d", index))
|
||||
}
|
||||
jobID := fmt.Sprintf("%s-%d-%s", jobType, now.AsTime().UnixNano(), hex.EncodeToString(suffix))
|
||||
|
||||
jobSpec := &plugin_pb.JobSpec{
|
||||
JobId: jobID,
|
||||
JobType: jobType,
|
||||
Priority: plugin_pb.JobPriority_JOB_PRIORITY_NORMAL,
|
||||
CreatedAt: now,
|
||||
Labels: make(map[string]string),
|
||||
Parameters: make(map[string]*plugin_pb.ConfigValue),
|
||||
DedupeKey: "",
|
||||
}
|
||||
|
||||
if proposal != nil {
|
||||
jobSpec.Summary = proposal.Summary
|
||||
jobSpec.Detail = proposal.Detail
|
||||
if proposal.Priority != plugin_pb.JobPriority_JOB_PRIORITY_UNSPECIFIED {
|
||||
jobSpec.Priority = proposal.Priority
|
||||
}
|
||||
jobSpec.DedupeKey = proposal.DedupeKey
|
||||
jobSpec.Parameters = plugin.CloneConfigValueMap(proposal.Parameters)
|
||||
if proposal.Labels != nil {
|
||||
for k, v := range proposal.Labels {
|
||||
jobSpec.Labels[k] = v
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return jobSpec
|
||||
}
|
||||
|
||||
func applyDescriptorDefaultsToPersistedConfig(
|
||||
config *plugin_pb.PersistedJobTypeConfig,
|
||||
descriptor *plugin_pb.JobTypeDescriptor,
|
||||
|
||||
@@ -115,6 +115,32 @@ func TestExpirePluginJobAPI(t *testing.T) {
|
||||
})
|
||||
}
|
||||
|
||||
func TestBuildJobSpecFromProposalDoesNotReuseProposalID(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
proposal := &plugin_pb.JobProposal{
|
||||
ProposalId: "vacuum-2",
|
||||
DedupeKey: "vacuum:2",
|
||||
JobType: "vacuum",
|
||||
}
|
||||
|
||||
jobA := buildJobSpecFromProposal("vacuum", proposal, 0)
|
||||
jobB := buildJobSpecFromProposal("vacuum", proposal, 1)
|
||||
|
||||
if jobA.JobId == proposal.ProposalId {
|
||||
t.Fatalf("job id must not reuse proposal id: %s", jobA.JobId)
|
||||
}
|
||||
if jobB.JobId == proposal.ProposalId {
|
||||
t.Fatalf("job id must not reuse proposal id: %s", jobB.JobId)
|
||||
}
|
||||
if jobA.JobId == jobB.JobId {
|
||||
t.Fatalf("job ids must be unique across jobs: %s", jobA.JobId)
|
||||
}
|
||||
if jobA.DedupeKey != proposal.DedupeKey {
|
||||
t.Fatalf("dedupe key must be preserved: got=%s want=%s", jobA.DedupeKey, proposal.DedupeKey)
|
||||
}
|
||||
}
|
||||
|
||||
func TestApplyDescriptorDefaultsToPersistedConfigBackfillsAdminDefaults(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
|
||||
@@ -787,6 +787,15 @@ func (s *WorkerGrpcServer) RequestTaskLogsFromAllWorkers(taskID string, maxEntri
|
||||
return results, nil
|
||||
}
|
||||
|
||||
// convertTaskParameters converts task parameters to protobuf format
|
||||
func convertTaskParameters(params map[string]interface{}) map[string]string {
|
||||
result := make(map[string]string)
|
||||
for key, value := range params {
|
||||
result[key] = fmt.Sprintf("%v", value)
|
||||
}
|
||||
return result
|
||||
}
|
||||
|
||||
func findClientAddress(ctx context.Context) string {
|
||||
// fmt.Printf("FromContext %+v\n", ctx)
|
||||
pr, ok := peer.FromContext(ctx)
|
||||
|
||||
+14
-50
@@ -53,8 +53,6 @@ type SyncOptions struct {
|
||||
chunkConcurrency *int
|
||||
aDoDeleteFiles *bool
|
||||
bDoDeleteFiles *bool
|
||||
aSecurity *string
|
||||
bSecurity *string
|
||||
clientId int32
|
||||
clientEpoch atomic.Int32
|
||||
debug *bool
|
||||
@@ -115,8 +113,6 @@ func init() {
|
||||
syncOptions.metricsHttpPort = cmdFilerSynchronize.Flag.Int("metricsPort", 0, "metrics listen port")
|
||||
syncOptions.aDoDeleteFiles = cmdFilerSynchronize.Flag.Bool("a.doDeleteFiles", true, "delete and update files when synchronizing on filer A")
|
||||
syncOptions.bDoDeleteFiles = cmdFilerSynchronize.Flag.Bool("b.doDeleteFiles", true, "delete and update files when synchronizing on filer B")
|
||||
syncOptions.aSecurity = cmdFilerSynchronize.Flag.String("a.security", "", "security.toml file for filer A when clusters use different certificates")
|
||||
syncOptions.bSecurity = cmdFilerSynchronize.Flag.String("b.security", "", "security.toml file for filer B when clusters use different certificates")
|
||||
syncOptions.debug = cmdFilerSynchronize.Flag.Bool("debug", false, "serves runtime profiling data via pprof on the port specified by -debug.port")
|
||||
syncOptions.debugPort = cmdFilerSynchronize.Flag.Int("debug.port", 6060, "http port for debugging")
|
||||
syncOptions.clientId = util.RandomInt32()
|
||||
@@ -148,22 +144,6 @@ func runFilerSynchronize(cmd *Command, args []string) bool {
|
||||
util.LoadSecurityConfiguration()
|
||||
grpcDialOption := security.LoadClientTLS(util.GetViper(), "grpc.client")
|
||||
|
||||
// per-filer TLS when clusters use different certificates
|
||||
grpcDialOptionA := grpcDialOption
|
||||
grpcDialOptionB := grpcDialOption
|
||||
if *syncOptions.aSecurity != "" {
|
||||
var err error
|
||||
if grpcDialOptionA, err = security.LoadClientTLSFromFile(*syncOptions.aSecurity, "grpc.client"); err != nil {
|
||||
glog.Fatalf("load security config for filer A: %v", err)
|
||||
}
|
||||
}
|
||||
if *syncOptions.bSecurity != "" {
|
||||
var err error
|
||||
if grpcDialOptionB, err = security.LoadClientTLSFromFile(*syncOptions.bSecurity, "grpc.client"); err != nil {
|
||||
glog.Fatalf("load security config for filer B: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
grace.SetupProfiling(*syncCpuProfile, *syncMemProfile)
|
||||
|
||||
filerA := pb.ServerAddress(*syncOptions.filerA)
|
||||
@@ -173,13 +153,13 @@ func runFilerSynchronize(cmd *Command, args []string) bool {
|
||||
go statsCollect.StartMetricsServer(*syncOptions.metricsHttpIp, *syncOptions.metricsHttpPort)
|
||||
|
||||
// read a filer signature
|
||||
aFilerSignature, aFilerErr := replication.ReadFilerSignature(grpcDialOptionA, filerA)
|
||||
aFilerSignature, aFilerErr := replication.ReadFilerSignature(grpcDialOption, filerA)
|
||||
if aFilerErr != nil {
|
||||
glog.Errorf("get filer 'a' signature %d error from %s to %s: %v", aFilerSignature, *syncOptions.filerA, *syncOptions.filerB, aFilerErr)
|
||||
return true
|
||||
}
|
||||
// read b filer signature
|
||||
bFilerSignature, bFilerErr := replication.ReadFilerSignature(grpcDialOptionB, filerB)
|
||||
bFilerSignature, bFilerErr := replication.ReadFilerSignature(grpcDialOption, filerB)
|
||||
if bFilerErr != nil {
|
||||
glog.Errorf("get filer 'b' signature %d error from %s to %s: %v", bFilerSignature, *syncOptions.filerA, *syncOptions.filerB, bFilerErr)
|
||||
return true
|
||||
@@ -209,9 +189,9 @@ func runFilerSynchronize(cmd *Command, args []string) bool {
|
||||
go func() {
|
||||
// a->b
|
||||
// set synchronization start timestamp to offset
|
||||
initOffsetError := initOffsetFromTsMs(grpcDialOptionB, filerB, aFilerSignature, *syncOptions.aFromTsMs, getSignaturePrefixByPath(*syncOptions.aPath))
|
||||
initOffsetError := initOffsetFromTsMs(grpcDialOption, filerB, aFilerSignature, *syncOptions.bFromTsMs, getSignaturePrefixByPath(*syncOptions.aPath))
|
||||
if initOffsetError != nil {
|
||||
glog.Errorf("init offset from timestamp %d error from %s to %s: %v", *syncOptions.aFromTsMs, *syncOptions.filerA, *syncOptions.filerB, initOffsetError)
|
||||
glog.Errorf("init offset from timestamp %d error from %s to %s: %v", *syncOptions.bFromTsMs, *syncOptions.filerA, *syncOptions.filerB, initOffsetError)
|
||||
os.Exit(2)
|
||||
}
|
||||
for {
|
||||
@@ -219,12 +199,11 @@ func runFilerSynchronize(cmd *Command, args []string) bool {
|
||||
err := doSubscribeFilerMetaChanges(
|
||||
syncOptions.clientId,
|
||||
syncOptions.clientEpoch.Load(),
|
||||
grpcDialOptionA,
|
||||
grpcDialOption,
|
||||
filerA,
|
||||
*syncOptions.aPath,
|
||||
util.StringSplit(*syncOptions.aExcludePaths, ","),
|
||||
*syncOptions.aProxyByFiler,
|
||||
grpcDialOptionB,
|
||||
filerB,
|
||||
*syncOptions.bPath,
|
||||
*syncOptions.bReplication,
|
||||
@@ -249,9 +228,9 @@ func runFilerSynchronize(cmd *Command, args []string) bool {
|
||||
if !*syncOptions.isActivePassive {
|
||||
// b->a
|
||||
// set synchronization start timestamp to offset
|
||||
initOffsetError := initOffsetFromTsMs(grpcDialOptionA, filerA, bFilerSignature, *syncOptions.bFromTsMs, getSignaturePrefixByPath(*syncOptions.bPath))
|
||||
initOffsetError := initOffsetFromTsMs(grpcDialOption, filerA, bFilerSignature, *syncOptions.aFromTsMs, getSignaturePrefixByPath(*syncOptions.bPath))
|
||||
if initOffsetError != nil {
|
||||
glog.Errorf("init offset from timestamp %d error from %s to %s: %v", *syncOptions.bFromTsMs, *syncOptions.filerB, *syncOptions.filerA, initOffsetError)
|
||||
glog.Errorf("init offset from timestamp %d error from %s to %s: %v", *syncOptions.aFromTsMs, *syncOptions.filerB, *syncOptions.filerA, initOffsetError)
|
||||
os.Exit(2)
|
||||
}
|
||||
go func() {
|
||||
@@ -260,12 +239,11 @@ func runFilerSynchronize(cmd *Command, args []string) bool {
|
||||
err := doSubscribeFilerMetaChanges(
|
||||
syncOptions.clientId,
|
||||
syncOptions.clientEpoch.Load(),
|
||||
grpcDialOptionB,
|
||||
grpcDialOption,
|
||||
filerB,
|
||||
*syncOptions.bPath,
|
||||
util.StringSplit(*syncOptions.bExcludePaths, ","),
|
||||
*syncOptions.bProxyByFiler,
|
||||
grpcDialOptionA,
|
||||
filerA,
|
||||
*syncOptions.aPath,
|
||||
*syncOptions.aReplication,
|
||||
@@ -307,12 +285,12 @@ func initOffsetFromTsMs(grpcDialOption grpc.DialOption, targetFiler pb.ServerAdd
|
||||
return nil
|
||||
}
|
||||
|
||||
func doSubscribeFilerMetaChanges(clientId int32, clientEpoch int32, sourceGrpcDialOption grpc.DialOption, sourceFiler pb.ServerAddress, sourcePath string, sourceExcludePaths []string, sourceReadChunkFromFiler bool, targetGrpcDialOption grpc.DialOption, targetFiler pb.ServerAddress, targetPath string,
|
||||
func doSubscribeFilerMetaChanges(clientId int32, clientEpoch int32, grpcDialOption grpc.DialOption, sourceFiler pb.ServerAddress, sourcePath string, sourceExcludePaths []string, sourceReadChunkFromFiler bool, targetFiler pb.ServerAddress, targetPath string,
|
||||
replicationStr, collection string, ttlSec int, sinkWriteChunkByFiler bool, diskType string, debug bool, concurrency int, chunkConcurrency int, doDeleteFiles bool, sourceFilerSignature int32, targetFilerSignature int32, statePtr *atomic.Pointer[syncState]) error {
|
||||
|
||||
// if first time, start from now
|
||||
// if has previously synced, resume from that point of time
|
||||
sourceFilerOffsetTsNs, err := getOffset(targetGrpcDialOption, targetFiler, getSignaturePrefixByPath(sourcePath), sourceFilerSignature)
|
||||
sourceFilerOffsetTsNs, err := getOffset(grpcDialOption, targetFiler, getSignaturePrefixByPath(sourcePath), sourceFilerSignature)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -322,9 +300,8 @@ func doSubscribeFilerMetaChanges(clientId int32, clientEpoch int32, sourceGrpcDi
|
||||
// create filer sink
|
||||
filerSource := &source.FilerSource{}
|
||||
filerSource.DoInitialize(sourceFiler.ToHttpAddress(), sourceFiler.ToGrpcAddress(), sourcePath, sourceReadChunkFromFiler)
|
||||
filerSource.SetGrpcDialOption(sourceGrpcDialOption)
|
||||
filerSink := &filersink.FilerSink{}
|
||||
filerSink.DoInitialize(targetFiler.ToHttpAddress(), targetFiler.ToGrpcAddress(), targetPath, replicationStr, collection, ttlSec, diskType, targetGrpcDialOption, sinkWriteChunkByFiler)
|
||||
filerSink.DoInitialize(targetFiler.ToHttpAddress(), targetFiler.ToGrpcAddress(), targetPath, replicationStr, collection, ttlSec, diskType, grpcDialOption, sinkWriteChunkByFiler)
|
||||
filerSink.SetChunkConcurrency(chunkConcurrency)
|
||||
filerSink.SetSourceFiler(filerSource)
|
||||
|
||||
@@ -351,7 +328,7 @@ func doSubscribeFilerMetaChanges(clientId int32, clientEpoch int32, sourceGrpcDi
|
||||
if statePtr != nil {
|
||||
statePtr.Store(&syncState{
|
||||
processor: processor,
|
||||
grpcDialOption: targetGrpcDialOption,
|
||||
grpcDialOption: grpcDialOption,
|
||||
targetFiler: targetFiler,
|
||||
sourcePath: sourcePath,
|
||||
sourceFilerSignature: sourceFilerSignature,
|
||||
@@ -359,7 +336,6 @@ 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)
|
||||
@@ -373,21 +349,9 @@ 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)
|
||||
return setOffset(grpcDialOption, targetFiler, getSignaturePrefixByPath(sourcePath), sourceFilerSignature, offsetTsNs)
|
||||
})
|
||||
|
||||
prefix := sourcePath
|
||||
@@ -408,7 +372,7 @@ func doSubscribeFilerMetaChanges(clientId int32, clientEpoch int32, sourceGrpcDi
|
||||
EventErrorType: pb.RetryForeverOnError,
|
||||
}
|
||||
|
||||
return pb.FollowMetadata(sourceFiler, sourceGrpcDialOption, metadataFollowOption, processEventFnWithOffset)
|
||||
return pb.FollowMetadata(sourceFiler, grpcDialOption, metadataFollowOption, processEventFnWithOffset)
|
||||
|
||||
}
|
||||
|
||||
|
||||
@@ -2,6 +2,7 @@ package engine
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/binary"
|
||||
"fmt"
|
||||
"io"
|
||||
"strings"
|
||||
@@ -17,6 +18,7 @@ import (
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/master_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/mq_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/schema_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/credentials/insecure"
|
||||
jsonpb "google.golang.org/protobuf/encoding/protojson"
|
||||
@@ -507,3 +509,77 @@ func (c *BrokerClient) GetUnflushedMessages(ctx context.Context, namespace, topi
|
||||
|
||||
return logEntries, nil
|
||||
}
|
||||
|
||||
// getEarliestBufferStart finds the earliest buffer_start index from disk files in the partition
|
||||
//
|
||||
// This method handles three scenarios for seamless broker querying:
|
||||
// 1. Live log files exist: Uses their buffer_start metadata (most recent boundaries)
|
||||
// 2. Only Parquet files exist: Uses Parquet buffer_start metadata (preserved from archived sources)
|
||||
// 3. Mixed files: Uses earliest buffer_start from all sources for comprehensive coverage
|
||||
//
|
||||
// This ensures continuous real-time querying capability even after log file compaction/archival
|
||||
func (c *BrokerClient) getEarliestBufferStart(ctx context.Context, partitionPath string) (int64, error) {
|
||||
filerClient, err := c.GetFilerClient()
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("failed to get filer client: %v", err)
|
||||
}
|
||||
|
||||
var earliestBufferIndex int64 = -1 // -1 means no buffer_start found
|
||||
var logFileCount, parquetFileCount int
|
||||
var bufferStartSources []string // Track which files provide buffer_start
|
||||
|
||||
err = filer_pb.ReadDirAllEntries(ctx, filerClient, util.FullPath(partitionPath), "", func(entry *filer_pb.Entry, isLast bool) error {
|
||||
// Skip directories
|
||||
if entry.IsDirectory {
|
||||
return nil
|
||||
}
|
||||
|
||||
// Count file types for scenario detection
|
||||
if strings.HasSuffix(entry.Name, ".parquet") {
|
||||
parquetFileCount++
|
||||
} else {
|
||||
logFileCount++
|
||||
}
|
||||
|
||||
// Extract buffer_start from file extended attributes (both log files and parquet files)
|
||||
bufferStart := c.getBufferStartFromEntry(entry)
|
||||
if bufferStart != nil && bufferStart.StartIndex > 0 {
|
||||
if earliestBufferIndex == -1 || bufferStart.StartIndex < earliestBufferIndex {
|
||||
earliestBufferIndex = bufferStart.StartIndex
|
||||
}
|
||||
bufferStartSources = append(bufferStartSources, entry.Name)
|
||||
}
|
||||
|
||||
return nil
|
||||
})
|
||||
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("failed to scan partition directory: %v", err)
|
||||
}
|
||||
|
||||
if earliestBufferIndex == -1 {
|
||||
return 0, fmt.Errorf("no buffer_start metadata found in partition")
|
||||
}
|
||||
|
||||
return earliestBufferIndex, nil
|
||||
}
|
||||
|
||||
// getBufferStartFromEntry extracts LogBufferStart from file entry metadata
|
||||
// Only supports binary format (used by both log files and Parquet files)
|
||||
func (c *BrokerClient) getBufferStartFromEntry(entry *filer_pb.Entry) *LogBufferStart {
|
||||
if entry.Extended == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
if startData, exists := entry.Extended["buffer_start"]; exists {
|
||||
// Only support binary format
|
||||
if len(startData) == 8 {
|
||||
startIndex := int64(binary.BigEndian.Uint64(startData))
|
||||
if startIndex > 0 {
|
||||
return &LogBufferStart{StartIndex: startIndex}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -170,6 +170,27 @@ func (e *SQLEngine) convertRawValueToSchemaValue(rawValue interface{}) *schema_p
|
||||
}
|
||||
}
|
||||
|
||||
// convertJSONValueToSchemaValue converts JSON values to schema_pb.Value
|
||||
func (e *SQLEngine) convertJSONValueToSchemaValue(jsonValue interface{}) *schema_pb.Value {
|
||||
switch v := jsonValue.(type) {
|
||||
case string:
|
||||
return &schema_pb.Value{Kind: &schema_pb.Value_StringValue{StringValue: v}}
|
||||
case float64:
|
||||
// JSON numbers are always float64, try to detect if it's actually an integer
|
||||
if v == float64(int64(v)) {
|
||||
return &schema_pb.Value{Kind: &schema_pb.Value_Int64Value{Int64Value: int64(v)}}
|
||||
}
|
||||
return &schema_pb.Value{Kind: &schema_pb.Value_DoubleValue{DoubleValue: v}}
|
||||
case bool:
|
||||
return &schema_pb.Value{Kind: &schema_pb.Value_BoolValue{BoolValue: v}}
|
||||
case nil:
|
||||
return nil
|
||||
default:
|
||||
// Convert other types to string
|
||||
return &schema_pb.Value{Kind: &schema_pb.Value_StringValue{StringValue: fmt.Sprintf("%v", v)}}
|
||||
}
|
||||
}
|
||||
|
||||
// Helper functions for aggregation processing
|
||||
|
||||
// isNullValue checks if a schema_pb.Value is null or empty
|
||||
|
||||
@@ -2175,6 +2175,361 @@ func (e *SQLEngine) executeRegularSelectWithHybridScanner(ctx context.Context, h
|
||||
return e.ConvertToSQLResultWithExpressions(hybridScanner, results, stmt.SelectExprs), nil
|
||||
}
|
||||
|
||||
// executeSelectStatementWithBrokerStats handles SELECT queries with broker buffer statistics capture
|
||||
// This is used by EXPLAIN queries to capture complete data source information including broker memory
|
||||
func (e *SQLEngine) executeSelectStatementWithBrokerStats(ctx context.Context, stmt *SelectStatement, plan *QueryExecutionPlan) (*QueryResult, error) {
|
||||
// Parse FROM clause to get table (topic) information
|
||||
if len(stmt.From) != 1 {
|
||||
err := fmt.Errorf("SELECT supports single table queries only")
|
||||
return &QueryResult{Error: err}, err
|
||||
}
|
||||
|
||||
// Extract table reference
|
||||
var database, tableName string
|
||||
switch table := stmt.From[0].(type) {
|
||||
case *AliasedTableExpr:
|
||||
switch tableExpr := table.Expr.(type) {
|
||||
case TableName:
|
||||
tableName = tableExpr.Name.String()
|
||||
if tableExpr.Qualifier != nil && tableExpr.Qualifier.String() != "" {
|
||||
database = tableExpr.Qualifier.String()
|
||||
}
|
||||
default:
|
||||
err := fmt.Errorf("unsupported table expression: %T", tableExpr)
|
||||
return &QueryResult{Error: err}, err
|
||||
}
|
||||
default:
|
||||
err := fmt.Errorf("unsupported FROM clause: %T", table)
|
||||
return &QueryResult{Error: err}, err
|
||||
}
|
||||
|
||||
// Use current database context if not specified
|
||||
if database == "" {
|
||||
database = e.catalog.GetCurrentDatabase()
|
||||
if database == "" {
|
||||
database = "default"
|
||||
}
|
||||
}
|
||||
|
||||
// Auto-discover and register topic if not already in catalog
|
||||
if _, err := e.catalog.GetTableInfo(database, tableName); err != nil {
|
||||
// Topic not in catalog, try to discover and register it
|
||||
if regErr := e.discoverAndRegisterTopic(ctx, database, tableName); regErr != nil {
|
||||
// Return error immediately for non-existent topics instead of falling back to sample data
|
||||
return &QueryResult{Error: regErr}, regErr
|
||||
}
|
||||
}
|
||||
|
||||
// Create HybridMessageScanner for the topic (reads both live logs + Parquet files)
|
||||
// Get filerClient from broker connection (works with both real and mock brokers)
|
||||
var filerClient filer_pb.FilerClient
|
||||
var filerClientErr error
|
||||
filerClient, filerClientErr = e.catalog.brokerClient.GetFilerClient()
|
||||
if filerClientErr != nil {
|
||||
// Return error if filer client is not available for topic access
|
||||
return &QueryResult{Error: filerClientErr}, filerClientErr
|
||||
}
|
||||
|
||||
hybridScanner, err := NewHybridMessageScanner(filerClient, e.catalog.brokerClient, database, tableName, e)
|
||||
if err != nil {
|
||||
// Handle quiet topics gracefully: topics exist but have no active schema/brokers
|
||||
if IsNoSchemaError(err) {
|
||||
// Return empty result for quiet topics (normal in production environments)
|
||||
return &QueryResult{
|
||||
Columns: []string{},
|
||||
Rows: [][]sqltypes.Value{},
|
||||
Database: database,
|
||||
Table: tableName,
|
||||
}, nil
|
||||
}
|
||||
// Return error for other access issues (truly non-existent topics, etc.)
|
||||
topicErr := fmt.Errorf("failed to access topic %s.%s: %v", database, tableName, err)
|
||||
return &QueryResult{Error: topicErr}, topicErr
|
||||
}
|
||||
|
||||
// Parse SELECT columns and detect aggregation functions
|
||||
var columns []string
|
||||
var aggregations []AggregationSpec
|
||||
selectAll := false
|
||||
hasAggregations := false
|
||||
_ = hasAggregations // Used later in aggregation routing
|
||||
// Track required base columns for arithmetic expressions
|
||||
baseColumnsSet := make(map[string]bool)
|
||||
|
||||
for _, selectExpr := range stmt.SelectExprs {
|
||||
switch expr := selectExpr.(type) {
|
||||
case *StarExpr:
|
||||
selectAll = true
|
||||
case *AliasedExpr:
|
||||
switch col := expr.Expr.(type) {
|
||||
case *ColName:
|
||||
colName := col.Name.String()
|
||||
columns = append(columns, colName)
|
||||
baseColumnsSet[colName] = true
|
||||
case *ArithmeticExpr:
|
||||
// Handle arithmetic expressions like id+user_id and string concatenation like name||suffix
|
||||
columns = append(columns, e.getArithmeticExpressionAlias(col))
|
||||
// Extract base columns needed for this arithmetic expression
|
||||
e.extractBaseColumns(col, baseColumnsSet)
|
||||
case *SQLVal:
|
||||
// Handle string/numeric literals like 'good', 123, etc.
|
||||
columns = append(columns, e.getSQLValAlias(col))
|
||||
case *FuncExpr:
|
||||
// Distinguish between aggregation functions and string functions
|
||||
funcName := strings.ToUpper(col.Name.String())
|
||||
if e.isAggregationFunction(funcName) {
|
||||
// Handle aggregation functions
|
||||
aggSpec, err := e.parseAggregationFunction(col, expr)
|
||||
if err != nil {
|
||||
return &QueryResult{Error: err}, err
|
||||
}
|
||||
aggregations = append(aggregations, *aggSpec)
|
||||
hasAggregations = true
|
||||
} else if e.isStringFunction(funcName) {
|
||||
// Handle string functions like UPPER, LENGTH, etc.
|
||||
columns = append(columns, e.getStringFunctionAlias(col))
|
||||
// Extract base columns needed for this string function
|
||||
e.extractBaseColumnsFromFunction(col, baseColumnsSet)
|
||||
} else if e.isDateTimeFunction(funcName) {
|
||||
// Handle datetime functions like CURRENT_DATE, NOW, EXTRACT, DATE_TRUNC
|
||||
columns = append(columns, e.getDateTimeFunctionAlias(col))
|
||||
// Extract base columns needed for this datetime function
|
||||
e.extractBaseColumnsFromFunction(col, baseColumnsSet)
|
||||
} else {
|
||||
return &QueryResult{Error: fmt.Errorf("unsupported function: %s", funcName)}, fmt.Errorf("unsupported function: %s", funcName)
|
||||
}
|
||||
default:
|
||||
err := fmt.Errorf("unsupported SELECT expression: %T", col)
|
||||
return &QueryResult{Error: err}, err
|
||||
}
|
||||
default:
|
||||
err := fmt.Errorf("unsupported SELECT expression: %T", expr)
|
||||
return &QueryResult{Error: err}, err
|
||||
}
|
||||
}
|
||||
|
||||
// If we have aggregations, use aggregation query path
|
||||
if hasAggregations {
|
||||
return e.executeAggregationQuery(ctx, hybridScanner, aggregations, stmt)
|
||||
}
|
||||
|
||||
// Parse WHERE clause for predicate pushdown
|
||||
var predicate func(*schema_pb.RecordValue) bool
|
||||
if stmt.Where != nil {
|
||||
predicate, err = e.buildPredicateWithContext(stmt.Where.Expr, stmt.SelectExprs)
|
||||
if err != nil {
|
||||
return &QueryResult{Error: err}, err
|
||||
}
|
||||
}
|
||||
|
||||
// Parse LIMIT and OFFSET clauses
|
||||
// Use -1 to distinguish "no LIMIT" from "LIMIT 0"
|
||||
limit := -1
|
||||
offset := 0
|
||||
if stmt.Limit != nil && stmt.Limit.Rowcount != nil {
|
||||
switch limitExpr := stmt.Limit.Rowcount.(type) {
|
||||
case *SQLVal:
|
||||
if limitExpr.Type == IntVal {
|
||||
var parseErr error
|
||||
limit64, parseErr := strconv.ParseInt(string(limitExpr.Val), 10, 64)
|
||||
if parseErr != nil {
|
||||
return &QueryResult{Error: parseErr}, parseErr
|
||||
}
|
||||
if limit64 > math.MaxInt32 || limit64 < 0 {
|
||||
return &QueryResult{Error: fmt.Errorf("LIMIT value %d is out of valid range", limit64)}, fmt.Errorf("LIMIT value %d is out of valid range", limit64)
|
||||
}
|
||||
limit = int(limit64)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Parse OFFSET clause if present
|
||||
if stmt.Limit != nil && stmt.Limit.Offset != nil {
|
||||
switch offsetExpr := stmt.Limit.Offset.(type) {
|
||||
case *SQLVal:
|
||||
if offsetExpr.Type == IntVal {
|
||||
var parseErr error
|
||||
offset64, parseErr := strconv.ParseInt(string(offsetExpr.Val), 10, 64)
|
||||
if parseErr != nil {
|
||||
return &QueryResult{Error: parseErr}, parseErr
|
||||
}
|
||||
if offset64 > math.MaxInt32 || offset64 < 0 {
|
||||
return &QueryResult{Error: fmt.Errorf("OFFSET value %d is out of valid range", offset64)}, fmt.Errorf("OFFSET value %d is out of valid range", offset64)
|
||||
}
|
||||
offset = int(offset64)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Build hybrid scan options
|
||||
// Extract time filters from WHERE clause to optimize scanning
|
||||
startTimeNs, stopTimeNs := int64(0), int64(0)
|
||||
if stmt.Where != nil {
|
||||
startTimeNs, stopTimeNs = e.extractTimeFilters(stmt.Where.Expr)
|
||||
}
|
||||
|
||||
hybridScanOptions := HybridScanOptions{
|
||||
StartTimeNs: startTimeNs, // Extracted from WHERE clause time comparisons
|
||||
StopTimeNs: stopTimeNs, // Extracted from WHERE clause time comparisons
|
||||
Limit: limit,
|
||||
Offset: offset,
|
||||
Predicate: predicate,
|
||||
}
|
||||
|
||||
if !selectAll {
|
||||
// Convert baseColumnsSet to slice for hybrid scan options
|
||||
baseColumns := make([]string, 0, len(baseColumnsSet))
|
||||
for columnName := range baseColumnsSet {
|
||||
baseColumns = append(baseColumns, columnName)
|
||||
}
|
||||
// Use base columns (not expression aliases) for data retrieval
|
||||
if len(baseColumns) > 0 {
|
||||
hybridScanOptions.Columns = baseColumns
|
||||
} else {
|
||||
// If no base columns found (shouldn't happen), use original columns
|
||||
hybridScanOptions.Columns = columns
|
||||
}
|
||||
}
|
||||
|
||||
// Execute the hybrid scan with stats capture for EXPLAIN
|
||||
var results []HybridScanResult
|
||||
if plan != nil {
|
||||
// EXPLAIN mode - capture broker buffer stats
|
||||
var stats *HybridScanStats
|
||||
results, stats, err = hybridScanner.ScanWithStats(ctx, hybridScanOptions)
|
||||
if err != nil {
|
||||
return &QueryResult{Error: err}, err
|
||||
}
|
||||
|
||||
// Populate plan with broker buffer information
|
||||
if stats != nil {
|
||||
plan.BrokerBufferQueried = stats.BrokerBufferQueried
|
||||
plan.BrokerBufferMessages = stats.BrokerBufferMessages
|
||||
plan.BufferStartIndex = stats.BufferStartIndex
|
||||
|
||||
// Add broker_buffer to data sources if buffer was queried
|
||||
if stats.BrokerBufferQueried {
|
||||
// Check if broker_buffer is already in data sources
|
||||
hasBrokerBuffer := false
|
||||
for _, source := range plan.DataSources {
|
||||
if source == "broker_buffer" {
|
||||
hasBrokerBuffer = true
|
||||
break
|
||||
}
|
||||
}
|
||||
if !hasBrokerBuffer {
|
||||
plan.DataSources = append(plan.DataSources, "broker_buffer")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Populate execution plan details with source file information for Data Sources Tree
|
||||
if partitions, discoverErr := e.discoverTopicPartitions(database, tableName); discoverErr == nil {
|
||||
// Add partition paths to execution plan details
|
||||
plan.Details["partition_paths"] = partitions
|
||||
// Persist time filter details for downstream pruning/diagnostics
|
||||
plan.Details[PlanDetailStartTimeNs] = startTimeNs
|
||||
plan.Details[PlanDetailStopTimeNs] = stopTimeNs
|
||||
|
||||
// Collect actual file information for each partition
|
||||
var parquetFiles []string
|
||||
var liveLogFiles []string
|
||||
parquetSources := make(map[string]bool)
|
||||
|
||||
var parquetReadErrors []string
|
||||
var liveLogListErrors []string
|
||||
for _, partitionPath := range partitions {
|
||||
// Get parquet files for this partition
|
||||
if parquetStats, err := hybridScanner.ReadParquetStatistics(partitionPath); err == nil {
|
||||
// Prune files by time range with debug logging
|
||||
filteredStats := pruneParquetFilesByTime(ctx, parquetStats, hybridScanner, startTimeNs, stopTimeNs)
|
||||
|
||||
// Further prune by column statistics from WHERE clause
|
||||
if stmt.Where != nil {
|
||||
beforeColumnPrune := len(filteredStats)
|
||||
filteredStats = e.pruneParquetFilesByColumnStats(ctx, filteredStats, stmt.Where.Expr)
|
||||
columnPrunedCount := beforeColumnPrune - len(filteredStats)
|
||||
|
||||
if columnPrunedCount > 0 {
|
||||
// Track column statistics optimization
|
||||
if !contains(plan.OptimizationsUsed, "column_statistics_pruning") {
|
||||
plan.OptimizationsUsed = append(plan.OptimizationsUsed, "column_statistics_pruning")
|
||||
}
|
||||
}
|
||||
}
|
||||
for _, stats := range filteredStats {
|
||||
parquetFiles = append(parquetFiles, fmt.Sprintf("%s/%s", partitionPath, stats.FileName))
|
||||
}
|
||||
} else {
|
||||
parquetReadErrors = append(parquetReadErrors, fmt.Sprintf("%s: %v", partitionPath, err))
|
||||
}
|
||||
|
||||
// Merge accurate parquet sources from metadata
|
||||
if sources, err := e.getParquetSourceFilesFromMetadata(partitionPath); err == nil {
|
||||
for src := range sources {
|
||||
parquetSources[src] = true
|
||||
}
|
||||
}
|
||||
|
||||
// Get live log files for this partition
|
||||
if liveFiles, err := e.collectLiveLogFileNames(hybridScanner.filerClient, partitionPath); err == nil {
|
||||
for _, fileName := range liveFiles {
|
||||
// Exclude live log files that have been converted to parquet (deduplicated)
|
||||
if parquetSources[fileName] {
|
||||
continue
|
||||
}
|
||||
liveLogFiles = append(liveLogFiles, fmt.Sprintf("%s/%s", partitionPath, fileName))
|
||||
}
|
||||
} else {
|
||||
liveLogListErrors = append(liveLogListErrors, fmt.Sprintf("%s: %v", partitionPath, err))
|
||||
}
|
||||
}
|
||||
|
||||
if len(parquetFiles) > 0 {
|
||||
plan.Details["parquet_files"] = parquetFiles
|
||||
}
|
||||
if len(liveLogFiles) > 0 {
|
||||
plan.Details["live_log_files"] = liveLogFiles
|
||||
}
|
||||
if len(parquetReadErrors) > 0 {
|
||||
plan.Details["error_parquet_statistics"] = parquetReadErrors
|
||||
}
|
||||
if len(liveLogListErrors) > 0 {
|
||||
plan.Details["error_live_log_listing"] = liveLogListErrors
|
||||
}
|
||||
|
||||
// Update scan statistics for execution plan display
|
||||
plan.PartitionsScanned = len(partitions)
|
||||
plan.ParquetFilesScanned = len(parquetFiles)
|
||||
plan.LiveLogFilesScanned = len(liveLogFiles)
|
||||
} else {
|
||||
// Handle partition discovery error
|
||||
plan.Details["error_partition_discovery"] = discoverErr.Error()
|
||||
}
|
||||
} else {
|
||||
// Normal mode - just get results
|
||||
results, err = hybridScanner.Scan(ctx, hybridScanOptions)
|
||||
if err != nil {
|
||||
return &QueryResult{Error: err}, err
|
||||
}
|
||||
}
|
||||
|
||||
// Convert to SQL result format
|
||||
if selectAll {
|
||||
if len(columns) > 0 {
|
||||
// SELECT *, specific_columns - include both auto-discovered and explicit columns
|
||||
return hybridScanner.ConvertToSQLResultWithMixedColumns(results, columns), nil
|
||||
} else {
|
||||
// SELECT * only - let converter determine all columns (excludes system columns)
|
||||
columns = nil
|
||||
return hybridScanner.ConvertToSQLResult(results, columns), nil
|
||||
}
|
||||
}
|
||||
|
||||
// Handle custom column expressions (including arithmetic)
|
||||
return e.ConvertToSQLResultWithExpressions(hybridScanner, results, stmt.SelectExprs), nil
|
||||
}
|
||||
|
||||
// extractTimeFilters extracts time range filters from WHERE clause for optimization
|
||||
// This allows push-down of time-based queries to improve scan performance
|
||||
// Returns (startTimeNs, stopTimeNs) where 0 means unbounded
|
||||
|
||||
@@ -1254,6 +1254,69 @@ func TestSQLEngine_LogBufferDeduplication_ServerRestartScenario(t *testing.T) {
|
||||
// prevent false positive duplicates across server restarts
|
||||
}
|
||||
|
||||
func TestBrokerClient_BinaryBufferStartFormat(t *testing.T) {
|
||||
// Test scenario: getBufferStartFromEntry should only support binary format
|
||||
// This tests the standardized binary format for buffer_start metadata
|
||||
realBrokerClient := &BrokerClient{}
|
||||
|
||||
// Test binary format (used by both log files and Parquet files)
|
||||
binaryEntry := &filer_pb.Entry{
|
||||
Name: "2025-01-07-14-30-45",
|
||||
IsDirectory: false,
|
||||
Extended: map[string][]byte{
|
||||
"buffer_start": func() []byte {
|
||||
// Binary format: 8-byte BigEndian
|
||||
buf := make([]byte, 8)
|
||||
binary.BigEndian.PutUint64(buf, uint64(2000001))
|
||||
return buf
|
||||
}(),
|
||||
},
|
||||
}
|
||||
|
||||
bufferStart := realBrokerClient.getBufferStartFromEntry(binaryEntry)
|
||||
assert.NotNil(t, bufferStart)
|
||||
assert.Equal(t, int64(2000001), bufferStart.StartIndex, "Should parse binary buffer_start metadata")
|
||||
|
||||
// Test Parquet file (same binary format)
|
||||
parquetEntry := &filer_pb.Entry{
|
||||
Name: "2025-01-07-14-30.parquet",
|
||||
IsDirectory: false,
|
||||
Extended: map[string][]byte{
|
||||
"buffer_start": func() []byte {
|
||||
buf := make([]byte, 8)
|
||||
binary.BigEndian.PutUint64(buf, uint64(1500001))
|
||||
return buf
|
||||
}(),
|
||||
},
|
||||
}
|
||||
|
||||
bufferStart = realBrokerClient.getBufferStartFromEntry(parquetEntry)
|
||||
assert.NotNil(t, bufferStart)
|
||||
assert.Equal(t, int64(1500001), bufferStart.StartIndex, "Should parse binary buffer_start from Parquet file")
|
||||
|
||||
// Test missing metadata
|
||||
emptyEntry := &filer_pb.Entry{
|
||||
Name: "no-metadata",
|
||||
IsDirectory: false,
|
||||
Extended: nil,
|
||||
}
|
||||
|
||||
bufferStart = realBrokerClient.getBufferStartFromEntry(emptyEntry)
|
||||
assert.Nil(t, bufferStart, "Should return nil for entry without buffer_start metadata")
|
||||
|
||||
// Test invalid format (wrong size)
|
||||
invalidEntry := &filer_pb.Entry{
|
||||
Name: "invalid-metadata",
|
||||
IsDirectory: false,
|
||||
Extended: map[string][]byte{
|
||||
"buffer_start": []byte("invalid"),
|
||||
},
|
||||
}
|
||||
|
||||
bufferStart = realBrokerClient.getBufferStartFromEntry(invalidEntry)
|
||||
assert.Nil(t, bufferStart, "Should return nil for invalid buffer_start metadata")
|
||||
}
|
||||
|
||||
// TestGetSQLValAlias tests the getSQLValAlias function, particularly for SQL injection prevention
|
||||
func TestGetSQLValAlias(t *testing.T) {
|
||||
engine := &SQLEngine{}
|
||||
|
||||
@@ -320,6 +320,12 @@ func (hms *HybridMessageScanner) ScanWithStats(ctx context.Context, options Hybr
|
||||
return results, stats, nil
|
||||
}
|
||||
|
||||
// scanUnflushedData queries brokers for unflushed in-memory data using buffer_start deduplication
|
||||
func (hms *HybridMessageScanner) scanUnflushedData(ctx context.Context, partition topic.Partition, options HybridScanOptions) ([]HybridScanResult, error) {
|
||||
results, _, err := hms.scanUnflushedDataWithStats(ctx, partition, options)
|
||||
return results, err
|
||||
}
|
||||
|
||||
// scanUnflushedDataWithStats queries brokers for unflushed data and returns statistics
|
||||
func (hms *HybridMessageScanner) scanUnflushedDataWithStats(ctx context.Context, partition topic.Partition, options HybridScanOptions) ([]HybridScanResult, *HybridScanStats, error) {
|
||||
var results []HybridScanResult
|
||||
@@ -430,6 +436,27 @@ func (hms *HybridMessageScanner) scanUnflushedDataWithStats(ctx context.Context,
|
||||
return results, stats, nil
|
||||
}
|
||||
|
||||
// convertDataMessageToRecord converts mq_pb.DataMessage to schema_pb.RecordValue
|
||||
func (hms *HybridMessageScanner) convertDataMessageToRecord(msg *mq_pb.DataMessage) (*schema_pb.RecordValue, string, error) {
|
||||
// Parse the message data as RecordValue
|
||||
recordValue := &schema_pb.RecordValue{}
|
||||
if err := proto.Unmarshal(msg.Value, recordValue); err != nil {
|
||||
return nil, "", fmt.Errorf("failed to unmarshal message data: %v", err)
|
||||
}
|
||||
|
||||
// Add system columns
|
||||
if recordValue.Fields == nil {
|
||||
recordValue.Fields = make(map[string]*schema_pb.Value)
|
||||
}
|
||||
|
||||
// Add timestamp
|
||||
recordValue.Fields[SW_COLUMN_NAME_TIMESTAMP] = &schema_pb.Value{
|
||||
Kind: &schema_pb.Value_Int64Value{Int64Value: msg.TsNs},
|
||||
}
|
||||
|
||||
return recordValue, string(msg.Key), nil
|
||||
}
|
||||
|
||||
// discoverTopicPartitions discovers the actual partitions for this topic by scanning the filesystem
|
||||
// This finds real partition directories like v2025-09-01-07-16-34/0000-0630/
|
||||
func (hms *HybridMessageScanner) discoverTopicPartitions(ctx context.Context) ([]topic.Partition, error) {
|
||||
@@ -494,6 +521,15 @@ func (hms *HybridMessageScanner) discoverTopicPartitions(ctx context.Context) ([
|
||||
return allPartitions, nil
|
||||
}
|
||||
|
||||
// scanPartitionHybrid scans a specific partition using the hybrid approach
|
||||
// This is where the magic happens - seamlessly reading ALL data sources:
|
||||
// 1. Unflushed in-memory data from brokers (REAL-TIME)
|
||||
// 2. Live logs + Parquet files from disk (FLUSHED/ARCHIVED)
|
||||
func (hms *HybridMessageScanner) scanPartitionHybrid(ctx context.Context, partition topic.Partition, options HybridScanOptions) ([]HybridScanResult, error) {
|
||||
results, _, err := hms.scanPartitionHybridWithStats(ctx, partition, options)
|
||||
return results, err
|
||||
}
|
||||
|
||||
// scanPartitionHybridWithStats scans a specific partition using streaming merge for memory efficiency
|
||||
// PERFORMANCE IMPROVEMENT: Uses heap-based streaming merge instead of collecting all data and sorting
|
||||
// - Memory usage: O(k) where k = number of data sources, instead of O(n) where n = total records
|
||||
@@ -611,6 +647,23 @@ func (hms *HybridMessageScanner) countLiveLogFiles(partition topic.Partition) (i
|
||||
return fileCount, nil
|
||||
}
|
||||
|
||||
// isControlEntry checks if a log entry is a control entry without actual data
|
||||
// Based on MQ system analysis, control entries are:
|
||||
// 1. DataMessages with populated Ctrl field (publisher close signals)
|
||||
// 2. Entries with empty keys (as filtered by subscriber)
|
||||
// NOTE: Messages with empty data but valid keys (like NOOP messages) are NOT control entries
|
||||
func (hms *HybridMessageScanner) isControlEntry(logEntry *filer_pb.LogEntry) bool {
|
||||
// Pre-decode DataMessage if needed
|
||||
var dataMessage *mq_pb.DataMessage
|
||||
if len(logEntry.Data) > 0 {
|
||||
dataMessage = &mq_pb.DataMessage{}
|
||||
if err := proto.Unmarshal(logEntry.Data, dataMessage); err != nil {
|
||||
dataMessage = nil // Failed to decode, treat as raw data
|
||||
}
|
||||
}
|
||||
return hms.isControlEntryWithDecoded(logEntry, dataMessage)
|
||||
}
|
||||
|
||||
// isControlEntryWithDecoded checks if a log entry is a control entry using pre-decoded DataMessage
|
||||
// This avoids duplicate protobuf unmarshaling when the DataMessage is already decoded
|
||||
func (hms *HybridMessageScanner) isControlEntryWithDecoded(logEntry *filer_pb.LogEntry, dataMessage *mq_pb.DataMessage) bool {
|
||||
@@ -629,6 +682,26 @@ func (hms *HybridMessageScanner) isControlEntryWithDecoded(logEntry *filer_pb.Lo
|
||||
return false
|
||||
}
|
||||
|
||||
// isNullOrEmpty checks if a schema_pb.Value is null or empty
|
||||
func isNullOrEmpty(value *schema_pb.Value) bool {
|
||||
if value == nil {
|
||||
return true
|
||||
}
|
||||
|
||||
switch v := value.Kind.(type) {
|
||||
case *schema_pb.Value_StringValue:
|
||||
return v.StringValue == ""
|
||||
case *schema_pb.Value_BytesValue:
|
||||
return len(v.BytesValue) == 0
|
||||
case *schema_pb.Value_ListValue:
|
||||
return v.ListValue == nil || len(v.ListValue.Values) == 0
|
||||
case nil:
|
||||
return true // No kind set means null
|
||||
default:
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
// isSchemaless checks if the scanner is configured for a schema-less topic
|
||||
// Schema-less topics only have system fields: _ts_ns, _key, and _value
|
||||
func (hms *HybridMessageScanner) isSchemaless() bool {
|
||||
@@ -663,6 +736,61 @@ func (hms *HybridMessageScanner) isSchemaless() bool {
|
||||
return hasValue && dataFieldCount == 1
|
||||
}
|
||||
|
||||
// convertLogEntryToRecordValue converts a filer_pb.LogEntry to schema_pb.RecordValue
|
||||
// This handles both:
|
||||
// 1. Live log entries (raw message format)
|
||||
// 2. Parquet entries (already in schema_pb.RecordValue format)
|
||||
// 3. Schema-less topics (raw bytes in _value field)
|
||||
func (hms *HybridMessageScanner) convertLogEntryToRecordValue(logEntry *filer_pb.LogEntry) (*schema_pb.RecordValue, string, error) {
|
||||
// For schema-less topics, put raw data directly into _value field
|
||||
if hms.isSchemaless() {
|
||||
recordValue := &schema_pb.RecordValue{
|
||||
Fields: make(map[string]*schema_pb.Value),
|
||||
}
|
||||
recordValue.Fields[SW_COLUMN_NAME_TIMESTAMP] = &schema_pb.Value{
|
||||
Kind: &schema_pb.Value_Int64Value{Int64Value: logEntry.TsNs},
|
||||
}
|
||||
recordValue.Fields[SW_COLUMN_NAME_KEY] = &schema_pb.Value{
|
||||
Kind: &schema_pb.Value_BytesValue{BytesValue: logEntry.Key},
|
||||
}
|
||||
recordValue.Fields[SW_COLUMN_NAME_VALUE] = &schema_pb.Value{
|
||||
Kind: &schema_pb.Value_BytesValue{BytesValue: logEntry.Data},
|
||||
}
|
||||
return recordValue, "live_log", nil
|
||||
}
|
||||
|
||||
// Try to unmarshal as RecordValue first (Parquet format)
|
||||
recordValue := &schema_pb.RecordValue{}
|
||||
if err := proto.Unmarshal(logEntry.Data, recordValue); err == nil {
|
||||
// This is an archived message from Parquet files
|
||||
// FIX: Add system columns from LogEntry to RecordValue
|
||||
if recordValue.Fields == nil {
|
||||
recordValue.Fields = make(map[string]*schema_pb.Value)
|
||||
}
|
||||
|
||||
// Add system columns from LogEntry
|
||||
recordValue.Fields[SW_COLUMN_NAME_TIMESTAMP] = &schema_pb.Value{
|
||||
Kind: &schema_pb.Value_Int64Value{Int64Value: logEntry.TsNs},
|
||||
}
|
||||
recordValue.Fields[SW_COLUMN_NAME_KEY] = &schema_pb.Value{
|
||||
Kind: &schema_pb.Value_BytesValue{BytesValue: logEntry.Key},
|
||||
}
|
||||
|
||||
return recordValue, "parquet_archive", nil
|
||||
}
|
||||
|
||||
// If not a RecordValue, this is raw live message data - parse with schema
|
||||
return hms.parseRawMessageWithSchema(logEntry)
|
||||
}
|
||||
|
||||
// min returns the minimum of two integers
|
||||
func min(a, b int) int {
|
||||
if a < b {
|
||||
return a
|
||||
}
|
||||
return b
|
||||
}
|
||||
|
||||
// parseRawMessageWithSchema parses raw live message data using the topic's schema
|
||||
// This provides proper type conversion and field mapping instead of treating everything as strings
|
||||
func (hms *HybridMessageScanner) parseRawMessageWithSchema(logEntry *filer_pb.LogEntry) (*schema_pb.RecordValue, string, error) {
|
||||
|
||||
@@ -6,6 +6,8 @@ import (
|
||||
"math/big"
|
||||
"time"
|
||||
|
||||
"github.com/parquet-go/parquet-go"
|
||||
"github.com/seaweedfs/seaweedfs/weed/filer"
|
||||
"github.com/seaweedfs/seaweedfs/weed/mq/schema"
|
||||
"github.com/seaweedfs/seaweedfs/weed/mq/topic"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
@@ -170,6 +172,90 @@ func (ps *ParquetScanner) scanPartition(ctx context.Context, partition topic.Par
|
||||
return results, nil
|
||||
}
|
||||
|
||||
// scanParquetFile scans a single Parquet file (real implementation)
|
||||
func (ps *ParquetScanner) scanParquetFile(ctx context.Context, entry *filer_pb.Entry, options ScanOptions) ([]ScanResult, error) {
|
||||
var results []ScanResult
|
||||
|
||||
// Create reader for the Parquet file (same pattern as logstore)
|
||||
lookupFileIdFn := filer.LookupFn(ps.filerClient)
|
||||
fileSize := filer.FileSize(entry)
|
||||
visibleIntervals, _ := filer.NonOverlappingVisibleIntervals(ctx, lookupFileIdFn, entry.Chunks, 0, int64(fileSize))
|
||||
chunkViews := filer.ViewFromVisibleIntervals(visibleIntervals, 0, int64(fileSize))
|
||||
readerCache := filer.NewReaderCache(32, ps.chunkCache, lookupFileIdFn)
|
||||
readerAt := filer.NewChunkReaderAtFromClient(ctx, readerCache, chunkViews, int64(fileSize), filer.DefaultPrefetchCount)
|
||||
|
||||
// Create Parquet reader
|
||||
parquetReader := parquet.NewReader(readerAt)
|
||||
defer parquetReader.Close()
|
||||
|
||||
rows := make([]parquet.Row, 128) // Read in batches like logstore
|
||||
|
||||
for {
|
||||
rowCount, readErr := parquetReader.ReadRows(rows)
|
||||
|
||||
// Process rows even if EOF
|
||||
for i := 0; i < rowCount; i++ {
|
||||
// Convert Parquet row to schema value
|
||||
recordValue, err := schema.ToRecordValue(ps.recordSchema, ps.parquetLevels, rows[i])
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to convert row: %v", err)
|
||||
}
|
||||
|
||||
// Extract system columns
|
||||
timestamp := recordValue.Fields[SW_COLUMN_NAME_TIMESTAMP].GetInt64Value()
|
||||
key := recordValue.Fields[SW_COLUMN_NAME_KEY].GetBytesValue()
|
||||
|
||||
// Apply time filtering
|
||||
if options.StartTimeNs > 0 && timestamp < options.StartTimeNs {
|
||||
continue
|
||||
}
|
||||
if options.StopTimeNs > 0 && timestamp >= options.StopTimeNs {
|
||||
break // Assume data is time-ordered
|
||||
}
|
||||
|
||||
// Apply predicate filtering (WHERE clause)
|
||||
if options.Predicate != nil && !options.Predicate(recordValue) {
|
||||
continue
|
||||
}
|
||||
|
||||
// Apply column projection
|
||||
values := make(map[string]*schema_pb.Value)
|
||||
if len(options.Columns) == 0 {
|
||||
// Select all columns (excluding system columns from user view)
|
||||
for name, value := range recordValue.Fields {
|
||||
if name != SW_COLUMN_NAME_TIMESTAMP && name != SW_COLUMN_NAME_KEY {
|
||||
values[name] = value
|
||||
}
|
||||
}
|
||||
} else {
|
||||
// Select specified columns only
|
||||
for _, columnName := range options.Columns {
|
||||
if value, exists := recordValue.Fields[columnName]; exists {
|
||||
values[columnName] = value
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
results = append(results, ScanResult{
|
||||
Values: values,
|
||||
Timestamp: timestamp,
|
||||
Key: key,
|
||||
})
|
||||
|
||||
// Apply row limit
|
||||
if options.Limit > 0 && len(results) >= options.Limit {
|
||||
return results, nil
|
||||
}
|
||||
}
|
||||
|
||||
if readErr != nil {
|
||||
break // EOF or error
|
||||
}
|
||||
}
|
||||
|
||||
return results, nil
|
||||
}
|
||||
|
||||
// generateSampleData creates sample data for testing when no real Parquet files exist
|
||||
func (ps *ParquetScanner) generateSampleData(options ScanOptions) []ScanResult {
|
||||
now := time.Now().UnixNano()
|
||||
|
||||
@@ -50,6 +50,14 @@ func (e *SQLEngine) getSystemColumnDisplayName(columnName string) string {
|
||||
}
|
||||
}
|
||||
|
||||
// isSystemColumnDisplayName checks if a column name is a system column display name
|
||||
func (e *SQLEngine) isSystemColumnDisplayName(columnName string) bool {
|
||||
lowerName := strings.ToLower(columnName)
|
||||
return lowerName == SW_DISPLAY_NAME_TIMESTAMP ||
|
||||
lowerName == SW_COLUMN_NAME_KEY ||
|
||||
lowerName == SW_COLUMN_NAME_SOURCE
|
||||
}
|
||||
|
||||
// getSystemColumnInternalName returns the internal name for a system column display name
|
||||
func (e *SQLEngine) getSystemColumnInternalName(displayName string) string {
|
||||
lowerName := strings.ToLower(displayName)
|
||||
|
||||
@@ -241,14 +241,6 @@ 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
|
||||
@@ -290,11 +282,6 @@ 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{
|
||||
@@ -337,21 +324,11 @@ 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,7 +4,6 @@ import (
|
||||
"context"
|
||||
"fmt"
|
||||
"math"
|
||||
"sync"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/wdclient"
|
||||
@@ -21,19 +20,6 @@ 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
|
||||
@@ -49,7 +35,6 @@ type FilerSink struct {
|
||||
isIncremental bool
|
||||
executor *util.LimitedConcurrentExecutor
|
||||
signature int32
|
||||
activeTransfers sync.Map // chunkFileId -> *ChunkTransferStatus
|
||||
}
|
||||
|
||||
func init() {
|
||||
@@ -116,25 +101,6 @@ 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()
|
||||
|
||||
@@ -50,10 +50,6 @@ func (fs *FilerSource) DoInitialize(address, grpcAddress string, dir string, rea
|
||||
return nil
|
||||
}
|
||||
|
||||
func (fs *FilerSource) SetGrpcDialOption(option grpc.DialOption) {
|
||||
fs.grpcDialOption = option
|
||||
}
|
||||
|
||||
func (fs *FilerSource) LookupFileId(ctx context.Context, part string) (fileUrls []string, err error) {
|
||||
|
||||
vid2Locations := make(map[string]*filer_pb.Locations)
|
||||
|
||||
@@ -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 (never log key material)
|
||||
glog.V(4).Infof("SSE-C MD5 validation: provided='%s', expected='%s'", keyMD5, expectedMD5)
|
||||
// Debug logging for MD5 validation
|
||||
glog.V(4).Infof("SSE-C MD5 validation: provided='%s', expected='%s', keyBytes=%x", keyMD5, expectedMD5, keyBytes)
|
||||
|
||||
if keyMD5 != expectedMD5 {
|
||||
glog.Errorf("SSE-C MD5 mismatch: provided='%s', expected='%s'", keyMD5, expectedMD5)
|
||||
|
||||
@@ -7,13 +7,10 @@ import (
|
||||
"fmt"
|
||||
"net"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"slices"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/spf13/viper"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/glog"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
"google.golang.org/grpc"
|
||||
@@ -142,23 +139,6 @@ func LoadServerTLS(config *util.ViperProxy, component string) (grpc.ServerOption
|
||||
return grpc.Creds(ta), nil
|
||||
}
|
||||
|
||||
func LoadClientTLSFromFile(configFile string, component string) (grpc.DialOption, error) {
|
||||
v := viper.New()
|
||||
v.SetConfigFile(configFile)
|
||||
if err := v.ReadInConfig(); err != nil {
|
||||
return nil, fmt.Errorf("failed to read security config %s: %v", configFile, err)
|
||||
}
|
||||
// Resolve relative PEM paths against the config file's directory.
|
||||
configDir := filepath.Dir(configFile)
|
||||
for _, key := range []string{"grpc.ca", component + ".cert", component + ".key"} {
|
||||
p := v.GetString(key)
|
||||
if p != "" && !filepath.IsAbs(p) {
|
||||
v.Set(key, filepath.Join(configDir, p))
|
||||
}
|
||||
}
|
||||
return LoadClientTLS(&util.ViperProxy{Viper: v}, component), nil
|
||||
}
|
||||
|
||||
func LoadClientTLS(config *util.ViperProxy, component string) grpc.DialOption {
|
||||
if config == nil {
|
||||
return grpc.WithTransportCredentials(insecure.NewCredentials())
|
||||
|
||||
@@ -2,12 +2,14 @@ package weed_server
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
|
||||
"github.com/seaweedfs/raft"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/operation"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/master_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/volume_server_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/topology"
|
||||
)
|
||||
|
||||
func (ms *MasterServer) CollectionList(ctx context.Context, req *master_pb.CollectionListRequest) (*master_pb.CollectionListResponse, error) {
|
||||
@@ -35,6 +37,10 @@ func (ms *MasterServer) CollectionDelete(ctx context.Context, req *master_pb.Col
|
||||
|
||||
resp := &master_pb.CollectionDeleteResponse{}
|
||||
|
||||
if err := ms.ensureCollectionDeleteSafe(req.Name); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
err := ms.doDeleteNormalCollection(req.Name)
|
||||
|
||||
if err != nil {
|
||||
@@ -50,6 +56,19 @@ func (ms *MasterServer) CollectionDelete(ctx context.Context, req *master_pb.Col
|
||||
return resp, nil
|
||||
}
|
||||
|
||||
func (ms *MasterServer) ensureCollectionDeleteSafe(collectionName string) error {
|
||||
duplicates := ms.Topo.FindDuplicateVolumeIds(collectionName)
|
||||
if len(duplicates) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
return fmt.Errorf(
|
||||
"refusing to delete collection %q: duplicate volume IDs exist across collections (%s). Run `volume.check.duplicates` to inspect and resolve them first",
|
||||
collectionName,
|
||||
topology.FormatDuplicateVolumeIds(duplicates),
|
||||
)
|
||||
}
|
||||
|
||||
func (ms *MasterServer) doDeleteNormalCollection(collectionName string) error {
|
||||
|
||||
collection, ok := ms.Topo.FindCollection(collectionName)
|
||||
|
||||
@@ -4,6 +4,11 @@ import (
|
||||
"testing"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/cluster"
|
||||
"github.com/seaweedfs/seaweedfs/weed/sequence"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/needle"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/super_block"
|
||||
"github.com/seaweedfs/seaweedfs/weed/topology"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
@@ -35,3 +40,41 @@ func TestInitialLockRingUpdateSkipsNonFilers(t *testing.T) {
|
||||
|
||||
assert.Nil(t, ms.initialLockRingUpdate(cluster.BrokerType, "group-a"))
|
||||
}
|
||||
|
||||
func TestEnsureCollectionDeleteSafeRejectsDuplicateVolumeIds(t *testing.T) {
|
||||
topo := topology.NewTopology("weedfs", sequence.NewMemorySequencer(), 32*1024, 5, false)
|
||||
dc := topo.GetOrCreateDataCenter("dc1")
|
||||
rack := dc.GetOrCreateRack("rack1")
|
||||
maxVolumeCounts := map[string]uint32{"": 25}
|
||||
|
||||
dn1 := rack.GetOrCreateDataNode("127.0.0.1", 34534, 0, "127.0.0.1", "", maxVolumeCounts)
|
||||
dn2 := rack.GetOrCreateDataNode("127.0.0.2", 34535, 0, "127.0.0.2", "", maxVolumeCounts)
|
||||
|
||||
replicaPlacement := &super_block.ReplicaPlacement{}
|
||||
collectionA := storage.VolumeInfo{
|
||||
Id: needle.VolumeId(100),
|
||||
Collection: "collection-a",
|
||||
Version: needle.GetCurrentVersion(),
|
||||
ReplicaPlacement: replicaPlacement,
|
||||
Ttl: needle.EMPTY_TTL,
|
||||
}
|
||||
collectionB := storage.VolumeInfo{
|
||||
Id: needle.VolumeId(100),
|
||||
Collection: "collection-b",
|
||||
Version: needle.GetCurrentVersion(),
|
||||
ReplicaPlacement: replicaPlacement,
|
||||
Ttl: needle.EMPTY_TTL,
|
||||
}
|
||||
|
||||
dn1.UpdateVolumes([]storage.VolumeInfo{collectionA})
|
||||
dn2.UpdateVolumes([]storage.VolumeInfo{collectionB})
|
||||
topo.RegisterVolumeLayout(collectionA, dn1)
|
||||
topo.RegisterVolumeLayout(collectionB, dn2)
|
||||
|
||||
ms := &MasterServer{Topo: topo}
|
||||
|
||||
err := ms.ensureCollectionDeleteSafe("collection-a")
|
||||
require.Error(t, err)
|
||||
assert.Contains(t, err.Error(), "refusing to delete collection")
|
||||
assert.Contains(t, err.Error(), "100:[collection-a, collection-b]")
|
||||
}
|
||||
|
||||
@@ -30,6 +30,10 @@ func (ms *MasterServer) collectionDeleteHandler(w http.ResponseWriter, r *http.R
|
||||
writeJsonError(w, r, http.StatusBadRequest, fmt.Errorf("collection %s does not exist", collectionName))
|
||||
return
|
||||
}
|
||||
if err := ms.ensureCollectionDeleteSafe(collectionName); err != nil {
|
||||
writeJsonError(w, r, http.StatusConflict, err)
|
||||
return
|
||||
}
|
||||
for _, server := range collection.ListVolumeServers() {
|
||||
err := operation.WithVolumeServerClient(false, server.ServerAddress(), ms.grpcDialOption, func(client volume_server_pb.VolumeServerClient) error {
|
||||
_, deleteErr := client.DeleteCollection(context.Background(), &volume_server_pb.DeleteCollectionRequest{
|
||||
|
||||
@@ -609,17 +609,6 @@ func (vs *VolumeServer) VolumeEcShardsToVolume(ctx context.Context, req *volume_
|
||||
}
|
||||
|
||||
dataBaseFileName, indexBaseFileName := v.DataBaseFileName(), v.IndexBaseFileName()
|
||||
if !util.FileExists(indexBaseFileName + ".ecx") {
|
||||
indexBaseFileName = dataBaseFileName
|
||||
}
|
||||
|
||||
// Merge .ecj deletions into .ecx so that HasLiveNeedles and FindDatFileSize
|
||||
// see the full set of deleted needles. Without this, needles deleted after the
|
||||
// last ecx rebuild would still appear live, causing the decoded .dat to include
|
||||
// data that should be skipped and HasLiveNeedles to return a false positive.
|
||||
if err := erasure_coding.RebuildEcxFile(indexBaseFileName); err != nil {
|
||||
return nil, fmt.Errorf("RebuildEcxFile %s: %v", indexBaseFileName, err)
|
||||
}
|
||||
|
||||
// If the EC index contains no live entries, decoding should be a no-op:
|
||||
// just allow the caller to purge EC shards and do not generate an empty normal volume.
|
||||
@@ -647,29 +636,6 @@ func (vs *VolumeServer) VolumeEcShardsToVolume(ctx context.Context, req *volume_
|
||||
return nil, fmt.Errorf("WriteIdxFileFromEcIndex %s: %v", v.IndexBaseFileName(), err)
|
||||
}
|
||||
|
||||
var volumeLocation *storage.DiskLocation
|
||||
for _, location := range vs.store.Locations {
|
||||
if candidate, found := location.FindEcVolume(needle.VolumeId(req.VolumeId)); found && candidate == v {
|
||||
volumeLocation = location
|
||||
break
|
||||
}
|
||||
}
|
||||
if volumeLocation == nil {
|
||||
return nil, fmt.Errorf("ec volume %d location not found for offline compaction", req.VolumeId)
|
||||
}
|
||||
|
||||
if err := vs.store.CompactVolumeFiles(
|
||||
needle.VolumeId(req.VolumeId),
|
||||
v.Collection,
|
||||
volumeLocation,
|
||||
vs.needleMapKind,
|
||||
vs.ldbTimout,
|
||||
0,
|
||||
vs.compactionBytePerSecond,
|
||||
); err != nil {
|
||||
glog.Errorf("CompactVolumeFiles %s: %v", dataBaseFileName, err)
|
||||
}
|
||||
|
||||
return &volume_server_pb.VolumeEcShardsToVolumeResponse{}, nil
|
||||
}
|
||||
|
||||
|
||||
@@ -169,14 +169,8 @@ func doEcDecode(commandEnv *CommandEnv, topoInfo *master_pb.TopologyInfo, collec
|
||||
return fmt.Errorf("generate normal volume %d on %s: %v", vid, targetNodeLocation, err)
|
||||
}
|
||||
|
||||
// mount the decoded volume after server-side offline compaction succeeded
|
||||
err = mountDecodedVolume(commandEnv.option.GrpcDialOption, targetNodeLocation, vid)
|
||||
if err != nil {
|
||||
return fmt.Errorf("mount decoded volume %d on %s: %v", vid, targetNodeLocation, err)
|
||||
}
|
||||
|
||||
// delete the previous ec shards
|
||||
err = unmountAndDeleteEcShardsWithPrefix("deleteDecodedEcShards", commandEnv.option.GrpcDialOption, collection, nodeToEcShardsInfo, vid)
|
||||
err = mountVolumeAndDeleteEcShards(commandEnv.option.GrpcDialOption, collection, targetNodeLocation, nodeToEcShardsInfo, vid)
|
||||
if err != nil {
|
||||
return fmt.Errorf("delete ec shards for volume %d: %v", vid, err)
|
||||
}
|
||||
@@ -225,16 +219,23 @@ func unmountAndDeleteEcShardsWithPrefix(prefix string, grpcDialOption grpc.DialO
|
||||
return ewg.Wait()
|
||||
}
|
||||
|
||||
func mountDecodedVolume(grpcDialOption grpc.DialOption, targetNodeLocation pb.ServerAddress, vid needle.VolumeId) error {
|
||||
return operation.WithVolumeServerClient(false, targetNodeLocation, grpcDialOption, func(volumeServerClient volume_server_pb.VolumeServerClient) error {
|
||||
func mountVolumeAndDeleteEcShards(grpcDialOption grpc.DialOption, collection string, targetNodeLocation pb.ServerAddress, nodeToShardsInfo map[pb.ServerAddress]*erasure_coding.ShardsInfo, vid needle.VolumeId) error {
|
||||
|
||||
// mount volume
|
||||
if err := operation.WithVolumeServerClient(false, targetNodeLocation, grpcDialOption, func(volumeServerClient volume_server_pb.VolumeServerClient) error {
|
||||
_, mountErr := volumeServerClient.VolumeMount(context.Background(), &volume_server_pb.VolumeMountRequest{
|
||||
VolumeId: uint32(vid),
|
||||
})
|
||||
return mountErr
|
||||
})
|
||||
}); err != nil {
|
||||
return fmt.Errorf("mountVolumeAndDeleteEcShards mount volume %d on %s: %v", vid, targetNodeLocation, err)
|
||||
}
|
||||
|
||||
return unmountAndDeleteEcShardsWithPrefix("mountVolumeAndDeleteEcShards", grpcDialOption, collection, nodeToShardsInfo, vid)
|
||||
}
|
||||
|
||||
func generateNormalVolume(grpcDialOption grpc.DialOption, vid needle.VolumeId, collection string, sourceVolumeServer pb.ServerAddress) error {
|
||||
|
||||
fmt.Printf("generateNormalVolume from ec volume %d on %s\n", vid, sourceVolumeServer)
|
||||
|
||||
err := operation.WithVolumeServerClient(false, sourceVolumeServer, grpcDialOption, func(volumeServerClient volume_server_pb.VolumeServerClient) error {
|
||||
@@ -279,18 +280,13 @@ func collectEcShards(commandEnv *CommandEnv, nodeToShardsInfo map[pb.ServerAddre
|
||||
}
|
||||
|
||||
needToCopyShardsInfo := si.Minus(existingShardsInfo).MinusParityShards()
|
||||
if needToCopyShardsInfo.Count() == 0 {
|
||||
continue
|
||||
}
|
||||
|
||||
err = operation.WithVolumeServerClient(false, targetNodeLocation, commandEnv.option.GrpcDialOption, func(volumeServerClient volume_server_pb.VolumeServerClient) error {
|
||||
|
||||
// Always collect .ecj from every shard location. Each server's .ecj
|
||||
// only contains deletions for needles whose data resides in shards
|
||||
// held by that server. Without merging all .ecj files, deletions
|
||||
// recorded on other servers would be lost during decode.
|
||||
if needToCopyShardsInfo.Count() > 0 {
|
||||
fmt.Printf("copy %d.%v %s => %s\n", vid, needToCopyShardsInfo.Ids(), loc, targetNodeLocation)
|
||||
} else {
|
||||
fmt.Printf("collect ecj %d %s => %s\n", vid, loc, targetNodeLocation)
|
||||
}
|
||||
fmt.Printf("copy %d.%v %s => %s\n", vid, needToCopyShardsInfo.Ids(), loc, targetNodeLocation)
|
||||
|
||||
_, copyErr := volumeServerClient.VolumeEcShardsCopy(context.Background(), &volume_server_pb.VolumeEcShardsCopyRequest{
|
||||
VolumeId: uint32(vid),
|
||||
@@ -298,23 +294,21 @@ func collectEcShards(commandEnv *CommandEnv, nodeToShardsInfo map[pb.ServerAddre
|
||||
ShardIds: needToCopyShardsInfo.IdsUint32(),
|
||||
CopyEcxFile: false,
|
||||
CopyEcjFile: true,
|
||||
CopyVifFile: needToCopyShardsInfo.Count() > 0,
|
||||
CopyVifFile: true,
|
||||
SourceDataNode: string(loc),
|
||||
})
|
||||
if copyErr != nil {
|
||||
return fmt.Errorf("copy %d.%v %s => %s : %v\n", vid, needToCopyShardsInfo.Ids(), loc, targetNodeLocation, copyErr)
|
||||
}
|
||||
|
||||
if needToCopyShardsInfo.Count() > 0 {
|
||||
fmt.Printf("mount %d.%v on %s\n", vid, needToCopyShardsInfo.Ids(), targetNodeLocation)
|
||||
_, mountErr := volumeServerClient.VolumeEcShardsMount(context.Background(), &volume_server_pb.VolumeEcShardsMountRequest{
|
||||
VolumeId: uint32(vid),
|
||||
Collection: collection,
|
||||
ShardIds: needToCopyShardsInfo.IdsUint32(),
|
||||
})
|
||||
if mountErr != nil {
|
||||
return fmt.Errorf("mount %d.%v on %s : %v\n", vid, needToCopyShardsInfo.Ids(), targetNodeLocation, mountErr)
|
||||
}
|
||||
fmt.Printf("mount %d.%v on %s\n", vid, needToCopyShardsInfo.Ids(), targetNodeLocation)
|
||||
_, mountErr := volumeServerClient.VolumeEcShardsMount(context.Background(), &volume_server_pb.VolumeEcShardsMountRequest{
|
||||
VolumeId: uint32(vid),
|
||||
Collection: collection,
|
||||
ShardIds: needToCopyShardsInfo.IdsUint32(),
|
||||
})
|
||||
if mountErr != nil {
|
||||
return fmt.Errorf("mount %d.%v on %s : %v\n", vid, needToCopyShardsInfo.Ids(), targetNodeLocation, mountErr)
|
||||
}
|
||||
|
||||
return nil
|
||||
|
||||
@@ -20,6 +20,7 @@ import (
|
||||
"github.com/seaweedfs/seaweedfs/weed/operation"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/master_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/topology"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
util_http "github.com/seaweedfs/seaweedfs/weed/util/http"
|
||||
)
|
||||
@@ -70,7 +71,13 @@ func (c *commandFsMergeVolumes) Do(args []string, commandEnv *CommandEnv, writer
|
||||
fromVolumeId := needle.VolumeId(*fromVolumeArg)
|
||||
toVolumeId := needle.VolumeId(*toVolumeArg)
|
||||
|
||||
c.reloadVolumesInfo(commandEnv.MasterClient)
|
||||
if err := c.ensureNoAmbiguousDuplicateVolumeIds(commandEnv, *collectionArg, fromVolumeId, toVolumeId); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err = c.reloadVolumesInfo(commandEnv.MasterClient, *collectionArg); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if fromVolumeId != 0 && toVolumeId != 0 {
|
||||
if fromVolumeId == toVolumeId {
|
||||
@@ -146,6 +153,98 @@ func (c *commandFsMergeVolumes) Do(args []string, commandEnv *CommandEnv, writer
|
||||
})
|
||||
}
|
||||
|
||||
func (c *commandFsMergeVolumes) ensureNoAmbiguousDuplicateVolumeIds(commandEnv *CommandEnv, collection string, fromVolumeId, toVolumeId needle.VolumeId) error {
|
||||
duplicateInfos, err := c.findDuplicateVolumeIds(commandEnv)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
ambiguous := filterAmbiguousDuplicateVolumeIds(duplicateInfos, collection, fromVolumeId, toVolumeId)
|
||||
|
||||
if len(ambiguous) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
return fmt.Errorf(
|
||||
"cannot run fs.mergeVolumes with duplicate volume IDs across collections (%s). Resolve the duplicates first",
|
||||
topology.FormatDuplicateVolumeIds(ambiguous),
|
||||
)
|
||||
}
|
||||
|
||||
func filterAmbiguousDuplicateVolumeIds(duplicateInfos map[needle.VolumeId][]string, collection string, fromVolumeId, toVolumeId needle.VolumeId) map[needle.VolumeId][]string {
|
||||
ambiguous := make(map[needle.VolumeId][]string)
|
||||
for vid, collections := range duplicateInfos {
|
||||
if fromVolumeId != 0 && toVolumeId != 0 && vid != fromVolumeId && vid != toVolumeId {
|
||||
continue
|
||||
}
|
||||
if fromVolumeId != 0 && toVolumeId == 0 && vid != fromVolumeId {
|
||||
continue
|
||||
}
|
||||
if toVolumeId != 0 && fromVolumeId == 0 && vid != toVolumeId {
|
||||
continue
|
||||
}
|
||||
if collection != "*" {
|
||||
matchesCollection := false
|
||||
for _, name := range collections {
|
||||
if name == collection {
|
||||
matchesCollection = true
|
||||
break
|
||||
}
|
||||
}
|
||||
if !matchesCollection {
|
||||
continue
|
||||
}
|
||||
}
|
||||
ambiguous[vid] = collections
|
||||
}
|
||||
return ambiguous
|
||||
}
|
||||
|
||||
func (c *commandFsMergeVolumes) findDuplicateVolumeIds(commandEnv *CommandEnv) (map[needle.VolumeId][]string, error) {
|
||||
volumeCollections := make(map[needle.VolumeId]map[string]struct{})
|
||||
|
||||
err := commandEnv.MasterClient.WithClient(false, func(client master_pb.SeaweedClient) error {
|
||||
volumes, err := client.VolumeList(context.Background(), &master_pb.VolumeListRequest{})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
for _, dc := range volumes.TopologyInfo.DataCenterInfos {
|
||||
for _, rack := range dc.RackInfos {
|
||||
for _, node := range rack.DataNodeInfos {
|
||||
for _, disk := range node.DiskInfos {
|
||||
for _, volume := range disk.VolumeInfos {
|
||||
vid := needle.VolumeId(volume.Id)
|
||||
if volumeCollections[vid] == nil {
|
||||
volumeCollections[vid] = make(map[string]struct{})
|
||||
}
|
||||
volumeCollections[vid][volume.Collection] = struct{}{}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
duplicates := make(map[needle.VolumeId][]string)
|
||||
for vid, collectionSet := range volumeCollections {
|
||||
if len(collectionSet) < 2 {
|
||||
continue
|
||||
}
|
||||
collections := make([]string, 0, len(collectionSet))
|
||||
for name := range collectionSet {
|
||||
collections = append(collections, name)
|
||||
}
|
||||
sort.Strings(collections)
|
||||
duplicates[vid] = collections
|
||||
}
|
||||
return duplicates, nil
|
||||
}
|
||||
|
||||
func (c *commandFsMergeVolumes) getVolumeInfoById(vid needle.VolumeId) (*master_pb.VolumeInformationMessage, error) {
|
||||
info := c.volumes[vid]
|
||||
var err error
|
||||
@@ -169,7 +268,7 @@ func (c *commandFsMergeVolumes) volumesAreCompatible(src needle.VolumeId, dest n
|
||||
srcInfo.ReplicaPlacement == destInfo.ReplicaPlacement), nil
|
||||
}
|
||||
|
||||
func (c *commandFsMergeVolumes) reloadVolumesInfo(masterClient *wdclient.MasterClient) error {
|
||||
func (c *commandFsMergeVolumes) reloadVolumesInfo(masterClient *wdclient.MasterClient, collection string) error {
|
||||
c.volumes = make(map[needle.VolumeId]*master_pb.VolumeInformationMessage)
|
||||
|
||||
return masterClient.WithClient(false, func(client master_pb.SeaweedClient) error {
|
||||
@@ -186,8 +285,24 @@ func (c *commandFsMergeVolumes) reloadVolumesInfo(masterClient *wdclient.MasterC
|
||||
for _, disk := range node.DiskInfos {
|
||||
for _, volume := range disk.VolumeInfos {
|
||||
vid := needle.VolumeId(volume.Id)
|
||||
if found := c.volumes[vid]; found == nil {
|
||||
existing := c.volumes[vid]
|
||||
if existing == nil {
|
||||
c.volumes[vid] = volume
|
||||
} else if existing.Collection != volume.Collection {
|
||||
// Duplicate volume ID across collections detected.
|
||||
// Prefer the volume matching the requested collection filter.
|
||||
if collection != "*" && volume.Collection == collection {
|
||||
fmt.Printf("WARNING: duplicate volume id %d across collections %q and %q, using %q\n",
|
||||
vid, existing.Collection, volume.Collection, collection)
|
||||
c.volumes[vid] = volume
|
||||
} else if collection != "*" && existing.Collection == collection {
|
||||
// Already have the right one
|
||||
fmt.Printf("WARNING: duplicate volume id %d across collections %q and %q, using %q\n",
|
||||
vid, existing.Collection, volume.Collection, collection)
|
||||
} else {
|
||||
fmt.Printf("WARNING: duplicate volume id %d across collections %q and %q\n",
|
||||
vid, existing.Collection, volume.Collection)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,72 @@
|
||||
package shell
|
||||
|
||||
import (
|
||||
"reflect"
|
||||
"testing"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/needle"
|
||||
)
|
||||
|
||||
func TestFilterAmbiguousDuplicateVolumeIds(t *testing.T) {
|
||||
duplicates := map[needle.VolumeId][]string{
|
||||
needle.VolumeId(10): {"alpha", "beta"},
|
||||
needle.VolumeId(20): {"beta", "gamma"},
|
||||
}
|
||||
|
||||
tests := []struct {
|
||||
name string
|
||||
collection string
|
||||
from needle.VolumeId
|
||||
to needle.VolumeId
|
||||
expected map[needle.VolumeId][]string
|
||||
}{
|
||||
{
|
||||
name: "all collections include all duplicates",
|
||||
collection: "*",
|
||||
expected: map[needle.VolumeId][]string{
|
||||
needle.VolumeId(10): {"alpha", "beta"},
|
||||
needle.VolumeId(20): {"beta", "gamma"},
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "collection scope keeps matching duplicates",
|
||||
collection: "alpha",
|
||||
expected: map[needle.VolumeId][]string{
|
||||
needle.VolumeId(10): {"alpha", "beta"},
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "from volume narrows duplicate set",
|
||||
collection: "*",
|
||||
from: needle.VolumeId(20),
|
||||
expected: map[needle.VolumeId][]string{
|
||||
needle.VolumeId(20): {"beta", "gamma"},
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "explicit pair checks both ends",
|
||||
collection: "*",
|
||||
from: needle.VolumeId(10),
|
||||
to: needle.VolumeId(20),
|
||||
expected: map[needle.VolumeId][]string{
|
||||
needle.VolumeId(10): {"alpha", "beta"},
|
||||
needle.VolumeId(20): {"beta", "gamma"},
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "collection plus explicit volume ignores unrelated duplicate",
|
||||
collection: "alpha",
|
||||
to: needle.VolumeId(20),
|
||||
expected: map[needle.VolumeId][]string{},
|
||||
},
|
||||
}
|
||||
|
||||
for _, tc := range tests {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
got := filterAmbiguousDuplicateVolumeIds(duplicates, tc.collection, tc.from, tc.to)
|
||||
if !reflect.DeepEqual(got, tc.expected) {
|
||||
t.Fatalf("unexpected duplicates: got=%v want=%v", got, tc.expected)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -337,9 +337,8 @@ func (c *commandVolumeFixReplication) fixUnderReplicatedVolumes(commandEnv *Comm
|
||||
}
|
||||
for _, vid := range volumeIds {
|
||||
for i := 0; i < retryCount+1; i++ {
|
||||
var copied bool
|
||||
if copied, err = c.fixOneUnderReplicatedVolume(commandEnv, writer, applyChanges, volumeReplicas, vid, allLocations); err == nil {
|
||||
if applyChanges && copied {
|
||||
if err = c.fixOneUnderReplicatedVolume(commandEnv, writer, applyChanges, volumeReplicas, vid, allLocations); err == nil {
|
||||
if applyChanges {
|
||||
fixedVolumes[strconv.FormatUint(uint64(vid), 10)] = len(volumeReplicas[vid])
|
||||
}
|
||||
break
|
||||
@@ -351,7 +350,7 @@ func (c *commandVolumeFixReplication) fixUnderReplicatedVolumes(commandEnv *Comm
|
||||
return fixedVolumes, nil
|
||||
}
|
||||
|
||||
func (c *commandVolumeFixReplication) fixOneUnderReplicatedVolume(commandEnv *CommandEnv, writer io.Writer, applyChanges bool, volumeReplicas map[uint32][]*VolumeReplica, vid uint32, allLocations []location) (bool, error) {
|
||||
func (c *commandVolumeFixReplication) fixOneUnderReplicatedVolume(commandEnv *CommandEnv, writer io.Writer, applyChanges bool, volumeReplicas map[uint32][]*VolumeReplica, vid uint32, allLocations []location) error {
|
||||
replicas := volumeReplicas[vid]
|
||||
replica := pickOneReplicaToCopyFrom(replicas)
|
||||
replicaPlacement, _ := super_block.NewReplicaPlacementFromByte(byte(replica.info.ReplicaPlacement))
|
||||
@@ -371,7 +370,7 @@ func (c *commandVolumeFixReplication) fixOneUnderReplicatedVolume(commandEnv *Co
|
||||
var err error
|
||||
matched, err = filepath.Match(*c.collectionPattern, replica.info.Collection)
|
||||
if err != nil {
|
||||
return false, fmt.Errorf("match pattern %s with collection %s: %v", *c.collectionPattern, replica.info.Collection, err)
|
||||
return fmt.Errorf("match pattern %s with collection %s: %v", *c.collectionPattern, replica.info.Collection, err)
|
||||
}
|
||||
}
|
||||
if !matched {
|
||||
@@ -387,7 +386,7 @@ func (c *commandVolumeFixReplication) fixOneUnderReplicatedVolume(commandEnv *Co
|
||||
if !applyChanges {
|
||||
// adjust volume count
|
||||
addVolumeCount(dst.dataNode.DiskInfos[replica.info.DiskType], 1)
|
||||
return true, nil
|
||||
break
|
||||
}
|
||||
|
||||
err := operation.WithVolumeServerClient(false, pb.NewServerAddressFromDataNode(dst.dataNode), commandEnv.option.GrpcDialOption, func(volumeServerClient volume_server_pb.VolumeServerClient) error {
|
||||
@@ -416,19 +415,19 @@ func (c *commandVolumeFixReplication) fixOneUnderReplicatedVolume(commandEnv *Co
|
||||
})
|
||||
|
||||
if err != nil {
|
||||
return false, err
|
||||
return err
|
||||
}
|
||||
|
||||
// adjust volume count
|
||||
addVolumeCount(dst.dataNode.DiskInfos[replica.info.DiskType], 1)
|
||||
return true, nil
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
if !foundNewLocation && !hasSkippedCollection {
|
||||
fmt.Fprintf(writer, "failed to place volume %d replica as %s, existing:%+v\n", replica.info.Id, replicaPlacement, len(replicas))
|
||||
}
|
||||
return false, nil
|
||||
return nil
|
||||
}
|
||||
|
||||
func addVolumeCount(info *master_pb.DiskInfo, count int) {
|
||||
|
||||
@@ -174,218 +174,6 @@ func TestWriteIdxFileFromEcIndex_ProcessesEcjJournal(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestDecodeWithNonEmptyEcj_AllDeleted verifies the full decode pre-processing
|
||||
// when .ecj contains deletions for ALL live entries in .ecx.
|
||||
// After RebuildEcxFile merges .ecj into .ecx, HasLiveNeedles must return false
|
||||
// and WriteIdxFileFromEcIndex must produce an .idx where every entry is tombstoned.
|
||||
func TestDecodeWithNonEmptyEcj_AllDeleted(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
base := filepath.Join(dir, "test_1")
|
||||
|
||||
// .ecx: two live entries
|
||||
needle1 := makeNeedleMapEntry(types.NeedleId(1), types.ToOffset(64), types.Size(100))
|
||||
needle2 := makeNeedleMapEntry(types.NeedleId(2), types.ToOffset(128), types.Size(200))
|
||||
ecxData := append(needle1, needle2...)
|
||||
if err := os.WriteFile(base+".ecx", ecxData, 0644); err != nil {
|
||||
t.Fatalf("write ecx: %v", err)
|
||||
}
|
||||
|
||||
// .ecj: both needles deleted
|
||||
ecjData := make([]byte, 2*types.NeedleIdSize)
|
||||
types.NeedleIdToBytes(ecjData[0:types.NeedleIdSize], types.NeedleId(1))
|
||||
types.NeedleIdToBytes(ecjData[types.NeedleIdSize:], types.NeedleId(2))
|
||||
if err := os.WriteFile(base+".ecj", ecjData, 0644); err != nil {
|
||||
t.Fatalf("write ecj: %v", err)
|
||||
}
|
||||
|
||||
// Before rebuild, ecx entries look live
|
||||
hasLive, err := erasure_coding.HasLiveNeedles(base)
|
||||
if err != nil {
|
||||
t.Fatalf("HasLiveNeedles before rebuild: %v", err)
|
||||
}
|
||||
if !hasLive {
|
||||
t.Fatal("expected live entries before rebuild")
|
||||
}
|
||||
|
||||
// Simulate what VolumeEcShardsToVolume now does: merge .ecj into .ecx
|
||||
if err := erasure_coding.RebuildEcxFile(base); err != nil {
|
||||
t.Fatalf("RebuildEcxFile: %v", err)
|
||||
}
|
||||
|
||||
// .ecj should be removed after rebuild
|
||||
if _, err := os.Stat(base + ".ecj"); !os.IsNotExist(err) {
|
||||
t.Fatal("expected .ecj to be removed after RebuildEcxFile")
|
||||
}
|
||||
|
||||
// After rebuild, HasLiveNeedles must return false
|
||||
hasLive, err = erasure_coding.HasLiveNeedles(base)
|
||||
if err != nil {
|
||||
t.Fatalf("HasLiveNeedles after rebuild: %v", err)
|
||||
}
|
||||
if hasLive {
|
||||
t.Fatal("expected no live entries after rebuild merged all deletions")
|
||||
}
|
||||
|
||||
// WriteIdxFileFromEcIndex should still work (no .ecj to process)
|
||||
if err := erasure_coding.WriteIdxFileFromEcIndex(base); err != nil {
|
||||
t.Fatalf("WriteIdxFileFromEcIndex: %v", err)
|
||||
}
|
||||
|
||||
idxData, err := os.ReadFile(base + ".idx")
|
||||
if err != nil {
|
||||
t.Fatalf("read idx: %v", err)
|
||||
}
|
||||
|
||||
// .idx should have exactly 2 entries (copied from .ecx, both now tombstoned)
|
||||
entrySize := types.NeedleIdSize + types.OffsetSize + types.SizeSize
|
||||
if len(idxData) != 2*entrySize {
|
||||
t.Fatalf("idx file size: got %d, want %d", len(idxData), 2*entrySize)
|
||||
}
|
||||
|
||||
// Both entries must be tombstoned
|
||||
for i := 0; i < 2; i++ {
|
||||
entry := idxData[i*entrySize : (i+1)*entrySize]
|
||||
size := types.BytesToSize(entry[types.NeedleIdSize+types.OffsetSize:])
|
||||
if !size.IsDeleted() {
|
||||
t.Fatalf("entry %d: expected tombstone, got size %d", i+1, size)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestDecodeWithNonEmptyEcj_PartiallyDeleted verifies decode pre-processing
|
||||
// when .ecj deletes only some entries. After RebuildEcxFile, HasLiveNeedles
|
||||
// must still return true for the surviving entries, and WriteIdxFileFromEcIndex
|
||||
// must produce an .idx that correctly distinguishes live from deleted needles.
|
||||
func TestDecodeWithNonEmptyEcj_PartiallyDeleted(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
base := filepath.Join(dir, "test_1")
|
||||
|
||||
// .ecx: three live entries
|
||||
needle1 := makeNeedleMapEntry(types.NeedleId(1), types.ToOffset(64), types.Size(100))
|
||||
needle2 := makeNeedleMapEntry(types.NeedleId(2), types.ToOffset(128), types.Size(200))
|
||||
needle3 := makeNeedleMapEntry(types.NeedleId(3), types.ToOffset(256), types.Size(300))
|
||||
ecxData := append(append(needle1, needle2...), needle3...)
|
||||
if err := os.WriteFile(base+".ecx", ecxData, 0644); err != nil {
|
||||
t.Fatalf("write ecx: %v", err)
|
||||
}
|
||||
|
||||
// .ecj: only needle 2 is deleted
|
||||
ecjData := make([]byte, types.NeedleIdSize)
|
||||
types.NeedleIdToBytes(ecjData, types.NeedleId(2))
|
||||
if err := os.WriteFile(base+".ecj", ecjData, 0644); err != nil {
|
||||
t.Fatalf("write ecj: %v", err)
|
||||
}
|
||||
|
||||
// Merge .ecj into .ecx
|
||||
if err := erasure_coding.RebuildEcxFile(base); err != nil {
|
||||
t.Fatalf("RebuildEcxFile: %v", err)
|
||||
}
|
||||
|
||||
// HasLiveNeedles must still return true (needles 1 and 3 survive)
|
||||
hasLive, err := erasure_coding.HasLiveNeedles(base)
|
||||
if err != nil {
|
||||
t.Fatalf("HasLiveNeedles: %v", err)
|
||||
}
|
||||
if !hasLive {
|
||||
t.Fatal("expected live entries after partial deletion")
|
||||
}
|
||||
|
||||
// WriteIdxFileFromEcIndex
|
||||
if err := erasure_coding.WriteIdxFileFromEcIndex(base); err != nil {
|
||||
t.Fatalf("WriteIdxFileFromEcIndex: %v", err)
|
||||
}
|
||||
|
||||
idxData, err := os.ReadFile(base + ".idx")
|
||||
if err != nil {
|
||||
t.Fatalf("read idx: %v", err)
|
||||
}
|
||||
|
||||
entrySize := types.NeedleIdSize + types.OffsetSize + types.SizeSize
|
||||
if len(idxData) != 3*entrySize {
|
||||
t.Fatalf("idx file size: got %d, want %d", len(idxData), 3*entrySize)
|
||||
}
|
||||
|
||||
// Verify each entry
|
||||
for i := 0; i < 3; i++ {
|
||||
entry := idxData[i*entrySize : (i+1)*entrySize]
|
||||
key := types.BytesToNeedleId(entry[0:types.NeedleIdSize])
|
||||
size := types.BytesToSize(entry[types.NeedleIdSize+types.OffsetSize:])
|
||||
|
||||
switch key {
|
||||
case types.NeedleId(1), types.NeedleId(3):
|
||||
if size.IsDeleted() {
|
||||
t.Fatalf("needle %d: should be live, got tombstone", key)
|
||||
}
|
||||
case types.NeedleId(2):
|
||||
if !size.IsDeleted() {
|
||||
t.Fatalf("needle %d: should be tombstoned, got size %d", key, size)
|
||||
}
|
||||
default:
|
||||
t.Fatalf("unexpected needle id %d", key)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestDecodeWithEmptyEcj verifies that the decode flow is a no-op when
|
||||
// .ecj exists but is empty (no deletions recorded).
|
||||
func TestDecodeWithEmptyEcj(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
base := filepath.Join(dir, "test_1")
|
||||
|
||||
// .ecx: one live entry
|
||||
needle1 := makeNeedleMapEntry(types.NeedleId(1), types.ToOffset(64), types.Size(100))
|
||||
if err := os.WriteFile(base+".ecx", needle1, 0644); err != nil {
|
||||
t.Fatalf("write ecx: %v", err)
|
||||
}
|
||||
|
||||
// .ecj: empty
|
||||
if err := os.WriteFile(base+".ecj", []byte{}, 0644); err != nil {
|
||||
t.Fatalf("write ecj: %v", err)
|
||||
}
|
||||
|
||||
// RebuildEcxFile with empty .ecj should not change anything
|
||||
if err := erasure_coding.RebuildEcxFile(base); err != nil {
|
||||
t.Fatalf("RebuildEcxFile: %v", err)
|
||||
}
|
||||
|
||||
// HasLiveNeedles must still return true
|
||||
hasLive, err := erasure_coding.HasLiveNeedles(base)
|
||||
if err != nil {
|
||||
t.Fatalf("HasLiveNeedles: %v", err)
|
||||
}
|
||||
if !hasLive {
|
||||
t.Fatal("expected live entries with empty .ecj")
|
||||
}
|
||||
}
|
||||
|
||||
// TestDecodeWithNoEcjFile verifies that the decode flow works when no .ecj
|
||||
// file exists at all.
|
||||
func TestDecodeWithNoEcjFile(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
base := filepath.Join(dir, "test_1")
|
||||
|
||||
// .ecx: one live entry
|
||||
needle1 := makeNeedleMapEntry(types.NeedleId(1), types.ToOffset(64), types.Size(100))
|
||||
if err := os.WriteFile(base+".ecx", needle1, 0644); err != nil {
|
||||
t.Fatalf("write ecx: %v", err)
|
||||
}
|
||||
|
||||
// No .ecj file
|
||||
|
||||
// RebuildEcxFile should be a no-op
|
||||
if err := erasure_coding.RebuildEcxFile(base); err != nil {
|
||||
t.Fatalf("RebuildEcxFile: %v", err)
|
||||
}
|
||||
|
||||
hasLive, err := erasure_coding.HasLiveNeedles(base)
|
||||
if err != nil {
|
||||
t.Fatalf("HasLiveNeedles: %v", err)
|
||||
}
|
||||
if !hasLive {
|
||||
t.Fatal("expected live entries without .ecj file")
|
||||
}
|
||||
}
|
||||
|
||||
// TestEcxFileDeletionVisibleAfterSync verifies that deletions made to .ecx
|
||||
// via MarkNeedleDeleted are visible to other readers after Sync().
|
||||
// This is a regression test for issue #7751.
|
||||
|
||||
@@ -19,9 +19,27 @@ func (s *Store) CheckCompactVolume(volumeId needle.VolumeId) (float64, error) {
|
||||
|
||||
func (s *Store) CompactVolume(vid needle.VolumeId, preallocate int64, compactionBytePerSecond int64, progressFn ProgressFunc) error {
|
||||
if v := s.findVolume(vid); v != nil {
|
||||
if err := ensureCompactVolumeSpace(v, preallocate); err != nil {
|
||||
return err
|
||||
// Get current volume size for space calculation
|
||||
volumeSize, indexSize, _ := v.FileStat()
|
||||
|
||||
// Calculate space needed for compaction:
|
||||
// 1. Space for the new compacted volume (approximately same as current volume size)
|
||||
// 2. Use the larger of preallocate or estimated volume size
|
||||
estimatedCompactSize := int64(volumeSize + indexSize)
|
||||
spaceNeeded := preallocate
|
||||
if estimatedCompactSize > preallocate {
|
||||
spaceNeeded = estimatedCompactSize
|
||||
}
|
||||
|
||||
diskStatus := stats.NewDiskStatus(v.dir)
|
||||
if int64(diskStatus.Free) < spaceNeeded {
|
||||
return fmt.Errorf("insufficient free space for compaction: need %d bytes (volume: %d, index: %d, buffer: 10%%), but only %d bytes available",
|
||||
spaceNeeded, volumeSize, indexSize, diskStatus.Free)
|
||||
}
|
||||
|
||||
glog.V(1).Infof("volume %d compaction space check: volume=%d, index=%d, space_needed=%d, free_space=%d",
|
||||
vid, volumeSize, indexSize, spaceNeeded, diskStatus.Free)
|
||||
|
||||
return v.CompactByIndex(&CompactOptions{
|
||||
PreallocateBytes: preallocate,
|
||||
MaxBytesPerSecond: compactionBytePerSecond,
|
||||
@@ -53,71 +71,3 @@ func (s *Store) CommitCleanupVolume(vid needle.VolumeId) error {
|
||||
}
|
||||
return fmt.Errorf("volume id %d is not found during cleaning up", vid)
|
||||
}
|
||||
|
||||
func ensureCompactVolumeSpace(v *Volume, preallocate int64) error {
|
||||
// Get current volume size for space calculation
|
||||
volumeSize, indexSize, _ := v.FileStat()
|
||||
|
||||
// Calculate space needed for compaction:
|
||||
// 1. Space for the new compacted volume (approximately same as current volume size)
|
||||
// 2. Use the larger of preallocate or estimated volume size
|
||||
estimatedCompactSize := int64(volumeSize + indexSize)
|
||||
spaceNeeded := preallocate
|
||||
if estimatedCompactSize > preallocate {
|
||||
spaceNeeded = estimatedCompactSize
|
||||
}
|
||||
|
||||
diskStatus := stats.NewDiskStatus(v.dir)
|
||||
if int64(diskStatus.Free) < spaceNeeded {
|
||||
return fmt.Errorf("insufficient free space for compaction: need %d bytes (volume: %d, index: %d), but only %d bytes available",
|
||||
spaceNeeded, volumeSize, indexSize, diskStatus.Free)
|
||||
}
|
||||
|
||||
glog.V(1).Infof("volume %d compaction space check: volume=%d, index=%d, space_needed=%d, free_space=%d",
|
||||
v.Id, volumeSize, indexSize, spaceNeeded, diskStatus.Free)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *Store) CompactVolumeFiles(vid needle.VolumeId, collection string, location *DiskLocation, needleMapKind NeedleMapKind, ldbTimeout int64, preallocate int64, compactionBytePerSecond int64) (err error) {
|
||||
if location == nil {
|
||||
return fmt.Errorf("volume %d compaction location is nil", vid)
|
||||
}
|
||||
|
||||
tempVolume, err := loadVolumeWithoutWorker(location.Directory, location.IdxDirectory, collection, vid, needleMapKind, ldbTimeout)
|
||||
if err != nil {
|
||||
return fmt.Errorf("load volume %d for offline compaction: %w", vid, err)
|
||||
}
|
||||
tempVolume.location = location
|
||||
|
||||
defer func() {
|
||||
if tempVolume.tmpNm != nil {
|
||||
tempVolume.tmpNm.Close()
|
||||
tempVolume.tmpNm = nil
|
||||
}
|
||||
tempVolume.doClose()
|
||||
}()
|
||||
|
||||
if err := ensureCompactVolumeSpace(tempVolume, preallocate); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := tempVolume.CompactByIndex(&CompactOptions{
|
||||
PreallocateBytes: preallocate,
|
||||
MaxBytesPerSecond: compactionBytePerSecond,
|
||||
}); err != nil {
|
||||
if cleanupErr := tempVolume.cleanupCompact(); cleanupErr != nil {
|
||||
return fmt.Errorf("compact volume %d: %v (cleanup failed: %v)", vid, err, cleanupErr)
|
||||
}
|
||||
return fmt.Errorf("compact volume %d: %w", vid, err)
|
||||
}
|
||||
|
||||
if err := tempVolume.CommitCompact(); err != nil {
|
||||
if cleanupErr := tempVolume.cleanupCompact(); cleanupErr != nil {
|
||||
return fmt.Errorf("commit compact volume %d: %v (cleanup failed: %v)", vid, err, cleanupErr)
|
||||
}
|
||||
return fmt.Errorf("commit compact volume %d: %w", vid, err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -24,20 +24,6 @@ func loadVolumeWithoutIndex(dirname string, collection string, id needle.VolumeI
|
||||
return
|
||||
}
|
||||
|
||||
func loadVolumeWithoutWorker(dirname string, dirIdx string, collection string, id needle.VolumeId, needleMapKind NeedleMapKind, ldbTimeout int64) (v *Volume, err error) {
|
||||
v = &Volume{
|
||||
dir: dirname,
|
||||
dirIdx: dirIdx,
|
||||
Collection: collection,
|
||||
Id: id,
|
||||
needleMapKind: needleMapKind,
|
||||
ldbTimeout: ldbTimeout,
|
||||
}
|
||||
v.SuperBlock = super_block.SuperBlock{}
|
||||
err = v.load(true, false, needleMapKind, 0, needle.GetCurrentVersion())
|
||||
return
|
||||
}
|
||||
|
||||
func (v *Volume) load(alsoLoadIndex bool, createDatIfMissing bool, needleMapKind NeedleMapKind, preallocate int64, ver needle.Version) (err error) {
|
||||
alreadyHasSuperBlock := false
|
||||
|
||||
|
||||
@@ -2,8 +2,6 @@ package storage
|
||||
|
||||
import (
|
||||
"math/rand"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"reflect"
|
||||
"testing"
|
||||
"time"
|
||||
@@ -11,7 +9,6 @@ import (
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/needle"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/super_block"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/types"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
)
|
||||
|
||||
/*
|
||||
@@ -149,90 +146,6 @@ func testCompactionByIndex(t *testing.T, needleMapKind NeedleMapKind) {
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
func TestCompactVolumeFilesOffline(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
location := NewDiskLocation(dir, 10, util.MinFreeSpace{}, dir, "", nil)
|
||||
defer location.Close()
|
||||
|
||||
v, err := NewVolume(dir, dir, "", 1, NeedleMapInMemory, &super_block.ReplicaPlacement{}, &needle.TTL{}, 0, needle.GetCurrentVersion(), 0, 0)
|
||||
if err != nil {
|
||||
t.Fatalf("volume creation: %v", err)
|
||||
}
|
||||
|
||||
infos := make([]*needleInfo, 32)
|
||||
for i := 1; i <= 32; i++ {
|
||||
doSomeWritesDeletes(i, v, t, infos)
|
||||
}
|
||||
v.Close()
|
||||
|
||||
store := &Store{}
|
||||
if err := store.CompactVolumeFiles(needle.VolumeId(1), "", location, NeedleMapInMemory, 0, 0, 0); err != nil {
|
||||
t.Fatalf("CompactVolumeFiles: %v", err)
|
||||
}
|
||||
|
||||
reloaded, err := NewVolume(dir, dir, "", 1, NeedleMapInMemory, nil, nil, 0, needle.GetCurrentVersion(), 0, 0)
|
||||
if err != nil {
|
||||
t.Fatalf("volume reload: %v", err)
|
||||
}
|
||||
defer reloaded.Close()
|
||||
|
||||
if _, err := os.Stat(filepath.Join(dir, "1.cpd")); !os.IsNotExist(err) {
|
||||
t.Fatalf("expected no .cpd after successful offline compaction, got err=%v", err)
|
||||
}
|
||||
if _, err := os.Stat(filepath.Join(dir, "1.cpx")); !os.IsNotExist(err) {
|
||||
t.Fatalf("expected no .cpx after successful offline compaction, got err=%v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCleanupCompactRemovesTempFiles(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
location := NewDiskLocation(dir, 10, util.MinFreeSpace{}, dir, "", nil)
|
||||
defer location.Close()
|
||||
|
||||
v, err := NewVolume(dir, dir, "", 1, NeedleMapInMemory, &super_block.ReplicaPlacement{}, &needle.TTL{}, 0, needle.GetCurrentVersion(), 0, 0)
|
||||
if err != nil {
|
||||
t.Fatalf("volume creation: %v", err)
|
||||
}
|
||||
|
||||
infos := make([]*needleInfo, 16)
|
||||
for i := 1; i <= 16; i++ {
|
||||
doSomeWritesDeletes(i, v, t, infos)
|
||||
}
|
||||
v.Close()
|
||||
|
||||
if err := os.WriteFile(filepath.Join(dir, "1.cpx"), []byte("broken"), 0o644); err != nil {
|
||||
t.Fatalf("write broken cpx: %v", err)
|
||||
}
|
||||
if err := os.WriteFile(filepath.Join(dir, "1.cpd"), []byte("temp"), 0o644); err != nil {
|
||||
t.Fatalf("write cpd: %v", err)
|
||||
}
|
||||
if err := os.Mkdir(filepath.Join(dir, "1.cpldb"), 0o755); err != nil {
|
||||
t.Fatalf("mkdir cpldb: %v", err)
|
||||
}
|
||||
|
||||
tempVolume, err := loadVolumeWithoutWorker(dir, dir, "", needle.VolumeId(1), NeedleMapInMemory, 0)
|
||||
if err != nil {
|
||||
t.Fatalf("loadVolumeWithoutWorker: %v", err)
|
||||
}
|
||||
tempVolume.location = location
|
||||
defer tempVolume.doClose()
|
||||
|
||||
if err := tempVolume.cleanupCompact(); err != nil {
|
||||
t.Fatalf("cleanupCompact: %v", err)
|
||||
}
|
||||
|
||||
if _, err := os.Stat(filepath.Join(dir, "1.cpd")); !os.IsNotExist(err) {
|
||||
t.Fatalf("expected cleanup to remove .cpd, got err=%v", err)
|
||||
}
|
||||
if _, err := os.Stat(filepath.Join(dir, "1.cpx")); !os.IsNotExist(err) {
|
||||
t.Fatalf("expected cleanup to remove .cpx, got err=%v", err)
|
||||
}
|
||||
if _, err := os.Stat(filepath.Join(dir, "1.cpldb")); !os.IsNotExist(err) {
|
||||
t.Fatalf("expected cleanup to remove .cpldb, got err=%v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func doSomeWritesDeletes(i int, v *Volume, t *testing.T, infos []*needleInfo) {
|
||||
n := newRandomNeedle(uint64(i))
|
||||
_, size, _, err := v.writeNeedle2(n, true, false)
|
||||
|
||||
@@ -6,6 +6,7 @@ import (
|
||||
"fmt"
|
||||
"math/rand/v2"
|
||||
"slices"
|
||||
"strings"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
@@ -366,6 +367,74 @@ func (t *Topology) ListCollections(includeNormalVolumes, includeEcVolumes bool)
|
||||
return ret
|
||||
}
|
||||
|
||||
func (t *Topology) FindDuplicateVolumeIds(collectionName string) map[needle.VolumeId][]string {
|
||||
duplicates := make(map[needle.VolumeId][]string)
|
||||
collectionsByVid := make(map[needle.VolumeId]map[string]struct{})
|
||||
|
||||
t.collectionMap.RLock()
|
||||
defer t.collectionMap.RUnlock()
|
||||
|
||||
for _, item := range t.collectionMap.Items() {
|
||||
collection := item.(*Collection)
|
||||
for _, layout := range collection.GetAllVolumeLayouts() {
|
||||
layout.accessLock.RLock()
|
||||
for vid := range layout.vid2location {
|
||||
if collectionsByVid[vid] == nil {
|
||||
collectionsByVid[vid] = make(map[string]struct{})
|
||||
}
|
||||
collectionsByVid[vid][collection.Name] = struct{}{}
|
||||
}
|
||||
layout.accessLock.RUnlock()
|
||||
}
|
||||
}
|
||||
|
||||
for vid, collectionSet := range collectionsByVid {
|
||||
if len(collectionSet) < 2 {
|
||||
continue
|
||||
}
|
||||
if collectionName != "" {
|
||||
if _, found := collectionSet[collectionName]; !found {
|
||||
continue
|
||||
}
|
||||
}
|
||||
duplicateCollections := make([]string, 0, len(collectionSet))
|
||||
for name := range collectionSet {
|
||||
duplicateCollections = append(duplicateCollections, name)
|
||||
}
|
||||
slices.Sort(duplicateCollections)
|
||||
duplicates[vid] = duplicateCollections
|
||||
}
|
||||
|
||||
return duplicates
|
||||
}
|
||||
|
||||
func FormatDuplicateVolumeIds(duplicates map[needle.VolumeId][]string) string {
|
||||
if len(duplicates) == 0 {
|
||||
return ""
|
||||
}
|
||||
|
||||
vids := make([]needle.VolumeId, 0, len(duplicates))
|
||||
for vid := range duplicates {
|
||||
vids = append(vids, vid)
|
||||
}
|
||||
slices.Sort(vids)
|
||||
|
||||
lines := make([]string, 0, len(vids))
|
||||
for _, vid := range vids {
|
||||
displayCollections := make([]string, 0, len(duplicates[vid]))
|
||||
for _, name := range duplicates[vid] {
|
||||
if name == "" {
|
||||
displayCollections = append(displayCollections, "(default)")
|
||||
} else {
|
||||
displayCollections = append(displayCollections, name)
|
||||
}
|
||||
}
|
||||
lines = append(lines, fmt.Sprintf("%d:[%s]", vid, strings.Join(displayCollections, ", ")))
|
||||
}
|
||||
|
||||
return strings.Join(lines, "; ")
|
||||
}
|
||||
|
||||
func (t *Topology) FindCollection(collectionName string) (*Collection, bool) {
|
||||
c, hasCollection := t.collectionMap.Find(collectionName)
|
||||
if !hasCollection {
|
||||
|
||||
@@ -14,6 +14,75 @@ import (
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestFindDuplicateVolumeIds(t *testing.T) {
|
||||
topo := NewTopology("weedfs", sequence.NewMemorySequencer(), 32*1024, 5, false)
|
||||
|
||||
dc := topo.GetOrCreateDataCenter("dc1")
|
||||
rack := dc.GetOrCreateRack("rack1")
|
||||
maxVolumeCounts := map[string]uint32{"": 25}
|
||||
|
||||
dn1 := rack.GetOrCreateDataNode("127.0.0.1", 34534, 0, "127.0.0.1", "", maxVolumeCounts)
|
||||
dn2 := rack.GetOrCreateDataNode("127.0.0.2", 34535, 0, "127.0.0.2", "", maxVolumeCounts)
|
||||
|
||||
replicaPlacement := &super_block.ReplicaPlacement{}
|
||||
|
||||
collectionA := storage.VolumeInfo{
|
||||
Id: needle.VolumeId(100),
|
||||
Collection: "collection-a",
|
||||
DiskType: "",
|
||||
ReadOnly: false,
|
||||
Version: needle.GetCurrentVersion(),
|
||||
ReplicaPlacement: replicaPlacement,
|
||||
Ttl: needle.EMPTY_TTL,
|
||||
}
|
||||
collectionB := storage.VolumeInfo{
|
||||
Id: needle.VolumeId(100),
|
||||
Collection: "collection-b",
|
||||
DiskType: "",
|
||||
ReadOnly: false,
|
||||
Version: needle.GetCurrentVersion(),
|
||||
ReplicaPlacement: replicaPlacement,
|
||||
Ttl: needle.EMPTY_TTL,
|
||||
}
|
||||
unique := storage.VolumeInfo{
|
||||
Id: needle.VolumeId(200),
|
||||
Collection: "collection-a",
|
||||
DiskType: "",
|
||||
ReadOnly: false,
|
||||
Version: needle.GetCurrentVersion(),
|
||||
ReplicaPlacement: replicaPlacement,
|
||||
Ttl: needle.EMPTY_TTL,
|
||||
}
|
||||
|
||||
dn1.UpdateVolumes([]storage.VolumeInfo{collectionA, unique})
|
||||
dn2.UpdateVolumes([]storage.VolumeInfo{collectionB})
|
||||
topo.RegisterVolumeLayout(collectionA, dn1)
|
||||
topo.RegisterVolumeLayout(unique, dn1)
|
||||
topo.RegisterVolumeLayout(collectionB, dn2)
|
||||
|
||||
duplicates := topo.FindDuplicateVolumeIds("")
|
||||
got, ok := duplicates[needle.VolumeId(100)]
|
||||
if !ok {
|
||||
t.Fatalf("expected duplicate volume 100, got %+v", duplicates)
|
||||
}
|
||||
expected := []string{"collection-a", "collection-b"}
|
||||
if !reflect.DeepEqual(got, expected) {
|
||||
t.Fatalf("unexpected duplicate collections: got=%v want=%v", got, expected)
|
||||
}
|
||||
if _, exists := duplicates[needle.VolumeId(200)]; exists {
|
||||
t.Fatalf("did not expect unique volume 200 in duplicates: %+v", duplicates)
|
||||
}
|
||||
|
||||
scoped := topo.FindDuplicateVolumeIds("collection-a")
|
||||
if _, ok := scoped[needle.VolumeId(100)]; !ok {
|
||||
t.Fatalf("expected scoped duplicate for collection-a, got %+v", scoped)
|
||||
}
|
||||
scopedMissing := topo.FindDuplicateVolumeIds("collection-c")
|
||||
if len(scopedMissing) != 0 {
|
||||
t.Fatalf("expected no duplicates for unrelated collection, got %+v", scopedMissing)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRemoveDataCenter(t *testing.T) {
|
||||
topo := setup(topologyLayout)
|
||||
topo.UnlinkChildNode(NodeId("dc2"))
|
||||
|
||||
@@ -9,7 +9,7 @@ import (
|
||||
|
||||
var (
|
||||
MAJOR_VERSION = int32(4)
|
||||
MINOR_VERSION = int32(18)
|
||||
MINOR_VERSION = int32(17)
|
||||
VERSION_NUMBER = fmt.Sprintf("%d.%02d", MAJOR_VERSION, MINOR_VERSION)
|
||||
VERSION = util.SizeLimit + " " + VERSION_NUMBER
|
||||
COMMIT = ""
|
||||
|
||||
Reference in New Issue
Block a user