mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-02 04:32:42 +00:00
Compare commits
57
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
95af3f3649 | ||
|
|
7e6db15a88 | ||
|
|
e66983e9bd | ||
|
|
b010af88fb | ||
|
|
94cefd6f4c | ||
|
|
47156eb8ce | ||
|
|
1ce0174b23 | ||
|
|
120aa956a0 | ||
|
|
d6ff6ed6d4 | ||
|
|
532bf863ad | ||
|
|
23e4497b21 | ||
|
|
b0447d2479 | ||
|
|
fbff2cb39a | ||
|
|
d2a6066181 | ||
|
|
7e6e0261ab | ||
|
|
c7c7be42ed | ||
|
|
2e65966c06 | ||
|
|
61befd10fc | ||
|
|
70ddbee370 | ||
|
|
14c863dbff | ||
|
|
6bb9d8bac2 | ||
|
|
cc80ad3643 | ||
|
|
9009e38f7b | ||
|
|
b9fbb85af2 | ||
|
|
47d3001572 | ||
|
|
a12dd5f8d3 | ||
|
|
8e614486a3 | ||
|
|
a5864c3eb6 | ||
|
|
6302809442 | ||
|
|
27a80f7607 | ||
|
|
ec429e0361 | ||
|
|
90e82b15ce | ||
|
|
a3e1ee1653 | ||
|
|
2ab30900d4 | ||
|
|
62ee14fa61 | ||
|
|
ab95a6ef15 | ||
|
|
24965fd489 | ||
|
|
ed23e290fc | ||
|
|
9b57fb6961 | ||
|
|
1bb40b6bc5 | ||
|
|
34e342da63 | ||
|
|
4835d34438 | ||
|
|
5814729def | ||
|
|
37bf9b5ebf | ||
|
|
19201df6d7 | ||
|
|
4d61cbdeed | ||
|
|
3ce883624e | ||
|
|
de974c05d5 | ||
|
|
7768fda023 | ||
|
|
548b3d9a38 | ||
|
|
a7f50d23b5 | ||
|
|
6ce4d7eded | ||
|
|
3bd20e6a10 | ||
|
|
d402573ea8 | ||
|
|
63d08e8a91 | ||
|
|
880c2e1dab | ||
|
|
7beab85c21 |
@@ -9,6 +9,7 @@ on:
|
||||
- 'weed/storage/**'
|
||||
- 'weed/pb/volume_server.proto'
|
||||
- 'weed/pb/volume_server_pb/**'
|
||||
- 'rust/volume_server/**'
|
||||
- '.github/workflows/volume-server-integration-tests.yml'
|
||||
push:
|
||||
branches: [ master, main ]
|
||||
@@ -18,6 +19,7 @@ on:
|
||||
- 'weed/storage/**'
|
||||
- 'weed/pb/volume_server.proto'
|
||||
- 'weed/pb/volume_server_pb/**'
|
||||
- 'rust/volume_server/**'
|
||||
- '.github/workflows/volume-server-integration-tests.yml'
|
||||
|
||||
concurrency:
|
||||
@@ -120,3 +122,55 @@ jobs:
|
||||
echo "## Volume Server Integration Test Summary (${{ matrix.test-type }} - Shard ${{ matrix.shard }})" >> "$GITHUB_STEP_SUMMARY"
|
||||
echo "- Suite: test/volume_server/${{ matrix.test-type }} (Pattern: ${TEST_PATTERN})" >> "$GITHUB_STEP_SUMMARY"
|
||||
echo "- Command: go test -v -count=1 -timeout=${{ env.TEST_TIMEOUT }} ./test/volume_server/${{ matrix.test-type }}/... -run \"${TEST_PATTERN}\"" >> "$GITHUB_STEP_SUMMARY"
|
||||
|
||||
volume-server-rust-smoke:
|
||||
name: Volume Server Rust Smoke (${{ matrix.mode }})
|
||||
runs-on: ubuntu-22.04
|
||||
timeout-minutes: 40
|
||||
strategy:
|
||||
fail-fast: false
|
||||
matrix:
|
||||
mode: [exec, proxy, native]
|
||||
|
||||
steps:
|
||||
- name: Checkout code
|
||||
uses: actions/checkout@v6
|
||||
|
||||
- name: Set up Go ${{ env.GO_VERSION }}
|
||||
uses: actions/setup-go@v6
|
||||
with:
|
||||
go-version: ${{ env.GO_VERSION }}
|
||||
|
||||
- name: Set up Rust toolchain
|
||||
uses: dtolnay/rust-toolchain@stable
|
||||
|
||||
- name: Build SeaweedFS binary
|
||||
run: |
|
||||
cd weed
|
||||
go build -o weed .
|
||||
chmod +x weed
|
||||
./weed version
|
||||
|
||||
- name: Run Rust-mode HTTP smoke tests
|
||||
env:
|
||||
WEED_BINARY: ${{ github.workspace }}/weed/weed
|
||||
VOLUME_SERVER_IMPL: rust
|
||||
VOLUME_SERVER_RUST_MODE: ${{ matrix.mode }}
|
||||
run: |
|
||||
go test -v -count=1 -timeout=25m ./test/volume_server/http/... -run "TestAdminStatusAndHealthz|TestUploadReadRangeHeadDeleteRoundTrip"
|
||||
|
||||
- name: Run Rust-mode gRPC smoke tests
|
||||
env:
|
||||
WEED_BINARY: ${{ github.workspace }}/weed/weed
|
||||
VOLUME_SERVER_IMPL: rust
|
||||
VOLUME_SERVER_RUST_MODE: ${{ matrix.mode }}
|
||||
run: |
|
||||
go test -v -count=1 -timeout=25m ./test/volume_server/grpc/... -run "TestStateAndStatusRPCs|TestVolumeSyncStatusAndReadVolumeFileStatus"
|
||||
|
||||
- name: Rust smoke summary
|
||||
if: always()
|
||||
run: |
|
||||
echo "## Volume Server Rust Smoke Summary (${{ matrix.mode }})" >> "$GITHUB_STEP_SUMMARY"
|
||||
echo "- Mode: Go master + Rust volume server launcher (VOLUME_SERVER_IMPL=rust, VOLUME_SERVER_RUST_MODE=${{ matrix.mode }})" >> "$GITHUB_STEP_SUMMARY"
|
||||
echo "- HTTP Command: go test -v -count=1 -timeout=25m ./test/volume_server/http/... -run \"TestAdminStatusAndHealthz|TestUploadReadRangeHeadDeleteRoundTrip\"" >> "$GITHUB_STEP_SUMMARY"
|
||||
echo "- gRPC Command: go test -v -count=1 -timeout=25m ./test/volume_server/grpc/... -run \"TestStateAndStatusRPCs|TestVolumeSyncStatusAndReadVolumeFileStatus\"" >> "$GITHUB_STEP_SUMMARY"
|
||||
|
||||
Generated
+7
@@ -0,0 +1,7 @@
|
||||
# This file is automatically @generated by Cargo.
|
||||
# It is not intended for manual editing.
|
||||
version = 4
|
||||
|
||||
[[package]]
|
||||
name = "weed-volume-rs"
|
||||
version = "0.1.0"
|
||||
@@ -0,0 +1,6 @@
|
||||
[package]
|
||||
name = "weed-volume-rs"
|
||||
version = "0.1.0"
|
||||
edition = "2021"
|
||||
|
||||
[dependencies]
|
||||
@@ -0,0 +1,245 @@
|
||||
# Rust Volume Server Parity Implementation Plan
|
||||
|
||||
## Objective
|
||||
Implement a native Rust volume server that replicates Go volume-server behavior for HTTP and gRPC APIs, so it can become a drop-in replacement validated by existing integration suites.
|
||||
|
||||
## Current Focus (2026-02-16)
|
||||
- Program focus is now Rust implementation parity, not broad test expansion.
|
||||
- `test/volume_server` is treated as the parity gate.
|
||||
- Existing Rust launcher modes (`exec`, `proxy`) are transition tools; they are not the final target.
|
||||
|
||||
## Current Status
|
||||
- Rust crate and launcher are in place.
|
||||
- Integration harness can run:
|
||||
- Go master + Go volume (default)
|
||||
- Go master + Rust launcher (`VOLUME_SERVER_IMPL=rust`)
|
||||
- Rust launcher `proxy` mode has full-suite integration pass while delegating backend handlers to Go.
|
||||
- Rust launcher `native` mode is wired as the default Rust entrypoint and currently bootstraps via supervised Go backend delegation.
|
||||
- Native Rust HTTP handlers now serve control/surface paths in `native` mode:
|
||||
- `/status`, `/healthz`
|
||||
- `OPTIONS` admin/public behavior (method allow-list + CORS preflight headers)
|
||||
- `/ui/index.html` with config-driven JWT/access-ui gating
|
||||
- `/favicon.ico` and `/seaweedfsstatic/*`
|
||||
- public non-read methods (`POST`/`PUT`/`DELETE`/unsupported verbs) as `200` no-op passthrough parity
|
||||
- admin unsupported verbs as `400` parity
|
||||
- Native Rust control-surface parity now also includes:
|
||||
- `/healthz` status mirroring from backend state transitions (e.g. leave/stopping => `503`)
|
||||
- absolute-form HTTP request-target normalization before native route matching
|
||||
- Native Rust HTTP data-path prevalidation now includes:
|
||||
- early malformed vid/fid rejection (`400`) for GET/HEAD/POST/PUT fid-route shapes before delegation
|
||||
- slash-form fid-route parsing for `/{vid}/{fid}` and `/{vid}/{fid}/{filename}` with reserved-path exclusions (`/status`, `/healthz`, `/ui/index.html`, `/stats/*`, static assets)
|
||||
- write-path error prevalidation for fid routes:
|
||||
- malformed multipart form-data requests without a boundary => `400`
|
||||
- requests containing `Content-MD5` header => `400` mismatch parity branch
|
||||
- request bodies exceeding configured `-fileSizeLimitMB` (by `Content-Length`) => `400`
|
||||
- Native Rust API/storage handlers are not implemented yet.
|
||||
|
||||
## Parity Exit Criteria
|
||||
1. Native mode passes:
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=native go test -count=1 ./test/volume_server/http`
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=native go test -count=1 ./test/volume_server/grpc`
|
||||
2. CI runs native Rust mode integration coverage (at least smoke, then expanded shards).
|
||||
3. Rust mode defaults to native behavior for integration harness.
|
||||
4. Go-backend delegation is removed (or retained only as explicit fallback mode).
|
||||
|
||||
## Architecture Workstreams
|
||||
|
||||
### A. Runtime and Configuration Parity
|
||||
- [x] Add `native` runtime mode in `weed-volume-rs` (bootstrap delegation path).
|
||||
- [ ] Parse and honor volume-server CLI/config flags used by integration harness:
|
||||
- [ ] network/bind ports (`-ip`, `-port`, `-port.grpc`, `-port.public`)
|
||||
- [ ] master target/config dir/read mode/throttling/JWT-related config
|
||||
- [ ] size/timeout controls and maintenance state defaults
|
||||
- [ ] Implement graceful lifecycle behavior (signals, shutdown, readiness).
|
||||
|
||||
### B. Native HTTP Surface
|
||||
- [ ] Admin/control endpoints:
|
||||
- [x] `GET /status` (native Rust in `native` mode)
|
||||
- [x] `GET /healthz` (native Rust in `native` mode)
|
||||
- [x] `OPTIONS` admin/public method+CORS control behavior (native Rust in `native` mode)
|
||||
- [x] static/UI endpoints currently exercised
|
||||
- [ ] Data read path parity:
|
||||
- [ ] fid parsing/path variants
|
||||
- [ ] conditional headers (`If-Modified-Since`, `If-None-Match`)
|
||||
- [ ] range handling (single/multi/invalid)
|
||||
- [ ] deleted reads, auth checks, read-mode branches
|
||||
- [ ] chunk-manifest and compression/image transformation branches
|
||||
- [ ] Data write/delete parity:
|
||||
- [ ] write success/unchanged/error paths
|
||||
- [ ] replication and file-size-limit paths
|
||||
- [ ] delete and chunk-manifest delete branches
|
||||
- [ ] Method/CORS/public-port parity for split admin/public behavior.
|
||||
|
||||
### C. Native gRPC Surface
|
||||
- [ ] Control-plane RPCs:
|
||||
- [ ] `GetState`, `SetState`, `VolumeServerStatus`, `Ping`, `VolumeServerLeave`
|
||||
- [ ] admin lifecycle: allocate/mount/unmount/delete/configure/readonly/writable
|
||||
- [ ] Data RPCs:
|
||||
- [ ] `ReadNeedleBlob`, `ReadNeedleMeta`, `WriteNeedleBlob`
|
||||
- [ ] `BatchDelete`, `ReadAllNeedles`
|
||||
- [ ] sync/copy/receive and status endpoints
|
||||
- [ ] Stream RPCs:
|
||||
- [ ] tail sender/receiver
|
||||
- [ ] vacuum streams
|
||||
- [ ] query streams
|
||||
- [ ] Advanced families:
|
||||
- [ ] erasure coding RPC set
|
||||
- [ ] tiering/remote fetch
|
||||
- [ ] scrub/query mode matrix
|
||||
|
||||
### D. Storage Compatibility Layer
|
||||
- [ ] Implement volume data/index handling compatible with Go on-disk format.
|
||||
- [ ] Preserve cookie/checksum/timestamp semantics used by tests.
|
||||
- [ ] Match read/write/delete consistency and error mapping behavior.
|
||||
- [ ] Ensure EC metadata/data-path compatibility with existing files.
|
||||
|
||||
### E. Operational Hardening
|
||||
- [ ] Deterministic startup/readiness and shutdown semantics.
|
||||
- [ ] Log/error parity sufficient for debugging and CI triage.
|
||||
- [ ] Concurrency/timeout behavior alignment for throttling and streams.
|
||||
- [ ] Performance baseline checks vs Go for key flows.
|
||||
|
||||
## Milestone Plan
|
||||
|
||||
### M0 (Completed): Harness + Launcher Transition
|
||||
- [x] Rust launcher integrated into harness.
|
||||
- [x] Proxy mode full-suite validation with Go backend delegation.
|
||||
|
||||
### M1: Native Skeleton (Control Plane First)
|
||||
- [ ] `native` mode boots and serves:
|
||||
- [x] `/status`, `/healthz`
|
||||
- [ ] `GetState`, `SetState`, `VolumeServerStatus`, `Ping`, `VolumeServerLeave`
|
||||
- Gate:
|
||||
- [x] targeted HTTP/grpc control tests pass in `native` mode (delegated backend path).
|
||||
|
||||
### M2: Native Core Data Paths
|
||||
- [ ] Native HTTP read/write/delete baseline parity.
|
||||
- [ ] Native gRPC data baseline parity (`Read/WriteNeedle*`, `BatchDelete`, `ReadAllNeedles`).
|
||||
- Gate:
|
||||
- core HTTP and gRPC data suites pass in `native` mode.
|
||||
|
||||
### M3: Native Stream + Copy/Sync
|
||||
- [ ] Tail/copy/receive/sync paths in native mode.
|
||||
- Gate:
|
||||
- stream/copy families pass in `native` mode.
|
||||
|
||||
### M4: Native Advanced Feature Families
|
||||
- [ ] EC, tiering, scrub/query advanced branches.
|
||||
- Gate:
|
||||
- full `/test/volume_server/http` and `/test/volume_server/grpc` pass in `native` mode.
|
||||
|
||||
### M5: CI/Cutover
|
||||
- [x] Add/expand native-mode CI jobs (smoke matrix includes `native`).
|
||||
- [x] Make native mode default for Rust integration runs.
|
||||
- [ ] Keep `exec`/`proxy` only as explicit fallback modes during rollout.
|
||||
|
||||
## Immediate Next Steps
|
||||
1. Implement minimal native gRPC state/ping RPC handlers (`GetState`, `SetState`, `VolumeServerStatus`, `Ping`, `VolumeServerLeave`).
|
||||
2. Expand native HTTP parity from control/static surface into data-path handlers (read/write/delete variants) while preserving current native control behavior.
|
||||
3. Keep rerunning native-mode integration suites as each delegated branch is replaced.
|
||||
4. Add mismatch triage notes for each API moved from delegation to native implementation.
|
||||
|
||||
## Risk Register
|
||||
- On-disk format mismatch risk:
|
||||
- Mitigation: implement format-level compatibility tests early (idx/dat/needle encoding).
|
||||
- Behavioral drift in edge branches:
|
||||
- Mitigation: use integration suite failures as primary truth; only add tests for newly discovered untracked branches.
|
||||
- Stream/concurrency semantic mismatch:
|
||||
- Mitigation: stabilize with focused interruption/timeout parity tests.
|
||||
|
||||
## Progress Log
|
||||
- Date: 2026-02-15
|
||||
- Change: Added Rust launcher integration (`exec`) and harness wiring.
|
||||
- Validation: Rust launcher mode passed smoke and full integration suites while delegating to Go backend.
|
||||
- Commits: `7beab85c2`, `880c2e1da`, `63d08e8a9`, `d402573ea`, `3bd20e6a1`, `6ce4d7ede`
|
||||
|
||||
- Date: 2026-02-15
|
||||
- Change: Added Rust proxy supervisor mode and validated full integration suite.
|
||||
- Validation:
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=proxy go test -count=1 ./test/volume_server/http`
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=proxy go test -count=1 ./test/volume_server/grpc`
|
||||
- Commits: `a7f50d23b`, `548b3d9a3`
|
||||
|
||||
- Date: 2026-02-16
|
||||
- Change: Re-focused plan from test expansion to native Rust implementation parity.
|
||||
- Validation basis: latest Rust proxy full-suite pass keeps regression baseline stable while native implementation starts.
|
||||
- Commits: `14c863dbf`
|
||||
|
||||
- Date: 2026-02-16
|
||||
- Change: Added native Rust launcher mode bootstrap, set Rust launcher default mode to `native`, and expanded CI Rust smoke matrix to include `native`.
|
||||
- Validation:
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=native go test -count=1 ./test/volume_server/http`
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=native go test -count=1 ./test/volume_server/grpc`
|
||||
- `env VOLUME_SERVER_IMPL=rust go test -count=1 ./test/volume_server/http -run '^TestAdminStatusAndHealthz$'`
|
||||
- `env VOLUME_SERVER_IMPL=rust go test -count=1 ./test/volume_server/grpc -run '^TestStateAndStatusRPCs$'`
|
||||
- Commits: `70ddbee37`, `61befd10f`, `2e65966c0`
|
||||
|
||||
- Date: 2026-02-16
|
||||
- Change: Implemented first native Rust HTTP handlers in `native` mode for `/status` and `/healthz` on the admin listener while preserving proxy delegation for other APIs.
|
||||
- Validation:
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=native go test -count=1 ./test/volume_server/http -run '^TestAdminStatusAndHealthz$'`
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=native go test -count=1 ./test/volume_server/http`
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=native go test -count=1 ./test/volume_server/grpc`
|
||||
- Commits: `7e6e0261a`
|
||||
|
||||
- Date: 2026-02-16
|
||||
- Change: Implemented native Rust `OPTIONS` handling for admin and public listeners in `native` mode, including method allow-list and origin-driven CORS response parity used by integration tests.
|
||||
- Validation:
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=native go test -count=1 ./test/volume_server/http -run '^TestOptionsMethodsByPort$|^TestOptionsWithOriginIncludesCorsHeaders$'`
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=native go test -count=1 ./test/volume_server/http`
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=native go test -count=1 ./test/volume_server/grpc`
|
||||
- Commits: `fbff2cb39`
|
||||
|
||||
- Date: 2026-02-16
|
||||
- Change: Added native Rust UI/static/public-surface handlers in `native` mode:
|
||||
- `/ui/index.html` with config-driven JWT/access-ui gating parity
|
||||
- `/favicon.ico` and `/seaweedfsstatic/*` static endpoint parity
|
||||
- public non-read methods as `200` no-op passthrough parity
|
||||
- admin unsupported methods as `400` parity
|
||||
- Validation:
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=native go test -count=1 ./test/volume_server/http -run '^TestAdminStatusAndHealthz$|^TestUiIndexNotExposedWhenJwtSigningEnabled$|^TestUiIndexExposedWhenJwtSigningEnabledAndAccessUITrue$|^TestStaticAssetEndpoints$|^TestStaticAssetEndpointsOnPublicPort$'`
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=native go test -count=1 ./test/volume_server/http -run '^TestOptionsMethodsByPort$|^TestOptionsWithOriginIncludesCorsHeaders$|^TestPublicPortReadOnlyMethodBehavior$|^TestCorsAndUnsupportedMethodBehavior$|^TestUnsupportedMethodTraceParity$|^TestUnsupportedMethodPropfindParity$|^TestUnsupportedMethodConnectParity$|^TestUnsupportedMethodMkcolParity$'`
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=native go test -count=1 ./test/volume_server/http`
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=native go test -count=1 ./test/volume_server/grpc`
|
||||
- Commits: `23e4497b2`
|
||||
|
||||
- Date: 2026-02-16
|
||||
- Change: Hardened native Rust HTTP control-path parity:
|
||||
- `/healthz` now mirrors backend health status (`200`/`503`) in native mode
|
||||
- absolute-form request-target parsing (e.g. `GET http://host/path HTTP/1.1`) now normalizes to route path parity
|
||||
- Validation:
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=native go test -count=1 ./test/volume_server/http -run '^TestAdminStatusAndHealthz$'`
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=native go test -count=1 ./test/volume_server/http/...`
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=native go test -count=1 ./test/volume_server/grpc/...`
|
||||
- Commits: `d6ff6ed6d`
|
||||
|
||||
- Date: 2026-02-16
|
||||
- Change: Added native malformed fid-route validation in `native` mode:
|
||||
- GET/HEAD/POST/PUT requests matching fid URL shapes now return Rust-native `400 Bad Request` for invalid vid/fid tokens
|
||||
- valid fid-route requests still delegate to backend handlers for storage/data execution paths
|
||||
- Validation:
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=native go test -count=1 ./test/volume_server/http -run '^TestInvalidReadPathReturnsBadRequest$|^TestWriteInvalidVidAndFidReturnBadRequest$|^TestMalformedVidFidPathReturnsBadRequest$|^TestDownloadLimitInvalidVidWhileOverLimitReturnsBadRequest$'`
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=native go test -count=1 ./test/volume_server/http/...`
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=native go test -count=1 ./test/volume_server/grpc/...`
|
||||
- Commits: `1ce0174b2`
|
||||
|
||||
- Date: 2026-02-16
|
||||
- Change: Added native Rust write-path prevalidation branches for fid routes in `native` mode:
|
||||
- multipart form uploads without boundary now return native `400`
|
||||
- `Content-MD5` requests now return native `400` mismatch branch
|
||||
- `Content-Length` over configured `-fileSizeLimitMB` now returns native `400` limit branch
|
||||
- Validation:
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=native go test -count=1 ./test/volume_server/http -run '^TestUploadReadRangeHeadDeleteRoundTrip$|^TestWriteMalformedMultipartAndMD5Mismatch$|^TestWriteRejectsPayloadOverFileSizeLimit$'`
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=native go test -count=1 ./test/volume_server/http/...`
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=native go test -count=1 ./test/volume_server/grpc/...`
|
||||
- Commits: `94cefd6f4`
|
||||
|
||||
- Date: 2026-02-16
|
||||
- Change: Broadened native fid-route parser to cover slash-form URLs:
|
||||
- `/{vid}/{fid}` and `/{vid}/{fid}/{filename}` now flow through native malformed vid/fid validation
|
||||
- added explicit non-fid exclusions for control/stats/static paths to preserve route parity
|
||||
- Validation:
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=native go test -count=1 ./test/volume_server/http -run '^TestAdminStatusAndHealthz$|^TestReadPathShapesAndIfModifiedSince$|^TestMalformedVidFidPathReturnsBadRequest$'`
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=native go test -count=1 ./test/volume_server/http/...`
|
||||
- `env VOLUME_SERVER_IMPL=rust VOLUME_SERVER_RUST_MODE=native go test -count=1 ./test/volume_server/grpc/...`
|
||||
- Commits: `e66983e9b`
|
||||
File diff suppressed because it is too large
Load Diff
@@ -3,6 +3,12 @@
|
||||
## Goal
|
||||
Create a Go integration test suite under `test/volume_server` that validates **drop-in behavior parity** for the Volume Server HTTP and gRPC APIs, so a Rust rewrite can be verified against the current Go behavior.
|
||||
|
||||
## Current Program Focus (2026-02-16)
|
||||
- Primary execution focus has shifted to implementing native Rust volume-server parity.
|
||||
- This integration suite is now the parity gate for Rust implementation work.
|
||||
- New tests should be added only when native Rust implementation reveals uncovered Go behavior that is not yet captured.
|
||||
- Rust implementation roadmap lives in `rust/volume_server/DEV_PLAN.md`.
|
||||
|
||||
## Hard Requirements
|
||||
- Tests live under `test/volume_server`.
|
||||
- Tests are written in Go.
|
||||
@@ -1127,3 +1133,202 @@ Update this section during implementation:
|
||||
- Profiles covered: P1.
|
||||
- Gaps introduced/remaining: `499` cancellation and transport-interruption branches remain pending.
|
||||
- Commit: `d1e5f390a`
|
||||
|
||||
- Date: 2026-02-15
|
||||
- Change: Added Rust volume-server migration bootstrap wiring.
|
||||
- APIs covered: Harness now supports Go master + selectable volume binary (`VOLUME_SERVER_IMPL=rust` or `VOLUME_SERVER_BINARY`), with Rust-mode smoke coverage for representative HTTP/gRPC integration tests.
|
||||
- Profiles covered: P1 smoke and default Go matrix unchanged.
|
||||
- Gaps introduced/remaining: native Rust handler implementations are tracked in `rust/volume_server/DEV_PLAN.md`; current Rust mode is a compatibility launcher phase.
|
||||
- Commit: `7beab85c2`, `880c2e1da`, `63d08e8a9`, `d402573ea`
|
||||
|
||||
- Date: 2026-02-15
|
||||
- Change: Validated Rust-mode compatibility path with full integration suites.
|
||||
- APIs covered: full `/test/volume_server/http` and `/test/volume_server/grpc` suites executed successfully with `VOLUME_SERVER_IMPL=rust`.
|
||||
- Profiles covered: existing integration matrix in Rust launcher mode.
|
||||
- Gaps introduced/remaining: native Rust endpoint and storage/RPC implementation remains pending; current mode delegates execution to Go volume server.
|
||||
- Commit: `6ce4d7ede`
|
||||
|
||||
- Date: 2026-02-15
|
||||
- Change: Added Rust proxy-supervision launcher mode and validated all volume-server integration suites in proxy mode.
|
||||
- APIs covered: full `/test/volume_server/http` and `/test/volume_server/grpc` suites executed successfully with `VOLUME_SERVER_IMPL=rust` and `VOLUME_SERVER_RUST_MODE=proxy`.
|
||||
- Profiles covered: existing integration matrix in Rust proxy launcher mode.
|
||||
- Gaps introduced/remaining: native Rust endpoint/storage/RPC logic remains pending; current proxy mode still delegates backend handlers to Go.
|
||||
- Commit: `a7f50d23b`
|
||||
|
||||
- Date: 2026-02-15
|
||||
- Change: Expanded GitHub CI Rust smoke job to run both launcher modes (`exec` and `proxy`).
|
||||
- APIs covered: representative HTTP/gRPC smoke flows now validated in CI for both Rust launcher execution paths.
|
||||
- Profiles covered: CI smoke profile for P1 representative calls.
|
||||
- Gaps introduced/remaining: CI still runs Rust mode as smoke scope; full-suite Rust-mode CI remains optional due runtime cost.
|
||||
- Commit: `548b3d9a3`
|
||||
|
||||
- Date: 2026-02-15
|
||||
- Change: Added JWT UI override integration coverage with explicit `access.ui=true` profile support.
|
||||
- APIs covered: `/ui/index.html` behavior under JWT profile now verifies both default auth-gated path (`401`) and explicit UI-enable override path (`200` with rendered UI).
|
||||
- Profiles covered: P3 and P3+`access.ui=true`.
|
||||
- Gaps introduced/remaining: UI exposure behavior under JWT profiles is now covered for both secured and override branches.
|
||||
- Commit: `de974c05d`
|
||||
|
||||
- Date: 2026-02-15
|
||||
- Change: Added oversized upload integration coverage using explicit `-fileSizeLimitMB` profile configuration.
|
||||
- APIs covered: HTTP write path now verifies payloads larger than configured file-size limit are rejected (`400`) with limit-related error context.
|
||||
- Profiles covered: P1-derived profile with `FileSizeLimitMB=1`.
|
||||
- Gaps introduced/remaining: file-size limit rejection branch is now covered; replicate-write failure and cancellation (`499`) branches remain pending.
|
||||
- Commit: `4d61cbdee`
|
||||
|
||||
- Date: 2026-02-15
|
||||
- Change: Added gRPC EC-only metadata-read unsupported-path coverage.
|
||||
- APIs covered: `ReadNeedleMeta` now verifies the explicit unsupported branch when a volume is unmounted and only EC shards are mounted.
|
||||
- Profiles covered: P1 with EC shard generation/mount lifecycle setup.
|
||||
- Gaps introduced/remaining: low-level read metadata fault-injection and transport interruption branches remain pending.
|
||||
- Commit: `37bf9b5eb`
|
||||
|
||||
- Date: 2026-02-15
|
||||
- Change: Added deterministic replicated-write failure coverage for unmet replication requirements.
|
||||
- APIs covered: HTTP write path now verifies replication error handling (`500`) and no-local-commit outcome (`404` on follow-up read) when volume replication is configured (`001`) but replica set cannot be satisfied.
|
||||
- Profiles covered: P1-derived single-node profile with per-volume replication override.
|
||||
- Gaps introduced/remaining: replicated-write failure branch is now covered; cancellation (`499`) and deeper transport fault-injection branches remain pending.
|
||||
- Commit: `4835d3443`
|
||||
|
||||
- Date: 2026-02-15
|
||||
- Change: Added EC-backed `BatchDelete` positive-path integration coverage.
|
||||
- APIs covered: `BatchDelete` now validates successful EC-needle deletion (`202`) with post-delete `VolumeEcShardRead` deleted-marker verification (`IsDeleted=true`).
|
||||
- Profiles covered: P1 with EC shard generate/mount lifecycle; includes test-local gRPC offset proxy for EC internal dial-path parity.
|
||||
- Gaps introduced/remaining: EC `BatchDelete` success branch is now covered; deeper distributed failure-injection permutations remain pending.
|
||||
- Commit: `1bb40b6bc`
|
||||
|
||||
- Date: 2026-02-15
|
||||
- Change: Expanded HTTP conditional-header matrix for `HEAD` requests with combined header scenarios.
|
||||
- APIs covered: `HEAD` now verifies `If-Modified-Since` precedence over mismatched `If-None-Match`, and `If-None-Match` precedence when IMS is stale but ETag matches.
|
||||
- Profiles covered: P1.
|
||||
- Gaps introduced/remaining: combined conditional-header precedence paths for `HEAD`/`GET` are now explicitly covered.
|
||||
- Commit: `ed23e290f`
|
||||
|
||||
- Date: 2026-02-15
|
||||
- Change: Added multi-volume success-stream coverage for `ReadAllNeedles`.
|
||||
- APIs covered: `ReadAllNeedles` now verifies happy-path streaming across two existing volumes in one request, with per-volume payload validation.
|
||||
- Profiles covered: P1.
|
||||
- Gaps introduced/remaining: multi-volume stream boundary coverage is now present for successful reads; transport interruption branches remain pending.
|
||||
- Commit: `ab95a6ef1`
|
||||
|
||||
- Date: 2026-02-15
|
||||
- Change: Expanded split-port unsupported HTTP method matrix with `MKCOL` coverage.
|
||||
- APIs covered: admin/public parity for `MKCOL` (`400` on admin, passthrough `200` on public) with post-call data-integrity verification.
|
||||
- Profiles covered: P2.
|
||||
- Gaps introduced/remaining: unsupported-method sampling now covers `PATCH`, `TRACE`, `PROPFIND`, `CONNECT`, and `MKCOL`; cancellation/transport fault paths remain the main unaddressed area.
|
||||
- Commit: `2ab30900d`
|
||||
|
||||
- Date: 2026-02-15
|
||||
- Change: Hardened framework port allocation to keep volume ports within `10000..55535`.
|
||||
- APIs covered: test harness now guarantees valid `admin+10000` gRPC offset ranges, eliminating EC batch-delete flake scenarios that rely on offset dialing.
|
||||
- Profiles covered: all profiles using shared framework port allocation.
|
||||
- Gaps introduced/remaining: no new API gaps; transport cancellation/fault-injection branches remain the primary uncovered area.
|
||||
- Commit: `90e82b15c`
|
||||
|
||||
- Date: 2026-02-15
|
||||
- Change: Re-validated full Rust proxy-mode integration suite after latest HTTP/gRPC coverage additions and framework hardening.
|
||||
- APIs covered: full `/test/volume_server/http` and `/test/volume_server/grpc` packages.
|
||||
- Profiles covered: existing matrix in Rust proxy launcher mode.
|
||||
- Gaps introduced/remaining: native Rust endpoint/storage/RPC implementation remains pending (current Rust mode still delegates to Go backend handlers).
|
||||
- Commit: `6bb9d8bac`
|
||||
|
||||
- Date: 2026-02-15
|
||||
- Change: Added tail sender stream-cancellation interruption coverage.
|
||||
- APIs covered: `VolumeTailSender` now verifies client-side context cancellation behavior (`codes.Canceled`) after stream start/heartbeat.
|
||||
- Profiles covered: P1.
|
||||
- Gaps introduced/remaining: tail transport interruption coverage is now partially closed; receiver-side interruption and deeper network fault injection remain pending.
|
||||
- Commit: `27a80f760`
|
||||
|
||||
- Date: 2026-02-15
|
||||
- Change: Added receiver-side tail stream interruption coverage for unavailable source node.
|
||||
- APIs covered: `VolumeTailReceiver` now verifies source connect/dial failure branch when `SourceVolumeServer` is unreachable.
|
||||
- Profiles covered: P1.
|
||||
- Gaps introduced/remaining: both sender and receiver transport interruption branches are now covered at integration level; deeper injected network faults remain pending.
|
||||
- Commit: `a5864c3eb`
|
||||
|
||||
- Date: 2026-02-15
|
||||
- Change: Expanded query CSV matrix with explicit CSV-payload parity coverage.
|
||||
- APIs covered: `Query` now verifies current CSV-input behavior (no streamed rows / immediate EOF) even when source needle content is valid CSV text.
|
||||
- Profiles covered: P1.
|
||||
- Gaps introduced/remaining: CSV parsing/selection implementation still absent in current server behavior; integration parity now locks the current no-output semantics for both JSON and CSV payload shapes.
|
||||
- Commit: `a12dd5f8d`
|
||||
|
||||
- Date: 2026-02-15
|
||||
- Change: Expanded ping unreachable-target matrix with volume-server target coverage.
|
||||
- APIs covered: `Ping` now explicitly verifies unreachable-target error wrapping for `target_type=volumeServer`, alongside existing master/filer variants.
|
||||
- Profiles covered: P1.
|
||||
- Gaps introduced/remaining: ping target-type matrix is covered for success and unreachable branches across master/filer/volume-server.
|
||||
- Commit: `b9fbb85af`
|
||||
|
||||
- Date: 2026-02-15
|
||||
- Change: Expanded deleted-read HTTP matrix with `HEAD` parity coverage.
|
||||
- APIs covered: `HEAD` with `readDeleted=true` now validates current behavior (`200`, empty body, payload-size `Content-Length`) after delete.
|
||||
- Profiles covered: P1.
|
||||
- Gaps introduced/remaining: deleted-read parity now covers both `GET` and `HEAD` semantics on local-volume path.
|
||||
- Commit: `cc80ad364`
|
||||
|
||||
- Date: 2026-02-16
|
||||
- Change: Shifted planning priority to native Rust implementation parity.
|
||||
- APIs covered: no new API additions in this plan entry; integration suite remains the validation gate.
|
||||
- Profiles covered: unchanged.
|
||||
- Gaps introduced/remaining: primary remaining gap is native Rust handler/storage/RPC implementation replacing Go backend delegation.
|
||||
- Commit: `14c863dbf`
|
||||
|
||||
- Date: 2026-02-16
|
||||
- Change: Added Rust launcher `native` mode bootstrap, made it the default Rust launcher mode, and expanded CI Rust smoke matrix coverage to include `native`.
|
||||
- APIs covered: full `/test/volume_server/http` and `/test/volume_server/grpc` packages re-validated in `VOLUME_SERVER_RUST_MODE=native`; default Rust launcher path (`VOLUME_SERVER_IMPL=rust`) smoke-validated for HTTP and gRPC control tests.
|
||||
- Profiles covered: existing HTTP/gRPC matrix in native launcher mode; default-mode smoke checks on P1 control flows.
|
||||
- Gaps introduced/remaining: native Rust API/storage/RPC handlers still pending; current `native` mode remains a delegation bootstrap while parity replacement proceeds.
|
||||
- Commit: `70ddbee37`, `61befd10f`, `2e65966c0`
|
||||
|
||||
- Date: 2026-02-16
|
||||
- Change: Implemented first native Rust HTTP control handlers in `native` mode for `/status` and `/healthz` (admin listener path), with proxy fallback retained for remaining endpoints.
|
||||
- APIs covered: `/status` and `/healthz` now served by Rust-native code path in `VOLUME_SERVER_RUST_MODE=native`; full HTTP and gRPC integration suites re-validated.
|
||||
- Profiles covered: full existing HTTP/gRPC matrix under native mode.
|
||||
- Gaps introduced/remaining: gRPC control/data APIs and most HTTP data/static/auth paths remain delegated and require incremental native replacement.
|
||||
- Commit: `7e6e0261a`
|
||||
|
||||
- Date: 2026-02-16
|
||||
- Change: Implemented native Rust HTTP `OPTIONS` handling for admin/public listeners in `VOLUME_SERVER_RUST_MODE=native`.
|
||||
- APIs covered: admin/public `OPTIONS` method allow-list and CORS origin handling now served from Rust-native path; full HTTP and gRPC packages re-validated.
|
||||
- Profiles covered: P2 split-port CORS/OPTIONS profiles plus full existing matrix in native mode.
|
||||
- Gaps introduced/remaining: static/UI/auth and data-path HTTP handlers, plus gRPC control/data/stream handlers, remain delegated and pending native replacement.
|
||||
- Commit: `fbff2cb39`
|
||||
|
||||
- Date: 2026-02-16
|
||||
- Change: Implemented broader native Rust HTTP surface in `native` mode:
|
||||
- `/ui/index.html` auth-gating parity (JWT + access.ui handling)
|
||||
- static asset paths (`/favicon.ico`, `/seaweedfsstatic/*`)
|
||||
- public non-read method no-op parity (`200`) and admin unsupported-method rejection parity (`400`)
|
||||
- APIs covered: UI/static/admin/public control-surface behavior now served by Rust-native path; full `/test/volume_server/http` and `/test/volume_server/grpc` packages re-validated.
|
||||
- Profiles covered: P1/P2/P3 variants touched by UI/static/CORS/method behavior, plus full matrix in native mode.
|
||||
- Gaps introduced/remaining: HTTP data-path handlers and all gRPC handlers remain delegated and still require native replacement.
|
||||
- Commit: `23e4497b2`
|
||||
|
||||
- Date: 2026-02-16
|
||||
- Change: Improved native Rust control-path parity for health/status and request-target parsing.
|
||||
- APIs covered: native `/healthz` now mirrors backend service health transitions (`200`/`503`) and native route matching now handles absolute-form HTTP request targets before path dispatch.
|
||||
- Profiles covered: full existing HTTP/gRPC integration matrix in native mode (`VOLUME_SERVER_IMPL=rust`, `VOLUME_SERVER_RUST_MODE=native`).
|
||||
- Gaps introduced/remaining: core HTTP data handlers and all gRPC RPC handlers are still delegated and remain the primary native implementation gap.
|
||||
- Commit: `d6ff6ed6d`
|
||||
|
||||
- Date: 2026-02-16
|
||||
- Change: Added Rust-native malformed fid-route validation ahead of delegated data handlers.
|
||||
- APIs covered: GET/HEAD/POST/PUT fid-shaped paths now return native `400` for invalid vid/fid tokens (including invalid-read and invalid-write path variants) while valid routes continue through delegated storage handlers.
|
||||
- Profiles covered: full existing HTTP/gRPC integration matrix in native mode, plus targeted invalid-path parity runs.
|
||||
- Gaps introduced/remaining: data-path success branches and all gRPC handler bodies remain delegated; native replacement work continues on those core execution paths.
|
||||
- Commit: `1ce0174b2`
|
||||
|
||||
- Date: 2026-02-16
|
||||
- Change: Added Rust-native write error prevalidation for delegated fid routes.
|
||||
- APIs covered: native path now rejects malformed multipart boundary requests, `Content-MD5` writes, and over-limit `Content-Length` writes (using `-fileSizeLimitMB`) with `400` parity behavior before delegation.
|
||||
- Profiles covered: P1 default and P1 with explicit file-size limit profile, plus full existing HTTP/gRPC native-mode matrix revalidation.
|
||||
- Gaps introduced/remaining: positive write/read/delete data handlers and all gRPC handler bodies are still delegated; native implementation continues incrementally.
|
||||
- Commit: `94cefd6f4`
|
||||
|
||||
- Date: 2026-02-16
|
||||
- Change: Broadened Rust-native malformed-fid handling to slash-form read/write URL shapes with reserved-route exclusions.
|
||||
- APIs covered: `/{vid}/{fid}` and `/{vid}/{fid}/{filename}` malformed cases now hit native `400` validation path, while `/status`, `/healthz`, `/ui/index.html`, `/stats/*`, and static asset routes remain excluded from fid parsing.
|
||||
- Profiles covered: full existing HTTP/gRPC native-mode matrix, plus targeted admin/read-path malformed-slash validation run.
|
||||
- Gaps introduced/remaining: successful data-path read/write/delete handlers and all gRPC method bodies remain delegated and are still the primary native implementation gap.
|
||||
- Commit: `e66983e9b`
|
||||
|
||||
@@ -1,7 +1,10 @@
|
||||
.PHONY: test-volume-server test-volume-server-short
|
||||
.PHONY: test-volume-server test-volume-server-short test-volume-server-rust-smoke
|
||||
|
||||
test-volume-server:
|
||||
go test ./test/volume_server/... -v
|
||||
|
||||
test-volume-server-short:
|
||||
go test ./test/volume_server/... -short -v
|
||||
|
||||
test-volume-server-rust-smoke:
|
||||
VOLUME_SERVER_IMPL=rust go test ./test/volume_server/http ./test/volume_server/grpc -run 'TestAdminStatusAndHealthz|TestUploadReadRangeHeadDeleteRoundTrip|TestStateAndStatusRPCs|TestVolumeSyncStatusAndReadVolumeFileStatus' -v
|
||||
|
||||
@@ -15,6 +15,10 @@ If a `weed` binary is not found, the harness will build one automatically.
|
||||
## Optional environment variables
|
||||
|
||||
- `WEED_BINARY`: explicit path to the `weed` executable (disables auto-build).
|
||||
- `VOLUME_SERVER_IMPL`: select volume-server implementation (`go` default, `rust` for Rust launcher mode).
|
||||
- `VOLUME_SERVER_BINARY`: explicit path to the volume-server executable (overrides `VOLUME_SERVER_IMPL`).
|
||||
- `VOLUME_SERVER_RUST_MODE`: Rust launcher mode (`native` default, `exec` for direct Go delegation, `proxy` for Rust front proxy + Go backend).
|
||||
- `VOLUME_SERVER_RUST_REBUILD=1`: force rebuild of Rust volume-server binary in Rust mode.
|
||||
- `VOLUME_SERVER_IT_KEEP_LOGS=1`: keep temporary test directories and process logs.
|
||||
|
||||
## Current scope (Phase 0)
|
||||
@@ -24,4 +28,4 @@ If a `weed` binary is not found, the harness will build one automatically.
|
||||
- Initial HTTP admin endpoint checks
|
||||
- Initial gRPC state/status checks
|
||||
|
||||
More API coverage is tracked in `/Users/chris/dev/seaweedfs2/test/volume_server/DEV_PLAN.md`.
|
||||
More API coverage is tracked in `test/volume_server/DEV_PLAN.md`.
|
||||
|
||||
@@ -32,11 +32,12 @@ type Cluster struct {
|
||||
testingTB testing.TB
|
||||
profile matrix.Profile
|
||||
|
||||
weedBinary string
|
||||
baseDir string
|
||||
configDir string
|
||||
logsDir string
|
||||
keepLogs bool
|
||||
weedBinary string
|
||||
volumeBinary string
|
||||
baseDir string
|
||||
configDir string
|
||||
logsDir string
|
||||
keepLogs bool
|
||||
|
||||
masterPort int
|
||||
masterGrpcPort int
|
||||
@@ -54,9 +55,9 @@ type Cluster struct {
|
||||
func StartSingleVolumeCluster(t testing.TB, profile matrix.Profile) *Cluster {
|
||||
t.Helper()
|
||||
|
||||
weedBinary, err := FindOrBuildWeedBinary()
|
||||
weedBinary, volumeBinary, err := FindOrBuildServerBinaries()
|
||||
if err != nil {
|
||||
t.Fatalf("resolve weed binary: %v", err)
|
||||
t.Fatalf("resolve server binaries: %v", err)
|
||||
}
|
||||
|
||||
baseDir, keepLogs, err := newWorkDir()
|
||||
@@ -92,6 +93,7 @@ func StartSingleVolumeCluster(t testing.TB, profile matrix.Profile) *Cluster {
|
||||
testingTB: t,
|
||||
profile: profile,
|
||||
weedBinary: weedBinary,
|
||||
volumeBinary: volumeBinary,
|
||||
baseDir: baseDir,
|
||||
configDir: configDir,
|
||||
logsDir: logsDir,
|
||||
@@ -206,8 +208,11 @@ func (c *Cluster) startVolume(dataDir string) error {
|
||||
if c.profile.InflightDownloadTimeout > 0 {
|
||||
args = append(args, "-inflightDownloadDataTimeout="+c.profile.InflightDownloadTimeout.String())
|
||||
}
|
||||
if c.profile.FileSizeLimitMB > 0 {
|
||||
args = append(args, "-fileSizeLimitMB="+strconv.Itoa(c.profile.FileSizeLimitMB))
|
||||
}
|
||||
|
||||
c.volumeCmd = exec.Command(c.weedBinary, args...)
|
||||
c.volumeCmd = exec.Command(c.volumeBinary, args...)
|
||||
c.volumeCmd.Dir = c.baseDir
|
||||
c.volumeCmd.Stdout = logFile
|
||||
c.volumeCmd.Stderr = logFile
|
||||
@@ -264,18 +269,46 @@ func stopProcess(cmd *exec.Cmd) {
|
||||
}
|
||||
|
||||
func allocatePorts(count int) ([]int, error) {
|
||||
const minPort = 10000
|
||||
const maxPort = 55535
|
||||
const host = "127.0.0.1"
|
||||
rangeSize := maxPort - minPort + 1
|
||||
|
||||
listeners := make([]net.Listener, 0, count)
|
||||
ports := make([]int, 0, count)
|
||||
seen := make(map[int]struct{}, count)
|
||||
|
||||
closeAll := func() {
|
||||
for _, ll := range listeners {
|
||||
_ = ll.Close()
|
||||
}
|
||||
}
|
||||
|
||||
startOffset := int(time.Now().UnixNano() % int64(rangeSize))
|
||||
for i := 0; i < count; i++ {
|
||||
l, err := net.Listen("tcp", "127.0.0.1:0")
|
||||
if err != nil {
|
||||
for _, ll := range listeners {
|
||||
_ = ll.Close()
|
||||
found := false
|
||||
for offset := 0; offset < rangeSize; offset++ {
|
||||
port := minPort + (startOffset+offset)%rangeSize
|
||||
if _, exists := seen[port]; exists {
|
||||
continue
|
||||
}
|
||||
return nil, err
|
||||
|
||||
l, err := net.Listen("tcp", net.JoinHostPort(host, strconv.Itoa(port)))
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
|
||||
seen[port] = struct{}{}
|
||||
listeners = append(listeners, l)
|
||||
ports = append(ports, port)
|
||||
found = true
|
||||
startOffset = (startOffset + offset + 1) % rangeSize
|
||||
break
|
||||
}
|
||||
if !found {
|
||||
closeAll()
|
||||
return nil, fmt.Errorf("unable to allocate %d ports within range [%d,%d]", count, minPort, maxPort)
|
||||
}
|
||||
listeners = append(listeners, l)
|
||||
ports = append(ports, l.Addr().(*net.TCPAddr).Port)
|
||||
}
|
||||
for _, l := range listeners {
|
||||
_ = l.Close()
|
||||
@@ -326,6 +359,13 @@ func writeSecurityConfig(configDir string, profile matrix.Profile) error {
|
||||
b.WriteString("\"\n")
|
||||
b.WriteString("expires_after_seconds = 60\n")
|
||||
}
|
||||
if profile.AccessUI {
|
||||
if b.Len() > 0 {
|
||||
b.WriteString("\n")
|
||||
}
|
||||
b.WriteString("[access]\n")
|
||||
b.WriteString("ui = true\n")
|
||||
}
|
||||
if b.Len() == 0 {
|
||||
b.WriteString("# optional security config generated for integration tests\n")
|
||||
}
|
||||
@@ -341,17 +381,12 @@ func FindOrBuildWeedBinary() (string, error) {
|
||||
return "", fmt.Errorf("WEED_BINARY is set but not executable: %s", fromEnv)
|
||||
}
|
||||
|
||||
repoRoot := ""
|
||||
if _, file, _, ok := runtime.Caller(0); ok {
|
||||
repoRoot = filepath.Clean(filepath.Join(filepath.Dir(file), "..", "..", ".."))
|
||||
candidate := filepath.Join(repoRoot, "weed", "weed")
|
||||
if isExecutableFile(candidate) {
|
||||
return candidate, nil
|
||||
}
|
||||
repoRoot, err := detectRepoRoot()
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
if repoRoot == "" {
|
||||
return "", errors.New("unable to detect repository root")
|
||||
if candidate := filepath.Join(repoRoot, "weed", "weed"); isExecutableFile(candidate) {
|
||||
return candidate, nil
|
||||
}
|
||||
|
||||
binDir := filepath.Join(os.TempDir(), "seaweedfs_volume_server_it_bin")
|
||||
@@ -377,6 +412,82 @@ func FindOrBuildWeedBinary() (string, error) {
|
||||
return binPath, nil
|
||||
}
|
||||
|
||||
// FindOrBuildServerBinaries returns master and volume executables.
|
||||
// Master always runs from the Go weed binary; volume can be switched via env.
|
||||
func FindOrBuildServerBinaries() (masterBinary string, volumeBinary string, err error) {
|
||||
masterBinary, err = FindOrBuildWeedBinary()
|
||||
if err != nil {
|
||||
return "", "", err
|
||||
}
|
||||
volumeBinary, err = FindOrBuildVolumeServerBinary(masterBinary)
|
||||
if err != nil {
|
||||
return "", "", err
|
||||
}
|
||||
return masterBinary, volumeBinary, nil
|
||||
}
|
||||
|
||||
// FindOrBuildVolumeServerBinary resolves the executable used for volume-server processes.
|
||||
//
|
||||
// Behavior:
|
||||
// - `VOLUME_SERVER_BINARY=/path/to/bin`: use explicit executable path.
|
||||
// - `VOLUME_SERVER_IMPL=rust`: build/use Rust volume server launcher.
|
||||
// - default: use the same Go `weed` binary.
|
||||
func FindOrBuildVolumeServerBinary(defaultBinary string) (string, error) {
|
||||
if fromEnv := os.Getenv("VOLUME_SERVER_BINARY"); fromEnv != "" {
|
||||
if isExecutableFile(fromEnv) {
|
||||
return fromEnv, nil
|
||||
}
|
||||
return "", fmt.Errorf("VOLUME_SERVER_BINARY is set but not executable: %s", fromEnv)
|
||||
}
|
||||
|
||||
impl := strings.ToLower(strings.TrimSpace(os.Getenv("VOLUME_SERVER_IMPL")))
|
||||
if impl == "" || impl == "go" {
|
||||
return defaultBinary, nil
|
||||
}
|
||||
if impl != "rust" {
|
||||
return "", fmt.Errorf("unsupported VOLUME_SERVER_IMPL %q (supported: go, rust)", impl)
|
||||
}
|
||||
|
||||
repoRoot, err := detectRepoRoot()
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
return FindOrBuildRustVolumeServerBinary(repoRoot)
|
||||
}
|
||||
|
||||
// FindOrBuildRustVolumeServerBinary builds the Rust volume server launcher when needed.
|
||||
func FindOrBuildRustVolumeServerBinary(repoRoot string) (string, error) {
|
||||
manifestPath := filepath.Join(repoRoot, "rust", "volume_server", "Cargo.toml")
|
||||
if _, err := os.Stat(manifestPath); err != nil {
|
||||
return "", fmt.Errorf("rust volume server manifest not found at %s: %w", manifestPath, err)
|
||||
}
|
||||
|
||||
targetDir := filepath.Join(os.TempDir(), "seaweedfs_volume_server_it_rust_target")
|
||||
binPath := filepath.Join(targetDir, "release", "weed-volume-rs")
|
||||
if isExecutableFile(binPath) && os.Getenv("VOLUME_SERVER_RUST_REBUILD") != "1" {
|
||||
return binPath, nil
|
||||
}
|
||||
|
||||
cmd := exec.Command("cargo", "build", "--release", "--manifest-path", manifestPath, "--target-dir", targetDir)
|
||||
var out bytes.Buffer
|
||||
cmd.Stdout = &out
|
||||
cmd.Stderr = &out
|
||||
if err := cmd.Run(); err != nil {
|
||||
return "", fmt.Errorf("build rust volume server binary: %w\n%s", err, out.String())
|
||||
}
|
||||
if !isExecutableFile(binPath) {
|
||||
return "", fmt.Errorf("built rust volume server binary is not executable: %s", binPath)
|
||||
}
|
||||
return binPath, nil
|
||||
}
|
||||
|
||||
func detectRepoRoot() (string, error) {
|
||||
if _, file, _, ok := runtime.Caller(0); ok {
|
||||
return filepath.Clean(filepath.Join(filepath.Dir(file), "..", "..", "..")), nil
|
||||
}
|
||||
return "", errors.New("unable to detect repository root")
|
||||
}
|
||||
|
||||
func isExecutableFile(path string) bool {
|
||||
info, err := os.Stat(path)
|
||||
if err != nil || info.IsDir() {
|
||||
|
||||
@@ -17,11 +17,12 @@ type DualVolumeCluster struct {
|
||||
testingTB testing.TB
|
||||
profile matrix.Profile
|
||||
|
||||
weedBinary string
|
||||
baseDir string
|
||||
configDir string
|
||||
logsDir string
|
||||
keepLogs bool
|
||||
weedBinary string
|
||||
volumeBinary string
|
||||
baseDir string
|
||||
configDir string
|
||||
logsDir string
|
||||
keepLogs bool
|
||||
|
||||
masterPort int
|
||||
masterGrpcPort int
|
||||
@@ -33,7 +34,7 @@ type DualVolumeCluster struct {
|
||||
volumeGrpcPort1 int
|
||||
volumePubPort1 int
|
||||
|
||||
masterCmd *exec.Cmd
|
||||
masterCmd *exec.Cmd
|
||||
volumeCmd0 *exec.Cmd
|
||||
volumeCmd1 *exec.Cmd
|
||||
|
||||
@@ -43,9 +44,9 @@ type DualVolumeCluster struct {
|
||||
func StartDualVolumeCluster(t testing.TB, profile matrix.Profile) *DualVolumeCluster {
|
||||
t.Helper()
|
||||
|
||||
weedBinary, err := FindOrBuildWeedBinary()
|
||||
weedBinary, volumeBinary, err := FindOrBuildServerBinaries()
|
||||
if err != nil {
|
||||
t.Fatalf("resolve weed binary: %v", err)
|
||||
t.Fatalf("resolve server binaries: %v", err)
|
||||
}
|
||||
|
||||
baseDir, keepLogs, err := newWorkDir()
|
||||
@@ -79,21 +80,22 @@ func StartDualVolumeCluster(t testing.TB, profile matrix.Profile) *DualVolumeClu
|
||||
}
|
||||
|
||||
c := &DualVolumeCluster{
|
||||
testingTB: t,
|
||||
profile: profile,
|
||||
weedBinary: weedBinary,
|
||||
baseDir: baseDir,
|
||||
configDir: configDir,
|
||||
logsDir: logsDir,
|
||||
keepLogs: keepLogs,
|
||||
masterPort: masterPort,
|
||||
masterGrpcPort: masterGrpcPort,
|
||||
volumePort0: ports[0],
|
||||
testingTB: t,
|
||||
profile: profile,
|
||||
weedBinary: weedBinary,
|
||||
volumeBinary: volumeBinary,
|
||||
baseDir: baseDir,
|
||||
configDir: configDir,
|
||||
logsDir: logsDir,
|
||||
keepLogs: keepLogs,
|
||||
masterPort: masterPort,
|
||||
masterGrpcPort: masterGrpcPort,
|
||||
volumePort0: ports[0],
|
||||
volumeGrpcPort0: ports[1],
|
||||
volumePubPort0: ports[0],
|
||||
volumePort1: ports[2],
|
||||
volumePubPort0: ports[0],
|
||||
volumePort1: ports[2],
|
||||
volumeGrpcPort1: ports[3],
|
||||
volumePubPort1: ports[2],
|
||||
volumePubPort1: ports[2],
|
||||
}
|
||||
if profile.SplitPublicPort {
|
||||
c.volumePubPort0 = ports[4]
|
||||
@@ -227,7 +229,7 @@ func (c *DualVolumeCluster) startVolume(index int, dataDir string) error {
|
||||
args = append(args, "-inflightDownloadDataTimeout="+c.profile.InflightDownloadTimeout.String())
|
||||
}
|
||||
|
||||
cmd := exec.Command(c.weedBinary, args...)
|
||||
cmd := exec.Command(c.volumeBinary, args...)
|
||||
cmd.Dir = c.baseDir
|
||||
cmd.Stdout = logFile
|
||||
cmd.Stderr = logFile
|
||||
|
||||
@@ -14,18 +14,26 @@ import (
|
||||
|
||||
func AllocateVolume(t testing.TB, client volume_server_pb.VolumeServerClient, volumeID uint32, collection string) {
|
||||
t.Helper()
|
||||
AllocateVolumeWithReplication(t, client, volumeID, collection, "000")
|
||||
}
|
||||
|
||||
func AllocateVolumeWithReplication(t testing.TB, client volume_server_pb.VolumeServerClient, volumeID uint32, collection, replication string) {
|
||||
t.Helper()
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||||
defer cancel()
|
||||
if replication == "" {
|
||||
replication = "000"
|
||||
}
|
||||
|
||||
_, err := client.AllocateVolume(ctx, &volume_server_pb.AllocateVolumeRequest{
|
||||
VolumeId: volumeID,
|
||||
Collection: collection,
|
||||
Replication: "000",
|
||||
Replication: replication,
|
||||
Version: uint32(needle.GetCurrentVersion()),
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("allocate volume %d: %v", volumeID, err)
|
||||
t.Fatalf("allocate volume %d (replication=%s): %v", volumeID, replication, err)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -388,6 +388,17 @@ func TestPingUnknownAndUnreachableTargetPaths(t *testing.T) {
|
||||
if !strings.Contains(err.Error(), "ping filer") {
|
||||
t.Fatalf("Ping filer unreachable error mismatch: %v", err)
|
||||
}
|
||||
|
||||
_, err = grpcClient.Ping(ctx, &volume_server_pb.PingRequest{
|
||||
TargetType: cluster.VolumeServerType,
|
||||
Target: "127.0.0.1:1.1",
|
||||
})
|
||||
if err == nil {
|
||||
t.Fatalf("Ping volume server target should fail when target is unreachable")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "ping volumeServer") {
|
||||
t.Fatalf("Ping volume server unreachable error mismatch: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestPingMasterTargetSuccess(t *testing.T) {
|
||||
|
||||
@@ -3,7 +3,10 @@ package volume_server_grpc_test
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"io"
|
||||
"net"
|
||||
"net/http"
|
||||
"strconv"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
@@ -262,3 +265,141 @@ func TestBatchDeleteRejectsChunkManifestNeedles(t *testing.T) {
|
||||
t.Fatalf("chunk manifest should not be deleted by BatchDelete reject path, got %d", readResp.StatusCode)
|
||||
}
|
||||
}
|
||||
|
||||
func TestBatchDeleteEcNeedleSuccess(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("skipping integration test in short mode")
|
||||
}
|
||||
|
||||
cluster := framework.StartSingleVolumeCluster(t, matrix.P1())
|
||||
conn, client := framework.DialVolumeServer(t, cluster.VolumeGRPCAddress())
|
||||
defer conn.Close()
|
||||
|
||||
if stopProxy := maybeStartGrpcOffsetProxy(t, cluster.VolumeAdminAddress(), cluster.VolumeGRPCAddress()); stopProxy != nil {
|
||||
t.Cleanup(stopProxy)
|
||||
}
|
||||
|
||||
const volumeID = uint32(34)
|
||||
const needleID = uint64(930001)
|
||||
const cookie = uint32(0x6677BBCC)
|
||||
framework.AllocateVolume(t, client, volumeID, "")
|
||||
|
||||
httpClient := framework.NewHTTPClient()
|
||||
fid := framework.NewFileID(volumeID, needleID, cookie)
|
||||
uploadResp := framework.UploadBytes(t, httpClient, cluster.VolumeAdminURL(), fid, []byte("batch-delete-ec-success"))
|
||||
_ = framework.ReadAllAndClose(t, uploadResp)
|
||||
if uploadResp.StatusCode != http.StatusCreated {
|
||||
t.Fatalf("upload expected 201, got %d", uploadResp.StatusCode)
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
||||
defer cancel()
|
||||
|
||||
_, err := client.VolumeEcShardsGenerate(ctx, &volume_server_pb.VolumeEcShardsGenerateRequest{
|
||||
VolumeId: volumeID,
|
||||
Collection: "",
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("VolumeEcShardsGenerate failed: %v", err)
|
||||
}
|
||||
|
||||
_, 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 failed: %v", err)
|
||||
}
|
||||
|
||||
deleteResp, err := client.BatchDelete(ctx, &volume_server_pb.BatchDeleteRequest{FileIds: []string{fid}})
|
||||
if err != nil {
|
||||
t.Fatalf("BatchDelete EC needle should return response, got grpc error: %v", err)
|
||||
}
|
||||
if len(deleteResp.GetResults()) != 1 {
|
||||
t.Fatalf("BatchDelete EC needle expected one result, got %d", len(deleteResp.GetResults()))
|
||||
}
|
||||
if deleteResp.GetResults()[0].GetStatus() != http.StatusAccepted {
|
||||
t.Fatalf("BatchDelete EC needle expected status 202, got %d error=%q", deleteResp.GetResults()[0].GetStatus(), deleteResp.GetResults()[0].GetError())
|
||||
}
|
||||
if deleteResp.GetResults()[0].GetSize() == 0 {
|
||||
t.Fatalf("BatchDelete EC needle expected non-zero deleted size")
|
||||
}
|
||||
|
||||
deletedStream, err := client.VolumeEcShardRead(ctx, &volume_server_pb.VolumeEcShardReadRequest{
|
||||
VolumeId: volumeID,
|
||||
ShardId: 0,
|
||||
FileKey: needleID,
|
||||
Offset: 0,
|
||||
Size: 1,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("VolumeEcShardRead deleted-check start failed: %v", err)
|
||||
}
|
||||
deletedMsg, err := deletedStream.Recv()
|
||||
if err != nil {
|
||||
t.Fatalf("VolumeEcShardRead deleted-check recv failed: %v", err)
|
||||
}
|
||||
if !deletedMsg.GetIsDeleted() {
|
||||
t.Fatalf("VolumeEcShardRead expected IsDeleted=true after EC batch delete")
|
||||
}
|
||||
}
|
||||
|
||||
func maybeStartGrpcOffsetProxy(t testing.TB, adminAddr, actualGrpcAddr string) func() {
|
||||
t.Helper()
|
||||
|
||||
host, portText, err := net.SplitHostPort(adminAddr)
|
||||
if err != nil {
|
||||
t.Fatalf("split admin address %q: %v", adminAddr, err)
|
||||
}
|
||||
port, err := strconv.Atoi(portText)
|
||||
if err != nil {
|
||||
t.Fatalf("parse admin port %q: %v", portText, err)
|
||||
}
|
||||
|
||||
expectedGrpcAddr := net.JoinHostPort(host, strconv.Itoa(port+10000))
|
||||
if expectedGrpcAddr == actualGrpcAddr {
|
||||
return nil
|
||||
}
|
||||
|
||||
listener, err := net.Listen("tcp", expectedGrpcAddr)
|
||||
if err != nil {
|
||||
t.Fatalf("listen grpc-offset proxy %s -> %s: %v", expectedGrpcAddr, actualGrpcAddr, err)
|
||||
}
|
||||
|
||||
stopped := make(chan struct{})
|
||||
go func() {
|
||||
for {
|
||||
conn, acceptErr := listener.Accept()
|
||||
if acceptErr != nil {
|
||||
select {
|
||||
case <-stopped:
|
||||
return
|
||||
default:
|
||||
continue
|
||||
}
|
||||
}
|
||||
go proxyBidirectional(conn, actualGrpcAddr)
|
||||
}
|
||||
}()
|
||||
|
||||
return func() {
|
||||
close(stopped)
|
||||
_ = listener.Close()
|
||||
}
|
||||
}
|
||||
|
||||
func proxyBidirectional(inbound net.Conn, targetAddr string) {
|
||||
outbound, err := net.Dial("tcp", targetAddr)
|
||||
if err != nil {
|
||||
_ = inbound.Close()
|
||||
return
|
||||
}
|
||||
|
||||
go func() {
|
||||
_, _ = io.Copy(outbound, inbound)
|
||||
_ = outbound.Close()
|
||||
}()
|
||||
_, _ = io.Copy(inbound, outbound)
|
||||
_ = inbound.Close()
|
||||
}
|
||||
|
||||
@@ -2,6 +2,7 @@ package volume_server_grpc_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net/http"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
@@ -144,3 +145,66 @@ func TestReadNeedleBlobAndMetaInvalidOffsets(t *testing.T) {
|
||||
t.Fatalf("ReadNeedleMeta should fail for invalid offset")
|
||||
}
|
||||
}
|
||||
|
||||
func TestReadNeedleMetaUnsupportedForEcMountedVolume(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("skipping integration test in short mode")
|
||||
}
|
||||
|
||||
clusterHarness := framework.StartSingleVolumeCluster(t, matrix.P1())
|
||||
conn, grpcClient := framework.DialVolumeServer(t, clusterHarness.VolumeGRPCAddress())
|
||||
defer conn.Close()
|
||||
|
||||
const volumeID = uint32(93)
|
||||
const needleID = uint64(880002)
|
||||
const cookie = uint32(0x11AA22BB)
|
||||
framework.AllocateVolume(t, grpcClient, volumeID, "")
|
||||
|
||||
httpClient := framework.NewHTTPClient()
|
||||
fid := framework.NewFileID(volumeID, needleID, cookie)
|
||||
uploadResp := framework.UploadBytes(t, httpClient, clusterHarness.VolumeAdminURL(), fid, []byte("ec-meta-unsupported"))
|
||||
_ = framework.ReadAllAndClose(t, uploadResp)
|
||||
if uploadResp.StatusCode != http.StatusCreated {
|
||||
t.Fatalf("upload expected 201, got %d", uploadResp.StatusCode)
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
||||
defer cancel()
|
||||
|
||||
_, err := grpcClient.VolumeEcShardsGenerate(ctx, &volume_server_pb.VolumeEcShardsGenerateRequest{
|
||||
VolumeId: volumeID,
|
||||
Collection: "",
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("VolumeEcShardsGenerate failed: %v", err)
|
||||
}
|
||||
|
||||
_, err = grpcClient.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 failed: %v", err)
|
||||
}
|
||||
|
||||
_, err = grpcClient.VolumeUnmount(ctx, &volume_server_pb.VolumeUnmountRequest{
|
||||
VolumeId: volumeID,
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("VolumeUnmount failed: %v", err)
|
||||
}
|
||||
|
||||
_, err = grpcClient.ReadNeedleMeta(ctx, &volume_server_pb.ReadNeedleMetaRequest{
|
||||
VolumeId: volumeID,
|
||||
NeedleId: needleID,
|
||||
Offset: 0,
|
||||
Size: 128,
|
||||
})
|
||||
if err == nil {
|
||||
t.Fatalf("ReadNeedleMeta should fail when only EC shards are mounted")
|
||||
}
|
||||
if !strings.Contains(strings.ToLower(err.Error()), "ec shards is not supported") {
|
||||
t.Fatalf("ReadNeedleMeta EC-only error mismatch: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -3,6 +3,7 @@ package volume_server_grpc_test
|
||||
import (
|
||||
"context"
|
||||
"io"
|
||||
"net/http"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
@@ -230,6 +231,71 @@ func TestReadAllNeedlesExistingThenMissingVolumeAbortsStream(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestReadAllNeedlesStreamsAcrossMultipleVolumes(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("skipping integration test in short mode")
|
||||
}
|
||||
|
||||
clusterHarness := framework.StartSingleVolumeCluster(t, matrix.P1())
|
||||
conn, grpcClient := framework.DialVolumeServer(t, clusterHarness.VolumeGRPCAddress())
|
||||
defer conn.Close()
|
||||
|
||||
const volumeIDA = uint32(86)
|
||||
const volumeIDB = uint32(87)
|
||||
const needleIDA = uint64(445561)
|
||||
const needleIDB = uint64(445562)
|
||||
framework.AllocateVolume(t, grpcClient, volumeIDA, "")
|
||||
framework.AllocateVolume(t, grpcClient, volumeIDB, "")
|
||||
|
||||
client := framework.NewHTTPClient()
|
||||
fidA := framework.NewFileID(volumeIDA, needleIDA, 0xAA11CC22)
|
||||
fidB := framework.NewFileID(volumeIDB, needleIDB, 0xBB22DD33)
|
||||
payloadA := "read-all-multi-volume-a"
|
||||
payloadB := "read-all-multi-volume-b"
|
||||
|
||||
uploadA := framework.UploadBytes(t, client, clusterHarness.VolumeAdminURL(), fidA, []byte(payloadA))
|
||||
_ = framework.ReadAllAndClose(t, uploadA)
|
||||
if uploadA.StatusCode != http.StatusCreated {
|
||||
t.Fatalf("upload A expected 201, got %d", uploadA.StatusCode)
|
||||
}
|
||||
uploadB := framework.UploadBytes(t, client, clusterHarness.VolumeAdminURL(), fidB, []byte(payloadB))
|
||||
_ = framework.ReadAllAndClose(t, uploadB)
|
||||
if uploadB.StatusCode != http.StatusCreated {
|
||||
t.Fatalf("upload B expected 201, got %d", uploadB.StatusCode)
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||||
defer cancel()
|
||||
|
||||
stream, err := grpcClient.ReadAllNeedles(ctx, &volume_server_pb.ReadAllNeedlesRequest{
|
||||
VolumeIds: []uint32{volumeIDA, volumeIDB},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("ReadAllNeedles start failed: %v", err)
|
||||
}
|
||||
|
||||
seen := map[uint64]string{}
|
||||
for {
|
||||
msg, recvErr := stream.Recv()
|
||||
if recvErr == io.EOF {
|
||||
break
|
||||
}
|
||||
if recvErr != nil {
|
||||
t.Fatalf("ReadAllNeedles recv failed: %v", recvErr)
|
||||
}
|
||||
if msg.GetNeedleId() == needleIDA || msg.GetNeedleId() == needleIDB {
|
||||
seen[msg.GetNeedleId()] = string(msg.GetNeedleBlob())
|
||||
}
|
||||
}
|
||||
|
||||
if got := seen[needleIDA]; got != payloadA {
|
||||
t.Fatalf("ReadAllNeedles missing/mismatched payload for volume A needle: got %q want %q", got, payloadA)
|
||||
}
|
||||
if got := seen[needleIDB]; got != payloadB {
|
||||
t.Fatalf("ReadAllNeedles missing/mismatched payload for volume B needle: got %q want %q", got, payloadB)
|
||||
}
|
||||
}
|
||||
|
||||
func copyFileBytes(t testing.TB, grpcClient volume_server_pb.VolumeServerClient, req *volume_server_pb.CopyFileRequest) []byte {
|
||||
t.Helper()
|
||||
|
||||
|
||||
@@ -280,6 +280,54 @@ func TestQueryJsonSuccessAndCsvNoOutput(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestQueryCsvInputWithCsvPayloadStillReturnsNoRows(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("skipping integration test in short mode")
|
||||
}
|
||||
|
||||
clusterHarness := framework.StartSingleVolumeCluster(t, matrix.P1())
|
||||
conn, grpcClient := framework.DialVolumeServer(t, clusterHarness.VolumeGRPCAddress())
|
||||
defer conn.Close()
|
||||
|
||||
const volumeID = uint32(67)
|
||||
const needleID = uint64(777004)
|
||||
const cookie = uint32(0xDCDCADAD)
|
||||
framework.AllocateVolume(t, grpcClient, volumeID, "")
|
||||
|
||||
csvPayload := []byte("name,score\nalice,12\nbob,3\n")
|
||||
httpClient := framework.NewHTTPClient()
|
||||
fid := framework.NewFileID(volumeID, needleID, cookie)
|
||||
uploadResp := framework.UploadBytes(t, httpClient, clusterHarness.VolumeAdminURL(), fid, csvPayload)
|
||||
_ = framework.ReadAllAndClose(t, uploadResp)
|
||||
if uploadResp.StatusCode != 201 {
|
||||
t.Fatalf("upload expected 201, got %d", uploadResp.StatusCode)
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||||
defer cancel()
|
||||
|
||||
csvStream, err := grpcClient.Query(ctx, &volume_server_pb.QueryRequest{
|
||||
FromFileIds: []string{fid},
|
||||
Selections: []string{"score"},
|
||||
Filter: &volume_server_pb.QueryRequest_Filter{
|
||||
Field: "score",
|
||||
Operand: ">",
|
||||
Value: "10",
|
||||
},
|
||||
InputSerialization: &volume_server_pb.QueryRequest_InputSerialization{
|
||||
CsvInput: &volume_server_pb.QueryRequest_InputSerialization_CSVInput{},
|
||||
},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("Query csv start failed: %v", err)
|
||||
}
|
||||
|
||||
_, err = csvStream.Recv()
|
||||
if err != io.EOF {
|
||||
t.Fatalf("Query csv with csv payload expected EOF with no rows, got: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestQueryJsonNoMatchReturnsEmptyStripe(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("skipping integration test in short mode")
|
||||
|
||||
@@ -12,6 +12,8 @@ import (
|
||||
"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"
|
||||
"google.golang.org/grpc/codes"
|
||||
"google.golang.org/grpc/status"
|
||||
)
|
||||
|
||||
func TestVolumeTailSenderMissingVolume(t *testing.T) {
|
||||
@@ -91,6 +93,36 @@ func TestVolumeTailReceiverMissingVolume(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestVolumeTailReceiverSourceUnavailable(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("skipping integration test in short mode")
|
||||
}
|
||||
|
||||
clusterHarness := framework.StartSingleVolumeCluster(t, matrix.P1())
|
||||
conn, grpcClient := framework.DialVolumeServer(t, clusterHarness.VolumeGRPCAddress())
|
||||
defer conn.Close()
|
||||
|
||||
const volumeID = uint32(89)
|
||||
framework.AllocateVolume(t, grpcClient, volumeID, "")
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
|
||||
defer cancel()
|
||||
|
||||
_, err := grpcClient.VolumeTailReceiver(ctx, &volume_server_pb.VolumeTailReceiverRequest{
|
||||
VolumeId: volumeID,
|
||||
SourceVolumeServer: "127.0.0.1:19999.29999",
|
||||
SinceNs: 0,
|
||||
IdleTimeoutSeconds: 1,
|
||||
})
|
||||
if err == nil {
|
||||
t.Fatalf("VolumeTailReceiver should fail when source volume server is unavailable")
|
||||
}
|
||||
lowered := strings.ToLower(err.Error())
|
||||
if !strings.Contains(lowered, "dial") && !strings.Contains(lowered, "unavailable") && !strings.Contains(lowered, "connection refused") {
|
||||
t.Fatalf("VolumeTailReceiver source-unavailable error mismatch: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestVolumeTailReceiverReplicatesSourceUpdates(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("skipping integration test in short mode")
|
||||
@@ -204,3 +236,46 @@ func TestVolumeTailSenderLargeNeedleChunking(t *testing.T) {
|
||||
t.Fatalf("VolumeTailSender expected a final data chunk marked IsLastChunk=true")
|
||||
}
|
||||
}
|
||||
|
||||
func TestVolumeTailSenderStreamCancellation(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("skipping integration test in short mode")
|
||||
}
|
||||
|
||||
clusterHarness := framework.StartSingleVolumeCluster(t, matrix.P1())
|
||||
conn, grpcClient := framework.DialVolumeServer(t, clusterHarness.VolumeGRPCAddress())
|
||||
defer conn.Close()
|
||||
|
||||
const volumeID = uint32(74)
|
||||
framework.AllocateVolume(t, grpcClient, volumeID, "")
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
stream, err := grpcClient.VolumeTailSender(ctx, &volume_server_pb.VolumeTailSenderRequest{
|
||||
VolumeId: volumeID,
|
||||
SinceNs: 0,
|
||||
IdleTimeoutSeconds: 30,
|
||||
})
|
||||
if err != nil {
|
||||
cancel()
|
||||
t.Fatalf("VolumeTailSender start failed: %v", err)
|
||||
}
|
||||
|
||||
firstMsg, err := stream.Recv()
|
||||
if err != nil {
|
||||
cancel()
|
||||
t.Fatalf("VolumeTailSender first recv failed: %v", err)
|
||||
}
|
||||
if !firstMsg.GetIsLastChunk() {
|
||||
cancel()
|
||||
t.Fatalf("expected heartbeat first message before cancellation")
|
||||
}
|
||||
|
||||
cancel()
|
||||
_, err = stream.Recv()
|
||||
if err == nil {
|
||||
t.Fatalf("VolumeTailSender recv after cancellation expected error")
|
||||
}
|
||||
if status.Code(err) != codes.Canceled {
|
||||
t.Fatalf("VolumeTailSender cancellation expected grpc code Canceled, got %v (%v)", status.Code(err), err)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -164,6 +164,26 @@ func TestUiIndexNotExposedWhenJwtSigningEnabled(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestUiIndexExposedWhenJwtSigningEnabledAndAccessUITrue(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("skipping integration test in short mode")
|
||||
}
|
||||
|
||||
profile := matrix.P3()
|
||||
profile.AccessUI = true
|
||||
cluster := framework.StartSingleVolumeCluster(t, profile)
|
||||
client := framework.NewHTTPClient()
|
||||
|
||||
resp := framework.DoRequest(t, client, mustNewRequest(t, http.MethodGet, cluster.VolumeAdminURL()+"/ui/index.html"))
|
||||
body := framework.ReadAllAndClose(t, resp)
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
t.Fatalf("expected /ui/index.html to be exposed when access.ui=true under JWT profile, got %d body=%s", resp.StatusCode, string(body))
|
||||
}
|
||||
if !strings.Contains(strings.ToLower(string(body)), "volume") {
|
||||
t.Fatalf("ui page does not look like volume status page")
|
||||
}
|
||||
}
|
||||
|
||||
func mustNewRequest(t testing.TB, method, url string) *http.Request {
|
||||
t.Helper()
|
||||
req, err := http.NewRequest(method, url, nil)
|
||||
|
||||
@@ -251,6 +251,50 @@ func TestUnsupportedMethodConnectParity(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestUnsupportedMethodMkcolParity(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("skipping integration test in short mode")
|
||||
}
|
||||
|
||||
clusterHarness := framework.StartSingleVolumeCluster(t, matrix.P2())
|
||||
conn, grpcClient := framework.DialVolumeServer(t, clusterHarness.VolumeGRPCAddress())
|
||||
defer conn.Close()
|
||||
|
||||
const volumeID = uint32(88)
|
||||
framework.AllocateVolume(t, grpcClient, volumeID, "")
|
||||
|
||||
fid := framework.NewFileID(volumeID, 124004, 0x06060606)
|
||||
client := framework.NewHTTPClient()
|
||||
uploadResp := framework.UploadBytes(t, client, clusterHarness.VolumeAdminURL(), fid, []byte("mkcol-method-check"))
|
||||
_ = framework.ReadAllAndClose(t, uploadResp)
|
||||
if uploadResp.StatusCode != http.StatusCreated {
|
||||
t.Fatalf("upload expected 201, got %d", uploadResp.StatusCode)
|
||||
}
|
||||
|
||||
adminReq := mustNewRequest(t, "MKCOL", clusterHarness.VolumeAdminURL()+"/"+fid)
|
||||
adminResp := framework.DoRequest(t, client, adminReq)
|
||||
_ = framework.ReadAllAndClose(t, adminResp)
|
||||
if adminResp.StatusCode != http.StatusBadRequest {
|
||||
t.Fatalf("admin MKCOL expected 400, got %d", adminResp.StatusCode)
|
||||
}
|
||||
|
||||
publicReq := mustNewRequest(t, "MKCOL", clusterHarness.VolumePublicURL()+"/"+fid)
|
||||
publicResp := framework.DoRequest(t, client, publicReq)
|
||||
_ = framework.ReadAllAndClose(t, publicResp)
|
||||
if publicResp.StatusCode != http.StatusOK {
|
||||
t.Fatalf("public MKCOL expected passthrough 200, got %d", publicResp.StatusCode)
|
||||
}
|
||||
|
||||
verifyResp := framework.ReadBytes(t, client, clusterHarness.VolumeAdminURL(), fid)
|
||||
verifyBody := framework.ReadAllAndClose(t, verifyResp)
|
||||
if verifyResp.StatusCode != http.StatusOK {
|
||||
t.Fatalf("verify GET expected 200, got %d", verifyResp.StatusCode)
|
||||
}
|
||||
if string(verifyBody) != "mkcol-method-check" {
|
||||
t.Fatalf("MKCOL should not mutate data, got %q", string(verifyBody))
|
||||
}
|
||||
}
|
||||
|
||||
func TestPublicPortHeadReadParity(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("skipping integration test in short mode")
|
||||
|
||||
@@ -2,6 +2,7 @@ package volume_server_http_test
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"strconv"
|
||||
"testing"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/test/volume_server/framework"
|
||||
@@ -51,4 +52,17 @@ func TestReadDeletedQueryReturnsDeletedNeedleData(t *testing.T) {
|
||||
if string(readDeletedBody) != string(payload) {
|
||||
t.Fatalf("readDeleted body mismatch: got %q want %q", string(readDeletedBody), string(payload))
|
||||
}
|
||||
|
||||
headReadDeletedReq := mustNewRequest(t, http.MethodHead, clusterHarness.VolumeAdminURL()+"/"+fid+"?readDeleted=true")
|
||||
headReadDeletedResp := framework.DoRequest(t, client, headReadDeletedReq)
|
||||
headReadDeletedBody := framework.ReadAllAndClose(t, headReadDeletedResp)
|
||||
if headReadDeletedResp.StatusCode != http.StatusOK {
|
||||
t.Fatalf("HEAD with readDeleted=true expected 200, got %d", headReadDeletedResp.StatusCode)
|
||||
}
|
||||
if len(headReadDeletedBody) != 0 {
|
||||
t.Fatalf("HEAD with readDeleted=true expected empty body, got %d bytes", len(headReadDeletedBody))
|
||||
}
|
||||
if got := headReadDeletedResp.Header.Get("Content-Length"); got != strconv.Itoa(len(payload)) {
|
||||
t.Fatalf("HEAD with readDeleted=true content-length mismatch: got %q want %d", got, len(payload))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -164,6 +164,10 @@ func TestConditionalHeaderPrecedenceAndInvalidIfModifiedSince(t *testing.T) {
|
||||
if lastModified == "" {
|
||||
t.Fatalf("baseline read expected Last-Modified header")
|
||||
}
|
||||
etag := baselineResp.Header.Get("ETag")
|
||||
if etag == "" {
|
||||
t.Fatalf("baseline read expected ETag header")
|
||||
}
|
||||
|
||||
precedenceReq := mustNewRequest(t, http.MethodGet, clusterHarness.VolumeAdminURL()+"/"+fid)
|
||||
precedenceReq.Header.Set("If-Modified-Since", lastModified)
|
||||
@@ -188,4 +192,28 @@ func TestConditionalHeaderPrecedenceAndInvalidIfModifiedSince(t *testing.T) {
|
||||
if string(invalidIMSBody) != string(payload) {
|
||||
t.Fatalf("invalid If-Modified-Since fallback body mismatch: got %q want %q", string(invalidIMSBody), string(payload))
|
||||
}
|
||||
|
||||
headIMSPrecedenceReq := mustNewRequest(t, http.MethodHead, clusterHarness.VolumeAdminURL()+"/"+fid)
|
||||
headIMSPrecedenceReq.Header.Set("If-Modified-Since", lastModified)
|
||||
headIMSPrecedenceReq.Header.Set("If-None-Match", "\"definitely-different-etag\"")
|
||||
headIMSPrecedenceResp := framework.DoRequest(t, client, headIMSPrecedenceReq)
|
||||
headIMSPrecedenceBody := framework.ReadAllAndClose(t, headIMSPrecedenceResp)
|
||||
if headIMSPrecedenceResp.StatusCode != http.StatusNotModified {
|
||||
t.Fatalf("HEAD conditional precedence expected 304 when If-Modified-Since matches, got %d", headIMSPrecedenceResp.StatusCode)
|
||||
}
|
||||
if len(headIMSPrecedenceBody) != 0 {
|
||||
t.Fatalf("HEAD conditional precedence expected empty body, got %d bytes", len(headIMSPrecedenceBody))
|
||||
}
|
||||
|
||||
headETagPrecedenceReq := mustNewRequest(t, http.MethodHead, clusterHarness.VolumeAdminURL()+"/"+fid)
|
||||
headETagPrecedenceReq.Header.Set("If-Modified-Since", "Thu, 01 Jan 1970 00:00:00 GMT")
|
||||
headETagPrecedenceReq.Header.Set("If-None-Match", etag)
|
||||
headETagPrecedenceResp := framework.DoRequest(t, client, headETagPrecedenceReq)
|
||||
headETagPrecedenceBody := framework.ReadAllAndClose(t, headETagPrecedenceResp)
|
||||
if headETagPrecedenceResp.StatusCode != http.StatusNotModified {
|
||||
t.Fatalf("HEAD conditional precedence expected 304 when If-None-Match matches, got %d", headETagPrecedenceResp.StatusCode)
|
||||
}
|
||||
if len(headETagPrecedenceBody) != 0 {
|
||||
t.Fatalf("HEAD conditional etag precedence expected empty body, got %d bytes", len(headETagPrecedenceBody))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package volume_server_http_test
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"net/http"
|
||||
"strings"
|
||||
"testing"
|
||||
@@ -72,3 +73,64 @@ func TestWriteMalformedMultipartAndMD5Mismatch(t *testing.T) {
|
||||
t.Fatalf("content-md5 mismatch response should mention Content-MD5, got %q", string(md5MismatchBody))
|
||||
}
|
||||
}
|
||||
|
||||
func TestWriteRejectsPayloadOverFileSizeLimit(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("skipping integration test in short mode")
|
||||
}
|
||||
|
||||
profile := matrix.P1()
|
||||
profile.FileSizeLimitMB = 1
|
||||
clusterHarness := framework.StartSingleVolumeCluster(t, profile)
|
||||
conn, grpcClient := framework.DialVolumeServer(t, clusterHarness.VolumeGRPCAddress())
|
||||
defer conn.Close()
|
||||
|
||||
const volumeID = uint32(99)
|
||||
framework.AllocateVolume(t, grpcClient, volumeID, "")
|
||||
|
||||
client := framework.NewHTTPClient()
|
||||
fid := framework.NewFileID(volumeID, 772002, 0x2A3B4C5D)
|
||||
oversizedPayload := bytes.Repeat([]byte("z"), 1024*1024+1)
|
||||
|
||||
oversizedReq := newUploadRequest(t, clusterHarness.VolumeAdminURL()+"/"+fid, oversizedPayload)
|
||||
oversizedResp := framework.DoRequest(t, client, oversizedReq)
|
||||
oversizedBody := framework.ReadAllAndClose(t, oversizedResp)
|
||||
if oversizedResp.StatusCode != http.StatusBadRequest {
|
||||
t.Fatalf("oversized write expected 400, got %d", oversizedResp.StatusCode)
|
||||
}
|
||||
if !strings.Contains(strings.ToLower(string(oversizedBody)), "limited") {
|
||||
t.Fatalf("oversized write response should mention limit, got %q", string(oversizedBody))
|
||||
}
|
||||
}
|
||||
|
||||
func TestReplicatedWriteFailsWhenReplicaRequirementsNotMet(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("skipping integration test in short mode")
|
||||
}
|
||||
|
||||
clusterHarness := framework.StartSingleVolumeCluster(t, matrix.P1())
|
||||
conn, grpcClient := framework.DialVolumeServer(t, clusterHarness.VolumeGRPCAddress())
|
||||
defer conn.Close()
|
||||
|
||||
const volumeID = uint32(109)
|
||||
framework.AllocateVolumeWithReplication(t, grpcClient, volumeID, "", "001")
|
||||
|
||||
client := framework.NewHTTPClient()
|
||||
fid := framework.NewFileID(volumeID, 772003, 0x3A4B5C6D)
|
||||
payload := []byte("replicated-write-failure-path")
|
||||
|
||||
writeResp := framework.UploadBytes(t, client, clusterHarness.VolumeAdminURL(), fid, payload)
|
||||
writeBody := framework.ReadAllAndClose(t, writeResp)
|
||||
if writeResp.StatusCode != http.StatusInternalServerError {
|
||||
t.Fatalf("replicated write with unmet replication requirements expected 500, got %d body=%s", writeResp.StatusCode, string(writeBody))
|
||||
}
|
||||
if !strings.Contains(strings.ToLower(string(writeBody)), "replica") {
|
||||
t.Fatalf("replicated write failure response should mention replica write failure, got %q", string(writeBody))
|
||||
}
|
||||
|
||||
readResp := framework.ReadBytes(t, client, clusterHarness.VolumeAdminURL(), fid)
|
||||
_ = framework.ReadAllAndClose(t, readResp)
|
||||
if readResp.StatusCode != http.StatusNotFound {
|
||||
t.Fatalf("local read after failed replicate write expected 404 (not committed locally), got %d", readResp.StatusCode)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -12,12 +12,14 @@ type Profile struct {
|
||||
EnableJWT bool
|
||||
JWTSigningKey string
|
||||
JWTReadKey string
|
||||
AccessUI bool
|
||||
EnableMaintain bool
|
||||
|
||||
ConcurrentUploadLimitMB int
|
||||
ConcurrentDownloadLimitMB int
|
||||
InflightUploadTimeout time.Duration
|
||||
InflightDownloadTimeout time.Duration
|
||||
FileSizeLimitMB int
|
||||
|
||||
ReplicatedLayout bool
|
||||
HasErasureCoding bool
|
||||
|
||||
Reference in New Issue
Block a user