Files
seaweedfs/seaweed-volume/tests/http_integration.rs
T
c1ccbcda13 volume server: the write queue answers uploads with the needle's real ETag (#11414)
With SEAWEED_WRITE_QUEUE=1 every upload came back with ETag "00000000".
The upload handler built the needle with Needle::default(), so its
checksum was CRC(0), and handed a clone of it to the queue. The CRC was
only computed in the write path, on the worker's clone, and WriteResult
carries no checksum back, so n.etag() in the handler formatted the zero
checksum. The direct path writes through &mut n and was correct.

Compute the checksum in the handler while building the needle, the way
Go's CreateNeedleFromRequest does, over the same bytes the write path
hashes (the stored data, gzipped or not). The ETag and the has-name flag
are read before the write, so the needle is moved into the queue instead
of cloned, which also drops a full payload copy per queued upload.

Co-authored-by: Claude Fable 5.1 <noreply@anthropic.com>
2026-09-24 07:08:36 +08:00

1357 lines
44 KiB
Rust

//! Integration tests for the volume server HTTP handlers.
//!
//! Uses axum's Router with tower::ServiceExt::oneshot to test
//! end-to-end without starting a real TCP server.
use std::sync::{Arc, RwLock};
use axum::body::Body;
use axum::extract::connect_info::ConnectInfo;
use axum::http::{Request, StatusCode};
use tower::ServiceExt; // for `oneshot`
use seaweed_volume::security::{Guard, SigningKey};
use seaweed_volume::server::volume_server::{
VolumeServerState, build_admin_router, build_admin_router_with_ui, build_metrics_router,
build_public_router,
};
use seaweed_volume::storage::needle_map::NeedleMapKind;
use seaweed_volume::storage::store::Store;
use seaweed_volume::storage::types::{DiskType, VolumeId};
use seaweed_volume::storage::volume::VolumeSpec;
use tempfile::TempDir;
/// Create a test VolumeServerState with a temp directory, a single disk
/// location, and one pre-created volume (VolumeId 1).
fn test_state() -> (Arc<VolumeServerState>, TempDir) {
test_state_with_guard(Vec::new(), Vec::new())
}
fn test_state_with_signing_key(signing_key: Vec<u8>) -> (Arc<VolumeServerState>, TempDir) {
test_state_with_guard(Vec::new(), signing_key)
}
fn test_state_with_whitelist(whitelist: Vec<String>) -> (Arc<VolumeServerState>, TempDir) {
test_state_with_guard(whitelist, Vec::new())
}
fn test_state_with_guard(
whitelist: Vec<String>,
signing_key: Vec<u8>,
) -> (Arc<VolumeServerState>, TempDir) {
build_test_state(whitelist, signing_key, None, String::new())
}
/// Build a test state with volume 1 created using the given replica placement
/// (e.g. "001" for a two-copy volume) and master URL. An empty master URL plus
/// a None placement is the single-copy default used by most tests.
fn build_test_state(
whitelist: Vec<String>,
signing_key: Vec<u8>,
replication: Option<&str>,
master_url: String,
) -> (Arc<VolumeServerState>, TempDir) {
let tmp = TempDir::new().expect("failed to create temp dir");
let dir = tmp.path().to_str().unwrap();
let replica_placement = replication.map(|s| {
seaweed_volume::storage::super_block::ReplicaPlacement::from_string(s)
.expect("invalid replica placement")
});
let mut store = Store::new(NeedleMapKind::InMemory);
store
.add_location(
dir,
dir,
10,
DiskType::HardDrive,
seaweed_volume::config::MinFreeSpace::Percent(1.0),
Vec::new(),
)
.expect("failed to add location");
store
.add_volume(
VolumeId(1),
DiskType::HardDrive,
&VolumeSpec {
replica_placement,
..Default::default()
},
)
.expect("failed to create volume");
let guard = Guard::new(
&whitelist,
SigningKey(signing_key),
0,
SigningKey(vec![]),
0,
);
let state = Arc::new(VolumeServerState {
store: RwLock::new(store),
guard: RwLock::new(guard),
is_stopping: RwLock::new(false),
maintenance: std::sync::atomic::AtomicBool::new(false),
state_version: std::sync::atomic::AtomicU32::new(0),
concurrent_upload_limit: 0,
concurrent_download_limit: 0,
inflight_upload_data_timeout: std::time::Duration::from_secs(60),
inflight_download_data_timeout: std::time::Duration::from_secs(60),
inflight_upload_bytes: std::sync::atomic::AtomicI64::new(0),
inflight_download_bytes: std::sync::atomic::AtomicI64::new(0),
upload_notify: tokio::sync::Notify::new(),
download_notify: tokio::sync::Notify::new(),
data_center: String::new(),
rack: String::new(),
file_size_limit_bytes: 0,
maintenance_byte_per_second: 0,
is_heartbeating: std::sync::atomic::AtomicBool::new(true),
has_master: false,
pre_stop_seconds: 0,
volume_state_notify: tokio::sync::Notify::new(),
write_queue: std::sync::OnceLock::new(),
read_mode: seaweed_volume::config::ReadMode::Local,
allow_untrusted_remote_endpoints: false,
master_url,
master_urls: Vec::new(),
seed_master_set: std::collections::HashSet::new(),
current_master_url: tokio::sync::RwLock::new(String::new()),
self_url: String::new(),
http_client: reqwest::Client::new(),
outgoing_http_scheme: "http".to_string(),
outgoing_grpc_tls: None,
metrics_runtime: std::sync::RwLock::new(
seaweed_volume::server::volume_server::RuntimeMetricsConfig::default(),
),
metrics_notify: tokio::sync::Notify::new(),
fix_jpg_orientation: false,
has_slow_read: false,
read_buffer_size_bytes: 1024 * 1024,
security_file: String::new(),
cli_white_list: vec![],
state_file_path: String::new(),
});
(state, tmp)
}
/// Helper: read the entire response body as bytes.
async fn body_bytes(response: axum::response::Response) -> Vec<u8> {
let body = response.into_body();
axum::body::to_bytes(body, usize::MAX)
.await
.expect("failed to read body")
.to_vec()
}
fn with_remote_addr(request: Request<Body>, remote_addr: &str) -> Request<Body> {
let mut request = request;
let remote_addr = remote_addr
.parse::<std::net::SocketAddr>()
.expect("invalid socket address");
request.extensions_mut().insert(ConnectInfo(remote_addr));
request
}
// ============================================================================
// 1. GET /healthz returns 200 when server is running
// ============================================================================
#[tokio::test]
async fn healthz_returns_200_when_running() {
let (state, _tmp) = test_state();
let app = build_admin_router(state);
let response = app
.oneshot(
Request::builder()
.uri("/healthz")
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::OK);
}
// ============================================================================
// 2. GET /healthz returns 503 when is_stopping=true
// ============================================================================
#[tokio::test]
async fn healthz_returns_503_when_stopping() {
let (state, _tmp) = test_state();
*state.is_stopping.write().unwrap() = true;
let app = build_admin_router(state);
let response = app
.oneshot(
Request::builder()
.uri("/healthz")
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE);
}
// ============================================================================
// 3. GET /status returns JSON with version and volumes array
// ============================================================================
#[tokio::test]
async fn status_returns_json_with_version_and_volumes() {
let (state, _tmp) = test_state();
let app = build_admin_router(state);
let response = app
.oneshot(
Request::builder()
.uri("/status")
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::OK);
let body = body_bytes(response).await;
let json: serde_json::Value =
serde_json::from_slice(&body).expect("response is not valid JSON");
assert!(json.get("Version").is_some(), "missing 'Version' field");
assert!(json["Version"].is_string(), "'Version' should be a string");
assert!(json.get("Volumes").is_some(), "missing 'Volumes' field");
assert!(json["Volumes"].is_array(), "'Volumes' should be an array");
// We created one volume in test_state, so the array should have one entry
let volumes = json["Volumes"].as_array().unwrap();
assert_eq!(volumes.len(), 1, "expected 1 volume");
assert_eq!(volumes[0]["Id"], 1);
}
#[tokio::test]
async fn admin_router_does_not_expose_metrics() {
let (state, _tmp) = test_state();
let app = build_admin_router(state);
let response = app
.oneshot(
Request::builder()
.uri("/metrics")
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::BAD_REQUEST);
}
#[tokio::test]
async fn metrics_router_serves_metrics() {
let app = build_metrics_router();
let response = app
.oneshot(
Request::builder()
.uri("/metrics")
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::OK);
}
#[tokio::test]
async fn admin_router_rejects_non_whitelisted_uploads() {
let (state, _tmp) = test_state_with_whitelist(vec!["127.0.0.1".to_string()]);
let app = build_admin_router(state);
let response = app
.oneshot(with_remote_addr(
Request::builder()
.method("POST")
.uri("/1,000000000000000001")
.body(Body::from("blocked"))
.unwrap(),
"10.0.0.9:12345",
))
.await
.unwrap();
assert_eq!(response.status(), StatusCode::UNAUTHORIZED);
}
#[tokio::test]
async fn admin_router_rejects_non_whitelisted_deletes() {
let (state, _tmp) = test_state_with_whitelist(vec!["127.0.0.1".to_string()]);
let app = build_admin_router(state);
let response = app
.oneshot(with_remote_addr(
Request::builder()
.method("DELETE")
.uri("/1,000000000000000001")
.body(Body::empty())
.unwrap(),
"10.0.0.9:12345",
))
.await
.unwrap();
assert_eq!(response.status(), StatusCode::UNAUTHORIZED);
}
// Go's volume_server.go has /stats/* endpoints commented out (L130-134).
// Requests to /stats/counter fall through to the store handler which returns 400.
#[tokio::test]
async fn admin_router_does_not_expose_stats_routes() {
let (state, _tmp) = test_state_with_whitelist(vec!["127.0.0.1".to_string()]);
let app = build_admin_router_with_ui(state, true);
let response = app
.oneshot(with_remote_addr(
Request::builder()
.uri("/stats/counter")
.body(Body::empty())
.unwrap(),
"127.0.0.1:12345",
))
.await
.unwrap();
// Falls through to store handler → 400 (bad volume id)
assert_eq!(response.status(), StatusCode::BAD_REQUEST);
}
// ============================================================================
// 4. POST writes data, then GET reads it back
// ============================================================================
#[tokio::test]
async fn write_then_read_needle() {
let (state, _tmp) = test_state();
// The fid "01637037d6" encodes NeedleId=0x01, Cookie=0x637037d6
let uri = "/1,01637037d6";
let payload = b"hello, seaweedfs!";
// --- POST (write) ---
let app = build_admin_router(state.clone());
let response = app
.oneshot(
Request::builder()
.method("POST")
.uri(uri)
.body(Body::from(payload.to_vec()))
.unwrap(),
)
.await
.unwrap();
assert_eq!(
response.status(),
StatusCode::CREATED,
"POST should return 201 Created"
);
let body = body_bytes(response).await;
let json: serde_json::Value =
serde_json::from_slice(&body).expect("POST response is not valid JSON");
assert_eq!(json["size"], payload.len() as u64);
// --- GET (read back) ---
let app = build_admin_router(state.clone());
let response = app
.oneshot(Request::builder().uri(uri).body(Body::empty()).unwrap())
.await
.unwrap();
assert_eq!(response.status(), StatusCode::OK, "GET should return 200");
let body = body_bytes(response).await;
assert_eq!(body, payload, "GET body should match written data");
}
// A durable upload takes the same route through the handler; the flush is not
// observable from here, but a broken wiring would show up as a failed write.
#[tokio::test]
async fn write_with_fsync_then_read_needle() {
let (state, _tmp) = test_state();
let uri = "/1,01637037d6?fsync=true";
let payload = b"durable through the handler";
let app = build_admin_router(state.clone());
let response = app
.oneshot(
Request::builder()
.method("POST")
.uri(uri)
.body(Body::from(payload.to_vec()))
.unwrap(),
)
.await
.unwrap();
assert_eq!(
response.status(),
StatusCode::CREATED,
"a durable POST should return 201 Created"
);
let app = build_admin_router(state.clone());
let response = app
.oneshot(
Request::builder()
.uri("/1,01637037d6")
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::OK);
assert_eq!(body_bytes(response).await, payload);
}
// ============================================================================
// 5. DELETE deletes a needle, subsequent GET returns 404
// ============================================================================
#[tokio::test]
async fn delete_then_get_returns_404() {
let (state, _tmp) = test_state();
let uri = "/1,01637037d6";
let payload = b"to be deleted";
// Write the needle first
let app = build_admin_router(state.clone());
let response = app
.oneshot(
Request::builder()
.method("POST")
.uri(uri)
.body(Body::from(payload.to_vec()))
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::CREATED);
// Delete
let app = build_admin_router(state.clone());
let response = app
.oneshot(
Request::builder()
.method("DELETE")
.uri(uri)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(
response.status(),
StatusCode::ACCEPTED,
"DELETE should return 202 Accepted"
);
// GET should now return 404
let app = build_admin_router(state.clone());
let response = app
.oneshot(Request::builder().uri(uri).body(Body::empty()).unwrap())
.await
.unwrap();
assert_eq!(
response.status(),
StatusCode::NOT_FOUND,
"GET after DELETE should return 404"
);
}
// ============================================================================
// 6. HEAD returns headers without body
// ============================================================================
#[tokio::test]
async fn head_returns_headers_without_body() {
let (state, _tmp) = test_state();
let uri = "/1,01637037d6";
let payload = b"head test data";
// Write needle
let app = build_admin_router(state.clone());
let response = app
.oneshot(
Request::builder()
.method("POST")
.uri(uri)
.body(Body::from(payload.to_vec()))
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::CREATED);
// HEAD
let app = build_admin_router(state.clone());
let response = app
.oneshot(
Request::builder()
.method("HEAD")
.uri(uri)
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::OK, "HEAD should return 200");
// Content-Length header should be present
let content_length = response
.headers()
.get("content-length")
.expect("HEAD should include Content-Length header");
let len: usize = content_length
.to_str()
.unwrap()
.parse()
.expect("Content-Length should be a number");
assert_eq!(
len,
payload.len(),
"Content-Length should match payload size"
);
// Body should be empty for HEAD
let body = body_bytes(response).await;
assert!(body.is_empty(), "HEAD body should be empty");
}
// ============================================================================
// 7. Invalid URL path returns 400
// ============================================================================
#[tokio::test]
async fn invalid_url_path_returns_400() {
let (state, _tmp) = test_state();
let app = build_admin_router(state);
// "invalidpath" has no comma or slash separator so parse_url_path returns None
let response = app
.oneshot(
Request::builder()
.uri("/invalidpath")
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(
response.status(),
StatusCode::BAD_REQUEST,
"invalid URL path should return 400"
);
}
#[tokio::test]
async fn deep_invalid_url_path_returns_400() {
let (state, _tmp) = test_state();
let app = build_admin_router(state);
let response = app
.oneshot(
Request::builder()
.uri("/not/a/valid/volume/path")
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::BAD_REQUEST);
}
#[tokio::test]
async fn admin_root_get_returns_400() {
let (state, _tmp) = test_state();
let app = build_admin_router(state);
let response = app
.oneshot(Request::builder().uri("/").body(Body::empty()).unwrap())
.await
.unwrap();
assert_eq!(response.status(), StatusCode::BAD_REQUEST);
}
#[tokio::test]
async fn public_root_get_returns_400() {
let (state, _tmp) = test_state();
let app = build_public_router(state);
let response = app
.oneshot(Request::builder().uri("/").body(Body::empty()).unwrap())
.await
.unwrap();
assert_eq!(response.status(), StatusCode::BAD_REQUEST);
}
#[tokio::test]
async fn public_router_does_not_expose_healthz() {
let (state, _tmp) = test_state();
let app = build_public_router(state);
let response = app
.oneshot(
Request::builder()
.uri("/healthz")
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::BAD_REQUEST);
}
// Go's volume_server.go has /stats/* endpoints commented out (L130-134).
#[tokio::test]
async fn admin_router_stats_routes_not_registered() {
let (state, _tmp) = test_state();
let app = build_admin_router(state);
let response = app
.oneshot(
Request::builder()
.uri("/stats/counter")
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
// Falls through to store handler → 400 (bad volume id)
assert_eq!(response.status(), StatusCode::BAD_REQUEST);
}
#[tokio::test]
async fn admin_router_hides_ui_when_write_jwt_is_configured() {
let (state, _tmp) = test_state_with_signing_key(b"secret".to_vec());
let app = build_admin_router(state);
let response = app
.oneshot(
Request::builder()
.uri("/ui/index.html")
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::BAD_REQUEST);
}
#[tokio::test]
async fn admin_router_can_expose_ui_with_explicit_override() {
let (state, _tmp) = test_state_with_signing_key(b"secret".to_vec());
let app = build_admin_router_with_ui(state, true);
let response = app
.oneshot(
Request::builder()
.uri("/ui/index.html")
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::OK);
let body = body_bytes(response).await;
let html = String::from_utf8(body).unwrap();
assert!(html.contains("Disk Stats"));
assert!(html.contains("System Stats"));
assert!(html.contains("Volumes"));
}
#[tokio::test]
async fn admin_router_ui_override_ignores_read_jwt_checks() {
let (state, _tmp) = test_state_with_signing_key(b"write-secret".to_vec());
state.guard.write().unwrap().read_signing_key = SigningKey(b"read-secret".to_vec());
let app = build_admin_router_with_ui(state, true);
let response = app
.oneshot(
Request::builder()
.uri("/ui/index.html")
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::OK);
}
#[tokio::test]
async fn admin_router_serves_volume_ui_static_assets() {
let (state, _tmp) = test_state();
let app = build_admin_router(state);
let response = app
.oneshot(
Request::builder()
.uri("/seaweedfsstatic/bootstrap/3.3.1/css/bootstrap.min.css")
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::OK);
assert_eq!(
response
.headers()
.get("content-type")
.and_then(|value| value.to_str().ok()),
Some("text/css; charset=utf-8")
);
let body = body_bytes(response).await;
assert!(body.len() > 1000);
}
// ============================================================================
// Replicated (fan-out) writes
//
// A chunked S3 upload has the gateway write every replica holder directly with
// `type=replicate` instead of relaying through the primary volume. Each holder
// must accept the copy and store it locally without re-replicating. These tests
// pin that contract on the Rust volume server.
// ============================================================================
/// Raw-body `type=replicate` write is stored and reads back.
#[tokio::test]
async fn replicate_write_raw_body_is_stored() {
let (state, _tmp) = test_state();
let uri = "/1,01637037d6?type=replicate";
let payload = b"fan-out replica bytes";
let app = build_admin_router(state.clone());
let response = app
.oneshot(
Request::builder()
.method("POST")
.uri(uri)
.body(Body::from(payload.to_vec()))
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::CREATED);
let app = build_admin_router(state.clone());
let response = app
.oneshot(
Request::builder()
.uri("/1,01637037d6")
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::OK);
assert_eq!(body_bytes(response).await, payload);
}
/// Go reads fsync through r.FormValue, which decodes the query, so a
/// percent-encoded value has to reach the write path here too.
#[tokio::test]
async fn write_with_percent_encoded_fsync_is_accepted() {
let (state, _tmp) = test_state();
let payload = b"encoded durable payload";
let app = build_admin_router(state.clone());
let response = app
.oneshot(
Request::builder()
.method("POST")
.uri("/1,01637037d6?fsync=%74rue")
.body(Body::from(payload.to_vec()))
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::CREATED);
let app = build_admin_router(state.clone());
let response = app
.oneshot(
Request::builder()
.uri("/1,01637037d6")
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::OK);
assert_eq!(body_bytes(response).await, payload);
}
/// The fan-out query a Go primary sends for a durable write: `fsync=true` rides
/// along with `type=replicate`, and the replica has to honor it rather than ack
/// out of the page cache.
#[tokio::test]
async fn replicate_write_with_fsync_is_stored() {
let (state, _tmp) = test_state();
let uri = "/1,01637037d6?fsync=true&type=replicate";
let payload = b"durable replica bytes";
let app = build_admin_router(state.clone());
let response = app
.oneshot(
Request::builder()
.method("POST")
.uri(uri)
.body(Body::from(payload.to_vec()))
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::CREATED);
let app = build_admin_router(state.clone());
let response = app
.oneshot(
Request::builder()
.uri("/1,01637037d6")
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::OK);
assert_eq!(body_bytes(response).await, payload);
}
/// Multipart `type=replicate` write (the shape the Go gateway uploader sends)
/// is stored and reads back.
#[tokio::test]
async fn replicate_write_multipart_body_is_stored() {
let (state, _tmp) = test_state();
let boundary = "----replicateboundary";
let payload = b"fan-out replica bytes";
let mut body = Vec::new();
body.extend_from_slice(format!("--{}\r\n", boundary).as_bytes());
body.extend_from_slice(
b"Content-Disposition: form-data; name=\"file\"; filename=\"chunk\"\r\n",
);
body.extend_from_slice(b"Content-Type: application/octet-stream\r\n\r\n");
body.extend_from_slice(payload);
body.extend_from_slice(format!("\r\n--{}--\r\n", boundary).as_bytes());
let app = build_admin_router(state.clone());
let response = app
.oneshot(
Request::builder()
.method("POST")
.uri("/1,01637037d6?type=replicate")
.header(
"Content-Type",
format!("multipart/form-data; boundary={}", boundary),
)
.body(Body::from(body))
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::CREATED);
let app = build_admin_router(state.clone());
let response = app
.oneshot(
Request::builder()
.uri("/1,01637037d6")
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::OK);
assert_eq!(body_bytes(response).await, payload);
}
/// On a multi-copy volume, a `type=replicate` write must store locally without
/// re-replicating: it must not contact the master, even an unreachable one.
/// A plain write to the same volume does attempt replication and fails against
/// the dead master, which proves the `type=replicate` flag is what suppresses
/// the fan-out (not a disabled replication path).
#[tokio::test]
async fn replicate_write_does_not_re_replicate() {
// Port 0 is rejected by the socket layer immediately, so a stray lookup
// fails fast instead of hanging on a connect timeout in firewalled CI.
let dead_master = "127.0.0.1:0".to_string();
let (state, _tmp) = build_test_state(Vec::new(), Vec::new(), Some("001"), dead_master);
// type=replicate: stored locally, no master lookup -> 201.
let app = build_admin_router(state.clone());
let response = app
.oneshot(
Request::builder()
.method("POST")
.uri("/1,01637037d6?type=replicate")
.body(Body::from(b"replicated copy".to_vec()))
.unwrap(),
)
.await
.unwrap();
assert_eq!(
response.status(),
StatusCode::CREATED,
"replicate write must not depend on the master"
);
// plain write: tries to fan out to the dead master and fails.
let app = build_admin_router(state.clone());
let response = app
.oneshot(
Request::builder()
.method("POST")
.uri("/1,02637037d7")
.body(Body::from(b"primary copy".to_vec()))
.unwrap(),
)
.await
.unwrap();
assert_eq!(
response.status(),
StatusCode::INTERNAL_SERVER_ERROR,
"plain write to a multi-copy volume must attempt replication"
);
}
// ============================================================================
// Chunk-manifest expansion resolves chunks on EC volumes
//
// A chunked object whose data chunks live on an EC-encoded volume must expand
// by reconstruct-on-read from the shards, not by a local regular-volume lookup
// (which finds nothing once the volume is EC-encoded). Mirrors Go's
// ChunkedFileReader, which resolves every chunk through the master.
// ============================================================================
#[tokio::test]
async fn chunk_manifest_expands_chunk_stored_on_ec_volume() {
use seaweed_volume::storage::erasure_coding::ec_encoder::write_ec_files;
use seaweed_volume::storage::erasure_coding::ec_shard::ShardId;
use seaweed_volume::storage::needle::needle::{FileId, Needle};
use seaweed_volume::storage::types::{Cookie, NeedleId};
use seaweed_volume::storage::volume::{Volume, VolumeSpec};
let (state, tmp) = test_state();
let dir = tmp.path().to_str().unwrap();
// A chunk large enough to be worth EC-encoding; its bytes are the payload we
// expect the manifest GET to return.
let chunk_data: Vec<u8> = (0..4096u32).map(|i| (i % 251) as u8).collect();
let chunk_nid = NeedleId(0x2a);
let chunk_cookie = Cookie(0x1234abcd);
// Build regular volume 2 holding the chunk, then EC-encode it and mount all
// 14 shards locally so reconstruct-on-read is a pure local read.
{
let mut v = Volume::new(
dir,
dir,
VolumeId(2),
NeedleMapKind::InMemory,
&VolumeSpec::default(),
)
.unwrap();
let mut n = Needle {
id: chunk_nid,
cookie: chunk_cookie,
data: chunk_data.clone(),
data_size: chunk_data.len() as u32,
..Needle::default()
};
v.write_needle(&mut n, true, false).unwrap();
v.sync_to_disk().unwrap();
v.close();
}
write_ec_files(dir, dir, "", VolumeId(2), 10, 4).unwrap();
// Volume 2 was built standalone (never registered in the store), so it only
// exists as EC shards — the chunk resolves through the EC path, as it would
// after ec.encode retired the regular volume.
{
let mut store = state.store.write().unwrap();
let shard_ids: Vec<ShardId> = (0..14).collect();
store.mount_ec_shards(VolumeId(2), "", &shard_ids).unwrap();
}
// Write the chunk-manifest needle to regular volume 1.
let chunk_fid = FileId::new(VolumeId(2), chunk_nid, chunk_cookie).to_string();
let manifest = format!(
r#"{{"name":"big.bin","mime":"application/octet-stream","size":{},"chunks":[{{"fid":"{}","offset":0,"size":{}}}]}}"#,
chunk_data.len(),
chunk_fid,
chunk_data.len()
);
let manifest_nid = NeedleId(0x7);
let manifest_cookie = Cookie(0x55667788);
{
let mut store = state.store.write().unwrap();
let mut n = Needle {
id: manifest_nid,
cookie: manifest_cookie,
data: manifest.into_bytes(),
..Needle::default()
};
n.data_size = n.data.len() as u32;
n.set_is_chunk_manifest();
store
.write_volume_needle(VolumeId(1), &mut n, false)
.unwrap();
}
// GET the manifest object; expect the reconstructed chunk bytes.
let manifest_fid = FileId::new(VolumeId(1), manifest_nid, manifest_cookie).to_string();
let app = build_admin_router(state.clone());
let response = app
.oneshot(
Request::builder()
.uri(format!("/{}", manifest_fid))
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::OK);
assert_eq!(body_bytes(response).await, chunk_data);
}
// ============================================================================
// HTTP DELETE on an EC volume whose shards are not all mounted locally
//
// The delete handler used to validate the cookie with the local-only
// `EcVolume::read_ec_shard_needle`, which errors "ec shard N not available
// locally" for any interval held by a peer. Every such error was mapped to 500
// and no `.ecj` tombstone was written, so on a standard 10+4 spread across 14
// servers no HTTP delete of an EC needle could ever succeed.
//
// A single node can exercise that path without peers: mount 13 of the 14
// shards, leaving out the one holding the needle's interval. The distributed
// reader seeds its Reed-Solomon buffers from locally mounted siblings (Phase 0
// in `store_ec.rs`), so with >= 10 survivors it reconstructs without any peer
// fan-out — while the local-only read still fails outright.
// ============================================================================
#[tokio::test]
async fn delete_on_ec_volume_succeeds_when_the_needles_shard_is_not_mounted() {
use seaweed_volume::storage::erasure_coding::ec_encoder::write_ec_files;
use seaweed_volume::storage::erasure_coding::ec_shard::ShardId;
use seaweed_volume::storage::needle::needle::{FileId, Needle};
use seaweed_volume::storage::types::{Cookie, NeedleId};
use seaweed_volume::storage::volume::{Volume, VolumeSpec};
let (state, tmp) = test_state();
let dir = tmp.path().to_str().unwrap();
let data: Vec<u8> = (0..4096u32).map(|i| (i % 251) as u8).collect();
let nid = NeedleId(0x5c);
let cookie = Cookie(0x0badc0de);
// Build regular volume 4, then EC-encode it. The volume is standalone and
// never registered in the store, so afterwards it exists only as shards.
{
let mut v = Volume::new(
dir,
dir,
VolumeId(4),
NeedleMapKind::InMemory,
&VolumeSpec::default(),
)
.unwrap();
let mut n = Needle {
id: nid,
cookie,
data: data.clone(),
data_size: data.len() as u32,
..Needle::default()
};
v.write_needle(&mut n, true, false).unwrap();
v.sync_to_disk().unwrap();
v.close();
}
write_ec_files(dir, dir, "", VolumeId(4), 10, 4).unwrap();
// Mount shards 1..=13 only. A needle at .dat offset 0 lives in shard 0's
// first small block, so the interval this delete needs is deliberately the
// one shard that is absent.
{
let mut store = state.store.write().unwrap();
let shard_ids: Vec<ShardId> = (1..14).collect();
store.mount_ec_shards(VolumeId(4), "", &shard_ids).unwrap();
}
let fid = FileId::new(VolumeId(4), nid, cookie).to_string();
let app = build_admin_router(state.clone());
let response = app
.oneshot(
Request::builder()
.method("DELETE")
.uri(format!("/{}", fid))
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(
response.status(),
StatusCode::ACCEPTED,
"DELETE of an EC needle must reconstruct through the distributed \
reader instead of failing 500 on a non-local shard"
);
// The tombstone must actually have landed: a later GET is a 404.
let app = build_admin_router(state.clone());
let response = app
.oneshot(
Request::builder()
.uri(format!("/{}", fid))
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(
response.status(),
StatusCode::NOT_FOUND,
"the delete must have been journalled, not just answered 202"
);
}
// ============================================================================
// Hostile response-header override params must not panic the handler
//
// The `response-*` query params are attacker-controlled and were inserted with
// `parse().unwrap()`. `%0A` decodes to a newline, `HeaderValue::from_str`
// rejects it, and the unwrap panicked the connection task — unauthenticated.
// The override must simply be skipped.
// ============================================================================
#[tokio::test]
async fn hostile_response_header_overrides_are_skipped_not_panicked() {
let (state, _tmp) = test_state();
let uri = "/1,01637037d6";
let app = build_admin_router(state.clone());
let response = app
.oneshot(
Request::builder()
.method("POST")
.uri(uri)
.body(Body::from(b"payload".to_vec()))
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::CREATED);
// One request per override param, each carrying a raw newline.
for param in [
"response-cache-control",
"response-content-encoding",
"response-expires",
"response-content-language",
"response-content-disposition",
"response-content-type",
] {
let app = build_admin_router(state.clone());
let response = app
.oneshot(
Request::builder()
.uri(format!("{}?{}=%0Aevil", uri, param))
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(
response.status(),
StatusCode::OK,
"{} with a newline must be ignored, not panic",
param
);
assert_eq!(body_bytes(response).await, b"payload".to_vec());
}
}
// ============================================================================
// Non-ASCII in the fid and in ?ttl= must be rejected, not panic
//
// `parse_needle_id_cookie` split the hex by BYTE offset and `TTL::read` took
// the unit as the last BYTE, so a multi-byte character split inside itself.
// Both are reachable unauthenticated from the request line / query string.
// ============================================================================
#[tokio::test]
async fn non_ascii_fid_and_ttl_are_rejected_not_panicked() {
let (state, _tmp) = test_state();
// A fid whose hex part is multi-byte UTF-8.
let app = build_admin_router(state.clone());
let response = app
.oneshot(
Request::builder()
.uri("/1,%C3%A9%C3%A9%C3%A9%C3%A9a")
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert!(
response.status().is_client_error() || response.status().is_server_error(),
"non-ASCII fid must produce an error status, got {}",
response.status()
);
// A TTL whose unit character is multi-byte. The upload path does
// `TTL::read(..).ok()`, so *any* unparseable TTL is simply dropped and the
// write succeeds — the point here is that a non-ASCII one now takes that
// same road instead of panicking. Assert it matches an ASCII-invalid TTL
// rather than inventing a stricter contract than the handler has.
let mut statuses = Vec::new();
// Distinct needle ids: reusing one id with a different cookie is a
// cookie-mismatch overwrite, which would mask what this test measures.
for (fid, ttl) in [("/1,03637037d7", "5%C3%A9"), ("/1,04637037d8", "5z")] {
let app = build_admin_router(state.clone());
let response = app
.oneshot(
Request::builder()
.method("POST")
.uri(format!("{}?ttl={}", fid, ttl))
.body(Body::from(b"x".to_vec()))
.unwrap(),
)
.await
.unwrap();
statuses.push(response.status());
}
assert_eq!(
statuses[0], statuses[1],
"a non-ASCII ttl must behave like any other invalid ttl, not panic"
);
}
// ============================================================================
// The write queue answers an upload with the needle's real ETag
//
// The queue worker computes the CRC on the needle it was handed, so the handler
// has to know the checksum before it submits. Without that every queued upload
// came back as "00000000". The direct path is the reference: same payload, same
// ETag, for a plain body and for one the handler gzips before storing.
// ============================================================================
#[tokio::test]
async fn write_queue_upload_returns_same_etag_as_direct_write() {
use seaweed_volume::server::write_queue::WriteQueue;
let (direct_state, _direct_tmp) = test_state();
let (queued_state, _queued_tmp) = test_state();
let wq = WriteQueue::new(queued_state.clone(), 128);
let _ = queued_state.write_queue.set(wq);
let compressible = "seaweedfs ".repeat(200).into_bytes();
let uploads: [(&str, &[u8]); 2] = [
("/1,01637037d6", b"hello, seaweedfs!"),
("/1/02637037d6/notes.txt", &compressible),
];
for (uri, payload) in uploads {
let mut etags = Vec::new();
for state in [&direct_state, &queued_state] {
let app = build_admin_router(state.clone());
let response = app
.oneshot(
Request::builder()
.method("POST")
.uri(uri)
.body(Body::from(payload.to_vec()))
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::CREATED);
let header = response
.headers()
.get("ETag")
.expect("upload response has no ETag")
.to_str()
.unwrap()
.to_string();
let body = body_bytes(response).await;
let json: serde_json::Value =
serde_json::from_slice(&body).expect("POST response is not valid JSON");
let etag = json["eTag"].as_str().unwrap().to_string();
assert_eq!(header, format!("\"{}\"", etag));
etags.push(etag);
}
assert_ne!(etags[0], "00000000", "{}: direct ETag is the zero CRC", uri);
assert_eq!(
etags[1], etags[0],
"{}: queued upload must return the direct path's ETag",
uri
);
}
// The second upload really was stored gzipped, so its ETag is the CRC of
// the compressed bytes on both paths.
let app = build_admin_router(queued_state.clone());
let response = app
.oneshot(
Request::builder()
.uri(uploads[1].0)
.header("Accept-Encoding", "gzip")
.body(Body::empty())
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::OK);
assert_eq!(response.headers()["Content-Encoding"], "gzip");
// Re-uploading the same bytes is the unchanged path: 204 with the same ETag.
let (uri, payload) = uploads[0];
let mut etags = Vec::new();
for state in [&direct_state, &queued_state] {
let app = build_admin_router(state.clone());
let response = app
.oneshot(
Request::builder()
.method("POST")
.uri(uri)
.body(Body::from(payload.to_vec()))
.unwrap(),
)
.await
.unwrap();
assert_eq!(response.status(), StatusCode::NO_CONTENT);
etags.push(response.headers()["ETag"].to_str().unwrap().to_string());
}
assert_eq!(etags[1], etags[0], "unchanged upload ETag differs");
}