Compare commits

...
Author SHA1 Message Date
Chris Lu 95af3f3649 docs(volume_server): use repo-relative paths in plan/readme 2026-05-07 22:29:14 -07:00
Chris Lu 7e6db15a88 docs(volume-server): record slash fid-path native validation progress 2026-02-16 05:06:57 -08:00
Chris Lu e66983e9bd feat(rust-volume-server): broaden native slash fid-path validation 2026-02-16 05:06:21 -08:00
Chris Lu b010af88fb docs(volume-server): log native write prevalidation parity progress 2026-02-16 04:37:41 -08:00
Chris Lu 94cefd6f4c feat(rust-volume-server): add native write prevalidation for multipart md5 and size limits 2026-02-16 04:37:11 -08:00
Chris Lu 47156eb8ce docs(volume-server): record native malformed-fid parity milestone 2026-02-16 04:07:51 -08:00
Chris Lu 1ce0174b23 feat(rust-volume-server): validate malformed fid routes in native http path 2026-02-16 04:07:22 -08:00
Chris Lu 120aa956a0 docs(volume-server): update native parity progress logs 2026-02-16 03:36:34 -08:00
Chris Lu d6ff6ed6d4 feat(rust-volume-server): mirror backend healthz and normalize absolute request targets 2026-02-16 03:36:06 -08:00
Chris Lu 532bf863ad docs(volume-server): log native ui static method parity progress 2026-02-16 03:00:10 -08:00
Chris Lu 23e4497b21 feat(rust-volume-server): add native ui static and public no-op handlers 2026-02-16 02:44:14 -08:00
Chris Lu b0447d2479 docs(volume-server): track native rust options parity milestone 2026-02-16 01:52:46 -08:00
Chris Lu fbff2cb39a feat(rust-volume-server): serve native OPTIONS on admin and public listeners 2026-02-16 01:52:02 -08:00
Chris Lu d2a6066181 docs(volume-server): log native rust status healthz milestone 2026-02-16 01:22:40 -08:00
Chris Lu 7e6e0261ab feat(rust-volume-server): serve native /status and /healthz in native mode 2026-02-16 01:22:06 -08:00
Chris Lu c7c7be42ed docs(volume-server): track native launcher bootstrap progress 2026-02-16 00:44:31 -08:00
Chris Lu 2e65966c06 feat(rust-volume-server): default rust launcher mode to native 2026-02-16 00:43:35 -08:00
Chris Lu 61befd10fc ci(volume-server): include native rust mode in smoke matrix 2026-02-16 00:15:45 -08:00
Chris Lu 70ddbee370 feat(rust-volume-server): add native mode bootstrap entrypoint 2026-02-16 00:15:33 -08:00
Chris Lu 14c863dbff docs(volume-server): refocus plan on native rust parity 2026-02-16 00:09:31 -08:00
Chris Lu 6bb9d8bac2 docs(volume_server): log head readDeleted parity coverage 2026-02-15 23:51:50 -08:00
Chris Lu cc80ad3643 test(volume_server/http): add head readDeleted parity coverage 2026-02-15 23:51:37 -08:00
Chris Lu 9009e38f7b docs(volume_server): log ping volume-server unreachable coverage 2026-02-15 23:50:08 -08:00
Chris Lu b9fbb85af2 test(volume_server/grpc): add ping unreachable volume-server target case 2026-02-15 23:49:53 -08:00
Chris Lu 47d3001572 docs(volume_server): log csv query payload parity coverage 2026-02-15 23:48:35 -08:00
Chris Lu a12dd5f8d3 test(volume_server/grpc): cover csv-query payload no-output parity 2026-02-15 23:48:22 -08:00
Chris Lu 8e614486a3 docs(volume_server): log tail-receiver interruption coverage 2026-02-15 23:47:06 -08:00
Chris Lu a5864c3eb6 test(volume_server/grpc): cover tail-receiver source-unavailable branch 2026-02-15 23:46:55 -08:00
Chris Lu 6302809442 docs(volume_server): log tail sender cancellation coverage 2026-02-15 18:45:38 -08:00
Chris Lu 27a80f7607 test(volume_server/grpc): add tail-sender cancellation interruption coverage 2026-02-15 18:45:23 -08:00
Chris Lu ec429e0361 docs(volume_server): log framework port-range hardening and rerun 2026-02-15 18:44:13 -08:00
Chris Lu 90e82b15ce test(volume_server/framework): allocate volume ports within safe grpc-offset range 2026-02-15 18:43:57 -08:00
Chris Lu a3e1ee1653 docs(volume_server): log mkcol method parity coverage 2026-02-15 18:13:51 -08:00
Chris Lu 2ab30900d4 test(volume_server/http): add mkcol unsupported-method parity 2026-02-15 18:13:40 -08:00
Chris Lu 62ee14fa61 docs(volume_server): log read-all-needles multi-volume coverage 2026-02-15 17:58:44 -08:00
Chris Lu ab95a6ef15 test(volume_server/grpc): cover read-all-needles multi-volume success 2026-02-15 17:58:34 -08:00
Chris Lu 24965fd489 docs(volume_server): log head conditional precedence coverage 2026-02-15 17:56:47 -08:00
Chris Lu ed23e290fc test(volume_server/http): expand head conditional precedence coverage 2026-02-15 17:56:36 -08:00
Chris Lu 9b57fb6961 docs(volume_server): log ec batch delete success coverage 2026-02-15 17:55:43 -08:00
Chris Lu 1bb40b6bc5 test(volume_server/grpc): add ec batch delete success coverage 2026-02-15 17:55:29 -08:00
Chris Lu 34e342da63 docs(volume_server): log replicated write failure coverage 2026-02-15 17:45:36 -08:00
Chris Lu 4835d34438 test(volume_server/http): cover replicated write failure when replication unmet 2026-02-15 17:45:25 -08:00
Chris Lu 5814729def docs(volume_server): log ec-only read meta coverage 2026-02-15 14:49:31 -08:00
Chris Lu 37bf9b5ebf test(volume_server/grpc): cover ec-only read needle meta unsupported path 2026-02-15 14:49:18 -08:00
Chris Lu 19201df6d7 docs(volume_server): log oversized upload limit coverage 2026-02-15 14:47:18 -08:00
Chris Lu 4d61cbdeed test(volume_server/http): cover oversized upload file-size limit rejection 2026-02-15 14:47:09 -08:00
Chris Lu 3ce883624e docs(volume_server): log jwt ui access override coverage 2026-02-15 14:45:40 -08:00
Chris Lu de974c05d5 test(volume_server/http): cover jwt ui access override behavior 2026-02-15 14:45:28 -08:00
Chris Lu 7768fda023 docs(volume_server): record proxy-mode validation and CI matrix 2026-02-15 14:29:28 -08:00
Chris Lu 548b3d9a38 ci(volume_server): run rust smoke tests in exec and proxy modes 2026-02-15 14:28:35 -08:00
Chris Lu a7f50d23b5 feat(rust/volume_server): add proxy supervision mode for integration parity 2026-02-15 14:28:12 -08:00
Chris Lu 6ce4d7eded docs(volume_server): record rust-mode full-suite validation 2026-02-15 12:11:31 -08:00
Chris Lu 3bd20e6a10 chore(rust/volume_server): add Cargo.lock 2026-02-15 11:55:42 -08:00
Chris Lu d402573ea8 docs(volume_server): document rust-mode harness and tracking 2026-02-15 11:55:33 -08:00
Chris Lu 63d08e8a91 ci(volume_server): add rust-mode integration smoke job 2026-02-15 11:55:28 -08:00
Chris Lu 880c2e1dab feat(rust/volume_server): add compatibility launcher and migration plan 2026-02-15 11:55:25 -08:00
Chris Lu 7beab85c21 test(volume_server/framework): support selectable volume server binary 2026-02-15 11:55:15 -08:00
23 changed files with 2507 additions and 51 deletions
@@ -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"
+7
View File
@@ -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"
+6
View File
@@ -0,0 +1,6 @@
[package]
name = "weed-volume-rs"
version = "0.1.0"
edition = "2021"
[dependencies]
+245
View File
@@ -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
+205
View File
@@ -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`
+4 -1
View File
@@ -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
+5 -1
View File
@@ -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`.
+136 -25
View File
@@ -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() {
+24 -22
View File
@@ -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
+10 -2
View File
@@ -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()
}
+64
View File
@@ -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")
+75
View File
@@ -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)
}
}
+20
View File
@@ -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