mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-11 09:05:50 +00:00
Compare commits
37
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
50f25bb5cd | ||
|
|
512912cbb8 | ||
|
|
8d6c5cbb58 | ||
|
|
f3151900e4 | ||
|
|
7aaa431bb4 | ||
|
|
8049fcc516 | ||
|
|
06cbd2acdf | ||
|
|
2ee6907c19 | ||
|
|
cc5b246973 | ||
|
|
36ae7e04b5 | ||
|
|
46c0e56bb8 | ||
|
|
baa65c3823 | ||
|
|
f4bfe60549 | ||
|
|
67a2810d2d | ||
|
|
80db692728 | ||
|
|
ae08e77979 | ||
|
|
28d1ef24ec | ||
|
|
edf7d2a074 | ||
|
|
10e7f0f2bc | ||
|
|
9cae95d749 | ||
|
|
e8a8449553 | ||
|
|
b37bbf541a | ||
|
|
10b0bdce02 | ||
|
|
e2c79af6ec | ||
|
|
388cc018ab | ||
|
|
41ff105f47 | ||
|
|
c390448906 | ||
|
|
e648c76bcf | ||
|
|
066f7c3a0d | ||
|
|
ae724ac9d5 | ||
|
|
2e64c0fe2a | ||
|
|
ef30d91b7d | ||
|
|
8aa5809824 | ||
|
|
39e76b8e94 | ||
|
|
2a7ec8d033 | ||
|
|
07cd741380 | ||
|
|
2264941a17 |
@@ -142,7 +142,7 @@ jobs:
|
||||
password: ${{ secrets.GHCR_TOKEN }}
|
||||
|
||||
- name: Build
|
||||
uses: docker/build-push-action@v7
|
||||
uses: docker/build-push-action@v7.1.0
|
||||
with:
|
||||
context: ./docker
|
||||
push: ${{ github.event_name != 'pull_request' }}
|
||||
|
||||
@@ -150,7 +150,7 @@ jobs:
|
||||
fi
|
||||
|
||||
- name: Build and push image
|
||||
uses: docker/build-push-action@v7
|
||||
uses: docker/build-push-action@v7.1.0
|
||||
with:
|
||||
context: ./docker
|
||||
push: ${{ github.event_name != 'pull_request' }}
|
||||
|
||||
@@ -232,7 +232,7 @@ jobs:
|
||||
username: ${{ secrets.GHCR_USERNAME }}
|
||||
password: ${{ secrets.GHCR_TOKEN }}
|
||||
- name: Build ${{ matrix.platform }} ${{ matrix.variant }}
|
||||
uses: docker/build-push-action@v7
|
||||
uses: docker/build-push-action@v7.1.0
|
||||
env:
|
||||
DOCKER_BUILDKIT: 1
|
||||
with:
|
||||
@@ -322,7 +322,7 @@ jobs:
|
||||
buildkitd-config: /tmp/buildkitd.toml
|
||||
- name: Build local scan image tarball
|
||||
if: needs.setup.outputs.publish != 'true'
|
||||
uses: docker/build-push-action@v7
|
||||
uses: docker/build-push-action@v7.1.0
|
||||
env:
|
||||
DOCKER_BUILDKIT: 1
|
||||
with:
|
||||
|
||||
@@ -57,7 +57,7 @@ jobs:
|
||||
fi
|
||||
-
|
||||
name: Build
|
||||
uses: docker/build-push-action@v7
|
||||
uses: docker/build-push-action@v7.1.0
|
||||
with:
|
||||
context: ./docker
|
||||
push: ${{ github.event_name != 'pull_request' }}
|
||||
|
||||
@@ -251,7 +251,7 @@ jobs:
|
||||
|
||||
- name: Build and push ${{ matrix.variant }}
|
||||
if: github.event_name != 'workflow_dispatch' || github.event.inputs.variant == 'all' || github.event.inputs.variant == matrix.variant
|
||||
uses: docker/build-push-action@v7
|
||||
uses: docker/build-push-action@v7.1.0
|
||||
env:
|
||||
DOCKER_BUILDKIT: 1
|
||||
with:
|
||||
|
||||
@@ -94,7 +94,7 @@ jobs:
|
||||
password: ${{ secrets.DOCKER_PASSWORD }}
|
||||
|
||||
- name: Build and push image
|
||||
uses: docker/build-push-action@d08e5c354a6adb9ed34480a06d141179aa583294 # v2
|
||||
uses: docker/build-push-action@bcafcacb16a39f128d818304e6c9c0c18556b85f # v2
|
||||
with:
|
||||
context: ./docker
|
||||
push: true
|
||||
|
||||
@@ -84,7 +84,7 @@ jobs:
|
||||
rm weed-volume-normal
|
||||
|
||||
- name: Upload dev release assets
|
||||
uses: softprops/action-gh-release@v2
|
||||
uses: softprops/action-gh-release@v3
|
||||
with:
|
||||
tag_name: dev
|
||||
prerelease: true
|
||||
@@ -155,7 +155,7 @@ jobs:
|
||||
rm weed-volume-normal
|
||||
|
||||
- name: Upload dev release assets
|
||||
uses: softprops/action-gh-release@v2
|
||||
uses: softprops/action-gh-release@v3
|
||||
with:
|
||||
tag_name: dev
|
||||
prerelease: true
|
||||
|
||||
@@ -89,7 +89,7 @@ jobs:
|
||||
|
||||
- name: Upload release assets
|
||||
if: startsWith(github.ref, 'refs/tags/')
|
||||
uses: softprops/action-gh-release@v2
|
||||
uses: softprops/action-gh-release@v3
|
||||
with:
|
||||
files: |
|
||||
weed-volume_large_disk_${{ matrix.asset_suffix }}.tar.gz
|
||||
@@ -166,7 +166,7 @@ jobs:
|
||||
|
||||
- name: Upload release assets
|
||||
if: startsWith(github.ref, 'refs/tags/')
|
||||
uses: softprops/action-gh-release@v2
|
||||
uses: softprops/action-gh-release@v3
|
||||
with:
|
||||
files: |
|
||||
weed-volume_large_disk_${{ matrix.asset_suffix }}.tar.gz
|
||||
@@ -235,7 +235,7 @@ jobs:
|
||||
|
||||
- name: Upload release assets
|
||||
if: startsWith(github.ref, 'refs/tags/')
|
||||
uses: softprops/action-gh-release@v2
|
||||
uses: softprops/action-gh-release@v3
|
||||
with:
|
||||
files: |
|
||||
weed-volume_large_disk_windows_amd64.zip
|
||||
|
||||
@@ -0,0 +1,56 @@
|
||||
name: "Vacuum Integration Tests"
|
||||
|
||||
on:
|
||||
push:
|
||||
branches: [ master ]
|
||||
pull_request:
|
||||
branches: [ master ]
|
||||
|
||||
permissions:
|
||||
contents: read
|
||||
|
||||
jobs:
|
||||
vacuum-integration-tests:
|
||||
name: Vacuum Integration Tests
|
||||
runs-on: ubuntu-22.04
|
||||
timeout-minutes: 15
|
||||
steps:
|
||||
- name: Set up Go 1.x
|
||||
uses: actions/setup-go@v6
|
||||
with:
|
||||
go-version: ^1.25
|
||||
id: go
|
||||
|
||||
- name: Check out code into the Go module directory
|
||||
uses: actions/checkout@v6
|
||||
|
||||
- name: Build weed binary
|
||||
run: |
|
||||
cd weed && go build -o weed .
|
||||
|
||||
- name: Run Vacuum Integration Tests
|
||||
working-directory: test/vacuum
|
||||
run: |
|
||||
go test -v -timeout 10m
|
||||
|
||||
- name: Collect server logs on failure
|
||||
if: failure()
|
||||
run: |
|
||||
echo "Collecting server logs from temp directories..."
|
||||
mkdir -p /tmp/vacuum-test-logs
|
||||
find /tmp -maxdepth 1 -type d -name "TestVacuum*" 2>/dev/null | while read dir; do
|
||||
if [ -d "$dir" ]; then
|
||||
echo "Found test directory: $dir"
|
||||
cp -r "$dir" /tmp/vacuum-test-logs/ 2>/dev/null || true
|
||||
fi
|
||||
done
|
||||
echo "Collected logs:"
|
||||
find /tmp/vacuum-test-logs -type f -name "*.log" 2>/dev/null || echo "No logs found"
|
||||
|
||||
- name: Archive logs
|
||||
if: failure()
|
||||
uses: actions/upload-artifact@v7
|
||||
with:
|
||||
name: vacuum-integration-test-logs
|
||||
path: /tmp/vacuum-test-logs/
|
||||
retention-days: 14
|
||||
@@ -92,14 +92,14 @@ require (
|
||||
gocloud.dev v0.45.0
|
||||
gocloud.dev/pubsub/natspubsub v0.45.0
|
||||
gocloud.dev/pubsub/rabbitpubsub v0.45.0
|
||||
golang.org/x/crypto v0.49.0
|
||||
golang.org/x/crypto v0.50.0
|
||||
golang.org/x/exp v0.0.0-20260218203240-3dfff04db8fa
|
||||
golang.org/x/image v0.38.0
|
||||
golang.org/x/net v0.52.0
|
||||
golang.org/x/net v0.53.0
|
||||
golang.org/x/oauth2 v0.36.0
|
||||
golang.org/x/sys v0.42.0
|
||||
golang.org/x/text v0.35.0 // indirect
|
||||
golang.org/x/tools v0.42.0 // indirect
|
||||
golang.org/x/sys v0.43.0
|
||||
golang.org/x/text v0.36.0 // indirect
|
||||
golang.org/x/tools v0.43.0 // indirect
|
||||
golang.org/x/xerrors v0.0.0-20240903120638-7835f813f4da // indirect
|
||||
google.golang.org/api v0.274.0
|
||||
google.golang.org/genproto v0.0.0-20260316180232-0b37fe3546d5 // indirect
|
||||
@@ -118,15 +118,15 @@ require (
|
||||
github.com/Jille/raft-grpc-transport v1.6.1
|
||||
github.com/ThreeDotsLabs/watermill v1.5.1
|
||||
github.com/a-h/templ v0.3.977
|
||||
github.com/apache/cassandra-gocql-driver/v2 v2.0.0
|
||||
github.com/apache/cassandra-gocql-driver/v2 v2.1.0
|
||||
github.com/apache/iceberg-go v0.5.0
|
||||
github.com/apple/foundationdb/bindings/go v0.0.0-20250911184653-27f7192f47c3
|
||||
github.com/arangodb/go-driver v1.6.9
|
||||
github.com/armon/go-metrics v0.4.1
|
||||
github.com/aws/aws-sdk-go-v2 v1.41.5
|
||||
github.com/aws/aws-sdk-go-v2/config v1.32.13
|
||||
github.com/aws/aws-sdk-go-v2/config v1.32.14
|
||||
github.com/aws/aws-sdk-go-v2/credentials v1.19.14
|
||||
github.com/aws/aws-sdk-go-v2/service/s3 v1.98.0
|
||||
github.com/aws/aws-sdk-go-v2/service/s3 v1.99.0
|
||||
github.com/cognusion/imaging v1.0.2
|
||||
github.com/fluent/fluent-logger-golang v1.10.1
|
||||
github.com/getsentry/sentry-go v0.44.1
|
||||
@@ -154,7 +154,7 @@ require (
|
||||
github.com/tikv/client-go/v2 v2.0.7
|
||||
github.com/xeipuuv/gojsonschema v1.2.0
|
||||
github.com/ydb-platform/ydb-go-sdk-auth-environ v0.5.1
|
||||
github.com/ydb-platform/ydb-go-sdk/v3 v3.125.3
|
||||
github.com/ydb-platform/ydb-go-sdk/v3 v3.134.0
|
||||
go.etcd.io/etcd/client/pkg/v3 v3.6.10
|
||||
go.uber.org/atomic v1.11.0
|
||||
golang.org/x/sync v0.20.0
|
||||
@@ -282,8 +282,8 @@ require (
|
||||
go.uber.org/mock v0.5.2 // indirect
|
||||
go.yaml.in/yaml/v2 v2.4.3 // indirect
|
||||
go.yaml.in/yaml/v3 v3.0.4 // indirect
|
||||
golang.org/x/mod v0.33.0 // indirect
|
||||
golang.org/x/telemetry v0.0.0-20260209163413-e7419c687ee4 // indirect
|
||||
golang.org/x/mod v0.34.0 // indirect
|
||||
golang.org/x/telemetry v0.0.0-20260311193753-579e4da9a98c // indirect
|
||||
gonum.org/v1/gonum v0.17.0 // indirect
|
||||
)
|
||||
|
||||
@@ -476,7 +476,7 @@ require (
|
||||
github.com/vmihailenco/tagparser/v2 v2.0.0 // indirect
|
||||
github.com/xanzy/ssh-agent v0.3.3 // indirect
|
||||
github.com/yandex-cloud/go-genproto v0.0.0-20211115083454-9ca41db5ed9e // indirect
|
||||
github.com/ydb-platform/ydb-go-genproto v0.0.0-20251125145508-6d7ef87db5cb // indirect
|
||||
github.com/ydb-platform/ydb-go-genproto v0.0.0-20260311095541-ebbf792c1180 // indirect
|
||||
github.com/ydb-platform/ydb-go-yc v0.12.1 // indirect
|
||||
github.com/ydb-platform/ydb-go-yc-metadata v0.6.1 // indirect
|
||||
github.com/yunify/qingstor-sdk-go/v3 v3.2.0 // indirect
|
||||
@@ -496,7 +496,7 @@ require (
|
||||
go.opentelemetry.io/otel/trace v1.43.0 // indirect
|
||||
go.uber.org/multierr v1.11.0 // indirect
|
||||
go.uber.org/zap v1.27.1 // indirect
|
||||
golang.org/x/term v0.41.0
|
||||
golang.org/x/term v0.42.0
|
||||
golang.org/x/time v0.15.0 // indirect
|
||||
google.golang.org/genproto/googleapis/api v0.0.0-20260316180232-0b37fe3546d5 // indirect
|
||||
google.golang.org/genproto/googleapis/rpc v0.0.0-20260319201613-d00831a3d3e7 // indirect
|
||||
|
||||
@@ -692,8 +692,8 @@ github.com/antlr4-go/antlr/v4 v4.13.1/go.mod h1:GKmUxMtwp6ZgGwZSva4eWPC5mS6vUAmO
|
||||
github.com/apache/arrow-go/v18 v18.5.2-0.20260220015023-a886a5722b87 h1:r/gg2gzUXiXoy72VU3jnODh8l/5rL0aslnTAbtmai/U=
|
||||
github.com/apache/arrow-go/v18 v18.5.2-0.20260220015023-a886a5722b87/go.mod h1:IJTMBTlHe7cDOhRh0ioGuEKBl5iTR6xPfl5BN4AgirU=
|
||||
github.com/apache/arrow/go/v10 v10.0.1/go.mod h1:YvhnlEePVnBS4+0z3fhPfUy7W1Ikj0Ih0vcRo/gZ1M0=
|
||||
github.com/apache/cassandra-gocql-driver/v2 v2.0.0 h1:Omnzb1Z/P90Dr2TbVNu54ICQL7TKVIIsJO231w484HU=
|
||||
github.com/apache/cassandra-gocql-driver/v2 v2.0.0/go.mod h1:QH/asJjB3mHvY6Dot6ZKMMpTcOrWJ8i9GhsvG1g0PK4=
|
||||
github.com/apache/cassandra-gocql-driver/v2 v2.1.0 h1:VEbbeJ2ift4deKMZ6Fs55Vs3fq/RrkjCcxCnqUxhwf8=
|
||||
github.com/apache/cassandra-gocql-driver/v2 v2.1.0/go.mod h1:QH/asJjB3mHvY6Dot6ZKMMpTcOrWJ8i9GhsvG1g0PK4=
|
||||
github.com/apache/iceberg-go v0.5.0 h1:wQj4CK5YiXZcB+tj19gWG+Jf1I6MiORQ/StSL/E5gGQ=
|
||||
github.com/apache/iceberg-go v0.5.0/go.mod h1:F/rdP1yZmnO4mQ0Qew2HTGdc+ZV57cRfxbbq/uJm1eM=
|
||||
github.com/apache/thrift v0.16.0/go.mod h1:PHK3hniurgQaNMZYaCLEqXKsYK8upmhPbmdP2FXSqgU=
|
||||
@@ -718,8 +718,8 @@ github.com/aws/aws-sdk-go-v2 v1.41.5 h1:dj5kopbwUsVUVFgO4Fi5BIT3t4WyqIDjGKCangnV
|
||||
github.com/aws/aws-sdk-go-v2 v1.41.5/go.mod h1:mwsPRE8ceUUpiTgF7QmQIJ7lgsKUPQOUl3o72QBrE1o=
|
||||
github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.8 h1:eBMB84YGghSocM7PsjmmPffTa+1FBUeNvGvFou6V/4o=
|
||||
github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.8/go.mod h1:lyw7GFp3qENLh7kwzf7iMzAxDn+NzjXEAGjKS2UOKqI=
|
||||
github.com/aws/aws-sdk-go-v2/config v1.32.13 h1:5KgbxMaS2coSWRrx9TX/QtWbqzgQkOdEa3sZPhBhCSg=
|
||||
github.com/aws/aws-sdk-go-v2/config v1.32.13/go.mod h1:8zz7wedqtCbw5e9Mi2doEwDyEgHcEE9YOJp6a8jdSMY=
|
||||
github.com/aws/aws-sdk-go-v2/config v1.32.14 h1:opVIRo/ZbbI8OIqSOKmpFaY7IwfFUOCCXBsUpJOwDdI=
|
||||
github.com/aws/aws-sdk-go-v2/config v1.32.14/go.mod h1:U4/V0uKxh0Tl5sxmCBZ3AecYny4UNlVmObYjKuuaiOo=
|
||||
github.com/aws/aws-sdk-go-v2/credentials v1.19.14 h1:n+UcGWAIZHkXzYt87uMFBv/l8THYELoX6gVcUvgl6fI=
|
||||
github.com/aws/aws-sdk-go-v2/credentials v1.19.14/go.mod h1:cJKuyWB59Mqi0jM3nFYQRmnHVQIcgoxjEMAbLkpr62w=
|
||||
github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.21 h1:NUS3K4BTDArQqNu2ih7yeDLaS3bmHD0YndtA6UP884g=
|
||||
@@ -742,8 +742,8 @@ github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.21 h1:c31//R3x
|
||||
github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.21/go.mod h1:r6+pf23ouCB718FUxaqzZdbpYFyDtehyZcmP5KL9FkA=
|
||||
github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.21 h1:ZlvrNcHSFFWURB8avufQq9gFsheUgjVD9536obIknfM=
|
||||
github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.21/go.mod h1:cv3TNhVrssKR0O/xxLJVRfd2oazSnZnkUeTf6ctUwfQ=
|
||||
github.com/aws/aws-sdk-go-v2/service/s3 v1.98.0 h1:foqo/ocQ7WqKwy3FojGtZQJo0FR4vto9qnz9VaumbCo=
|
||||
github.com/aws/aws-sdk-go-v2/service/s3 v1.98.0/go.mod h1:uoA43SdFwacedBfSgfFSjjCvYe8aYBS7EnU5GZ/YKMM=
|
||||
github.com/aws/aws-sdk-go-v2/service/s3 v1.99.0 h1:hlSuz394kV0vhv9drL5lhuEFbEOEP1VyQpy15qWh1Pk=
|
||||
github.com/aws/aws-sdk-go-v2/service/s3 v1.99.0/go.mod h1:uoA43SdFwacedBfSgfFSjjCvYe8aYBS7EnU5GZ/YKMM=
|
||||
github.com/aws/aws-sdk-go-v2/service/signin v1.0.9 h1:QKZH0S178gCmFEgst8hN0mCX1KxLgHBKKY/CLqwP8lg=
|
||||
github.com/aws/aws-sdk-go-v2/service/signin v1.0.9/go.mod h1:7yuQJoT+OoH8aqIxw9vwF+8KpvLZ8AWmvmUWHsGQZvI=
|
||||
github.com/aws/aws-sdk-go-v2/service/sns v1.39.7 h1:fovS7qGMT+BBSuifkySdVaMWxXTyaYT6qaBx/1y6Ij4=
|
||||
@@ -2058,14 +2058,14 @@ github.com/yandex-cloud/go-genproto v0.0.0-20211115083454-9ca41db5ed9e h1:9LPdmD
|
||||
github.com/yandex-cloud/go-genproto v0.0.0-20211115083454-9ca41db5ed9e/go.mod h1:HEUYX/p8966tMUHHT+TsS0hF/Ca/NYwqprC5WXSDMfE=
|
||||
github.com/ydb-platform/ydb-go-genproto v0.0.0-20221215182650-986f9d10542f/go.mod h1:Er+FePu1dNUieD+XTMDduGpQuCPssK5Q4BjF+IIXJ3I=
|
||||
github.com/ydb-platform/ydb-go-genproto v0.0.0-20230528143953-42c825ace222/go.mod h1:Er+FePu1dNUieD+XTMDduGpQuCPssK5Q4BjF+IIXJ3I=
|
||||
github.com/ydb-platform/ydb-go-genproto v0.0.0-20251125145508-6d7ef87db5cb h1:LZ6dhVfWzhicf/P5Xh7fA0Jd7rfGduxmB2QZpD+Lz9Q=
|
||||
github.com/ydb-platform/ydb-go-genproto v0.0.0-20251125145508-6d7ef87db5cb/go.mod h1:Er+FePu1dNUieD+XTMDduGpQuCPssK5Q4BjF+IIXJ3I=
|
||||
github.com/ydb-platform/ydb-go-genproto v0.0.0-20260311095541-ebbf792c1180 h1:avIdi8eGXjKbn1WLokNR1Ofnz1k8t7tJ88YQLD/iCi8=
|
||||
github.com/ydb-platform/ydb-go-genproto v0.0.0-20260311095541-ebbf792c1180/go.mod h1:Er+FePu1dNUieD+XTMDduGpQuCPssK5Q4BjF+IIXJ3I=
|
||||
github.com/ydb-platform/ydb-go-sdk-auth-environ v0.5.1 h1:XaRxeVrOyl3y6v9CiYMWaFdZ6zevvYe+TRxOR8ifa2s=
|
||||
github.com/ydb-platform/ydb-go-sdk-auth-environ v0.5.1/go.mod h1:9YzkhlIymWaJGX6KMU3vh5sOf3UKbCXkG/ZdjaI3zNM=
|
||||
github.com/ydb-platform/ydb-go-sdk/v3 v3.44.0/go.mod h1:oSLwnuilwIpaF5bJJMAofnGgzPJusoI3zWMNb8I+GnM=
|
||||
github.com/ydb-platform/ydb-go-sdk/v3 v3.47.3/go.mod h1:bWnOIcUHd7+Sl7DN+yhyY1H/I61z53GczvwJgXMgvj0=
|
||||
github.com/ydb-platform/ydb-go-sdk/v3 v3.125.3 h1:hTwpF+PdbuR7vcixN+4AC6yu4asaUIAljQCxL51eDII=
|
||||
github.com/ydb-platform/ydb-go-sdk/v3 v3.125.3/go.mod h1:stS1mQYjbJvwwYaYzKyFY9eMiuVXWWXQA6T+SpOLg9c=
|
||||
github.com/ydb-platform/ydb-go-sdk/v3 v3.134.0 h1:Voog4d56wPNT8JQvbgN6QPdZ1gKG36YM5615PJB/XpM=
|
||||
github.com/ydb-platform/ydb-go-sdk/v3 v3.134.0/go.mod h1:VYUUkRJkKuQPkIpgtZJj6+58Fa2g8ccAqdmaaK6HP5k=
|
||||
github.com/ydb-platform/ydb-go-yc v0.12.1 h1:qw3Fa+T81+Kpu5Io2vYHJOwcrYrVjgJlT6t/0dOXJrA=
|
||||
github.com/ydb-platform/ydb-go-yc v0.12.1/go.mod h1:t/ZA4ECdgPWjAb4jyDe8AzQZB5dhpGbi3iCahFaNwBY=
|
||||
github.com/ydb-platform/ydb-go-yc-metadata v0.6.1 h1:9E5q8Nsy2RiJMZDNVy0A3KUrIMBPakJ2VgloeWbcI84=
|
||||
@@ -2208,8 +2208,8 @@ golang.org/x/crypto v0.14.0/go.mod h1:MVFd36DqK4CsrnJYDkBA3VC4m2GkXAM0PvzMCn4JQf
|
||||
golang.org/x/crypto v0.19.0/go.mod h1:Iy9bg/ha4yyC70EfRS8jz+B6ybOBKMaSxLj6P6oBDfU=
|
||||
golang.org/x/crypto v0.23.0/go.mod h1:CKFgDieR+mRhux2Lsu27y0fO304Db0wZe70UKqHu0v8=
|
||||
golang.org/x/crypto v0.31.0/go.mod h1:kDsLvtWBEx7MV9tJOj9bnXsPbxwJQ6csT/x4KIN4Ssk=
|
||||
golang.org/x/crypto v0.49.0 h1:+Ng2ULVvLHnJ/ZFEq4KdcDd/cfjrrjjNSXNzxg0Y4U4=
|
||||
golang.org/x/crypto v0.49.0/go.mod h1:ErX4dUh2UM+CFYiXZRTcMpEcN8b/1gxEuv3nODoYtCA=
|
||||
golang.org/x/crypto v0.50.0 h1:zO47/JPrL6vsNkINmLoo/PH1gcxpls50DNogFvB5ZGI=
|
||||
golang.org/x/crypto v0.50.0/go.mod h1:3muZ7vA7PBCE6xgPX7nkzzjiUq87kRItoJQM1Yo8S+Q=
|
||||
golang.org/x/exp v0.0.0-20180321215751-8460e604b9de/go.mod h1:CJ0aWSM057203Lf6IL+f9T1iT9GByDxfZKAQTCR3kQA=
|
||||
golang.org/x/exp v0.0.0-20180807140117-3d87b88a115f/go.mod h1:CJ0aWSM057203Lf6IL+f9T1iT9GByDxfZKAQTCR3kQA=
|
||||
golang.org/x/exp v0.0.0-20190121172915-509febef88a4/go.mod h1:CJ0aWSM057203Lf6IL+f9T1iT9GByDxfZKAQTCR3kQA=
|
||||
@@ -2275,8 +2275,8 @@ golang.org/x/mod v0.12.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs=
|
||||
golang.org/x/mod v0.13.0/go.mod h1:hTbmBsO62+eylJbnUtE2MGJUyE7QWk4xUqPFrRgJ+7c=
|
||||
golang.org/x/mod v0.15.0/go.mod h1:hTbmBsO62+eylJbnUtE2MGJUyE7QWk4xUqPFrRgJ+7c=
|
||||
golang.org/x/mod v0.17.0/go.mod h1:hTbmBsO62+eylJbnUtE2MGJUyE7QWk4xUqPFrRgJ+7c=
|
||||
golang.org/x/mod v0.33.0 h1:tHFzIWbBifEmbwtGz65eaWyGiGZatSrT9prnU8DbVL8=
|
||||
golang.org/x/mod v0.33.0/go.mod h1:swjeQEj+6r7fODbD2cqrnje9PnziFuw4bmLbBZFrQ5w=
|
||||
golang.org/x/mod v0.34.0 h1:xIHgNUUnW6sYkcM5Jleh05DvLOtwc6RitGHbDk4akRI=
|
||||
golang.org/x/mod v0.34.0/go.mod h1:ykgH52iCZe79kzLLMhyCUzhMci+nQj+0XkbXpNYtVjY=
|
||||
golang.org/x/net v0.0.0-20180724234803-3673e40ba225/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4=
|
||||
golang.org/x/net v0.0.0-20180826012351-8a410e7b638d/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4=
|
||||
golang.org/x/net v0.0.0-20180906233101-161cd47e91fd/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4=
|
||||
@@ -2344,8 +2344,8 @@ golang.org/x/net v0.16.0/go.mod h1:NxSsAGuq816PNPmqtQdLE42eU2Fs7NoRIZrHJAlaCOE=
|
||||
golang.org/x/net v0.21.0/go.mod h1:bIjVDfnllIU7BJ2DNgfnXvpSvtn8VRwhlsaeUTyUS44=
|
||||
golang.org/x/net v0.25.0/go.mod h1:JkAGAh7GEvH74S6FOH42FLoXpXbE/aqXSrIQjXgsiwM=
|
||||
golang.org/x/net v0.33.0/go.mod h1:HXLR5J+9DxmrqMwG9qjGCxZ+zKXxBru04zlTvWlWuN4=
|
||||
golang.org/x/net v0.52.0 h1:He/TN1l0e4mmR3QqHMT2Xab3Aj3L9qjbhRm78/6jrW0=
|
||||
golang.org/x/net v0.52.0/go.mod h1:R1MAz7uMZxVMualyPXb+VaqGSa3LIaUqk0eEt3w36Sw=
|
||||
golang.org/x/net v0.53.0 h1:d+qAbo5L0orcWAr0a9JweQpjXF19LMXJE8Ey7hwOdUA=
|
||||
golang.org/x/net v0.53.0/go.mod h1:JvMuJH7rrdiCfbeHoo3fCQU24Lf5JJwT9W3sJFulfgs=
|
||||
golang.org/x/oauth2 v0.0.0-20180821212333-d2e6202438be/go.mod h1:N/0e6XlmueqKjAGxoOufVs8QHGRruUQn6yWY3a++T0U=
|
||||
golang.org/x/oauth2 v0.0.0-20190226205417-e64efc72b421/go.mod h1:gOpvHmFTYa4IltrdGE7lF6nIHvwfUNPOp7c8zoXwtLw=
|
||||
golang.org/x/oauth2 v0.0.0-20190604053449-0f29369cfe45/go.mod h1:gOpvHmFTYa4IltrdGE7lF6nIHvwfUNPOp7c8zoXwtLw=
|
||||
@@ -2503,11 +2503,11 @@ golang.org/x/sys v0.13.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||
golang.org/x/sys v0.17.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
|
||||
golang.org/x/sys v0.20.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
|
||||
golang.org/x/sys v0.28.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
|
||||
golang.org/x/sys v0.42.0 h1:omrd2nAlyT5ESRdCLYdm3+fMfNFE/+Rf4bDIQImRJeo=
|
||||
golang.org/x/sys v0.42.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
|
||||
golang.org/x/sys v0.43.0 h1:Rlag2XtaFTxp19wS8MXlJwTvoh8ArU6ezoyFsMyCTNI=
|
||||
golang.org/x/sys v0.43.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
|
||||
golang.org/x/telemetry v0.0.0-20240228155512-f48c80bd79b2/go.mod h1:TeRTkGYfJXctD9OcfyVLyj2J3IxLnKwHJR8f4D8a3YE=
|
||||
golang.org/x/telemetry v0.0.0-20260209163413-e7419c687ee4 h1:bTLqdHv7xrGlFbvf5/TXNxy/iUwwdkjhqQTJDjW7aj0=
|
||||
golang.org/x/telemetry v0.0.0-20260209163413-e7419c687ee4/go.mod h1:g5NllXBEermZrmR51cJDQxmJUHUOfRAaNyWBM+R+548=
|
||||
golang.org/x/telemetry v0.0.0-20260311193753-579e4da9a98c h1:6a8FdnNk6bTXBjR4AGKFgUKuo+7GnR3FX5L7CbveeZc=
|
||||
golang.org/x/telemetry v0.0.0-20260311193753-579e4da9a98c/go.mod h1:TpUTTEp9frx7rTdLpC9gFG9kdI7zVLFTFFlqaH2Cncw=
|
||||
golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo=
|
||||
golang.org/x/term v0.0.0-20210220032956-6a3ed077a48d/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo=
|
||||
golang.org/x/term v0.0.0-20210615171337-6886f2dfbf5b/go.mod h1:jbD1KX2456YbFQfuXm/mYQcufACuNUgVhRMnK/tPxf8=
|
||||
@@ -2523,8 +2523,8 @@ golang.org/x/term v0.13.0/go.mod h1:LTmsnFJwVN6bCy1rVCoS+qHT1HhALEFxKncY3WNNh4U=
|
||||
golang.org/x/term v0.17.0/go.mod h1:lLRBjIVuehSbZlaOtGMbcMncT+aqLLLmKrsjNrUguwk=
|
||||
golang.org/x/term v0.20.0/go.mod h1:8UkIAJTvZgivsXaD6/pH6U9ecQzZ45awqEOzuCvwpFY=
|
||||
golang.org/x/term v0.27.0/go.mod h1:iMsnZpn0cago0GOrHO2+Y7u7JPn5AylBrcoWkElMTSM=
|
||||
golang.org/x/term v0.41.0 h1:QCgPso/Q3RTJx2Th4bDLqML4W6iJiaXFq2/ftQF13YU=
|
||||
golang.org/x/term v0.41.0/go.mod h1:3pfBgksrReYfZ5lvYM0kSO0LIkAl4Yl2bXOkKP7Ec2A=
|
||||
golang.org/x/term v0.42.0 h1:UiKe+zDFmJobeJ5ggPwOshJIVt6/Ft0rcfrXZDLWAWY=
|
||||
golang.org/x/term v0.42.0/go.mod h1:Dq/D+snpsbazcBG5+F9Q1n2rXV8Ma+71xEjTRufARgY=
|
||||
golang.org/x/text v0.0.0-20170915032832-14c0d48ead0c/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
|
||||
golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
|
||||
golang.org/x/text v0.3.1-0.20180807135948-17ff2d5776d2/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
|
||||
@@ -2545,8 +2545,8 @@ golang.org/x/text v0.13.0/go.mod h1:TvPlkZtksWOMsz7fbANvkp4WM8x/WCo/om8BMLbz+aE=
|
||||
golang.org/x/text v0.14.0/go.mod h1:18ZOQIKpY8NJVqYksKHtTdi31H5itFRjB5/qKTNYzSU=
|
||||
golang.org/x/text v0.15.0/go.mod h1:18ZOQIKpY8NJVqYksKHtTdi31H5itFRjB5/qKTNYzSU=
|
||||
golang.org/x/text v0.21.0/go.mod h1:4IBbMaMmOPCJ8SecivzSH54+73PCFmPWxNTLm+vZkEQ=
|
||||
golang.org/x/text v0.35.0 h1:JOVx6vVDFokkpaq1AEptVzLTpDe9KGpj5tR4/X+ybL8=
|
||||
golang.org/x/text v0.35.0/go.mod h1:khi/HExzZJ2pGnjenulevKNX1W67CUy0AsXcNubPGCA=
|
||||
golang.org/x/text v0.36.0 h1:JfKh3XmcRPqZPKevfXVpI1wXPTqbkE5f7JA92a55Yxg=
|
||||
golang.org/x/text v0.36.0/go.mod h1:NIdBknypM8iqVmPiuco0Dh6P5Jcdk8lJL0CUebqK164=
|
||||
golang.org/x/time v0.0.0-20181108054448-85acf8d2951c/go.mod h1:tRJNPiyCQ0inRvYxbN9jk5I+vvW/OXSQhTDSoE431IQ=
|
||||
golang.org/x/time v0.0.0-20190308202827-9d24e82272b4/go.mod h1:tRJNPiyCQ0inRvYxbN9jk5I+vvW/OXSQhTDSoE431IQ=
|
||||
golang.org/x/time v0.0.0-20191024005414-555d28b269f0/go.mod h1:tRJNPiyCQ0inRvYxbN9jk5I+vvW/OXSQhTDSoE431IQ=
|
||||
@@ -2624,8 +2624,8 @@ golang.org/x/tools v0.7.0/go.mod h1:4pg6aUX35JBAogB10C9AtvVL+qowtN4pT3CGSQex14s=
|
||||
golang.org/x/tools v0.13.0/go.mod h1:HvlwmtVNQAhOuCjW7xxvovg8wbNq7LwfXh/k7wXUl58=
|
||||
golang.org/x/tools v0.14.0/go.mod h1:uYBEerGOWcJyEORxN+Ek8+TT266gXkNlHdJBwexUsBg=
|
||||
golang.org/x/tools v0.21.1-0.20240508182429-e35e4ccd0d2d/go.mod h1:aiJjzUbINMkxbQROHiO6hDPo2LHcIPhhQsa9DLh0yGk=
|
||||
golang.org/x/tools v0.42.0 h1:uNgphsn75Tdz5Ji2q36v/nsFSfR/9BRFvqhGBaJGd5k=
|
||||
golang.org/x/tools v0.42.0/go.mod h1:Ma6lCIwGZvHK6XtgbswSoWroEkhugApmsXyrUmBhfr0=
|
||||
golang.org/x/tools v0.43.0 h1:12BdW9CeB3Z+J/I/wj34VMl8X+fEXBxVR90JeMX5E7s=
|
||||
golang.org/x/tools v0.43.0/go.mod h1:uHkMso649BX2cZK6+RpuIPXS3ho2hZo4FVwfoy1vIk0=
|
||||
golang.org/x/tools/godoc v0.1.0-deprecated h1:o+aZ1BOj6Hsx/GBdJO/s815sqftjSnrZZwyYTHODvtk=
|
||||
golang.org/x/tools/godoc v0.1.0-deprecated/go.mod h1:qM63CriJ961IHWmnWa9CjZnBndniPt4a3CK0PVB9bIg=
|
||||
golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
apiVersion: v1
|
||||
description: SeaweedFS
|
||||
name: seaweedfs
|
||||
appVersion: "4.19"
|
||||
appVersion: "4.20"
|
||||
# Dev note: Trigger a helm chart release by `git tag -a helm-<version>`
|
||||
version: 4.19.0
|
||||
version: 4.20.0
|
||||
|
||||
@@ -49,6 +49,35 @@ CREATE TABLE IF NOT EXISTS `filemeta` (
|
||||
|
||||
Alternative database can also be configured (e.g. leveldb, postgres) following the instructions at `filer.extraEnvironmentVars`.
|
||||
|
||||
#### RocksDB variant
|
||||
|
||||
The `_large_disk_rocksdb` image tag ships with RocksDB pre-configured as the filer backend.
|
||||
To use this image with the Helm chart, override the image on all three components and disable
|
||||
the chart's default `WEED_LEVELDB2_ENABLED`, which would otherwise re-enable LevelDB2 and
|
||||
override the image's built-in RocksDB configuration:
|
||||
|
||||
```yaml
|
||||
# Replace <VERSION> with the desired seaweedfs version, e.g. 3.80_large_disk_rocksdb.
|
||||
master:
|
||||
imageOverride: chrislusf/seaweedfs:<VERSION>_large_disk_rocksdb
|
||||
|
||||
volume:
|
||||
imageOverride: chrislusf/seaweedfs:<VERSION>_large_disk_rocksdb
|
||||
|
||||
filer:
|
||||
enablePVC: true
|
||||
imageOverride: chrislusf/seaweedfs:<VERSION>_large_disk_rocksdb
|
||||
extraEnvironmentVars:
|
||||
WEED_LEVELDB2_ENABLED: "false"
|
||||
```
|
||||
|
||||
Notes:
|
||||
|
||||
* `master` and `volume` use the same image tag so that all components share a consistent
|
||||
SeaweedFS build; RocksDB itself is only used by the filer.
|
||||
* `filer.enablePVC: true` (or another form of persistent storage for the filer) is required
|
||||
so that the RocksDB metadata store survives pod restarts — otherwise metadata will be lost.
|
||||
|
||||
### Node Labels
|
||||
Kubernetes nodes can have labels which help to define which node(Host) will run which pod:
|
||||
|
||||
|
||||
@@ -82,7 +82,7 @@ spec:
|
||||
{{- end }}
|
||||
containers:
|
||||
- name: seaweedfs
|
||||
image: {{ template "admin.image" . }}
|
||||
image: {{ template "seaweedfs.admin.image" . }}
|
||||
imagePullPolicy: {{ default "IfNotPresent" .Values.global.seaweedfs.imagePullPolicy }}
|
||||
{{- $adminAuthEnabled := or .Values.admin.secret.existingSecret .Values.admin.secret.adminPassword }}
|
||||
{{- $urlPrefix := .Values.admin.urlPrefix }}
|
||||
@@ -242,7 +242,7 @@ spec:
|
||||
securityContext: {{- omit .Values.admin.containerSecurityContext "enabled" | toYaml | nindent 12 }}
|
||||
{{- end }}
|
||||
{{- if .Values.admin.sidecars }}
|
||||
{{- include "common.tplvalues.render" (dict "value" .Values.admin.sidecars "context" $) | nindent 8 }}
|
||||
{{- include "seaweedfs.tplvalues.render" (dict "value" .Values.admin.sidecars "context" $) | nindent 8 }}
|
||||
{{- end }}
|
||||
volumes:
|
||||
{{- if eq .Values.admin.data.type "hostPath" }}
|
||||
@@ -303,7 +303,7 @@ spec:
|
||||
nodeSelector:
|
||||
{{ tpl .Values.admin.nodeSelector . | indent 8 | trim }}
|
||||
{{- end }}
|
||||
{{- $pvc_exists := include "admin.pvc_exists" . -}}
|
||||
{{- $pvc_exists := include "seaweedfs.admin.pvc_exists" . -}}
|
||||
{{- if $pvc_exists }}
|
||||
volumeClaimTemplates:
|
||||
{{- if eq .Values.admin.data.type "persistentVolumeClaim" }}
|
||||
|
||||
@@ -77,7 +77,7 @@ spec:
|
||||
{{- end }}
|
||||
containers:
|
||||
- name: seaweedfs
|
||||
image: {{ template "master.image" . }}
|
||||
image: {{ template "seaweedfs.master.image" . }}
|
||||
imagePullPolicy: {{ default "IfNotPresent" .Values.global.seaweedfs.imagePullPolicy }}
|
||||
env:
|
||||
{{- /* Determine default cluster alias and the corresponding env var keys to avoid conflicts */}}
|
||||
@@ -418,7 +418,7 @@ spec:
|
||||
{{- omit .Values.allInOne.containerSecurityContext "enabled" | toYaml | nindent 12 }}
|
||||
{{- end }}
|
||||
{{- if .Values.allInOne.sidecars }}
|
||||
{{- include "common.tplvalues.render" (dict "value" .Values.allInOne.sidecars "context" $) | nindent 8 }}
|
||||
{{- include "seaweedfs.tplvalues.render" (dict "value" .Values.allInOne.sidecars "context" $) | nindent 8 }}
|
||||
{{- end }}
|
||||
volumes:
|
||||
- name: data
|
||||
|
||||
@@ -164,7 +164,7 @@ spec:
|
||||
securityContext: {{- omit .Values.cosi.containerSecurityContext "enabled" | toYaml | nindent 12 }}
|
||||
{{- end }}
|
||||
{{- if .Values.cosi.sidecars }}
|
||||
{{- include "common.tplvalues.render" (dict "value" .Values.cosi.sidecars "context" $) | nindent 8 }}
|
||||
{{- include "seaweedfs.tplvalues.render" (dict "value" .Values.cosi.sidecars "context" $) | nindent 8 }}
|
||||
{{- end }}
|
||||
volumes:
|
||||
- name: socket
|
||||
|
||||
@@ -86,7 +86,7 @@ spec:
|
||||
{{- end }}
|
||||
containers:
|
||||
- name: seaweedfs
|
||||
image: {{ template "filer.image" . }}
|
||||
image: {{ template "seaweedfs.filer.image" . }}
|
||||
imagePullPolicy: {{ default "IfNotPresent" .Values.global.seaweedfs.imagePullPolicy }}
|
||||
env:
|
||||
- name: POD_IP
|
||||
@@ -317,7 +317,7 @@ spec:
|
||||
securityContext: {{- omit .Values.filer.containerSecurityContext "enabled" | toYaml | nindent 12 }}
|
||||
{{- end }}
|
||||
{{- if .Values.filer.sidecars }}
|
||||
{{- include "common.tplvalues.render" (dict "value" .Values.filer.sidecars "context" $) | nindent 8 }}
|
||||
{{- include "seaweedfs.tplvalues.render" (dict "value" .Values.filer.sidecars "context" $) | nindent 8 }}
|
||||
{{- end }}
|
||||
volumes:
|
||||
{{- if eq .Values.filer.logs.type "hostPath" }}
|
||||
@@ -413,7 +413,7 @@ spec:
|
||||
storageClassName: {{ .Values.filer.storageClass }}
|
||||
{{- end }}
|
||||
{{- end }}
|
||||
{{- $pvc_exists := include "filer.pvc_exists" . -}}
|
||||
{{- $pvc_exists := include "seaweedfs.filer.pvc_exists" . -}}
|
||||
{{- if $pvc_exists }}
|
||||
volumeClaimTemplates:
|
||||
{{- if eq .Values.filer.data.type "persistentVolumeClaim" }}
|
||||
|
||||
@@ -80,7 +80,7 @@ spec:
|
||||
{{- end }}
|
||||
containers:
|
||||
- name: seaweedfs
|
||||
image: {{ template "master.image" . }}
|
||||
image: {{ template "seaweedfs.master.image" . }}
|
||||
imagePullPolicy: {{ default "IfNotPresent" .Values.global.seaweedfs.imagePullPolicy }}
|
||||
env:
|
||||
- name: POD_IP
|
||||
@@ -251,7 +251,7 @@ spec:
|
||||
securityContext: {{- omit .Values.master.containerSecurityContext "enabled" | toYaml | nindent 12 }}
|
||||
{{- end }}
|
||||
{{- if .Values.master.sidecars }}
|
||||
{{- include "common.tplvalues.render" (dict "value" .Values.master.sidecars "context" $) | nindent 8 }}
|
||||
{{- include "seaweedfs.tplvalues.render" (dict "value" .Values.master.sidecars "context" $) | nindent 8 }}
|
||||
{{- end }}
|
||||
volumes:
|
||||
{{- if eq .Values.master.logs.type "hostPath" }}
|
||||
@@ -312,7 +312,7 @@ spec:
|
||||
nodeSelector:
|
||||
{{ tpl .Values.master.nodeSelector . | indent 8 | trim }}
|
||||
{{- end }}
|
||||
{{- $pvc_exists := include "master.pvc_exists" . -}}
|
||||
{{- $pvc_exists := include "seaweedfs.master.pvc_exists" . -}}
|
||||
{{- if $pvc_exists }}
|
||||
volumeClaimTemplates:
|
||||
{{- if eq .Values.master.data.type "persistentVolumeClaim"}}
|
||||
|
||||
@@ -74,7 +74,7 @@ spec:
|
||||
{{- end }}
|
||||
containers:
|
||||
- name: seaweedfs
|
||||
image: {{ template "s3.image" . }}
|
||||
image: {{ template "seaweedfs.s3.image" . }}
|
||||
imagePullPolicy: {{ default "IfNotPresent" .Values.global.seaweedfs.imagePullPolicy }}
|
||||
env:
|
||||
- name: POD_IP
|
||||
@@ -226,7 +226,7 @@ spec:
|
||||
securityContext: {{- omit .Values.s3.containerSecurityContext "enabled" | toYaml | nindent 12 }}
|
||||
{{- end }}
|
||||
{{- if .Values.s3.sidecars }}
|
||||
{{- include "common.tplvalues.render" (dict "value" .Values.s3.sidecars "context" $) | nindent 8 }}
|
||||
{{- include "seaweedfs.tplvalues.render" (dict "value" .Values.s3.sidecars "context" $) | nindent 8 }}
|
||||
{{- end }}
|
||||
volumes:
|
||||
{{- if .Values.s3.enableAuth }}
|
||||
|
||||
@@ -15,15 +15,15 @@
|
||||
{{- $access_key_admin := $adminCreds.accessKey -}}
|
||||
{{- $secret_key_admin := $adminCreds.secretKey -}}
|
||||
{{- if not (and $access_key_admin $secret_key_admin) -}}
|
||||
{{- $access_key_admin = include "getOrGeneratePassword" (dict "namespace" .Release.Namespace "secretName" $secretName "key" "admin_access_key_id" "length" 20 "existingSecret" (ternary $existingSecret nil $reuse)) -}}
|
||||
{{- $secret_key_admin = include "getOrGeneratePassword" (dict "namespace" .Release.Namespace "secretName" $secretName "key" "admin_secret_access_key" "length" 40 "existingSecret" (ternary $existingSecret nil $reuse)) -}}
|
||||
{{- $access_key_admin = include "seaweedfs.getOrGeneratePassword" (dict "namespace" .Release.Namespace "secretName" $secretName "key" "admin_access_key_id" "length" 20 "existingSecret" (ternary $existingSecret nil $reuse)) -}}
|
||||
{{- $secret_key_admin = include "seaweedfs.getOrGeneratePassword" (dict "namespace" .Release.Namespace "secretName" $secretName "key" "admin_secret_access_key" "length" 40 "existingSecret" (ternary $existingSecret nil $reuse)) -}}
|
||||
{{- end -}}
|
||||
{{- $readCreds := $creds.read | default dict -}}
|
||||
{{- $access_key_read := $readCreds.accessKey -}}
|
||||
{{- $secret_key_read := $readCreds.secretKey -}}
|
||||
{{- if not (and $access_key_read $secret_key_read) -}}
|
||||
{{- $access_key_read = include "getOrGeneratePassword" (dict "namespace" .Release.Namespace "secretName" $secretName "key" "read_access_key_id" "length" 20 "existingSecret" (ternary $existingSecret nil $reuse)) -}}
|
||||
{{- $secret_key_read = include "getOrGeneratePassword" (dict "namespace" .Release.Namespace "secretName" $secretName "key" "read_secret_access_key" "length" 40 "existingSecret" (ternary $existingSecret nil $reuse)) -}}
|
||||
{{- $access_key_read = include "seaweedfs.getOrGeneratePassword" (dict "namespace" .Release.Namespace "secretName" $secretName "key" "read_access_key_id" "length" 20 "existingSecret" (ternary $existingSecret nil $reuse)) -}}
|
||||
{{- $secret_key_read = include "seaweedfs.getOrGeneratePassword" (dict "namespace" .Release.Namespace "secretName" $secretName "key" "read_secret_access_key" "length" 40 "existingSecret" (ternary $existingSecret nil $reuse)) -}}
|
||||
{{- end -}}
|
||||
apiVersion: v1
|
||||
kind: Secret
|
||||
|
||||
@@ -74,7 +74,7 @@ spec:
|
||||
{{- end }}
|
||||
containers:
|
||||
- name: seaweedfs
|
||||
image: {{ template "sftp.image" . }}
|
||||
image: {{ template "seaweedfs.sftp.image" . }}
|
||||
imagePullPolicy: {{ default "IfNotPresent" .Values.global.seaweedfs.imagePullPolicy }}
|
||||
env:
|
||||
- name: POD_IP
|
||||
@@ -233,7 +233,7 @@ spec:
|
||||
securityContext: {{- omit .Values.sftp.containerSecurityContext "enabled" | toYaml | nindent 12 }}
|
||||
{{- end }}
|
||||
{{- if .Values.sftp.sidecars }}
|
||||
{{- include "common.tplvalues.render" (dict "value" .Values.sftp.sidecars "context" $) | nindent 8 }}
|
||||
{{- include "seaweedfs.tplvalues.render" (dict "value" .Values.sftp.sidecars "context" $) | nindent 8 }}
|
||||
{{- end }}
|
||||
volumes:
|
||||
{{- if .Values.sftp.enableAuth }}
|
||||
|
||||
@@ -1,8 +1,8 @@
|
||||
{{- if or .Values.sftp.enabled .Values.allInOne.enabled }}
|
||||
{{- $secretName := printf "%s-sftp-secret" (include "seaweedfs.fullname" .) }}
|
||||
{{- $admin_pwd := include "getOrGeneratePassword" (dict "namespace" .Release.Namespace "secretName" $secretName "key" "admin_password" 20) -}}
|
||||
{{- $read_user_pwd := include "getOrGeneratePassword" (dict "namespace" .Release.Namespace "secretName" $secretName "key" "readonly_password" 20) -}}
|
||||
{{- $public_user_pwd := include "getOrGeneratePassword" (dict "namespace" .Release.Namespace "secretName" $secretName "key" "public_user_password" 20) -}}
|
||||
{{- $admin_pwd := include "seaweedfs.getOrGeneratePassword" (dict "namespace" .Release.Namespace "secretName" $secretName "key" "admin_password" "length" 20) -}}
|
||||
{{- $read_user_pwd := include "seaweedfs.getOrGeneratePassword" (dict "namespace" .Release.Namespace "secretName" $secretName "key" "readonly_password" "length" 20) -}}
|
||||
{{- $public_user_pwd := include "seaweedfs.getOrGeneratePassword" (dict "namespace" .Release.Namespace "secretName" $secretName "key" "public_user_password" "length" 20) -}}
|
||||
apiVersion: v1
|
||||
kind: Secret
|
||||
type: Opaque
|
||||
|
||||
@@ -72,77 +72,77 @@ Inject extra environment vars in the format key:value, if populated
|
||||
{{- end -}}
|
||||
|
||||
{{/* Return the proper filer image */}}
|
||||
{{- define "filer.image" -}}
|
||||
{{- define "seaweedfs.filer.image" -}}
|
||||
{{- if .Values.filer.imageOverride -}}
|
||||
{{- $imageOverride := .Values.filer.imageOverride -}}
|
||||
{{- printf "%s" $imageOverride -}}
|
||||
{{- else -}}
|
||||
{{- include "common.image" . }}
|
||||
{{- include "seaweedfs.image" . }}
|
||||
{{- end -}}
|
||||
{{- end -}}
|
||||
|
||||
{{/* Return the proper master image */}}
|
||||
{{- define "master.image" -}}
|
||||
{{- define "seaweedfs.master.image" -}}
|
||||
{{- if .Values.master.imageOverride -}}
|
||||
{{- $imageOverride := .Values.master.imageOverride -}}
|
||||
{{- printf "%s" $imageOverride -}}
|
||||
{{- else -}}
|
||||
{{- include "common.image" . }}
|
||||
{{- include "seaweedfs.image" . }}
|
||||
{{- end -}}
|
||||
{{- end -}}
|
||||
|
||||
{{/* Return the proper s3 image */}}
|
||||
{{- define "s3.image" -}}
|
||||
{{- define "seaweedfs.s3.image" -}}
|
||||
{{- if .Values.s3.imageOverride -}}
|
||||
{{- $imageOverride := .Values.s3.imageOverride -}}
|
||||
{{- printf "%s" $imageOverride -}}
|
||||
{{- else -}}
|
||||
{{- include "common.image" . }}
|
||||
{{- include "seaweedfs.image" . }}
|
||||
{{- end -}}
|
||||
{{- end -}}
|
||||
|
||||
{{/* Return the proper sftp image */}}
|
||||
{{- define "sftp.image" -}}
|
||||
{{- define "seaweedfs.sftp.image" -}}
|
||||
{{- if .Values.sftp.imageOverride -}}
|
||||
{{- $imageOverride := .Values.sftp.imageOverride -}}
|
||||
{{- printf "%s" $imageOverride -}}
|
||||
{{- else -}}
|
||||
{{- include "common.image" . }}
|
||||
{{- include "seaweedfs.image" . }}
|
||||
{{- end -}}
|
||||
{{- end -}}
|
||||
|
||||
{{/* Return the proper admin image */}}
|
||||
{{- define "admin.image" -}}
|
||||
{{- define "seaweedfs.admin.image" -}}
|
||||
{{- if .Values.admin.imageOverride -}}
|
||||
{{- $imageOverride := .Values.admin.imageOverride -}}
|
||||
{{- printf "%s" $imageOverride -}}
|
||||
{{- else -}}
|
||||
{{- include "common.image" . }}
|
||||
{{- include "seaweedfs.image" . }}
|
||||
{{- end -}}
|
||||
{{- end -}}
|
||||
|
||||
{{/* Return the proper worker image */}}
|
||||
{{- define "worker.image" -}}
|
||||
{{- define "seaweedfs.worker.image" -}}
|
||||
{{- if .Values.worker.imageOverride -}}
|
||||
{{- $imageOverride := .Values.worker.imageOverride -}}
|
||||
{{- printf "%s" $imageOverride -}}
|
||||
{{- else -}}
|
||||
{{- include "common.image" . }}
|
||||
{{- include "seaweedfs.image" . }}
|
||||
{{- end -}}
|
||||
{{- end -}}
|
||||
|
||||
{{/* Return the proper volume image */}}
|
||||
{{- define "volume.image" -}}
|
||||
{{- define "seaweedfs.volume.image" -}}
|
||||
{{- if .Values.volume.imageOverride -}}
|
||||
{{- $imageOverride := .Values.volume.imageOverride -}}
|
||||
{{- printf "%s" $imageOverride -}}
|
||||
{{- else -}}
|
||||
{{- include "common.image" . }}
|
||||
{{- include "seaweedfs.image" . }}
|
||||
{{- end -}}
|
||||
{{- end -}}
|
||||
|
||||
{{/* Computes the container image name for all components (if they are not overridden) */}}
|
||||
{{- define "common.image" -}}
|
||||
{{- define "seaweedfs.image" -}}
|
||||
{{- $registryName := default .Values.image.registry .Values.global.imageRegistry | toString -}}
|
||||
{{- $repositoryName := default .Values.image.repository .Values.global.seaweedfs.image.repository | toString -}}
|
||||
{{- $name := .Values.global.seaweedfs.image.name | toString -}}
|
||||
@@ -160,7 +160,7 @@ Inject extra environment vars in the format key:value, if populated
|
||||
{{- end -}}
|
||||
|
||||
{{/* check if any Volume PVC exists */}}
|
||||
{{- define "volume.pvc_exists" -}}
|
||||
{{- define "seaweedfs.volume.pvc_exists" -}}
|
||||
{{- if or (or (eq .Values.volume.data.type "persistentVolumeClaim") (and (eq .Values.volume.idx.type "persistentVolumeClaim") .Values.volume.dir_idx )) (eq .Values.volume.logs.type "persistentVolumeClaim") -}}
|
||||
{{- printf "true" -}}
|
||||
{{- else -}}
|
||||
@@ -169,7 +169,7 @@ Inject extra environment vars in the format key:value, if populated
|
||||
{{- end -}}
|
||||
|
||||
{{/* check if any Filer PVC exists */}}
|
||||
{{- define "filer.pvc_exists" -}}
|
||||
{{- define "seaweedfs.filer.pvc_exists" -}}
|
||||
{{- if or (eq .Values.filer.data.type "persistentVolumeClaim") (eq .Values.filer.logs.type "persistentVolumeClaim") -}}
|
||||
{{- printf "true" -}}
|
||||
{{- else -}}
|
||||
@@ -178,7 +178,7 @@ Inject extra environment vars in the format key:value, if populated
|
||||
{{- end -}}
|
||||
|
||||
{{/* check if any Master PVC exists */}}
|
||||
{{- define "master.pvc_exists" -}}
|
||||
{{- define "seaweedfs.master.pvc_exists" -}}
|
||||
{{- if or (eq .Values.master.data.type "persistentVolumeClaim") (eq .Values.master.logs.type "persistentVolumeClaim") -}}
|
||||
{{- printf "true" -}}
|
||||
{{- else -}}
|
||||
@@ -187,7 +187,7 @@ Inject extra environment vars in the format key:value, if populated
|
||||
{{- end -}}
|
||||
|
||||
{{/* check if any Admin PVC exists */}}
|
||||
{{- define "admin.pvc_exists" -}}
|
||||
{{- define "seaweedfs.admin.pvc_exists" -}}
|
||||
{{- if or (eq .Values.admin.data.type "persistentVolumeClaim") (eq .Values.admin.logs.type "persistentVolumeClaim") -}}
|
||||
{{- printf "true" -}}
|
||||
{{- else -}}
|
||||
@@ -196,7 +196,7 @@ Inject extra environment vars in the format key:value, if populated
|
||||
{{- end -}}
|
||||
|
||||
{{/* check if any InitContainers exist for Volumes */}}
|
||||
{{- define "volume.initContainers_exists" -}}
|
||||
{{- define "seaweedfs.volume.initContainers_exists" -}}
|
||||
{{- if or (not (empty .Values.volume.idx )) (not (empty .Values.volume.initContainers )) -}}
|
||||
{{- printf "true" -}}
|
||||
{{- else -}}
|
||||
@@ -225,10 +225,10 @@ imagePullSecrets:
|
||||
{{/*
|
||||
Renders a value that contains template perhaps with scope if the scope is present.
|
||||
Usage:
|
||||
{{ include "common.tplvalues.render" ( dict "value" .Values.path.to.the.Value "context" $ ) }}
|
||||
{{ include "common.tplvalues.render" ( dict "value" .Values.path.to.the.Value "context" $ "scope" $app ) }}
|
||||
{{ include "seaweedfs.tplvalues.render" ( dict "value" .Values.path.to.the.Value "context" $ ) }}
|
||||
{{ include "seaweedfs.tplvalues.render" ( dict "value" .Values.path.to.the.Value "context" $ "scope" $app ) }}
|
||||
*/}}
|
||||
{{- define "common.tplvalues.render" -}}
|
||||
{{- define "seaweedfs.tplvalues.render" -}}
|
||||
{{- $value := typeIs "string" .value | ternary .value (.value | toYaml) }}
|
||||
{{- if contains "{{" (toJson .value) }}
|
||||
{{- if .scope }}
|
||||
@@ -245,9 +245,9 @@ Usage:
|
||||
Converts a Kubernetes quantity like "256Mi" or "2G" to a float64 in base units,
|
||||
handling both binary (Ki, Mi, Gi) and decimal (m, k, M) suffixes; numeric inputs
|
||||
Usage:
|
||||
{{ include "common.resource-quantity" "10Gi" }}
|
||||
{{ include "seaweedfs.resource-quantity" "10Gi" }}
|
||||
*/}}
|
||||
{{- define "common.resource-quantity" -}}
|
||||
{{- define "seaweedfs.resource-quantity" -}}
|
||||
{{- $value := . -}}
|
||||
{{- $unit := 1.0 -}}
|
||||
{{- if typeIs "string" . -}}
|
||||
@@ -267,7 +267,7 @@ Usage:
|
||||
getOrGeneratePassword will check if a password exists in a secret and return it,
|
||||
or generate a new random password if it doesn't exist.
|
||||
*/}}
|
||||
{{- define "getOrGeneratePassword" -}}
|
||||
{{- define "seaweedfs.getOrGeneratePassword" -}}
|
||||
{{- $params := . -}}
|
||||
{{- $namespace := $params.namespace -}}
|
||||
{{- $secretName := $params.secretName -}}
|
||||
|
||||
@@ -68,7 +68,7 @@ spec:
|
||||
{{- include "seaweedfs.imagePullSecrets" $ | nindent 6 }}
|
||||
containers:
|
||||
- name: post-install-job
|
||||
image: {{ template "master.image" . }}
|
||||
image: {{ template "seaweedfs.master.image" . }}
|
||||
imagePullPolicy: {{ $.Values.global.seaweedfs.imagePullPolicy | default "IfNotPresent" }}
|
||||
env:
|
||||
- name: WEED_CLUSTER_DEFAULT
|
||||
|
||||
@@ -35,8 +35,8 @@
|
||||
{{- $pvcName := printf "%s-%s-%s-%d" $dir.name $seaweedfsName $volumeName $e }}
|
||||
{{- $currentPVC := (lookup "v1" "PersistentVolumeClaim" $.Release.Namespace $pvcName) }}
|
||||
{{- if and $currentPVC }}
|
||||
{{- $oldSize := include "common.resource-quantity" $currentPVC.spec.resources.requests.storage }}
|
||||
{{- $newSize := include "common.resource-quantity" $desiredSize }}
|
||||
{{- $oldSize := include "seaweedfs.resource-quantity" $currentPVC.spec.resources.requests.storage }}
|
||||
{{- $newSize := include "seaweedfs.resource-quantity" $desiredSize }}
|
||||
{{- if gt $newSize $oldSize }}
|
||||
{{- $commands = append $commands (printf "kubectl patch pvc %s-%s-%s-%d -p '{\"spec\":{\"resources\":{\"requests\":{\"storage\":\"%s\"}}}}'" $dir.name $seaweedfsName $volumeName $e $desiredSize) }}
|
||||
{{- end }}
|
||||
|
||||
@@ -71,12 +71,12 @@ spec:
|
||||
{{- end }}
|
||||
enableServiceLinks: false
|
||||
serviceAccountName: {{ $volume.serviceAccountName | default (include "seaweedfs.serviceAccountName" $) | quote }} # for deleting statefulset pods after migration
|
||||
{{- $initContainers_exists := include "volume.initContainers_exists" $ -}}
|
||||
{{- $initContainers_exists := include "seaweedfs.volume.initContainers_exists" $ -}}
|
||||
{{- if $initContainers_exists }}
|
||||
initContainers:
|
||||
{{- if $volume.idx }}
|
||||
- name: seaweedfs-vol-move-idx
|
||||
image: {{ template "volume.image" $ }}
|
||||
image: {{ template "seaweedfs.volume.image" $ }}
|
||||
imagePullPolicy: {{ $.Values.global.seaweedfs.imagePullPolicy | default "IfNotPresent" }}
|
||||
command: [ '/bin/sh', '-c' ]
|
||||
args: [ '{{range $dir := $volume.dataDirs }}if ls /{{$dir.name}}/*.idx >/dev/null 2>&1; then mv /{{$dir.name}}/*.idx /idx/ ; fi; {{end}}' ]
|
||||
@@ -104,7 +104,7 @@ spec:
|
||||
{{- end }}
|
||||
containers:
|
||||
- name: seaweedfs
|
||||
image: {{ template "volume.image" $ }}
|
||||
image: {{ template "seaweedfs.volume.image" $ }}
|
||||
imagePullPolicy: {{ default "IfNotPresent" $.Values.global.seaweedfs.imagePullPolicy }}
|
||||
env:
|
||||
- name: POD_NAME
|
||||
@@ -274,7 +274,7 @@ spec:
|
||||
securityContext: {{- omit $volume.containerSecurityContext "enabled" | toYaml | nindent 12 }}
|
||||
{{- end }}
|
||||
{{- if $volume.sidecars }}
|
||||
{{- include "common.tplvalues.render" (dict "value" (printf "{{ $volumeName := \"%s\" }}%s" $volumeName $volume.sidecars) "context" $) | nindent 8 }}
|
||||
{{- include "seaweedfs.tplvalues.render" (dict "value" (printf "{{ $volumeName := \"%s\" }}%s" $volumeName $volume.sidecars) "context" $) | nindent 8 }}
|
||||
{{- end }}
|
||||
volumes:
|
||||
|
||||
|
||||
@@ -77,7 +77,7 @@ spec:
|
||||
{{- end }}
|
||||
containers:
|
||||
- name: seaweedfs
|
||||
image: {{ template "worker.image" . }}
|
||||
image: {{ template "seaweedfs.worker.image" . }}
|
||||
imagePullPolicy: {{ default "IfNotPresent" .Values.global.seaweedfs.imagePullPolicy }}
|
||||
env:
|
||||
- name: POD_IP
|
||||
@@ -219,7 +219,7 @@ spec:
|
||||
securityContext: {{- omit .Values.worker.containerSecurityContext "enabled" | toYaml | nindent 12 }}
|
||||
{{- end }}
|
||||
{{- if .Values.worker.sidecars }}
|
||||
{{- include "common.tplvalues.render" (dict "value" .Values.worker.sidecars "context" $) | nindent 8 }}
|
||||
{{- include "seaweedfs.tplvalues.render" (dict "value" .Values.worker.sidecars "context" $) | nindent 8 }}
|
||||
{{- end }}
|
||||
volumes:
|
||||
{{- if eq .Values.worker.data.type "hostPath" }}
|
||||
|
||||
@@ -294,6 +294,7 @@ message AssignVolumeRequest {
|
||||
string rack = 7;
|
||||
string data_node = 9;
|
||||
string disk_type = 8;
|
||||
uint64 expected_data_size = 10; // hint for size-aware volume selection
|
||||
}
|
||||
|
||||
message AssignVolumeResponse {
|
||||
|
||||
@@ -231,6 +231,7 @@ message AssignRequest {
|
||||
uint32 memory_map_max_size_mb = 8;
|
||||
uint32 writable_volume_count = 9;
|
||||
string disk_type = 10;
|
||||
uint64 expected_data_size = 11; // hint for size-aware volume selection
|
||||
}
|
||||
|
||||
message VolumeGrowRequest {
|
||||
|
||||
@@ -39,11 +39,11 @@ type dlmTestCluster struct {
|
||||
filerGrpcPorts [2]int
|
||||
mountPoints [2]string
|
||||
|
||||
masterCmd *exec.Cmd
|
||||
volumeCmd *exec.Cmd
|
||||
filerCmds [2]*exec.Cmd
|
||||
mountCmds [2]*exec.Cmd
|
||||
logFiles []*os.File
|
||||
masterCmd *exec.Cmd
|
||||
volumeCmd *exec.Cmd
|
||||
filerCmds [2]*exec.Cmd
|
||||
mountCmds [2]*exec.Cmd
|
||||
logFiles []*os.File
|
||||
|
||||
cleanupOnce sync.Once
|
||||
}
|
||||
|
||||
@@ -11,7 +11,7 @@
|
||||
<properties>
|
||||
<maven.compiler.source>11</maven.compiler.source>
|
||||
<maven.compiler.target>11</maven.compiler.target>
|
||||
<kafka.version>3.9.1</kafka.version>
|
||||
<kafka.version>3.9.2</kafka.version>
|
||||
<confluent.version>7.6.0</confluent.version>
|
||||
</properties>
|
||||
|
||||
|
||||
@@ -52,8 +52,8 @@ type MasterCluster struct {
|
||||
|
||||
// clusterStatus is the JSON returned by /cluster/status.
|
||||
type clusterStatus struct {
|
||||
IsLeader bool `json:"IsLeader"`
|
||||
Leader string `json:"Leader"`
|
||||
IsLeader bool `json:"IsLeader"`
|
||||
Leader string `json:"Leader"`
|
||||
Peers []string `json:"Peers"`
|
||||
}
|
||||
|
||||
@@ -358,7 +358,6 @@ func (mc *MasterCluster) tailLog(i int) string {
|
||||
return strings.Join(lines, "\n")
|
||||
}
|
||||
|
||||
|
||||
func findOrBuildWeedBinary() (string, error) {
|
||||
if fromEnv := os.Getenv("WEED_BINARY"); fromEnv != "" {
|
||||
if isExecutableFile(fromEnv) {
|
||||
|
||||
@@ -13,8 +13,8 @@ RUN apt-get update && \
|
||||
&& apt-get clean \
|
||||
&& rm -rf /var/lib/apt/lists/*
|
||||
|
||||
ARG PJDFSTEST_REPO=https://github.com/pjd/pjdfstest.git
|
||||
ARG PJDFSTEST_REF=03eb25706d8dbf3611c3f820b45b7a5e09a36c06
|
||||
ARG PJDFSTEST_REPO=https://github.com/sanwan/pjdfstest.git
|
||||
ARG PJDFSTEST_REF=d25636a227606f8960e5179741d8f4ad7030ef41
|
||||
|
||||
RUN git clone "${PJDFSTEST_REPO}" /opt/pjdfstest && \
|
||||
cd /opt/pjdfstest && \
|
||||
|
||||
@@ -6,55 +6,17 @@
|
||||
# A failure in any test NOT listed here will cause the CI job to fail,
|
||||
# catching regressions immediately.
|
||||
|
||||
# ── Linux FUSE NAME_MAX=255 limitation ──────────────────────────────────
|
||||
# The Linux FUSE kernel module enforces NAME_MAX=255 at the VFS layer.
|
||||
# These tests create filenames >255 bytes which cannot be looked up via
|
||||
# normal syscalls (stat, chmod, etc.) after creation.
|
||||
tests/chmod/02.t
|
||||
tests/chmod/03.t
|
||||
tests/chown/02.t
|
||||
tests/chown/03.t
|
||||
tests/ftruncate/02.t
|
||||
tests/ftruncate/03.t
|
||||
tests/link/02.t
|
||||
tests/link/03.t
|
||||
tests/mkdir/02.t
|
||||
tests/mkdir/03.t
|
||||
tests/mkfifo/02.t
|
||||
tests/mkfifo/03.t
|
||||
tests/mknod/02.t
|
||||
tests/mknod/03.t
|
||||
tests/open/02.t
|
||||
tests/open/03.t
|
||||
tests/rename/01.t
|
||||
tests/rename/02.t
|
||||
tests/rmdir/02.t
|
||||
tests/rmdir/03.t
|
||||
tests/symlink/02.t
|
||||
tests/symlink/03.t
|
||||
tests/truncate/02.t
|
||||
tests/truncate/03.t
|
||||
tests/unlink/02.t
|
||||
tests/unlink/03.t
|
||||
|
||||
# ── Hard link nlink/ctime tracking (requires filer changes) ────────────
|
||||
# nlink counts are not correctly maintained across rename operations.
|
||||
tests/rename/23.t
|
||||
tests/rename/24.t
|
||||
|
||||
# ── Parent directory mtime/ctime on deferred file create ───────────────
|
||||
# When file creation is deferred (not flushed to filer immediately),
|
||||
# the parent directory mtime/ctime cannot be updated without invalidating
|
||||
# the just-cached child entry in the metadata cache.
|
||||
tests/open/00.t
|
||||
|
||||
# ── Directory rename permission edge case ──────────────────────────────
|
||||
# Cross-directory rename of a subdirectory with restricted permissions
|
||||
# causes cascading test failures within the test file.
|
||||
tests/rename/21.t
|
||||
|
||||
# ── rmdir after hard link unlink ───────────────────────────────────────
|
||||
# The filer may still report a directory as non-empty after all hard-linked
|
||||
# entries have been unlinked.
|
||||
tests/unlink/14.t
|
||||
# ── Hard link nlink count inconsistencies ────────────────────────────
|
||||
# link/00.t and unlink/00.t fail nlink assertions (e.g. expected nlink=2,
|
||||
# got nlink=3) after hard link creation/removal. This is a filer-side hard
|
||||
# link counter issue, not a mount mtime/ctime problem. The failures are
|
||||
# deterministic and surfaced by caching changes that affect the order in
|
||||
# which entries are loaded into the local meta cache.
|
||||
tests/link/00.t
|
||||
tests/unlink/00.t
|
||||
|
||||
|
||||
@@ -26,8 +26,8 @@ FILER_ADDR="127.0.0.1:${FILER_PORT}"
|
||||
|
||||
# Pin to an immutable upstream commit so CI is reproducible. Override via env
|
||||
# if you want to test against a different ref or fork.
|
||||
PJDFSTEST_REPO="${PJDFSTEST_REPO:-https://github.com/pjd/pjdfstest.git}"
|
||||
PJDFSTEST_REF="${PJDFSTEST_REF:-03eb25706d8dbf3611c3f820b45b7a5e09a36c06}"
|
||||
PJDFSTEST_REPO="${PJDFSTEST_REPO:-https://github.com/sanwan/pjdfstest.git}"
|
||||
PJDFSTEST_REF="${PJDFSTEST_REF:-d25636a227606f8960e5179741d8f4ad7030ef41}"
|
||||
PJDFSTEST_TESTS="${PJDFSTEST_TESTS:-tests/}"
|
||||
|
||||
mini_pid=""
|
||||
|
||||
@@ -592,7 +592,6 @@ func (c *distributedLockCluster) tailLog(name string) string {
|
||||
return strings.Join(lines, "\n")
|
||||
}
|
||||
|
||||
|
||||
func stopProcess(cmd *exec.Cmd) {
|
||||
if cmd == nil || cmd.Process == nil {
|
||||
return
|
||||
|
||||
@@ -39,19 +39,19 @@ const (
|
||||
|
||||
// TestCluster manages the weed mini instance for integration testing
|
||||
type TestCluster struct {
|
||||
dataDir string
|
||||
ctx context.Context
|
||||
cancel context.CancelFunc
|
||||
s3Client *s3.S3
|
||||
isRunning bool
|
||||
startOnce sync.Once
|
||||
wg sync.WaitGroup
|
||||
masterPort int
|
||||
volumePort int
|
||||
filerPort int
|
||||
s3Port int
|
||||
s3Endpoint string
|
||||
rustVolumeCmd *exec.Cmd
|
||||
dataDir string
|
||||
ctx context.Context
|
||||
cancel context.CancelFunc
|
||||
s3Client *s3.S3
|
||||
isRunning bool
|
||||
startOnce sync.Once
|
||||
wg sync.WaitGroup
|
||||
masterPort int
|
||||
volumePort int
|
||||
filerPort int
|
||||
s3Port int
|
||||
s3Endpoint string
|
||||
rustVolumeCmd *exec.Cmd
|
||||
}
|
||||
|
||||
// TestS3Integration demonstrates basic S3 operations against a running weed mini instance
|
||||
|
||||
@@ -317,4 +317,3 @@ func hasKey(policy map[string]interface{}, key string) bool {
|
||||
|
||||
return false
|
||||
}
|
||||
|
||||
|
||||
@@ -703,7 +703,6 @@ func uniqueName(prefix string) string {
|
||||
|
||||
// --- Test setup helpers ---
|
||||
|
||||
|
||||
func startMiniCluster(t *testing.T) (*TestCluster, error) {
|
||||
ports := testutil.MustAllocatePorts(t, 8)
|
||||
masterPort, masterGrpcPort := ports[0], ports[1]
|
||||
@@ -733,10 +732,12 @@ func startMiniCluster(t *testing.T) (*TestCluster, error) {
|
||||
err := os.WriteFile(securityToml, []byte("# Empty security config\n"), 0644)
|
||||
require.NoError(t, err)
|
||||
|
||||
// Configure credential store for IAM tests
|
||||
// Configure credential store for IAM tests.
|
||||
// Use filer_etc instead of memory because the memory store does not
|
||||
// persist groups or service accounts through LoadConfiguration/SaveConfiguration.
|
||||
credentialToml := filepath.Join(testDir, "credential.toml")
|
||||
credentialConfig := `
|
||||
[credential.memory]
|
||||
[credential.filer_etc]
|
||||
enabled = true
|
||||
`
|
||||
err = os.WriteFile(credentialToml, []byte(credentialConfig), 0644)
|
||||
@@ -806,7 +807,6 @@ enabled = true
|
||||
return cluster, nil
|
||||
}
|
||||
|
||||
|
||||
// startRustVolumeServer starts a Rust volume server that registers with the same master.
|
||||
func (c *TestCluster) startRustVolumeServer(t *testing.T) error {
|
||||
t.Helper()
|
||||
|
||||
@@ -0,0 +1,75 @@
|
||||
package policy
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"testing"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
// TestShellAccessKeyLifecycle exercises s3.accesskey.* commands end-to-end.
|
||||
func TestShellAccessKeyLifecycle(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("Skipping integration test in short mode")
|
||||
}
|
||||
cluster, err := startMiniCluster(t)
|
||||
require.NoError(t, err)
|
||||
defer cluster.Stop()
|
||||
|
||||
const weedCmd = "weed"
|
||||
master := string(pb.NewServerAddress("127.0.0.1", cluster.masterPort, cluster.masterGrpcPort))
|
||||
filer := string(pb.NewServerAddress("127.0.0.1", cluster.filerPort, cluster.filerGrpcPort))
|
||||
|
||||
userName := uniqueName("akuser")
|
||||
// Create user with explicit key so we know the initial value.
|
||||
initialAK := "INITIALAK1234567890X"
|
||||
initialSK := "initialsecret1234567890abcdefghijklmnop"
|
||||
execShell(t, weedCmd, master, filer,
|
||||
fmt.Sprintf("s3.user.create -name %s -access_key %s -secret_key %s", userName, initialAK, initialSK))
|
||||
defer execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.delete -name %s", userName))
|
||||
|
||||
t.Run("ListInitialKey", func(t *testing.T) {
|
||||
out := execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.accesskey.list -user %s", userName))
|
||||
requireContains(t, out, initialAK, "accesskey.list initial")
|
||||
})
|
||||
|
||||
var createdAK string
|
||||
t.Run("CreateAdditionalKey", func(t *testing.T) {
|
||||
out := execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.accesskey.create -user %s", userName))
|
||||
requireContains(t, out, "Access Key:", "accesskey.create output")
|
||||
requireContains(t, out, "Secret Key:", "accesskey.create output")
|
||||
createdAK = extractFieldAfter(out, "Access Key:")
|
||||
if createdAK == "" {
|
||||
t.Fatalf("failed to extract access key from create output:\n%s", out)
|
||||
}
|
||||
|
||||
out = execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.accesskey.list -user %s", userName))
|
||||
requireContains(t, out, initialAK, "list contains original")
|
||||
requireContains(t, out, createdAK, "list contains new key")
|
||||
})
|
||||
|
||||
t.Run("RotateKey", func(t *testing.T) {
|
||||
if createdAK == "" {
|
||||
t.Fatal("createdAK is empty; CreateAdditionalKey must run successfully first")
|
||||
}
|
||||
out := execShell(t, weedCmd, master, filer,
|
||||
fmt.Sprintf("s3.accesskey.rotate -user %s -access_key %s", userName, initialAK))
|
||||
requireContains(t, out, initialAK, "rotate shows old key")
|
||||
requireContains(t, out, "deleted", "rotate marks old key deleted")
|
||||
|
||||
out = execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.accesskey.list -user %s", userName))
|
||||
requireNotContains(t, out, initialAK, "old key removed")
|
||||
requireContains(t, out, createdAK, "other key still present")
|
||||
})
|
||||
|
||||
t.Run("DeleteKey", func(t *testing.T) {
|
||||
if createdAK == "" {
|
||||
t.Fatal("createdAK is empty; CreateAdditionalKey must run successfully first")
|
||||
}
|
||||
execShell(t, weedCmd, master, filer,
|
||||
fmt.Sprintf("s3.accesskey.delete -user %s -access_key %s", userName, createdAK))
|
||||
out := execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.accesskey.list -user %s", userName))
|
||||
requireNotContains(t, out, createdAK, "deleted key removed from list")
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,54 @@
|
||||
package policy
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"testing"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
// TestShellAnonymousAccess exercises s3.anonymous.* commands end-to-end.
|
||||
func TestShellAnonymousAccess(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("Skipping integration test in short mode")
|
||||
}
|
||||
cluster, err := startMiniCluster(t)
|
||||
require.NoError(t, err)
|
||||
defer cluster.Stop()
|
||||
|
||||
const weedCmd = "weed"
|
||||
master := string(pb.NewServerAddress("127.0.0.1", cluster.masterPort, cluster.masterGrpcPort))
|
||||
filer := string(pb.NewServerAddress("127.0.0.1", cluster.filerPort, cluster.filerGrpcPort))
|
||||
|
||||
bucketName := uniqueName("anon-bkt")
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.create -name %s", bucketName))
|
||||
defer execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.delete -name %s", bucketName))
|
||||
|
||||
t.Run("SetAndGet", func(t *testing.T) {
|
||||
out := execShell(t, weedCmd, master, filer,
|
||||
fmt.Sprintf("s3.anonymous.set -bucket %s -access Read,List", bucketName))
|
||||
requireContains(t, out, bucketName, "anonymous.set output")
|
||||
|
||||
out = execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.anonymous.get -bucket %s", bucketName))
|
||||
requireContains(t, out, bucketName, "anonymous.get bucket")
|
||||
requireContains(t, out, "Read", "anonymous.get read action")
|
||||
requireContains(t, out, "List", "anonymous.get list action")
|
||||
})
|
||||
|
||||
t.Run("List", func(t *testing.T) {
|
||||
out := execShell(t, weedCmd, master, filer, "s3.anonymous.list")
|
||||
requireContains(t, out, bucketName, "anonymous.list contains bucket")
|
||||
})
|
||||
|
||||
t.Run("SetNone", func(t *testing.T) {
|
||||
execShell(t, weedCmd, master, filer,
|
||||
fmt.Sprintf("s3.anonymous.set -bucket %s -access none", bucketName))
|
||||
|
||||
out := execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.anonymous.get -bucket %s", bucketName))
|
||||
requireContains(t, out, "none", "anonymous.get after set none")
|
||||
|
||||
out = execShell(t, weedCmd, master, filer, "s3.anonymous.list")
|
||||
requireNotContains(t, out, bucketName, "anonymous.list after clearing")
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,128 @@
|
||||
package policy
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"testing"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
// TestShellBucketLifecycle exercises s3.bucket.* commands end-to-end:
|
||||
// create/list/delete, owner, quota, versioning, lock, quota.enforce.
|
||||
func TestShellBucketLifecycle(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("Skipping integration test in short mode")
|
||||
}
|
||||
cluster, err := startMiniCluster(t)
|
||||
require.NoError(t, err)
|
||||
defer cluster.Stop()
|
||||
|
||||
const weedCmd = "weed"
|
||||
master := string(pb.NewServerAddress("127.0.0.1", cluster.masterPort, cluster.masterGrpcPort))
|
||||
filer := string(pb.NewServerAddress("127.0.0.1", cluster.filerPort, cluster.filerGrpcPort))
|
||||
|
||||
t.Run("CreateListDelete", func(t *testing.T) {
|
||||
bucketName := uniqueName("bkt")
|
||||
out := execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.create -name %s", bucketName))
|
||||
requireContains(t, out, bucketName, "bucket.create output")
|
||||
|
||||
out = execShell(t, weedCmd, master, filer, "s3.bucket.list")
|
||||
requireContains(t, out, bucketName, "bucket.list contains created")
|
||||
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.delete -name %s", bucketName))
|
||||
|
||||
out = execShell(t, weedCmd, master, filer, "s3.bucket.list")
|
||||
requireNotContains(t, out, bucketName, "bucket.list after delete")
|
||||
})
|
||||
|
||||
t.Run("Owner", func(t *testing.T) {
|
||||
bucketName := uniqueName("bkt-own")
|
||||
ownerName := uniqueName("owner")
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.create -name %s", ownerName))
|
||||
defer execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.delete -name %s", ownerName))
|
||||
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.create -name %s", bucketName))
|
||||
defer execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.delete -name %s", bucketName))
|
||||
|
||||
// Initially no owner.
|
||||
out := execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.owner -name %s", bucketName))
|
||||
requireContains(t, out, "none", "initial owner none")
|
||||
|
||||
// Set owner.
|
||||
execShell(t, weedCmd, master, filer,
|
||||
fmt.Sprintf("s3.bucket.owner -name %s -owner %s", bucketName, ownerName))
|
||||
out = execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.owner -name %s", bucketName))
|
||||
requireContains(t, out, ownerName, "owner set")
|
||||
|
||||
// Remove owner.
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.owner -name %s -delete", bucketName))
|
||||
out = execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.owner -name %s", bucketName))
|
||||
requireContains(t, out, "none", "owner removed")
|
||||
})
|
||||
|
||||
t.Run("Quota", func(t *testing.T) {
|
||||
bucketName := uniqueName("bkt-quota")
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.create -name %s", bucketName))
|
||||
defer execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.delete -name %s", bucketName))
|
||||
|
||||
execShell(t, weedCmd, master, filer,
|
||||
fmt.Sprintf("s3.bucket.quota -name %s -op=set -sizeMB=1024", bucketName))
|
||||
|
||||
out := execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.quota -name %s -op=get", bucketName))
|
||||
requireContains(t, out, "1024", "quota.get shows size")
|
||||
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.quota -name %s -op=disable", bucketName))
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.quota -name %s -op=enable", bucketName))
|
||||
|
||||
// Enforce should run on an empty bucket without error.
|
||||
execShell(t, weedCmd, master, filer, "s3.bucket.quota.enforce -apply")
|
||||
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.quota -name %s -op=remove", bucketName))
|
||||
})
|
||||
|
||||
t.Run("Versioning", func(t *testing.T) {
|
||||
bucketName := uniqueName("bkt-ver")
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.create -name %s", bucketName))
|
||||
defer execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.delete -name %s", bucketName))
|
||||
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.versioning -name %s -enable", bucketName))
|
||||
out := execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.versioning -name %s", bucketName))
|
||||
requireContains(t, out, "Enabled", "versioning enabled")
|
||||
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.versioning -name %s -suspend", bucketName))
|
||||
out = execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.versioning -name %s", bucketName))
|
||||
requireContains(t, out, "Suspended", "versioning suspended")
|
||||
})
|
||||
|
||||
t.Run("Lock", func(t *testing.T) {
|
||||
bucketName := uniqueName("bkt-lock")
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.create -name %s", bucketName))
|
||||
defer execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.delete -name %s", bucketName))
|
||||
|
||||
out := execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.lock -name %s", bucketName))
|
||||
requireContains(t, out, "Disabled", "lock initially disabled")
|
||||
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.lock -name %s -enable", bucketName))
|
||||
|
||||
out = execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.lock -name %s", bucketName))
|
||||
requireContains(t, out, "Enabled", "lock enabled")
|
||||
|
||||
// Versioning should have been auto-enabled.
|
||||
out = execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.versioning -name %s", bucketName))
|
||||
requireContains(t, out, "Enabled", "versioning auto-enabled by lock")
|
||||
})
|
||||
|
||||
t.Run("CreateWithLock", func(t *testing.T) {
|
||||
bucketName := uniqueName("bkt-wlock")
|
||||
out := execShell(t, weedCmd, master, filer,
|
||||
fmt.Sprintf("s3.bucket.create -name %s -withLock", bucketName))
|
||||
// Cleanup may fail if the bucket contains locked objects; we created none.
|
||||
defer execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.delete -name %s", bucketName))
|
||||
|
||||
requireContains(t, out, "Object Lock", "create -withLock output")
|
||||
|
||||
out = execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.lock -name %s", bucketName))
|
||||
requireContains(t, out, "Enabled", "lock enabled after create -withLock")
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,86 @@
|
||||
package policy
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
// TestShellConfigShow verifies s3.config.show outputs a summary of IAM config.
|
||||
func TestShellConfigShow(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("Skipping integration test in short mode")
|
||||
}
|
||||
cluster, err := startMiniCluster(t)
|
||||
require.NoError(t, err)
|
||||
defer cluster.Stop()
|
||||
|
||||
const weedCmd = "weed"
|
||||
master := string(pb.NewServerAddress("127.0.0.1", cluster.masterPort, cluster.masterGrpcPort))
|
||||
filer := string(pb.NewServerAddress("127.0.0.1", cluster.filerPort, cluster.filerGrpcPort))
|
||||
|
||||
userName := uniqueName("cfg-user")
|
||||
groupName := uniqueName("cfg-grp")
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.create -name %s", userName))
|
||||
defer execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.delete -name %s", userName))
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.group.create -name %s", groupName))
|
||||
defer execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.group.delete -name %s", groupName))
|
||||
|
||||
out := execShell(t, weedCmd, master, filer, "s3.config.show")
|
||||
requireContains(t, out, "S3 IAM Configuration Summary", "config.show header")
|
||||
requireContains(t, out, userName, "config.show contains user")
|
||||
requireContains(t, out, groupName, "config.show contains group")
|
||||
}
|
||||
|
||||
// TestShellIAMExportImport does a roundtrip: create resources, export, delete, import, verify.
|
||||
func TestShellIAMExportImport(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("Skipping integration test in short mode")
|
||||
}
|
||||
cluster, err := startMiniCluster(t)
|
||||
require.NoError(t, err)
|
||||
defer cluster.Stop()
|
||||
|
||||
const weedCmd = "weed"
|
||||
master := string(pb.NewServerAddress("127.0.0.1", cluster.masterPort, cluster.masterGrpcPort))
|
||||
filer := string(pb.NewServerAddress("127.0.0.1", cluster.filerPort, cluster.filerGrpcPort))
|
||||
|
||||
userName := uniqueName("exp-user")
|
||||
groupName := uniqueName("exp-grp")
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.create -name %s", userName))
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.group.create -name %s", groupName))
|
||||
|
||||
exportFile := filepath.Join(t.TempDir(), "iam_export.txt")
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.iam.export -file %s", exportFile))
|
||||
|
||||
data, err := os.ReadFile(exportFile)
|
||||
require.NoError(t, err)
|
||||
content := string(data)
|
||||
requireContains(t, content, userName, "export file contains user")
|
||||
requireContains(t, content, groupName, "export file contains group")
|
||||
|
||||
// Delete the resources.
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.delete -name %s", userName))
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.group.delete -name %s", groupName))
|
||||
|
||||
out := execShell(t, weedCmd, master, filer, "s3.user.list")
|
||||
requireNotContains(t, out, userName, "user gone before import")
|
||||
|
||||
// Import.
|
||||
out = execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.iam.import -file %s -apply", exportFile))
|
||||
requireContains(t, out, "Imported IAM configuration", "import output")
|
||||
|
||||
// Verify resources restored.
|
||||
out = execShell(t, weedCmd, master, filer, "s3.user.list")
|
||||
requireContains(t, out, userName, "user restored after import")
|
||||
out = execShell(t, weedCmd, master, filer, "s3.group.list")
|
||||
requireContains(t, out, groupName, "group restored after import")
|
||||
|
||||
// Cleanup.
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.delete -name %s", userName))
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.group.delete -name %s", groupName))
|
||||
}
|
||||
@@ -0,0 +1,62 @@
|
||||
package policy
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"testing"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
// TestShellGroupLifecycle exercises s3.group.* commands end-to-end.
|
||||
func TestShellGroupLifecycle(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("Skipping integration test in short mode")
|
||||
}
|
||||
cluster, err := startMiniCluster(t)
|
||||
require.NoError(t, err)
|
||||
defer cluster.Stop()
|
||||
|
||||
const weedCmd = "weed"
|
||||
master := string(pb.NewServerAddress("127.0.0.1", cluster.masterPort, cluster.masterGrpcPort))
|
||||
filer := string(pb.NewServerAddress("127.0.0.1", cluster.filerPort, cluster.filerGrpcPort))
|
||||
|
||||
groupName := uniqueName("grp")
|
||||
userName := uniqueName("grpuser")
|
||||
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.create -name %s", userName))
|
||||
defer execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.delete -name %s", userName))
|
||||
|
||||
t.Run("CreateShowList", func(t *testing.T) {
|
||||
out := execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.group.create -name %s", groupName))
|
||||
requireContains(t, out, groupName, "group.create output")
|
||||
|
||||
out = execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.group.show -name %s", groupName))
|
||||
requireContains(t, out, groupName, "group.show output")
|
||||
|
||||
out = execShell(t, weedCmd, master, filer, "s3.group.list")
|
||||
requireContains(t, out, groupName, "group.list")
|
||||
})
|
||||
|
||||
t.Run("AddRemoveUser", func(t *testing.T) {
|
||||
out := execShell(t, weedCmd, master, filer,
|
||||
fmt.Sprintf("s3.group.add-user -group %s -user %s", groupName, userName))
|
||||
requireContains(t, out, userName, "group.add-user output")
|
||||
|
||||
out = execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.group.show -name %s", groupName))
|
||||
requireContains(t, out, userName, "group.show after add")
|
||||
|
||||
out = execShell(t, weedCmd, master, filer,
|
||||
fmt.Sprintf("s3.group.remove-user -group %s -user %s", groupName, userName))
|
||||
requireContains(t, out, userName, "group.remove-user output")
|
||||
|
||||
out = execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.group.show -name %s", groupName))
|
||||
requireNotContains(t, out, fmt.Sprintf("\"%s\"", userName), "group.show after remove")
|
||||
})
|
||||
|
||||
t.Run("Delete", func(t *testing.T) {
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.group.delete -name %s", groupName))
|
||||
out := execShell(t, weedCmd, master, filer, "s3.group.list")
|
||||
requireNotContains(t, out, fmt.Sprintf("\"%s\"", groupName), "group.list after delete")
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,73 @@
|
||||
package policy
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
// requireContains fails the test if substr is not found in output.
|
||||
func requireContains(t *testing.T, output, substr, context string) {
|
||||
t.Helper()
|
||||
if !strings.Contains(output, substr) {
|
||||
t.Fatalf("%s: expected output to contain %q\n--- output ---\n%s\n--- end ---", context, substr, output)
|
||||
}
|
||||
}
|
||||
|
||||
// requireNotContains fails the test if substr IS found in output.
|
||||
func requireNotContains(t *testing.T, output, substr, context string) {
|
||||
t.Helper()
|
||||
if strings.Contains(output, substr) {
|
||||
t.Fatalf("%s: expected output to NOT contain %q\n--- output ---\n%s\n--- end ---", context, substr, output)
|
||||
}
|
||||
}
|
||||
|
||||
// extractFieldAfter returns the first occurrence of the value after a "Prefix: " line.
|
||||
// Example: extractFieldAfter(out, "Access Key:") -> "AKIAXXXX..."
|
||||
// Returns "" if not found.
|
||||
func extractFieldAfter(output, prefix string) string {
|
||||
for _, line := range strings.Split(output, "\n") {
|
||||
line = strings.TrimSpace(line)
|
||||
if strings.HasPrefix(line, prefix) {
|
||||
return strings.TrimSpace(strings.TrimPrefix(line, prefix))
|
||||
}
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// splitLines splits output into trimmed non-empty lines.
|
||||
func splitLines(output string) []string {
|
||||
var lines []string
|
||||
for _, line := range strings.Split(output, "\n") {
|
||||
if trimmed := strings.TrimSpace(line); trimmed != "" {
|
||||
lines = append(lines, trimmed)
|
||||
}
|
||||
}
|
||||
return lines
|
||||
}
|
||||
|
||||
// fieldsOf splits a line on whitespace.
|
||||
func fieldsOf(line string) []string {
|
||||
return strings.Fields(line)
|
||||
}
|
||||
|
||||
// extractServiceAccountID parses the tab-separated output of `s3.serviceaccount.list`
|
||||
// and returns the ID of the first row whose PARENT column matches parentUser.
|
||||
// The list output format is:
|
||||
//
|
||||
// ID PARENT STATUS DESCRIPTION
|
||||
// sa:user-yyy:a1b2c3d4e5f6... user-yyy enabled some desc
|
||||
func extractServiceAccountID(t *testing.T, listOutput, parentUser string) string {
|
||||
t.Helper()
|
||||
for _, line := range strings.Split(listOutput, "\n") {
|
||||
line = strings.TrimSpace(line)
|
||||
if line == "" || strings.HasPrefix(line, "ID") || strings.HasPrefix(line, "No service accounts") {
|
||||
continue
|
||||
}
|
||||
fields := strings.Fields(line)
|
||||
if len(fields) >= 2 && fields[1] == parentUser {
|
||||
return fields[0]
|
||||
}
|
||||
}
|
||||
t.Fatalf("could not find service account with parent=%q in output:\n%s", parentUser, listOutput)
|
||||
return ""
|
||||
}
|
||||
@@ -0,0 +1,67 @@
|
||||
package policy
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"os"
|
||||
"testing"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
// TestShellPolicyAttachDetach exercises s3.policy.attach and s3.policy.detach.
|
||||
func TestShellPolicyAttachDetach(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("Skipping integration test in short mode")
|
||||
}
|
||||
cluster, err := startMiniCluster(t)
|
||||
require.NoError(t, err)
|
||||
defer cluster.Stop()
|
||||
|
||||
const weedCmd = "weed"
|
||||
master := string(pb.NewServerAddress("127.0.0.1", cluster.masterPort, cluster.masterGrpcPort))
|
||||
filer := string(pb.NewServerAddress("127.0.0.1", cluster.filerPort, cluster.filerGrpcPort))
|
||||
|
||||
// Create a policy via file.
|
||||
policyJSON := `{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Action":"s3:GetObject","Resource":"*"}]}`
|
||||
tmpFile, err := os.CreateTemp("", "test_policy_*.json")
|
||||
require.NoError(t, err)
|
||||
defer os.Remove(tmpFile.Name())
|
||||
_, err = tmpFile.WriteString(policyJSON)
|
||||
require.NoError(t, err)
|
||||
require.NoError(t, tmpFile.Close())
|
||||
|
||||
policyName := uniqueName("attach-pol")
|
||||
userName := uniqueName("attach-user")
|
||||
|
||||
execShell(t, weedCmd, master, filer,
|
||||
fmt.Sprintf("s3.policy -put -name=%s -file=%s", policyName, tmpFile.Name()))
|
||||
defer execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.policy -delete -name=%s", policyName))
|
||||
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.create -name %s", userName))
|
||||
defer execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.delete -name %s", userName))
|
||||
|
||||
t.Run("AttachAndVerify", func(t *testing.T) {
|
||||
out := execShell(t, weedCmd, master, filer,
|
||||
fmt.Sprintf("s3.policy.attach -policy %s -user %s", policyName, userName))
|
||||
requireContains(t, out, policyName, "policy.attach output")
|
||||
|
||||
out = execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.show -name %s", userName))
|
||||
requireContains(t, out, policyName, "user.show after attach")
|
||||
})
|
||||
|
||||
t.Run("AttachIdempotent", func(t *testing.T) {
|
||||
// Should succeed without error per the command's idempotent design.
|
||||
execShell(t, weedCmd, master, filer,
|
||||
fmt.Sprintf("s3.policy.attach -policy %s -user %s", policyName, userName))
|
||||
})
|
||||
|
||||
t.Run("DetachAndVerify", func(t *testing.T) {
|
||||
out := execShell(t, weedCmd, master, filer,
|
||||
fmt.Sprintf("s3.policy.detach -policy %s -user %s", policyName, userName))
|
||||
requireContains(t, out, policyName, "policy.detach output")
|
||||
|
||||
out = execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.show -name %s", userName))
|
||||
requireNotContains(t, out, fmt.Sprintf("\"%s\"", policyName), "user.show after detach")
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,80 @@
|
||||
package policy
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"testing"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
// TestShellServiceAccountLifecycle exercises s3.serviceaccount.* commands end-to-end.
|
||||
func TestShellServiceAccountLifecycle(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("Skipping integration test in short mode")
|
||||
}
|
||||
cluster, err := startMiniCluster(t)
|
||||
require.NoError(t, err)
|
||||
defer cluster.Stop()
|
||||
|
||||
const weedCmd = "weed"
|
||||
master := string(pb.NewServerAddress("127.0.0.1", cluster.masterPort, cluster.masterGrpcPort))
|
||||
filer := string(pb.NewServerAddress("127.0.0.1", cluster.filerPort, cluster.filerGrpcPort))
|
||||
|
||||
userName := uniqueName("sauser")
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.create -name %s", userName))
|
||||
defer execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.delete -name %s", userName))
|
||||
|
||||
var saID string
|
||||
|
||||
t.Run("CreateAndList", func(t *testing.T) {
|
||||
description := "integration-test-sa"
|
||||
out := execShell(t, weedCmd, master, filer,
|
||||
fmt.Sprintf("s3.serviceaccount.create -user %s -description %s", userName, description))
|
||||
requireContains(t, out, "Created service account", "serviceaccount.create")
|
||||
requireContains(t, out, "Access Key:", "serviceaccount.create credentials")
|
||||
requireContains(t, out, "Secret Key:", "serviceaccount.create credentials")
|
||||
|
||||
out = execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.serviceaccount.list -user %s", userName))
|
||||
requireContains(t, out, userName, "serviceaccount.list parent column")
|
||||
saID = extractServiceAccountID(t, out, userName)
|
||||
})
|
||||
|
||||
t.Run("Show", func(t *testing.T) {
|
||||
if saID == "" {
|
||||
t.Skip("no saID extracted")
|
||||
}
|
||||
out := execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.serviceaccount.show -id %s", saID))
|
||||
requireContains(t, out, saID, "show contains id")
|
||||
requireContains(t, out, userName, "show contains parent")
|
||||
requireContains(t, out, "enabled", "show contains status")
|
||||
})
|
||||
|
||||
t.Run("CreateWithActions", func(t *testing.T) {
|
||||
out := execShell(t, weedCmd, master, filer,
|
||||
fmt.Sprintf("s3.serviceaccount.create -user %s -actions Read,List -expiry 24h", userName))
|
||||
requireContains(t, out, "Access Key:", "serviceaccount.create with options")
|
||||
|
||||
out = execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.serviceaccount.list -user %s", userName))
|
||||
// Should now have at least 2 service accounts — count lines that start with parent in field[1]
|
||||
count := 0
|
||||
for _, line := range splitLines(out) {
|
||||
fields := fieldsOf(line)
|
||||
if len(fields) >= 2 && fields[1] == userName {
|
||||
count++
|
||||
}
|
||||
}
|
||||
if count < 2 {
|
||||
t.Fatalf("expected at least 2 service accounts for %s, got %d\n%s", userName, count, out)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("Delete", func(t *testing.T) {
|
||||
if saID == "" {
|
||||
t.Skip("no saID extracted")
|
||||
}
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.serviceaccount.delete -id %s", saID))
|
||||
out := execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.serviceaccount.list -user %s", userName))
|
||||
requireNotContains(t, out, saID, "deleted sa removed from list")
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,110 @@
|
||||
package policy
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"testing"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
|
||||
// TestShellUserLifecycle exercises the s3.user.* commands end-to-end:
|
||||
// create, show, list, enable, disable, delete, and provision.
|
||||
func TestShellUserLifecycle(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("Skipping integration test in short mode")
|
||||
}
|
||||
cluster, err := startMiniCluster(t)
|
||||
require.NoError(t, err)
|
||||
defer cluster.Stop()
|
||||
|
||||
const weedCmd = "weed"
|
||||
master := string(pb.NewServerAddress("127.0.0.1", cluster.masterPort, cluster.masterGrpcPort))
|
||||
filer := string(pb.NewServerAddress("127.0.0.1", cluster.filerPort, cluster.filerGrpcPort))
|
||||
|
||||
t.Run("CreateShowListDelete", func(t *testing.T) {
|
||||
userName := uniqueName("user")
|
||||
|
||||
out := execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.create -name %s", userName))
|
||||
requireContains(t, out, userName, "user.create output")
|
||||
requireContains(t, out, "access_key", "user.create JSON")
|
||||
requireContains(t, out, "Secret Key:", "user.create stderr")
|
||||
|
||||
out = execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.show -name %s", userName))
|
||||
requireContains(t, out, userName, "user.show")
|
||||
|
||||
out = execShell(t, weedCmd, master, filer, "s3.user.list")
|
||||
requireContains(t, out, userName, "user.list")
|
||||
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.delete -name %s", userName))
|
||||
|
||||
out = execShell(t, weedCmd, master, filer, "s3.user.list")
|
||||
requireNotContains(t, out, userName, "user.list after delete")
|
||||
})
|
||||
|
||||
t.Run("CreateWithExplicitKeys", func(t *testing.T) {
|
||||
userName := uniqueName("user-expl")
|
||||
ak := "TESTAK1234567890ABCD"
|
||||
sk := "testsecretkey1234567890abcdefghijklmnopq"
|
||||
|
||||
out := execShell(t, weedCmd, master, filer,
|
||||
fmt.Sprintf("s3.user.create -name %s -access_key %s -secret_key %s", userName, ak, sk))
|
||||
requireContains(t, out, ak, "user.create with explicit keys")
|
||||
|
||||
out = execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.show -name %s", userName))
|
||||
requireContains(t, out, ak, "user.show reveals access key")
|
||||
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.delete -name %s", userName))
|
||||
})
|
||||
|
||||
t.Run("EnableDisable", func(t *testing.T) {
|
||||
userName := uniqueName("user-toggle")
|
||||
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.create -name %s", userName))
|
||||
|
||||
out := execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.disable -name %s", userName))
|
||||
requireContains(t, out, "disabled", "user.disable")
|
||||
|
||||
out = execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.show -name %s", userName))
|
||||
requireContains(t, out, "disabled", "user.show after disable")
|
||||
|
||||
out = execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.enable -name %s", userName))
|
||||
requireContains(t, out, "enabled", "user.enable")
|
||||
|
||||
out = execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.show -name %s", userName))
|
||||
requireContains(t, out, "enabled", "user.show after enable")
|
||||
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.delete -name %s", userName))
|
||||
})
|
||||
|
||||
t.Run("Provision", func(t *testing.T) {
|
||||
userName := uniqueName("prov-user")
|
||||
bucketName := uniqueName("prov-bkt")
|
||||
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.create -name %s", bucketName))
|
||||
|
||||
out := execShell(t, weedCmd, master, filer,
|
||||
fmt.Sprintf("s3.user.provision -name %s -bucket %s -role readwrite", userName, bucketName))
|
||||
requireContains(t, out, "Created policy", "provision output")
|
||||
requireContains(t, out, "Created user", "provision output")
|
||||
requireContains(t, out, "Access Key:", "provision credentials")
|
||||
requireContains(t, out, "Secret Key:", "provision credentials")
|
||||
|
||||
out = execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.show -name %s", userName))
|
||||
requireContains(t, out, userName, "user.show after provision")
|
||||
|
||||
// Second call with same user but different bucket/role should succeed without creating duplicate user.
|
||||
bucket2 := uniqueName("prov-bkt2")
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.create -name %s", bucket2))
|
||||
out = execShell(t, weedCmd, master, filer,
|
||||
fmt.Sprintf("s3.user.provision -name %s -bucket %s -role readonly", userName, bucket2))
|
||||
requireContains(t, out, "already exists", "provision on existing user")
|
||||
requireContains(t, out, "Created policy", "second policy created")
|
||||
requireContains(t, out, "Attached policy", "second policy attached to existing user")
|
||||
requireNotContains(t, out, "Access Key:", "no new credentials for existing user")
|
||||
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.user.delete -name %s", userName))
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.delete -name %s", bucketName))
|
||||
execShell(t, weedCmd, master, filer, fmt.Sprintf("s3.bucket.delete -name %s", bucket2))
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,402 @@
|
||||
package vacuum
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"net/http"
|
||||
"os"
|
||||
"os/exec"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/operation"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/volume_server_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/shell"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/needle"
|
||||
"github.com/stretchr/testify/require"
|
||||
"google.golang.org/grpc"
|
||||
)
|
||||
|
||||
type TestCluster struct {
|
||||
masterCmd *exec.Cmd
|
||||
volumeServers []*exec.Cmd
|
||||
}
|
||||
|
||||
func (c *TestCluster) Stop() {
|
||||
for _, cmd := range c.volumeServers {
|
||||
if cmd != nil && cmd.Process != nil {
|
||||
cmd.Process.Kill()
|
||||
cmd.Wait()
|
||||
}
|
||||
}
|
||||
if c.masterCmd != nil && c.masterCmd.Process != nil {
|
||||
c.masterCmd.Process.Kill()
|
||||
c.masterCmd.Wait()
|
||||
}
|
||||
}
|
||||
|
||||
func startCluster(ctx context.Context, dataDir string) (*TestCluster, error) {
|
||||
weedBinary := findWeedBinary()
|
||||
if weedBinary == "" {
|
||||
return nil, fmt.Errorf("weed binary not found - build with 'cd weed && go build' first")
|
||||
}
|
||||
|
||||
cluster := &TestCluster{}
|
||||
|
||||
masterDir := filepath.Join(dataDir, "master")
|
||||
os.MkdirAll(masterDir, 0755)
|
||||
|
||||
// Empty security.toml to disable JWT in tests
|
||||
os.WriteFile(filepath.Join(dataDir, "security.toml"), []byte("# test\n"), 0644)
|
||||
|
||||
// Start master
|
||||
masterCmd := exec.CommandContext(ctx, weedBinary, "master",
|
||||
"-port", "9333",
|
||||
"-mdir", masterDir,
|
||||
"-volumeSizeLimitMB", "10",
|
||||
"-ip", "127.0.0.1",
|
||||
)
|
||||
masterCmd.Dir = dataDir
|
||||
masterLog, _ := os.Create(filepath.Join(masterDir, "master.log"))
|
||||
masterCmd.Stdout = masterLog
|
||||
masterCmd.Stderr = masterLog
|
||||
if err := masterCmd.Start(); err != nil {
|
||||
return nil, fmt.Errorf("start master: %v", err)
|
||||
}
|
||||
cluster.masterCmd = masterCmd
|
||||
time.Sleep(2 * time.Second)
|
||||
|
||||
// Start 2 volume servers (enough for vacuum testing)
|
||||
for i := 0; i < 2; i++ {
|
||||
volumeDir := filepath.Join(dataDir, fmt.Sprintf("volume%d", i))
|
||||
os.MkdirAll(volumeDir, 0755)
|
||||
|
||||
port := fmt.Sprintf("808%d", i)
|
||||
volumeCmd := exec.CommandContext(ctx, weedBinary, "volume",
|
||||
"-port", port,
|
||||
"-dir", volumeDir,
|
||||
"-max", "10",
|
||||
"-master", "127.0.0.1:9333",
|
||||
"-ip", "127.0.0.1",
|
||||
)
|
||||
volumeCmd.Dir = dataDir
|
||||
volumeLog, _ := os.Create(filepath.Join(volumeDir, "volume.log"))
|
||||
volumeCmd.Stdout = volumeLog
|
||||
volumeCmd.Stderr = volumeLog
|
||||
if err := volumeCmd.Start(); err != nil {
|
||||
cluster.Stop()
|
||||
return nil, fmt.Errorf("start volume server %d: %v", i, err)
|
||||
}
|
||||
cluster.volumeServers = append(cluster.volumeServers, volumeCmd)
|
||||
}
|
||||
|
||||
time.Sleep(5 * time.Second)
|
||||
return cluster, nil
|
||||
}
|
||||
|
||||
func findWeedBinary() string {
|
||||
candidates := []string{
|
||||
"../../weed/weed",
|
||||
"../weed/weed",
|
||||
"./weed",
|
||||
}
|
||||
for _, c := range candidates {
|
||||
if _, err := os.Stat(c); err == nil {
|
||||
if abs, err := filepath.Abs(c); err == nil {
|
||||
return abs
|
||||
}
|
||||
return c
|
||||
}
|
||||
}
|
||||
if path, err := exec.LookPath("weed"); err == nil {
|
||||
return path
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
func waitForServer(address string, timeout time.Duration) error {
|
||||
start := time.Now()
|
||||
for time.Since(start) < timeout {
|
||||
if conn, err := net.DialTimeout("tcp", address, 1*time.Second); err == nil {
|
||||
conn.Close()
|
||||
return nil
|
||||
}
|
||||
time.Sleep(500 * time.Millisecond)
|
||||
}
|
||||
return fmt.Errorf("timeout waiting for server %s", address)
|
||||
}
|
||||
|
||||
func uploadData(masterAddr, collection string, data []byte) (string, needle.VolumeId, error) {
|
||||
assignResult, err := operation.Assign(context.Background(), func(ctx context.Context) pb.ServerAddress {
|
||||
return pb.ServerAddress(masterAddr)
|
||||
}, grpc.WithInsecure(), &operation.VolumeAssignRequest{
|
||||
Count: 1,
|
||||
Collection: collection,
|
||||
})
|
||||
if err != nil {
|
||||
return "", 0, fmt.Errorf("assign: %v", err)
|
||||
}
|
||||
|
||||
uploader, err := operation.NewUploader()
|
||||
if err != nil {
|
||||
return "", 0, fmt.Errorf("new uploader: %v", err)
|
||||
}
|
||||
|
||||
uploadResult, err, _ := uploader.Upload(context.Background(), bytes.NewReader(data), &operation.UploadOption{
|
||||
UploadUrl: "http://" + assignResult.Url + "/" + assignResult.Fid,
|
||||
Filename: "testfile.txt",
|
||||
MimeType: "text/plain",
|
||||
})
|
||||
if err != nil {
|
||||
return "", 0, fmt.Errorf("upload: %v", err)
|
||||
}
|
||||
if uploadResult.Error != "" {
|
||||
return "", 0, fmt.Errorf("upload error: %s", uploadResult.Error)
|
||||
}
|
||||
|
||||
fid, err := needle.ParseFileIdFromString(assignResult.Fid)
|
||||
if err != nil {
|
||||
return "", 0, err
|
||||
}
|
||||
return assignResult.Fid, fid.VolumeId, nil
|
||||
}
|
||||
|
||||
func deleteFile(masterAddr string, fid string) error {
|
||||
results := operation.DeleteFileIds(func(ctx context.Context) pb.ServerAddress {
|
||||
return pb.ServerAddress(masterAddr)
|
||||
}, false, grpc.WithInsecure(), []string{fid})
|
||||
for _, r := range results {
|
||||
if r.Error != "" {
|
||||
return fmt.Errorf("delete %s: %s", fid, r.Error)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func getGarbageRatio(volumeServerAddr string, volumeId uint32) (float64, error) {
|
||||
var ratio float64
|
||||
err := operation.WithVolumeServerClient(false, pb.ServerAddress(volumeServerAddr), grpc.WithInsecure(),
|
||||
func(client volume_server_pb.VolumeServerClient) error {
|
||||
resp, err := client.VacuumVolumeCheck(context.Background(), &volume_server_pb.VacuumVolumeCheckRequest{
|
||||
VolumeId: volumeId,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
ratio = resp.GarbageRatio
|
||||
return nil
|
||||
})
|
||||
return ratio, err
|
||||
}
|
||||
|
||||
// TestVacuumIntegration tests the full vacuum flow:
|
||||
// upload data → delete some → verify garbage → vacuum → verify cleanup
|
||||
func TestVacuumIntegration(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("Skipping integration test in short mode")
|
||||
}
|
||||
|
||||
testDir := t.TempDir()
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 120*time.Second)
|
||||
defer cancel()
|
||||
|
||||
cluster, err := startCluster(ctx, testDir)
|
||||
require.NoError(t, err)
|
||||
defer cluster.Stop()
|
||||
|
||||
require.NoError(t, waitForServer("127.0.0.1:9333", 30*time.Second))
|
||||
require.NoError(t, waitForServer("127.0.0.1:8080", 30*time.Second))
|
||||
require.NoError(t, waitForServer("127.0.0.1:8081", 30*time.Second))
|
||||
|
||||
masterAddr := "127.0.0.1:9333"
|
||||
collection := "vactest"
|
||||
|
||||
// Upload files large enough that deleting most creates significant garbage.
|
||||
// With volumeSizeLimitMB=10, we need several MB of garbage to exceed the
|
||||
// 10% threshold passed to vacuum.
|
||||
const fileSize = 500 * 1024 // 500 KB per file
|
||||
const totalFiles = 16
|
||||
const filesToDelete = 12 // delete 75% → ~6 MB garbage out of ~8 MB
|
||||
|
||||
var fids []string
|
||||
var payloads [][]byte
|
||||
var volumeId needle.VolumeId
|
||||
for i := 0; i < totalFiles; i++ {
|
||||
data := bytes.Repeat([]byte{byte('A' + i%26)}, fileSize)
|
||||
fid, vid, err := uploadData(masterAddr, collection, data)
|
||||
require.NoError(t, err, "upload %d", i)
|
||||
fids = append(fids, fid)
|
||||
payloads = append(payloads, data)
|
||||
volumeId = vid
|
||||
}
|
||||
t.Logf("Uploaded %d files (%d KB each) to volume %d", totalFiles, fileSize/1024, volumeId)
|
||||
|
||||
// Wait for heartbeat to report sizes
|
||||
time.Sleep(6 * time.Second)
|
||||
|
||||
// Delete most files to create garbage well above the threshold
|
||||
for i := 0; i < filesToDelete; i++ {
|
||||
err := deleteFile(masterAddr, fids[i])
|
||||
require.NoError(t, err, "delete %s", fids[i])
|
||||
}
|
||||
t.Logf("Deleted %d of %d files to create garbage", filesToDelete, totalFiles)
|
||||
|
||||
// Wait for heartbeat to report deletions
|
||||
time.Sleep(6 * time.Second)
|
||||
|
||||
// Verify garbage exists
|
||||
t.Run("verify_garbage_before_vacuum", func(t *testing.T) {
|
||||
for _, addr := range []string{"127.0.0.1:8080", "127.0.0.1:8081"} {
|
||||
ratio, err := getGarbageRatio(addr, uint32(volumeId))
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
t.Logf("Garbage ratio on %s: %.2f%%", addr, ratio*100)
|
||||
if ratio > 0.1 {
|
||||
return // sufficient garbage found
|
||||
}
|
||||
}
|
||||
t.Fatal("No server reported garbage > 10% — test data setup failed")
|
||||
})
|
||||
|
||||
// Execute vacuum via shell command
|
||||
t.Run("run_vacuum", func(t *testing.T) {
|
||||
options := &shell.ShellOptions{
|
||||
Masters: stringPtr(masterAddr),
|
||||
GrpcDialOption: grpc.WithInsecure(),
|
||||
FilerGroup: stringPtr("default"),
|
||||
}
|
||||
commandEnv := shell.NewCommandEnv(options)
|
||||
|
||||
shellCtx, shellCancel := context.WithTimeout(context.Background(), 60*time.Second)
|
||||
defer shellCancel()
|
||||
go commandEnv.MasterClient.KeepConnectedToMaster(shellCtx)
|
||||
commandEnv.MasterClient.WaitUntilConnected(shellCtx)
|
||||
time.Sleep(2 * time.Second)
|
||||
|
||||
// Acquire lock (required by shell commands)
|
||||
locked, unlock := tryLock(t, commandEnv, 30*time.Second)
|
||||
require.True(t, locked, "could not acquire shell lock")
|
||||
defer unlock()
|
||||
|
||||
// Find and execute vacuum command
|
||||
var output bytes.Buffer
|
||||
var found bool
|
||||
var err error
|
||||
for _, cmd := range shell.Commands {
|
||||
if cmd.Name() == "volume.vacuum" {
|
||||
err = cmd.Do(
|
||||
[]string{"-garbageThreshold", "0.1", "-collection", collection},
|
||||
commandEnv, &output,
|
||||
)
|
||||
found = true
|
||||
break
|
||||
}
|
||||
}
|
||||
require.True(t, found, "volume.vacuum command not found")
|
||||
t.Logf("Vacuum output: %s", output.String())
|
||||
require.NoError(t, err, "vacuum command failed")
|
||||
t.Log("Vacuum completed successfully")
|
||||
})
|
||||
|
||||
// Wait for vacuum effects to settle
|
||||
time.Sleep(6 * time.Second)
|
||||
|
||||
// Verify garbage was cleaned
|
||||
t.Run("verify_cleanup_after_vacuum", func(t *testing.T) {
|
||||
var volumeFound, cleanupVerified bool
|
||||
for _, addr := range []string{"127.0.0.1:8080", "127.0.0.1:8081"} {
|
||||
ratio, err := getGarbageRatio(addr, uint32(volumeId))
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
volumeFound = true
|
||||
t.Logf("Garbage ratio after vacuum on %s: %.2f%%", addr, ratio*100)
|
||||
if ratio < 0.05 {
|
||||
cleanupVerified = true
|
||||
}
|
||||
}
|
||||
if !volumeFound {
|
||||
t.Fatal("No server reported volume after vacuum")
|
||||
}
|
||||
if !cleanupVerified {
|
||||
t.Fatal("Garbage was not cleaned up after vacuum")
|
||||
}
|
||||
})
|
||||
|
||||
// Verify remaining files are still readable with correct contents
|
||||
t.Run("verify_remaining_data", func(t *testing.T) {
|
||||
for i := filesToDelete; i < totalFiles; i++ {
|
||||
fid := fids[i]
|
||||
expected := payloads[i]
|
||||
|
||||
// Read file via HTTP from volume server
|
||||
client := &http.Client{Timeout: 5 * time.Second}
|
||||
url := fmt.Sprintf("http://127.0.0.1:8080/%s", fid)
|
||||
resp, err := client.Get(url)
|
||||
if err != nil || resp.StatusCode == http.StatusNotFound {
|
||||
if resp != nil {
|
||||
resp.Body.Close()
|
||||
}
|
||||
url = fmt.Sprintf("http://127.0.0.1:8081/%s", fid)
|
||||
resp, err = client.Get(url)
|
||||
}
|
||||
require.NoError(t, err, "read fid %s", fid)
|
||||
body, err := io.ReadAll(resp.Body)
|
||||
resp.Body.Close()
|
||||
require.NoError(t, err, "read body of fid %s", fid)
|
||||
require.Equal(t, http.StatusOK, resp.StatusCode, "fid %s returned %d", fid, resp.StatusCode)
|
||||
require.Equal(t, len(expected), len(body), "fid %s size mismatch", fid)
|
||||
require.True(t, bytes.Equal(expected, body), "fid %s content mismatch", fid)
|
||||
t.Logf("File %s verified (%d bytes)", fid, len(body))
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
func stringPtr(s string) *string {
|
||||
return &s
|
||||
}
|
||||
|
||||
func tryLock(t *testing.T, commandEnv *shell.CommandEnv, timeout time.Duration) (locked bool, unlock func()) {
|
||||
t.Helper()
|
||||
type result struct {
|
||||
err error
|
||||
}
|
||||
done := make(chan result, 1)
|
||||
go func() {
|
||||
for _, cmd := range shell.Commands {
|
||||
if cmd.Name() == "lock" {
|
||||
var out bytes.Buffer
|
||||
done <- result{err: cmd.Do([]string{}, commandEnv, &out)}
|
||||
return
|
||||
}
|
||||
}
|
||||
done <- result{err: fmt.Errorf("lock command not found")}
|
||||
}()
|
||||
|
||||
select {
|
||||
case res := <-done:
|
||||
if res.err != nil {
|
||||
t.Logf("lock failed: %v", res.err)
|
||||
return false, nil
|
||||
}
|
||||
return true, func() {
|
||||
for _, cmd := range shell.Commands {
|
||||
if cmd.Name() == "unlock" {
|
||||
var out bytes.Buffer
|
||||
cmd.Do([]string{}, commandEnv, &out)
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
case <-time.After(timeout):
|
||||
t.Log("lock timed out")
|
||||
return false, nil
|
||||
}
|
||||
}
|
||||
@@ -272,7 +272,6 @@ func stopProcess(cmd *exec.Cmd) {
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
func newWorkDir() (dir string, keepLogs bool, err error) {
|
||||
keepLogs = os.Getenv("VOLUME_SERVER_IT_KEEP_LOGS") == "1"
|
||||
dir, err = os.MkdirTemp("", "seaweedfs_volume_server_it_")
|
||||
|
||||
@@ -778,8 +778,8 @@ func TestEcIndexConsistencyAfterEncode(t *testing.T) {
|
||||
needles := []testNeedle{
|
||||
{framework.NewFileID(volumeID, 1001, 0xAABB0001), []byte("small-needle-1")},
|
||||
{framework.NewFileID(volumeID, 1002, 0xAABB0002), make([]byte, 1024)}, // 1KB
|
||||
{framework.NewFileID(volumeID, 1003, 0xAABB0003), make([]byte, 64*1024)}, // 64KB
|
||||
{framework.NewFileID(volumeID, 1004, 0xAABB0004), make([]byte, 256*1024)}, // 256KB
|
||||
{framework.NewFileID(volumeID, 1003, 0xAABB0003), make([]byte, 64*1024)}, // 64KB
|
||||
{framework.NewFileID(volumeID, 1004, 0xAABB0004), make([]byte, 256*1024)}, // 256KB
|
||||
{framework.NewFileID(volumeID, 1005, 0xAABB0005), []byte("small-needle-2")},
|
||||
}
|
||||
|
||||
|
||||
@@ -459,9 +459,10 @@ func sortDurations(d []time.Duration) {
|
||||
// This reveals tail latency differences that short tests miss (GC pauses, lock contention, etc).
|
||||
//
|
||||
// Run:
|
||||
// go test -v -count=1 -timeout 600s -run TestSustainedP99 ./test/volume_server/loadtest/...
|
||||
// VOLUME_SERVER_IMPL=rust go test -v -count=1 -timeout 600s -run TestSustainedP99 ./test/volume_server/loadtest/...
|
||||
// LOADTEST_DURATION=120s VOLUME_SERVER_IMPL=rust go test -v -count=1 -timeout 600s -run TestSustainedP99 ./test/volume_server/loadtest/...
|
||||
//
|
||||
// go test -v -count=1 -timeout 600s -run TestSustainedP99 ./test/volume_server/loadtest/...
|
||||
// VOLUME_SERVER_IMPL=rust go test -v -count=1 -timeout 600s -run TestSustainedP99 ./test/volume_server/loadtest/...
|
||||
// LOADTEST_DURATION=120s VOLUME_SERVER_IMPL=rust go test -v -count=1 -timeout 600s -run TestSustainedP99 ./test/volume_server/loadtest/...
|
||||
func TestSustainedP99(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("skipping sustained load test in short mode")
|
||||
|
||||
@@ -289,8 +289,8 @@ func (cp *ConfigPersistence) LoadVacuumTaskConfig() (*VacuumTaskConfig, error) {
|
||||
|
||||
// Return default config if no valid config found
|
||||
return &VacuumTaskConfig{
|
||||
GarbageThreshold: 0.3,
|
||||
MinVolumeAgeHours: 24,
|
||||
GarbageThreshold: 0.3,
|
||||
MinVolumeAgeHours: 24,
|
||||
}, nil
|
||||
}
|
||||
|
||||
@@ -305,8 +305,8 @@ func (cp *ConfigPersistence) LoadVacuumTaskPolicy() (*worker_pb.TaskPolicy, erro
|
||||
CheckIntervalSeconds: 6 * 3600, // 6 hours in seconds
|
||||
TaskConfig: &worker_pb.TaskPolicy_VacuumConfig{
|
||||
VacuumConfig: &worker_pb.VacuumTaskConfig{
|
||||
GarbageThreshold: 0.3,
|
||||
MinVolumeAgeHours: 24,
|
||||
GarbageThreshold: 0.3,
|
||||
MinVolumeAgeHours: 24,
|
||||
},
|
||||
},
|
||||
}, nil
|
||||
@@ -325,8 +325,8 @@ func (cp *ConfigPersistence) LoadVacuumTaskPolicy() (*worker_pb.TaskPolicy, erro
|
||||
CheckIntervalSeconds: 6 * 3600, // 6 hours in seconds
|
||||
TaskConfig: &worker_pb.TaskPolicy_VacuumConfig{
|
||||
VacuumConfig: &worker_pb.VacuumTaskConfig{
|
||||
GarbageThreshold: 0.3,
|
||||
MinVolumeAgeHours: 24,
|
||||
GarbageThreshold: 0.3,
|
||||
MinVolumeAgeHours: 24,
|
||||
},
|
||||
},
|
||||
}, nil
|
||||
@@ -705,8 +705,8 @@ func buildPolicyFromTaskConfigs() *worker_pb.MaintenancePolicy {
|
||||
CheckIntervalSeconds: int32(vacuumConfig.ScanIntervalSeconds),
|
||||
TaskConfig: &worker_pb.TaskPolicy_VacuumConfig{
|
||||
VacuumConfig: &worker_pb.VacuumTaskConfig{
|
||||
GarbageThreshold: float64(vacuumConfig.GarbageThreshold),
|
||||
MinVolumeAgeHours: int32(vacuumConfig.MinVolumeAgeSeconds / 3600), // Convert seconds to hours
|
||||
GarbageThreshold: float64(vacuumConfig.GarbageThreshold),
|
||||
MinVolumeAgeHours: int32(vacuumConfig.MinVolumeAgeSeconds / 3600), // Convert seconds to hours
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
@@ -452,6 +452,13 @@ func (s *AdminServer) UpdatePluginJobTypeConfigAPI(w http.ResponseWriter, r *htt
|
||||
}
|
||||
config.UpdatedBy = username
|
||||
|
||||
// Reapply descriptor defaults so a save from an older form (or a UI
|
||||
// that omits new fields) cannot silently clear a baseline like
|
||||
// execution_timeout_seconds back to zero.
|
||||
if descriptor, err := s.LoadPluginJobTypeDescriptor(jobType); err == nil && descriptor != nil {
|
||||
applyDescriptorDefaultsToPersistedConfig(config, descriptor)
|
||||
}
|
||||
|
||||
if err := s.SavePluginJobTypeConfig(config); err != nil {
|
||||
writeJSONError(w, http.StatusInternalServerError, err.Error())
|
||||
return
|
||||
@@ -916,6 +923,9 @@ func applyDescriptorDefaultsToPersistedConfig(
|
||||
if runtime.JobTypeMaxRuntimeSeconds <= 0 {
|
||||
runtime.JobTypeMaxRuntimeSeconds = defaults.JobTypeMaxRuntimeSeconds
|
||||
}
|
||||
if runtime.ExecutionTimeoutSeconds <= 0 {
|
||||
runtime.ExecutionTimeoutSeconds = defaults.ExecutionTimeoutSeconds
|
||||
}
|
||||
if runtime.RetryBackoffSeconds <= 0 {
|
||||
runtime.RetryBackoffSeconds = defaults.RetryBackoffSeconds
|
||||
}
|
||||
|
||||
@@ -9,6 +9,7 @@ import (
|
||||
"mime/multipart"
|
||||
"net"
|
||||
"net/http"
|
||||
"net/url"
|
||||
"os"
|
||||
"path"
|
||||
"path/filepath"
|
||||
@@ -398,7 +399,7 @@ func (h *FileBrowserHandlers) uploadFileToFiler(filePath string, fileHeader *mul
|
||||
|
||||
// Create the upload URL - the httpClient will normalize to the correct scheme (http/https)
|
||||
// based on the https.client configuration in security.toml
|
||||
uploadURL := fmt.Sprintf("%s%s", filerHttpAddress, cleanFilePath)
|
||||
uploadURL := filerFileURL(filerHttpAddress, cleanFilePath)
|
||||
|
||||
// Normalize the URL scheme based on TLS configuration
|
||||
uploadURL, err = h.httpClient.NormalizeHttpScheme(uploadURL)
|
||||
@@ -503,14 +504,16 @@ func (h *FileBrowserHandlers) validateAndCleanFilePath(filePath string) (string,
|
||||
return "", fmt.Errorf("path traversal not allowed")
|
||||
}
|
||||
|
||||
// Additional validation: ensure path doesn't contain dangerous characters
|
||||
if strings.ContainsAny(cleanPath, "\x00\r\n") {
|
||||
return "", fmt.Errorf("path contains invalid characters")
|
||||
}
|
||||
|
||||
return cleanPath, nil
|
||||
}
|
||||
|
||||
// filerFileURL joins the filer HTTP address with a validated file path, URL-escaping
|
||||
// the path so that control characters and other bytes that are legal in S3 object keys
|
||||
// cannot inject into the HTTP request target.
|
||||
func filerFileURL(filerHttpAddress, cleanFilePath string) string {
|
||||
return filerHttpAddress + (&url.URL{Path: cleanFilePath}).EscapedPath()
|
||||
}
|
||||
|
||||
// fetchFileContent fetches file content from the filer and returns the content or an error.
|
||||
func (h *FileBrowserHandlers) fetchFileContent(filePath string, timeout time.Duration) (string, error) {
|
||||
filerAddress := h.adminServer.GetFilerAddress()
|
||||
@@ -529,7 +532,7 @@ func (h *FileBrowserHandlers) fetchFileContent(filePath string, timeout time.Dur
|
||||
}
|
||||
|
||||
// Create the file URL with proper scheme based on TLS configuration
|
||||
fileURL := fmt.Sprintf("%s%s", filerHttpAddress, cleanFilePath)
|
||||
fileURL := filerFileURL(filerHttpAddress, cleanFilePath)
|
||||
fileURL, err = h.httpClient.NormalizeHttpScheme(fileURL)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("failed to construct file URL: %w", err)
|
||||
@@ -597,7 +600,7 @@ func (h *FileBrowserHandlers) DownloadFile(w http.ResponseWriter, r *http.Reques
|
||||
}
|
||||
|
||||
// Create the download URL with proper scheme based on TLS configuration
|
||||
downloadURL := fmt.Sprintf("%s%s", filerHttpAddress, cleanFilePath)
|
||||
downloadURL := filerFileURL(filerHttpAddress, cleanFilePath)
|
||||
downloadURL, err = h.httpClient.NormalizeHttpScheme(downloadURL)
|
||||
if err != nil {
|
||||
writeJSONError(w, http.StatusInternalServerError, "Failed to construct download URL: "+err.Error())
|
||||
@@ -1043,7 +1046,7 @@ func (h *FileBrowserHandlers) isLikelyTextFile(filePath string, maxCheckSize int
|
||||
}
|
||||
|
||||
// Create the file URL with proper scheme based on TLS configuration
|
||||
fileURL := fmt.Sprintf("%s%s", filerHttpAddress, cleanFilePath)
|
||||
fileURL := filerFileURL(filerHttpAddress, cleanFilePath)
|
||||
fileURL, err = h.httpClient.NormalizeHttpScheme(fileURL)
|
||||
if err != nil {
|
||||
glog.Errorf("Failed to normalize URL scheme: %v", err)
|
||||
|
||||
@@ -0,0 +1,62 @@
|
||||
package handlers
|
||||
|
||||
import (
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestValidateAndCleanFilePath_AllowsControlChars(t *testing.T) {
|
||||
h := &FileBrowserHandlers{}
|
||||
|
||||
// S3 object keys may legally contain any UTF-8 bytes, including control
|
||||
// characters like \n, \r, and \x00. The admin UI must be able to browse
|
||||
// and manage such entries rather than silently stripping or rejecting them.
|
||||
cases := []struct {
|
||||
in string
|
||||
want string
|
||||
}{
|
||||
{"/buckets/profilebuilder/3testGB.zip\n ", "/buckets/profilebuilder/3testGB.zip\n "},
|
||||
{"/foo\rbar", "/foo\rbar"},
|
||||
{"/foo\x00bar", "/foo\x00bar"},
|
||||
{"/normal/path.txt", "/normal/path.txt"},
|
||||
// Missing leading slash should be added back.
|
||||
{"relative/path.txt", "/relative/path.txt"},
|
||||
// Duplicate slashes should be collapsed by path.Clean.
|
||||
{"/a//b", "/a/b"},
|
||||
}
|
||||
for _, tc := range cases {
|
||||
got, err := h.validateAndCleanFilePath(tc.in)
|
||||
if err != nil {
|
||||
t.Errorf("validateAndCleanFilePath(%q) unexpected error: %v", tc.in, err)
|
||||
continue
|
||||
}
|
||||
if got != tc.want {
|
||||
t.Errorf("validateAndCleanFilePath(%q) = %q, want %q", tc.in, got, tc.want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestValidateAndCleanFilePath_RejectsEmpty(t *testing.T) {
|
||||
h := &FileBrowserHandlers{}
|
||||
if _, err := h.validateAndCleanFilePath(""); err == nil {
|
||||
t.Errorf("expected empty path rejection")
|
||||
}
|
||||
}
|
||||
|
||||
func TestFilerFileURL_EscapesControlChars(t *testing.T) {
|
||||
cases := []struct {
|
||||
addr string
|
||||
path string
|
||||
want string
|
||||
}{
|
||||
{"http://127.0.0.1:8888", "/buckets/profilebuilder/3testGB.zip\n ", "http://127.0.0.1:8888/buckets/profilebuilder/3testGB.zip%0A%20"},
|
||||
{"http://127.0.0.1:8888", "/buckets/profilebuilder/file\rname", "http://127.0.0.1:8888/buckets/profilebuilder/file%0Dname"},
|
||||
{"http://127.0.0.1:8888", "/buckets/profilebuilder/file\x00name", "http://127.0.0.1:8888/buckets/profilebuilder/file%00name"},
|
||||
// Plain path round-trips unchanged.
|
||||
{"http://h:1", "/a/b.txt", "http://h:1/a/b.txt"},
|
||||
}
|
||||
for _, tc := range cases {
|
||||
if got := filerFileURL(tc.addr, tc.path); got != tc.want {
|
||||
t.Errorf("filerFileURL(%q, %q) = %q, want %q", tc.addr, tc.path, got, tc.want)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1044,6 +1044,7 @@ func (r *Plugin) ensureJobTypeConfigFromDescriptor(jobType string, descriptor *p
|
||||
RetryLimit: defaults.RetryLimit,
|
||||
RetryBackoffSeconds: defaults.RetryBackoffSeconds,
|
||||
JobTypeMaxRuntimeSeconds: defaults.JobTypeMaxRuntimeSeconds,
|
||||
ExecutionTimeoutSeconds: defaults.ExecutionTimeoutSeconds,
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -468,7 +468,7 @@ func (r *Plugin) loadSchedulerPolicy(jobType string) (schedulerPolicy, bool, err
|
||||
policy := schedulerPolicy{
|
||||
DetectionInterval: durationFromSeconds(adminRuntime.DetectionIntervalSeconds, defaultScheduledDetectionInterval),
|
||||
DetectionTimeout: durationFromSeconds(adminRuntime.DetectionTimeoutSeconds, defaultScheduledDetectionTimeout),
|
||||
ExecutionTimeout: defaultScheduledExecutionTimeout,
|
||||
ExecutionTimeout: durationFromSeconds(adminRuntime.ExecutionTimeoutSeconds, defaultScheduledExecutionTimeout),
|
||||
JobTypeMaxRuntime: durationFromSeconds(adminRuntime.JobTypeMaxRuntimeSeconds, defaultScheduledJobTypeMaxRuntime),
|
||||
RetryBackoff: durationFromSeconds(adminRuntime.RetryBackoffSeconds, defaultScheduledRetryBackoff),
|
||||
MaxResults: adminRuntime.MaxJobsPerDetection,
|
||||
@@ -503,12 +503,9 @@ func (r *Plugin) loadSchedulerPolicy(jobType string) (schedulerPolicy, bool, err
|
||||
policy.JobTypeMaxRuntime = defaultScheduledJobTypeMaxRuntime
|
||||
}
|
||||
|
||||
// Plugin protocol currently has only detection timeout in admin settings.
|
||||
execTimeout := time.Duration(adminRuntime.DetectionTimeoutSeconds*2) * time.Second
|
||||
if execTimeout < defaultScheduledExecutionTimeout {
|
||||
execTimeout = defaultScheduledExecutionTimeout
|
||||
if policy.ExecutionTimeout < defaultScheduledExecutionTimeout {
|
||||
policy.ExecutionTimeout = defaultScheduledExecutionTimeout
|
||||
}
|
||||
policy.ExecutionTimeout = execTimeout
|
||||
|
||||
return policy, true, nil
|
||||
}
|
||||
@@ -610,6 +607,37 @@ func deriveSchedulerAdminRuntime(
|
||||
) *plugin_pb.AdminRuntimeConfig {
|
||||
if cfg != nil && cfg.AdminRuntime != nil {
|
||||
adminConfig := *cfg.AdminRuntime
|
||||
// Overlay descriptor defaults for any zero numeric fields. Persisted
|
||||
// configs from older versions have no execution_timeout_seconds, and
|
||||
// without this overlay the scheduler would fall back to the 90s
|
||||
// default instead of the handler's declared baseline.
|
||||
if descriptor != nil && descriptor.AdminRuntimeDefaults != nil {
|
||||
defaults := descriptor.AdminRuntimeDefaults
|
||||
if adminConfig.DetectionIntervalSeconds <= 0 {
|
||||
adminConfig.DetectionIntervalSeconds = defaults.DetectionIntervalSeconds
|
||||
}
|
||||
if adminConfig.DetectionTimeoutSeconds <= 0 {
|
||||
adminConfig.DetectionTimeoutSeconds = defaults.DetectionTimeoutSeconds
|
||||
}
|
||||
if adminConfig.MaxJobsPerDetection <= 0 {
|
||||
adminConfig.MaxJobsPerDetection = defaults.MaxJobsPerDetection
|
||||
}
|
||||
if adminConfig.GlobalExecutionConcurrency <= 0 {
|
||||
adminConfig.GlobalExecutionConcurrency = defaults.GlobalExecutionConcurrency
|
||||
}
|
||||
if adminConfig.PerWorkerExecutionConcurrency <= 0 {
|
||||
adminConfig.PerWorkerExecutionConcurrency = defaults.PerWorkerExecutionConcurrency
|
||||
}
|
||||
if adminConfig.RetryBackoffSeconds <= 0 {
|
||||
adminConfig.RetryBackoffSeconds = defaults.RetryBackoffSeconds
|
||||
}
|
||||
if adminConfig.JobTypeMaxRuntimeSeconds <= 0 {
|
||||
adminConfig.JobTypeMaxRuntimeSeconds = defaults.JobTypeMaxRuntimeSeconds
|
||||
}
|
||||
if adminConfig.ExecutionTimeoutSeconds <= 0 {
|
||||
adminConfig.ExecutionTimeoutSeconds = defaults.ExecutionTimeoutSeconds
|
||||
}
|
||||
}
|
||||
return &adminConfig
|
||||
}
|
||||
|
||||
@@ -628,6 +656,7 @@ func deriveSchedulerAdminRuntime(
|
||||
RetryLimit: defaults.RetryLimit,
|
||||
RetryBackoffSeconds: defaults.RetryBackoffSeconds,
|
||||
JobTypeMaxRuntimeSeconds: defaults.JobTypeMaxRuntimeSeconds,
|
||||
ExecutionTimeoutSeconds: defaults.ExecutionTimeoutSeconds,
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -619,10 +619,10 @@ func TestRunLaneSchedulerIterationLockBehavior(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
tests := []struct {
|
||||
name string
|
||||
lane SchedulerLane
|
||||
jobType string
|
||||
wantLock bool
|
||||
name string
|
||||
lane SchedulerLane
|
||||
jobType string
|
||||
wantLock bool
|
||||
}{
|
||||
{"Default", LaneDefault, "vacuum", true},
|
||||
{"Iceberg", LaneIceberg, "iceberg_maintenance", false},
|
||||
|
||||
@@ -114,7 +114,7 @@ type schedulerLaneState struct {
|
||||
// Per-lane execution reservation pool. Each lane tracks how many
|
||||
// execution slots it has reserved on each worker independently,
|
||||
// so lanes cannot starve each other.
|
||||
execMu sync.Mutex
|
||||
execMu sync.Mutex
|
||||
execRes map[string]int
|
||||
}
|
||||
|
||||
|
||||
@@ -16,8 +16,8 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
<p class="mb-0 text-muted">Manage S3 API users and their access credentials</p>
|
||||
</div>
|
||||
<div class="d-flex gap-2">
|
||||
<button type="button" class="btn btn-primary"
|
||||
data-bs-toggle="modal"
|
||||
<button type="button" class="btn btn-primary"
|
||||
data-bs-toggle="modal"
|
||||
data-bs-target="#createUserModal">
|
||||
<i class="fas fa-plus me-1"></i>Create User
|
||||
</button>
|
||||
@@ -269,7 +269,7 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
<div class="mb-3">
|
||||
<label class="form-label">Bucket Scope</label>
|
||||
<small class="form-text text-muted d-block mb-2">Apply selected permissions to specific buckets or all buckets</small>
|
||||
|
||||
|
||||
<div class="form-check mb-2">
|
||||
<input class="form-check-input" type="radio" name="bucketScope" id="allBuckets" value="all" checked onchange="toggleBucketList()">
|
||||
<label class="form-check-label" for="allBuckets">
|
||||
@@ -282,7 +282,7 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
Specific Buckets
|
||||
</label>
|
||||
</div>
|
||||
|
||||
|
||||
<div id="bucketSelectionList" class="mt-2" style="display: none;">
|
||||
<select multiple class="form-select" id="selectedBuckets" size="5">
|
||||
<!-- Options loaded dynamically -->
|
||||
@@ -376,7 +376,7 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
<div class="mb-3">
|
||||
<label class="form-label">Bucket Scope</label>
|
||||
<small class="form-text text-muted d-block mb-2">Apply selected permissions to specific buckets or all buckets</small>
|
||||
|
||||
|
||||
<div class="form-check mb-2">
|
||||
<input class="form-check-input" type="radio" name="editBucketScope" id="editAllBuckets" value="all" checked onchange="toggleBucketList('edit')">
|
||||
<label class="form-check-label" for="editAllBuckets">
|
||||
@@ -389,7 +389,7 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
Specific Buckets
|
||||
</label>
|
||||
</div>
|
||||
|
||||
|
||||
<div id="editBucketSelectionList" class="mt-2" style="display: none;">
|
||||
<select multiple class="form-select" id="editSelectedBuckets" size="5">
|
||||
<!-- Options loaded dynamically -->
|
||||
@@ -505,15 +505,15 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
const STATUS_INACTIVE = 'Inactive';
|
||||
|
||||
document.addEventListener('DOMContentLoaded', function() {
|
||||
|
||||
|
||||
// Event delegation for user action buttons
|
||||
document.addEventListener('click', function(e) {
|
||||
const button = e.target.closest('[data-action]');
|
||||
if (!button) return;
|
||||
|
||||
|
||||
const action = button.getAttribute('data-action');
|
||||
const username = button.getAttribute('data-username');
|
||||
|
||||
|
||||
switch (action) {
|
||||
case 'show-user-details':
|
||||
showUserDetails(username);
|
||||
@@ -553,7 +553,7 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
|
||||
// Load policies for dropdowns
|
||||
loadPolicies();
|
||||
|
||||
|
||||
// Load buckets for bucket permissions
|
||||
loadBuckets();
|
||||
});
|
||||
@@ -627,20 +627,20 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
if (response.ok) {
|
||||
const data = await response.json();
|
||||
const policies = data.policies || [];
|
||||
|
||||
|
||||
const createSelect = document.getElementById('policies');
|
||||
const editSelect = document.getElementById('editPolicies');
|
||||
|
||||
|
||||
// Check if elements exist
|
||||
if (!createSelect || !editSelect) {
|
||||
console.warn('Policy select elements not found');
|
||||
return;
|
||||
}
|
||||
|
||||
|
||||
// Clear existing options
|
||||
createSelect.innerHTML = '';
|
||||
editSelect.innerHTML = '';
|
||||
|
||||
|
||||
if (policies && policies.length > 0) {
|
||||
policies.forEach(policy => {
|
||||
const option = document.createElement('option');
|
||||
@@ -671,7 +671,7 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
mode = mode || 'create';
|
||||
const adminCheckbox = document.getElementById(mode === 'edit' ? 'editBucketAdmin' : 'bucketAdmin');
|
||||
const permissionFields = document.getElementById(mode === 'edit' ? 'editBucketPermissionFields' : 'bucketPermissionFields');
|
||||
|
||||
|
||||
if (adminCheckbox && permissionFields) {
|
||||
permissionFields.style.display = adminCheckbox.checked ? 'none' : 'block';
|
||||
}
|
||||
@@ -682,7 +682,7 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
mode = mode || 'create';
|
||||
const specificRadio = document.getElementById(mode === 'edit' ? 'editSpecificBuckets' : 'specificBuckets');
|
||||
const bucketList = document.getElementById(mode === 'edit' ? 'editBucketSelectionList' : 'bucketSelectionList');
|
||||
|
||||
|
||||
if (specificRadio && bucketList) {
|
||||
bucketList.style.display = specificRadio.checked ? 'block' : 'none';
|
||||
}
|
||||
@@ -710,7 +710,7 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
function populateBucketSelections() {
|
||||
const createSelect = document.getElementById('selectedBuckets');
|
||||
const editSelect = document.getElementById('editSelectedBuckets');
|
||||
|
||||
|
||||
[createSelect, editSelect].forEach(select => {
|
||||
if (select) {
|
||||
select.innerHTML = '';
|
||||
@@ -732,17 +732,17 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
applyToAll: false,
|
||||
specificBuckets: []
|
||||
};
|
||||
|
||||
|
||||
// Check if user has Admin permission
|
||||
if (actions.includes('Admin')) {
|
||||
result.isAdmin = true;
|
||||
return result;
|
||||
}
|
||||
|
||||
|
||||
// Separate bucket-scoped from global actions
|
||||
const bucketActions = [];
|
||||
const globalBucketPerms = [];
|
||||
|
||||
|
||||
actions.forEach(action => {
|
||||
if (action.startsWith('s3tables:')) {
|
||||
const actionValue = action.slice('s3tables:'.length);
|
||||
@@ -767,7 +767,7 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
globalBucketPerms.push(action);
|
||||
}
|
||||
});
|
||||
|
||||
|
||||
// If we have global bucket permissions (no colon), they apply to all buckets
|
||||
if (globalBucketPerms.length > 0) {
|
||||
result.permissions = globalBucketPerms;
|
||||
@@ -776,12 +776,12 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
// Get unique permissions and buckets
|
||||
const perms = [...new Set(bucketActions.map(ba => ba.permission))];
|
||||
const buckets = [...new Set(bucketActions.map(ba => ba.bucketId))];
|
||||
|
||||
|
||||
result.permissions = perms;
|
||||
result.applyToAll = false;
|
||||
result.specificBuckets = buckets;
|
||||
}
|
||||
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
@@ -805,34 +805,34 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
mode = mode || 'create';
|
||||
const selectId = mode === 'edit' ? 'editActions' : 'actions';
|
||||
const permSelect = document.getElementById(selectId);
|
||||
|
||||
|
||||
if (!permSelect) return [];
|
||||
|
||||
|
||||
// Get selected permissions from the original multi-select
|
||||
const selectedPerms = Array.from(permSelect.selectedOptions).map(opt => opt.value);
|
||||
|
||||
|
||||
const hasAdmin = selectedPerms.includes('Admin');
|
||||
const hasS3TablesAdmin = selectedPerms.includes('S3TablesAdmin');
|
||||
|
||||
|
||||
if (selectedPerms.length === 0) {
|
||||
return [];
|
||||
}
|
||||
|
||||
|
||||
// Check if applying to all buckets or specific ones
|
||||
// Use querySelector to find the checked radio button by name group
|
||||
const scopeName = mode === 'edit' ? 'editBucketScope' : 'bucketScope';
|
||||
|
||||
|
||||
// Try multiple methods to find the checked radio
|
||||
let checkedRadio = document.querySelector(`input[name="${scopeName}"]:checked`);
|
||||
|
||||
|
||||
// Fallback: check both radio buttons explicitly
|
||||
if (!checkedRadio) {
|
||||
const allBucketsId = mode === 'edit' ? 'editAllBuckets' : 'allBuckets';
|
||||
const specificBucketsId = mode === 'edit' ? 'editSpecificBuckets' : 'specificBuckets';
|
||||
|
||||
|
||||
const allBucketsRadio = document.getElementById(allBucketsId);
|
||||
const specificBucketsRadio = document.getElementById(specificBucketsId);
|
||||
|
||||
|
||||
if (specificBucketsRadio && specificBucketsRadio.checked) {
|
||||
checkedRadio = specificBucketsRadio;
|
||||
} else if (allBucketsRadio && allBucketsRadio.checked) {
|
||||
@@ -867,18 +867,21 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
// Get selected specific buckets
|
||||
const bucketSelect = document.getElementById(mode === 'edit' ? 'editSelectedBuckets' : 'selectedBuckets');
|
||||
if (!bucketSelect) return null;
|
||||
|
||||
|
||||
const selectedBuckets = [...new Set(Array.from(bucketSelect.selectedOptions).map(opt => opt.value))];
|
||||
|
||||
|
||||
// Return null to signal validation failure if no buckets selected
|
||||
if (selectedBuckets.length === 0) {
|
||||
return null;
|
||||
}
|
||||
|
||||
|
||||
// Build bucket-scoped permissions
|
||||
const actions = [];
|
||||
if (hasAdmin) {
|
||||
actions.push('Admin');
|
||||
selectedBuckets.forEach((bucket) => {
|
||||
const bucketInfo = parseBucketOptionValue(bucket);
|
||||
actions.push("Admin:" + bucketInfo.name);
|
||||
});
|
||||
}
|
||||
if (hasS3TablesAdmin) {
|
||||
actions.push('s3tables:*');
|
||||
@@ -898,7 +901,7 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
|
||||
return [...new Set(actions)];
|
||||
}
|
||||
}
|
||||
@@ -929,11 +932,11 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
const response = await fetch(`/api/users/${encodedUsername}`);
|
||||
if (response.ok) {
|
||||
const user = await response.json();
|
||||
|
||||
|
||||
// Populate edit form
|
||||
document.getElementById('editUsername').value = username;
|
||||
document.getElementById('editEmail').value = user.email || '';
|
||||
|
||||
|
||||
// Set selected actions
|
||||
const actionsSelect = document.getElementById('editActions');
|
||||
Array.from(actionsSelect.options).forEach(option => {
|
||||
@@ -947,11 +950,11 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
option.selected = user.policy_names && user.policy_names.includes(option.value);
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
// Populate bucket permissions using original permissions dropdown
|
||||
if (user.actions && user.actions.length > 0) {
|
||||
const bucketPerms = parseBucketPermissions(user.actions);
|
||||
|
||||
|
||||
// Set permissions in the original multi-select
|
||||
const actionsSelect = document.getElementById('editActions');
|
||||
if (actionsSelect) {
|
||||
@@ -963,18 +966,18 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
// Set bucket scope (all or specific)
|
||||
const allBucketsRadio = document.getElementById('editAllBuckets');
|
||||
const specificBucketsRadio = document.getElementById('editSpecificBuckets');
|
||||
|
||||
|
||||
if (!bucketPerms.isAdmin) {
|
||||
if (bucketPerms.applyToAll) {
|
||||
if (allBucketsRadio) allBucketsRadio.checked = true;
|
||||
} else if (bucketPerms.specificBuckets.length > 0) {
|
||||
if (specificBucketsRadio) specificBucketsRadio.checked = true;
|
||||
toggleBucketList('edit');
|
||||
|
||||
|
||||
// Select specific buckets
|
||||
const bucketSelect = document.getElementById('editSelectedBuckets');
|
||||
if (bucketSelect) {
|
||||
@@ -985,7 +988,7 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
// Populate groups
|
||||
await populateEditUserGroups(username);
|
||||
|
||||
@@ -1029,7 +1032,7 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
const response = await fetch(`/api/users/${encodedUsername}`, {
|
||||
method: 'DELETE'
|
||||
});
|
||||
|
||||
|
||||
if (response.ok) {
|
||||
showSuccessMessage('User deleted successfully');
|
||||
setTimeout(() => window.location.reload(), 1000);
|
||||
@@ -1048,10 +1051,10 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
async function handleCreateUser() {
|
||||
const form = document.getElementById('createUserForm');
|
||||
const formData = new FormData(form);
|
||||
|
||||
|
||||
// Get permissions with bucket scope applied
|
||||
const allActions = buildBucketPermissions('create');
|
||||
|
||||
|
||||
if (allActions === null) {
|
||||
showAlert('Please select at least one bucket when using specific bucket permissions', 'error');
|
||||
return;
|
||||
@@ -1061,7 +1064,7 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
showAlert('At least one permission must be selected', 'error');
|
||||
return;
|
||||
}
|
||||
|
||||
|
||||
const userData = {
|
||||
username: formData.get('username'),
|
||||
email: formData.get('email'),
|
||||
@@ -1069,7 +1072,7 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
policy_names: Array.from(document.getElementById('policies').selectedOptions).map(option => option.value),
|
||||
generate_key: document.getElementById('generateKey').checked
|
||||
};
|
||||
|
||||
|
||||
try {
|
||||
const response = await fetch('/api/users', {
|
||||
method: 'POST',
|
||||
@@ -1078,16 +1081,16 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
},
|
||||
body: JSON.stringify(userData)
|
||||
});
|
||||
|
||||
|
||||
if (response.ok) {
|
||||
const result = await response.json();
|
||||
showSuccessMessage('User created successfully');
|
||||
|
||||
|
||||
// Show the created access key if generated
|
||||
if (result.user && result.user.access_key) {
|
||||
showNewAccessKeyModal(result.user);
|
||||
}
|
||||
|
||||
|
||||
// Close modal and refresh page
|
||||
const modal = bootstrap.Modal.getInstance(document.getElementById('createUserModal'));
|
||||
modal.hide();
|
||||
@@ -1210,28 +1213,28 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
showAlert('Username is required', 'error');
|
||||
return;
|
||||
}
|
||||
|
||||
|
||||
// Get permissions with bucket scope applied
|
||||
const allActions = buildBucketPermissions('edit');
|
||||
|
||||
|
||||
// Check for null (validation failure from buildBucketPermissions)
|
||||
if (allActions === null) {
|
||||
showAlert('Please select at least one bucket when using specific bucket permissions', 'error');
|
||||
return;
|
||||
}
|
||||
|
||||
|
||||
// Validate that permissions are not empty
|
||||
if (!allActions || allActions.length === 0) {
|
||||
showAlert('At least one permission must be selected', 'error');
|
||||
return;
|
||||
}
|
||||
|
||||
|
||||
const userData = {
|
||||
email: document.getElementById('editEmail').value,
|
||||
actions: allActions,
|
||||
policy_names: Array.from(document.getElementById('editPolicies').selectedOptions).map(option => option.value)
|
||||
};
|
||||
|
||||
|
||||
try {
|
||||
const encodedUsername = encodeURIComponent(username);
|
||||
const response = await fetch(`/api/users/${encodedUsername}`, {
|
||||
@@ -1241,10 +1244,10 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
},
|
||||
body: JSON.stringify(userData)
|
||||
});
|
||||
|
||||
|
||||
if (response.ok) {
|
||||
showSuccessMessage('User updated successfully');
|
||||
|
||||
|
||||
// Close modal and refresh page
|
||||
const modal = bootstrap.Modal.getInstance(document.getElementById('editUserModal'));
|
||||
modal.hide();
|
||||
@@ -1321,12 +1324,12 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
if (!user.access_keys || user.access_keys.length === 0) {
|
||||
return '<p class="text-muted">No access keys available</p>';
|
||||
}
|
||||
|
||||
|
||||
var keysHtml = '<div class="table-responsive">';
|
||||
keysHtml += '<table class="table table-sm">';
|
||||
keysHtml += '<thead><tr><th>Access Key</th><th>Status</th><th>Actions</th></tr></thead>';
|
||||
keysHtml += '<tbody>';
|
||||
|
||||
|
||||
user.access_keys.forEach(function(key) {
|
||||
keysHtml += '<tr>';
|
||||
keysHtml += '<td><code>' + escapeHtml(key.access_key) + '</code></td>';
|
||||
@@ -1348,11 +1351,11 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
keysHtml += '</td>';
|
||||
keysHtml += '</tr>';
|
||||
});
|
||||
|
||||
|
||||
keysHtml += '</tbody>';
|
||||
keysHtml += '</table>';
|
||||
keysHtml += '</div>';
|
||||
|
||||
|
||||
// Add delegated event listener for view secret buttons
|
||||
setTimeout(() => {
|
||||
document.querySelectorAll('.view-secret-btn').forEach(btn => {
|
||||
@@ -1363,7 +1366,7 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
});
|
||||
});
|
||||
}, 100);
|
||||
|
||||
|
||||
return keysHtml;
|
||||
}
|
||||
|
||||
@@ -1514,10 +1517,10 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
const response = await fetch(`/api/users/${encodedUsername}/access-keys/${encodedAccessKey}`, {
|
||||
method: 'DELETE'
|
||||
});
|
||||
|
||||
|
||||
if (response.ok) {
|
||||
showSuccessMessage('Access key deleted successfully');
|
||||
|
||||
|
||||
// Refresh access keys display
|
||||
refreshAccessKeysList(username);
|
||||
} else {
|
||||
@@ -1552,4 +1555,4 @@ templ ObjectStoreUsers(data dash.ObjectStoreUsersData) {
|
||||
}
|
||||
|
||||
// Helper functions for template
|
||||
|
||||
|
||||
|
||||
File diff suppressed because one or more lines are too long
@@ -266,6 +266,11 @@ templ Plugin(page string, initialJob string, lane string) {
|
||||
<label class="form-label" for="plugin-admin-detection-timeout">Detection Timeout (s)</label>
|
||||
<input type="number" class="form-control" id="plugin-admin-detection-timeout" min="0"/>
|
||||
</div>
|
||||
<div class="col-12">
|
||||
<label class="form-label" for="plugin-admin-execution-timeout">Execution Timeout (s)</label>
|
||||
<input type="number" class="form-control" id="plugin-admin-execution-timeout" min="0"/>
|
||||
<div class="form-text">Per-attempt deadline for one job. Size-aware tasks may extend this automatically.</div>
|
||||
</div>
|
||||
<div class="col-12">
|
||||
<label class="form-label" for="plugin-admin-max-runtime">Job Type Max Runtime (s)</label>
|
||||
<input type="number" class="form-control" id="plugin-admin-max-runtime" min="0"/>
|
||||
@@ -2519,6 +2524,7 @@ templ Plugin(page string, initialJob string, lane string) {
|
||||
document.getElementById('plugin-admin-enabled').checked = pickBool('enabled');
|
||||
document.getElementById('plugin-admin-detection-interval').value = String(pickNumber('detection_interval_seconds'));
|
||||
document.getElementById('plugin-admin-detection-timeout').value = String(pickNumber('detection_timeout_seconds'));
|
||||
document.getElementById('plugin-admin-execution-timeout').value = String(pickNumber('execution_timeout_seconds'));
|
||||
document.getElementById('plugin-admin-max-runtime').value = String(pickNumber('job_type_max_runtime_seconds'));
|
||||
document.getElementById('plugin-admin-max-results').value = String(pickNumber('max_jobs_per_detection'));
|
||||
document.getElementById('plugin-admin-global-exec').value = String(pickNumber('global_execution_concurrency'));
|
||||
@@ -2544,6 +2550,7 @@ templ Plugin(page string, initialJob string, lane string) {
|
||||
enabled: !!document.getElementById('plugin-admin-enabled').checked,
|
||||
detection_interval_seconds: getInt('plugin-admin-detection-interval'),
|
||||
detection_timeout_seconds: getInt('plugin-admin-detection-timeout'),
|
||||
execution_timeout_seconds: getInt('plugin-admin-execution-timeout'),
|
||||
job_type_max_runtime_seconds: getInt('plugin-admin-max-runtime'),
|
||||
max_jobs_per_detection: getInt('plugin-admin-max-results'),
|
||||
global_execution_concurrency: getInt('plugin-admin-global-exec'),
|
||||
|
||||
File diff suppressed because one or more lines are too long
@@ -20,10 +20,10 @@ var NoLockServerError = fmt.Errorf("no lock server found")
|
||||
type ReplicateFunc func(server pb.ServerAddress, key string, expiredAtNs int64, token string, owner string, generation int64, seq int64, isUnlock bool)
|
||||
|
||||
type DistributedLockManager struct {
|
||||
lockManager *LockManager
|
||||
LockRing *LockRing
|
||||
Host pb.ServerAddress
|
||||
ReplicateFn ReplicateFunc // set by filer server after creation
|
||||
lockManager *LockManager
|
||||
LockRing *LockRing
|
||||
Host pb.ServerAddress
|
||||
ReplicateFn ReplicateFunc // set by filer server after creation
|
||||
}
|
||||
|
||||
func NewDistributedLockManager(host pb.ServerAddress) *DistributedLockManager {
|
||||
|
||||
@@ -23,7 +23,7 @@ const DefaultVnodeCount = 50
|
||||
type HashRing struct {
|
||||
mu sync.RWMutex
|
||||
vnodeCount int
|
||||
sortedHashes []uint32 // sorted ring positions
|
||||
sortedHashes []uint32 // sorted ring positions
|
||||
vnodeToServer map[uint32]pb.ServerAddress // ring position → server
|
||||
servers map[pb.ServerAddress]struct{} // set of all servers
|
||||
}
|
||||
|
||||
@@ -18,18 +18,18 @@ var LockNotFound = fmt.Errorf("lock not found")
|
||||
|
||||
// LockManager local lock manager, used by distributed lock manager
|
||||
type LockManager struct {
|
||||
locks map[string]*Lock
|
||||
accessLock sync.RWMutex
|
||||
nextGeneration atomic.Int64
|
||||
locks map[string]*Lock
|
||||
accessLock sync.RWMutex
|
||||
nextGeneration atomic.Int64
|
||||
}
|
||||
type Lock struct {
|
||||
Token string
|
||||
ExpiredAtNs int64
|
||||
Key string // only used for moving locks
|
||||
Owner string
|
||||
IsBackup bool // true if this node holds the lock as a backup
|
||||
Generation int64 // monotonic fencing token, increments on fresh acquisition
|
||||
Seq int64 // per-lock sequence number, increments on every mutation (acquire/renew/unlock)
|
||||
IsBackup bool // true if this node holds the lock as a backup
|
||||
Generation int64 // monotonic fencing token, increments on fresh acquisition
|
||||
Seq int64 // per-lock sequence number, increments on every mutation (acquire/renew/unlock)
|
||||
}
|
||||
|
||||
func NewLockManager() *LockManager {
|
||||
|
||||
@@ -22,7 +22,7 @@ type LockRing struct {
|
||||
onTakeSnapshot func(snapshot []pb.ServerAddress)
|
||||
cleanupWg sync.WaitGroup
|
||||
Ring *HashRing // consistent hash ring
|
||||
version int64 // monotonic version from master, rejects stale updates
|
||||
version int64 // monotonic version from master, rejects stale updates
|
||||
}
|
||||
|
||||
func NewLockRing(snapshotInterval time.Duration) *LockRing {
|
||||
|
||||
@@ -17,13 +17,13 @@ const LockRingStabilizationInterval = 1 * time.Second
|
||||
// so filers receive a single consistent ring update instead of multiple
|
||||
// intermediate states.
|
||||
type LockRingManager struct {
|
||||
mu sync.Mutex
|
||||
members map[FilerGroupName]map[pb.ServerAddress]struct{}
|
||||
version map[FilerGroupName]int64
|
||||
lastBroadcast map[FilerGroupName]*master_pb.LockRingUpdate
|
||||
pendingTimer map[FilerGroupName]*time.Timer
|
||||
broadcastFn func(resp *master_pb.KeepConnectedResponse)
|
||||
stabilizeDelay time.Duration
|
||||
mu sync.Mutex
|
||||
members map[FilerGroupName]map[pb.ServerAddress]struct{}
|
||||
version map[FilerGroupName]int64
|
||||
lastBroadcast map[FilerGroupName]*master_pb.LockRingUpdate
|
||||
pendingTimer map[FilerGroupName]*time.Timer
|
||||
broadcastFn func(resp *master_pb.KeepConnectedResponse)
|
||||
stabilizeDelay time.Duration
|
||||
}
|
||||
|
||||
func NewLockRingManager(broadcastFn func(resp *master_pb.KeepConnectedResponse)) *LockRingManager {
|
||||
|
||||
+81
-10
@@ -11,6 +11,7 @@ import (
|
||||
"runtime"
|
||||
"sort"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/spf13/viper"
|
||||
@@ -31,6 +32,7 @@ import (
|
||||
weed_server "github.com/seaweedfs/seaweedfs/weed/server"
|
||||
stats_collect "github.com/seaweedfs/seaweedfs/weed/stats"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util/grace"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util/version"
|
||||
)
|
||||
|
||||
@@ -381,6 +383,15 @@ func (fo *FilerOptions) startFiler() {
|
||||
glog.Fatalf("Filer startup error: %v", nfs_err)
|
||||
}
|
||||
|
||||
// Ensure fs.Shutdown() runs exactly once, whether triggered by a signal hook
|
||||
// or by the main goroutine after Serve() returns (e.g., MiniCluster tests).
|
||||
var shutdownOnce sync.Once
|
||||
shutdownFiler := func() {
|
||||
shutdownOnce.Do(func() {
|
||||
fs.Shutdown()
|
||||
})
|
||||
}
|
||||
|
||||
if *fo.publicPort != 0 {
|
||||
publicListeningAddress := util.JoinHostPort(*fo.bindIp, *fo.publicPort)
|
||||
glog.V(0).Infoln("Start Seaweed filer server", version.Version(), "public at", publicListeningAddress)
|
||||
@@ -434,6 +445,24 @@ func (fo *FilerOptions) startFiler() {
|
||||
go grpcS.Serve(grpcL)
|
||||
pb.ServeGrpcOnLocalSocket(grpcS, grpcPort)
|
||||
|
||||
// Register graceful shutdown for gRPC server to wait for active RPCs
|
||||
grace.OnInterrupt(func() {
|
||||
glog.V(0).Infof("Gracefully stopping gRPC server")
|
||||
stopped := make(chan struct{})
|
||||
go func() {
|
||||
grpcS.GracefulStop()
|
||||
close(stopped)
|
||||
}()
|
||||
select {
|
||||
case <-stopped:
|
||||
glog.V(0).Infof("gRPC server stopped gracefully")
|
||||
case <-time.After(10 * time.Second):
|
||||
glog.V(0).Infof("gRPC server graceful stop timed out, forcing stop")
|
||||
grpcS.Stop()
|
||||
}
|
||||
})
|
||||
|
||||
var socketServer *http.Server
|
||||
if runtime.GOOS != "windows" {
|
||||
localSocket := *fo.localSocket
|
||||
if localSocket == "" {
|
||||
@@ -442,14 +471,12 @@ func (fo *FilerOptions) startFiler() {
|
||||
if err := os.Remove(localSocket); err != nil && !os.IsNotExist(err) {
|
||||
glog.Fatalf("Failed to remove %s, error: %s", localSocket, err.Error())
|
||||
}
|
||||
go func() {
|
||||
// start on local unix socket
|
||||
filerSocketListener, err := net.Listen("unix", localSocket)
|
||||
if err != nil {
|
||||
glog.Fatalf("Failed to listen on %s: %v", localSocket, err)
|
||||
}
|
||||
newHttpServer(defaultMux, nil).Serve(filerSocketListener)
|
||||
}()
|
||||
filerSocketListener, err := net.Listen("unix", localSocket)
|
||||
if err != nil {
|
||||
glog.Fatalf("Failed to listen on %s: %v", localSocket, err)
|
||||
}
|
||||
socketServer = newHttpServer(defaultMux, nil)
|
||||
go socketServer.Serve(filerSocketListener)
|
||||
}
|
||||
|
||||
if viper.GetString("https.filer.key") != "" {
|
||||
@@ -489,14 +516,34 @@ func (fo *FilerOptions) startFiler() {
|
||||
|
||||
security.FixTlsConfig(util.GetViper(), tlsConfig)
|
||||
|
||||
var localTLSServer *http.Server
|
||||
if filerLocalListener != nil {
|
||||
localTLSServer = newHttpServer(defaultMux, tlsConfig)
|
||||
go func() {
|
||||
if err := newHttpServer(defaultMux, tlsConfig).ServeTLS(filerLocalListener, "", ""); err != nil {
|
||||
if err := localTLSServer.ServeTLS(filerLocalListener, "", ""); err != nil {
|
||||
glog.Errorf("Filer Fail to serve: %v", err)
|
||||
}
|
||||
}()
|
||||
}
|
||||
httpS := newHttpServer(defaultMux, tlsConfig)
|
||||
|
||||
// Register shutdown hooks: stop all HTTP servers, then close filer database
|
||||
grace.OnInterrupt(func() {
|
||||
glog.V(0).Infof("Gracefully stopping all HTTP servers")
|
||||
shutdownCtx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
|
||||
defer cancel()
|
||||
if socketServer != nil {
|
||||
socketServer.Shutdown(shutdownCtx)
|
||||
}
|
||||
if localTLSServer != nil {
|
||||
localTLSServer.Shutdown(shutdownCtx)
|
||||
}
|
||||
if err := httpS.Shutdown(shutdownCtx); err != nil {
|
||||
glog.Warningf("HTTPS server shutdown: %v", err)
|
||||
}
|
||||
})
|
||||
grace.OnInterrupt(shutdownFiler)
|
||||
|
||||
if MiniClusterCtx != nil {
|
||||
ctx := MiniClusterCtx
|
||||
go func() {
|
||||
@@ -508,15 +555,37 @@ func (fo *FilerOptions) startFiler() {
|
||||
if err := httpS.ServeTLS(filerListener, "", ""); err != nil && err != http.ErrServerClosed {
|
||||
glog.Fatalf("Filer Fail to serve: %v", err)
|
||||
}
|
||||
// Close database after servers have stopped to prevent data corruption
|
||||
shutdownFiler()
|
||||
} else {
|
||||
var localHTTPServer *http.Server
|
||||
if filerLocalListener != nil {
|
||||
localHTTPServer = newHttpServer(defaultMux, nil)
|
||||
go func() {
|
||||
if err := newHttpServer(defaultMux, nil).Serve(filerLocalListener); err != nil {
|
||||
if err := localHTTPServer.Serve(filerLocalListener); err != nil {
|
||||
glog.Errorf("Filer Fail to serve: %v", err)
|
||||
}
|
||||
}()
|
||||
}
|
||||
httpS := newHttpServer(defaultMux, nil)
|
||||
|
||||
// Register shutdown hooks: stop all HTTP servers, then close filer database
|
||||
grace.OnInterrupt(func() {
|
||||
glog.V(0).Infof("Gracefully stopping all HTTP servers")
|
||||
shutdownCtx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
|
||||
defer cancel()
|
||||
if socketServer != nil {
|
||||
socketServer.Shutdown(shutdownCtx)
|
||||
}
|
||||
if localHTTPServer != nil {
|
||||
localHTTPServer.Shutdown(shutdownCtx)
|
||||
}
|
||||
if err := httpS.Shutdown(shutdownCtx); err != nil {
|
||||
glog.Warningf("HTTP server shutdown: %v", err)
|
||||
}
|
||||
})
|
||||
grace.OnInterrupt(shutdownFiler)
|
||||
|
||||
if MiniClusterCtx != nil {
|
||||
ctx := MiniClusterCtx
|
||||
go func() {
|
||||
@@ -528,5 +597,7 @@ func (fo *FilerOptions) startFiler() {
|
||||
if err := httpS.Serve(filerListener); err != nil && err != http.ErrServerClosed {
|
||||
glog.Fatalf("Filer Fail to serve: %v", err)
|
||||
}
|
||||
// Close database after servers have stopped to prevent data corruption
|
||||
shutdownFiler()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -11,11 +11,11 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/glog"
|
||||
"github.com/seaweedfs/seaweedfs/weed/operation"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/replication"
|
||||
"github.com/seaweedfs/seaweedfs/weed/replication/sink"
|
||||
"github.com/seaweedfs/seaweedfs/weed/operation"
|
||||
"github.com/seaweedfs/seaweedfs/weed/replication/sink/filersink"
|
||||
"github.com/seaweedfs/seaweedfs/weed/replication/source"
|
||||
"github.com/seaweedfs/seaweedfs/weed/security"
|
||||
|
||||
@@ -15,10 +15,10 @@ import (
|
||||
// tsMinHeap implements heap.Interface for int64 timestamps.
|
||||
type tsMinHeap []int64
|
||||
|
||||
func (h tsMinHeap) Len() int { return len(h) }
|
||||
func (h tsMinHeap) Less(i, j int) bool { return h[i] < h[j] }
|
||||
func (h tsMinHeap) Swap(i, j int) { h[i], h[j] = h[j], h[i] }
|
||||
func (h *tsMinHeap) Push(x any) { *h = append(*h, x.(int64)) }
|
||||
func (h tsMinHeap) Len() int { return len(h) }
|
||||
func (h tsMinHeap) Less(i, j int) bool { return h[i] < h[j] }
|
||||
func (h tsMinHeap) Swap(i, j int) { h[i], h[j] = h[j], h[i] }
|
||||
func (h *tsMinHeap) Push(x any) { *h = append(*h, x.(int64)) }
|
||||
func (h *tsMinHeap) Pop() any {
|
||||
old := *h
|
||||
n := len(old)
|
||||
@@ -34,11 +34,11 @@ type syncJobPaths struct {
|
||||
}
|
||||
|
||||
type MetadataProcessor struct {
|
||||
activeJobs map[int64]*syncJobPaths
|
||||
activeJobsLock sync.Mutex
|
||||
activeJobsCond *sync.Cond
|
||||
concurrencyLimit int
|
||||
fn pb.ProcessMetadataFunc
|
||||
activeJobs map[int64]*syncJobPaths
|
||||
activeJobsLock sync.Mutex
|
||||
activeJobsCond *sync.Cond
|
||||
concurrencyLimit int
|
||||
fn pb.ProcessMetadataFunc
|
||||
processedTsWatermark atomic.Int64
|
||||
|
||||
// Indexes for O(depth) conflict detection, replacing O(n) linear scan.
|
||||
|
||||
@@ -11,8 +11,8 @@ import (
|
||||
|
||||
func makeResp(dir, name string, isDir bool, tsNs int64, isNew bool) *filer_pb.SubscribeMetadataResponse {
|
||||
resp := &filer_pb.SubscribeMetadataResponse{
|
||||
Directory: dir,
|
||||
TsNs: tsNs,
|
||||
Directory: dir,
|
||||
TsNs: tsNs,
|
||||
EventNotification: &filer_pb.EventNotification{},
|
||||
}
|
||||
entry := &filer_pb.Entry{
|
||||
|
||||
+39
-31
@@ -6,38 +6,39 @@ import (
|
||||
)
|
||||
|
||||
type MountOptions struct {
|
||||
filer *string
|
||||
filerMountRootPath *string
|
||||
dir *string
|
||||
dirAutoCreate *bool
|
||||
collection *string
|
||||
collectionQuota *int
|
||||
replication *string
|
||||
diskType *string
|
||||
ttlSec *int
|
||||
chunkSizeLimitMB *int
|
||||
concurrentWriters *int
|
||||
concurrentReaders *int
|
||||
cacheMetaTtlSec *int
|
||||
cacheDirForRead *string
|
||||
cacheDirForWrite *string
|
||||
cacheSizeMBForRead *int64
|
||||
dataCenter *string
|
||||
allowOthers *bool
|
||||
defaultPermissions *bool
|
||||
umaskString *string
|
||||
nonempty *bool
|
||||
volumeServerAccess *string
|
||||
uidMap *string
|
||||
gidMap *string
|
||||
readOnly *bool
|
||||
filer *string
|
||||
filerMountRootPath *string
|
||||
dir *string
|
||||
dirAutoCreate *bool
|
||||
collection *string
|
||||
collectionQuota *int
|
||||
replication *string
|
||||
diskType *string
|
||||
ttlSec *int
|
||||
chunkSizeLimitMB *int
|
||||
concurrentWriters *int
|
||||
concurrentReaders *int
|
||||
cacheMetaTtlSec *int
|
||||
cacheDirForRead *string
|
||||
cacheDirForWrite *string
|
||||
cacheSizeMBForRead *int64
|
||||
dataCenter *string
|
||||
allowOthers *bool
|
||||
defaultPermissions *bool
|
||||
umaskString *string
|
||||
nonempty *bool
|
||||
volumeServerAccess *string
|
||||
uidMap *string
|
||||
gidMap *string
|
||||
readOnly *bool
|
||||
includeSystemEntries *bool
|
||||
debug *bool
|
||||
debugPort *int
|
||||
localSocket *string
|
||||
disableXAttr *bool
|
||||
extraOptions []string
|
||||
fuseCommandPid int
|
||||
debug *bool
|
||||
debugPort *int
|
||||
debugFuse *bool
|
||||
localSocket *string
|
||||
disableXAttr *bool
|
||||
extraOptions []string
|
||||
fuseCommandPid int
|
||||
|
||||
// Periodic metadata flush to protect against orphan chunk cleanup
|
||||
metadataFlushSeconds *int
|
||||
@@ -55,6 +56,9 @@ type MountOptions struct {
|
||||
// Distributed lock for cross-mount write coordination
|
||||
distributedLock *bool
|
||||
|
||||
// POSIX compliance options
|
||||
posixDirNlink *bool
|
||||
|
||||
// FUSE performance options
|
||||
writebackCache *bool
|
||||
asyncDio *bool
|
||||
@@ -106,6 +110,7 @@ func init() {
|
||||
mountOptions.includeSystemEntries = cmdMount.Flag.Bool("includeSystemEntries", false, "show filer system entries (e.g. /topics, /etc) in directory listings")
|
||||
mountOptions.debug = cmdMount.Flag.Bool("debug", false, "serves runtime profiling data, e.g., http://localhost:<debug.port>/debug/pprof/goroutine?debug=2")
|
||||
mountOptions.debugPort = cmdMount.Flag.Int("debug.port", 6061, "http port for debugging")
|
||||
mountOptions.debugFuse = cmdMount.Flag.Bool("debug.fuse", false, "log raw FUSE protocol requests and responses")
|
||||
mountOptions.localSocket = cmdMount.Flag.String("localSocket", "", "default to /tmp/seaweedfs-mount-<mount_dir_hash>.sock")
|
||||
mountOptions.disableXAttr = cmdMount.Flag.Bool("disableXAttr", false, "disable xattr")
|
||||
mountOptions.hasAutofs = cmdMount.Flag.Bool("autofs", false, "ignore autofs mounted on the same mountpoint (useful when systemd.automount and autofs is used)")
|
||||
@@ -131,6 +136,9 @@ func init() {
|
||||
// Distributed lock for cross-mount write coordination
|
||||
mountOptions.distributedLock = cmdMount.Flag.Bool("dlm", false, "enable distributed lock for cross-mount write coordination (only one mount can write a file at a time)")
|
||||
|
||||
// POSIX compliance options
|
||||
mountOptions.posixDirNlink = cmdMount.Flag.Bool("posix.dirNLink", false, "report POSIX-compliant directory nlink (2 + subdirectory count); costs one directory listing per stat")
|
||||
|
||||
// FUSE performance options
|
||||
mountOptions.writebackCache = cmdMount.Flag.Bool("writebackCache", false, "enable FUSE writeback cache for improved write performance (at risk of data loss on crash)")
|
||||
mountOptions.asyncDio = cmdMount.Flag.Bool("asyncDio", false, "enable async direct I/O for better concurrency")
|
||||
|
||||
@@ -251,7 +251,7 @@ func RunMount(option *MountOptions, umask os.FileMode) bool {
|
||||
Name: "seaweedfs",
|
||||
SingleThreaded: false,
|
||||
DisableXAttrs: *option.disableXAttr,
|
||||
Debug: *option.debug,
|
||||
Debug: *option.debugFuse,
|
||||
EnableLocks: true,
|
||||
ExplicitDataCacheControl: false,
|
||||
DirectMount: true,
|
||||
@@ -345,15 +345,16 @@ func RunMount(option *MountOptions, umask os.FileMode) bool {
|
||||
IsMacOs: runtime.GOOS == "darwin",
|
||||
MetadataFlushSeconds: *option.metadataFlushSeconds,
|
||||
// RDMA acceleration options
|
||||
RdmaEnabled: *option.rdmaEnabled,
|
||||
RdmaSidecarAddr: *option.rdmaSidecarAddr,
|
||||
RdmaFallback: *option.rdmaFallback,
|
||||
RdmaReadOnly: *option.rdmaReadOnly,
|
||||
RdmaMaxConcurrent: *option.rdmaMaxConcurrent,
|
||||
RdmaTimeoutMs: *option.rdmaTimeoutMs,
|
||||
DirIdleEvictSec: *option.dirIdleEvictSec,
|
||||
RdmaEnabled: *option.rdmaEnabled,
|
||||
RdmaSidecarAddr: *option.rdmaSidecarAddr,
|
||||
RdmaFallback: *option.rdmaFallback,
|
||||
RdmaReadOnly: *option.rdmaReadOnly,
|
||||
RdmaMaxConcurrent: *option.rdmaMaxConcurrent,
|
||||
RdmaTimeoutMs: *option.rdmaTimeoutMs,
|
||||
DirIdleEvictSec: *option.dirIdleEvictSec,
|
||||
EnableDistributedLock: option.distributedLock != nil && *option.distributedLock,
|
||||
WritebackCache: option.writebackCache != nil && *option.writebackCache,
|
||||
WritebackCache: option.writebackCache != nil && *option.writebackCache,
|
||||
PosixDirNlink: option.posixDirNlink != nil && *option.posixDirNlink,
|
||||
})
|
||||
|
||||
// create mount root
|
||||
|
||||
@@ -15,7 +15,7 @@ var (
|
||||
shellOptions shell.ShellOptions
|
||||
shellInitialFiler *string
|
||||
shellCluster *string
|
||||
shellDebug *bool
|
||||
shellDebug *bool
|
||||
)
|
||||
|
||||
func init() {
|
||||
|
||||
@@ -271,6 +271,40 @@ func (cm *CredentialManager) UpdatePolicy(ctx context.Context, name string, docu
|
||||
return cm.Store.PutPolicy(ctx, name, document)
|
||||
}
|
||||
|
||||
// PutUserInlinePolicy stores a per-user inline policy document.
|
||||
// Returns nil without error if the underlying store does not support inline policies.
|
||||
func (cm *CredentialManager) PutUserInlinePolicy(ctx context.Context, userName, policyName string, document policy_engine.PolicyDocument) error {
|
||||
if store, ok := cm.Store.(InlinePolicyStore); ok {
|
||||
return store.PutUserInlinePolicy(ctx, userName, policyName, document)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// GetUserInlinePolicy retrieves a per-user inline policy document.
|
||||
// Returns nil without error if the underlying store does not support inline policies.
|
||||
func (cm *CredentialManager) GetUserInlinePolicy(ctx context.Context, userName, policyName string) (*policy_engine.PolicyDocument, error) {
|
||||
if store, ok := cm.Store.(InlinePolicyStore); ok {
|
||||
return store.GetUserInlinePolicy(ctx, userName, policyName)
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
// DeleteUserInlinePolicy removes a per-user inline policy document.
|
||||
func (cm *CredentialManager) DeleteUserInlinePolicy(ctx context.Context, userName, policyName string) error {
|
||||
if store, ok := cm.Store.(InlinePolicyStore); ok {
|
||||
return store.DeleteUserInlinePolicy(ctx, userName, policyName)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// ListUserInlinePolicies returns the names of all inline policies for a user.
|
||||
func (cm *CredentialManager) ListUserInlinePolicies(ctx context.Context, userName string) ([]string, error) {
|
||||
if store, ok := cm.Store.(InlinePolicyStore); ok {
|
||||
return store.ListUserInlinePolicies(ctx, userName)
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
// LoadS3ConfigFile reads a static S3 identity config file and registers
|
||||
// the identities so they appear in LoadConfiguration and listing results.
|
||||
func (cm *CredentialManager) LoadS3ConfigFile(path string) error {
|
||||
|
||||
@@ -137,5 +137,16 @@ type PolicyManager interface {
|
||||
GetPolicy(ctx context.Context, name string) (*policy_engine.PolicyDocument, error)
|
||||
}
|
||||
|
||||
// InlinePolicyStore is an optional interface for credential stores that support
|
||||
// per-user inline policy storage. Stores that implement this interface preserve
|
||||
// the exact policy document submitted via PutUserPolicy, enabling lossless
|
||||
// round-trips through GetUserPolicy.
|
||||
type InlinePolicyStore interface {
|
||||
PutUserInlinePolicy(ctx context.Context, userName, policyName string, document policy_engine.PolicyDocument) error
|
||||
GetUserInlinePolicy(ctx context.Context, userName, policyName string) (*policy_engine.PolicyDocument, error)
|
||||
DeleteUserInlinePolicy(ctx context.Context, userName, policyName string) error
|
||||
ListUserInlinePolicies(ctx context.Context, userName string) ([]string, error)
|
||||
}
|
||||
|
||||
// Stores holds all available credential store implementations
|
||||
var Stores []CredentialStore
|
||||
|
||||
@@ -316,6 +316,69 @@ func (store *FilerEtcStore) GetPolicy(ctx context.Context, name string) (*policy
|
||||
return nil, nil // Policy not found
|
||||
}
|
||||
|
||||
// PutUserInlinePolicy stores a per-user inline policy document.
|
||||
func (store *FilerEtcStore) PutUserInlinePolicy(ctx context.Context, userName, policyName string, document policy_engine.PolicyDocument) error {
|
||||
store.policyMu.Lock()
|
||||
defer store.policyMu.Unlock()
|
||||
|
||||
policiesCollection, _, err := store.loadLegacyPoliciesCollection(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if policiesCollection.InlinePolicies[userName] == nil {
|
||||
policiesCollection.InlinePolicies[userName] = make(map[string]policy_engine.PolicyDocument)
|
||||
}
|
||||
policiesCollection.InlinePolicies[userName][policyName] = document
|
||||
return store.saveLegacyPoliciesCollection(ctx, policiesCollection)
|
||||
}
|
||||
|
||||
// GetUserInlinePolicy retrieves a per-user inline policy document.
|
||||
func (store *FilerEtcStore) GetUserInlinePolicy(ctx context.Context, userName, policyName string) (*policy_engine.PolicyDocument, error) {
|
||||
policiesCollection, _, err := store.loadLegacyPoliciesCollection(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if userPolicies := policiesCollection.InlinePolicies[userName]; userPolicies != nil {
|
||||
if doc, exists := userPolicies[policyName]; exists {
|
||||
return &doc, nil
|
||||
}
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
// DeleteUserInlinePolicy removes a per-user inline policy document.
|
||||
func (store *FilerEtcStore) DeleteUserInlinePolicy(ctx context.Context, userName, policyName string) error {
|
||||
store.policyMu.Lock()
|
||||
defer store.policyMu.Unlock()
|
||||
|
||||
policiesCollection, _, err := store.loadLegacyPoliciesCollection(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if userPolicies := policiesCollection.InlinePolicies[userName]; userPolicies != nil {
|
||||
delete(userPolicies, policyName)
|
||||
if len(userPolicies) == 0 {
|
||||
delete(policiesCollection.InlinePolicies, userName)
|
||||
}
|
||||
}
|
||||
return store.saveLegacyPoliciesCollection(ctx, policiesCollection)
|
||||
}
|
||||
|
||||
// ListUserInlinePolicies returns the names of all inline policies for a user.
|
||||
func (store *FilerEtcStore) ListUserInlinePolicies(ctx context.Context, userName string) ([]string, error) {
|
||||
policiesCollection, _, err := store.loadLegacyPoliciesCollection(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
userPolicies := policiesCollection.InlinePolicies[userName]
|
||||
names := make([]string, 0, len(userPolicies))
|
||||
for name := range userPolicies {
|
||||
names = append(names, name)
|
||||
}
|
||||
return names, nil
|
||||
}
|
||||
|
||||
// ListPolicyNames returns all managed policy names stored in the filer.
|
||||
func (store *FilerEtcStore) ListPolicyNames(ctx context.Context) ([]string, error) {
|
||||
names := make([]string, 0)
|
||||
|
||||
@@ -105,3 +105,71 @@ func (store *MemoryStore) DeletePolicy(ctx context.Context, name string) error {
|
||||
delete(store.policies, name)
|
||||
return nil
|
||||
}
|
||||
|
||||
// PutUserInlinePolicy stores a per-user inline policy document.
|
||||
func (store *MemoryStore) PutUserInlinePolicy(ctx context.Context, userName, policyName string, document policy_engine.PolicyDocument) error {
|
||||
store.mu.Lock()
|
||||
defer store.mu.Unlock()
|
||||
|
||||
if !store.initialized {
|
||||
return fmt.Errorf("store not initialized")
|
||||
}
|
||||
|
||||
if store.inlinePolicies[userName] == nil {
|
||||
store.inlinePolicies[userName] = make(map[string]policy_engine.PolicyDocument)
|
||||
}
|
||||
store.inlinePolicies[userName][policyName] = document
|
||||
return nil
|
||||
}
|
||||
|
||||
// GetUserInlinePolicy retrieves a per-user inline policy document.
|
||||
func (store *MemoryStore) GetUserInlinePolicy(ctx context.Context, userName, policyName string) (*policy_engine.PolicyDocument, error) {
|
||||
store.mu.RLock()
|
||||
defer store.mu.RUnlock()
|
||||
|
||||
if !store.initialized {
|
||||
return nil, fmt.Errorf("store not initialized")
|
||||
}
|
||||
|
||||
if userPolicies := store.inlinePolicies[userName]; userPolicies != nil {
|
||||
if doc, exists := userPolicies[policyName]; exists {
|
||||
return &doc, nil
|
||||
}
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
// DeleteUserInlinePolicy removes a per-user inline policy document.
|
||||
func (store *MemoryStore) DeleteUserInlinePolicy(ctx context.Context, userName, policyName string) error {
|
||||
store.mu.Lock()
|
||||
defer store.mu.Unlock()
|
||||
|
||||
if !store.initialized {
|
||||
return fmt.Errorf("store not initialized")
|
||||
}
|
||||
|
||||
if userPolicies := store.inlinePolicies[userName]; userPolicies != nil {
|
||||
delete(userPolicies, policyName)
|
||||
if len(userPolicies) == 0 {
|
||||
delete(store.inlinePolicies, userName)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// ListUserInlinePolicies returns the names of all inline policies for a user.
|
||||
func (store *MemoryStore) ListUserInlinePolicies(ctx context.Context, userName string) ([]string, error) {
|
||||
store.mu.RLock()
|
||||
defer store.mu.RUnlock()
|
||||
|
||||
if !store.initialized {
|
||||
return nil, fmt.Errorf("store not initialized")
|
||||
}
|
||||
|
||||
userPolicies := store.inlinePolicies[userName]
|
||||
names := make([]string, 0, len(userPolicies))
|
||||
for name := range userPolicies {
|
||||
names = append(names, name)
|
||||
}
|
||||
return names, nil
|
||||
}
|
||||
|
||||
@@ -18,12 +18,13 @@ func init() {
|
||||
// This is primarily intended for testing purposes
|
||||
type MemoryStore struct {
|
||||
mu sync.RWMutex
|
||||
users map[string]*iam_pb.Identity // username -> identity
|
||||
accessKeys map[string]string // access_key -> username
|
||||
serviceAccounts map[string]*iam_pb.ServiceAccount // id -> service_account
|
||||
serviceAccountAccessKeys map[string]string // access_key -> id
|
||||
policies map[string]policy_engine.PolicyDocument // policy_name -> policy_document
|
||||
groups map[string]*iam_pb.Group // group_name -> group
|
||||
users map[string]*iam_pb.Identity // username -> identity
|
||||
accessKeys map[string]string // access_key -> username
|
||||
serviceAccounts map[string]*iam_pb.ServiceAccount // id -> service_account
|
||||
serviceAccountAccessKeys map[string]string // access_key -> id
|
||||
policies map[string]policy_engine.PolicyDocument // policy_name -> policy_document
|
||||
inlinePolicies map[string]map[string]policy_engine.PolicyDocument // username -> policy_name -> document
|
||||
groups map[string]*iam_pb.Group // group_name -> group
|
||||
initialized bool
|
||||
}
|
||||
|
||||
@@ -44,6 +45,7 @@ func (store *MemoryStore) Initialize(configuration util.Configuration, prefix st
|
||||
store.serviceAccounts = make(map[string]*iam_pb.ServiceAccount)
|
||||
store.serviceAccountAccessKeys = make(map[string]string)
|
||||
store.policies = make(map[string]policy_engine.PolicyDocument)
|
||||
store.inlinePolicies = make(map[string]map[string]policy_engine.PolicyDocument)
|
||||
store.groups = make(map[string]*iam_pb.Group)
|
||||
store.initialized = true
|
||||
|
||||
@@ -59,6 +61,7 @@ func (store *MemoryStore) Shutdown() {
|
||||
store.serviceAccounts = nil
|
||||
store.serviceAccountAccessKeys = nil
|
||||
store.policies = nil
|
||||
store.inlinePolicies = nil
|
||||
store.groups = nil
|
||||
store.initialized = false
|
||||
}
|
||||
@@ -74,6 +77,7 @@ func (store *MemoryStore) Reset() {
|
||||
store.serviceAccounts = make(map[string]*iam_pb.ServiceAccount)
|
||||
store.serviceAccountAccessKeys = make(map[string]string)
|
||||
store.policies = make(map[string]policy_engine.PolicyDocument)
|
||||
store.inlinePolicies = make(map[string]map[string]policy_engine.PolicyDocument)
|
||||
store.groups = make(map[string]*iam_pb.Group)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -280,6 +280,34 @@ func (s *PropagatingCredentialStore) LoadInlinePolicies(ctx context.Context) (ma
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
func (s *PropagatingCredentialStore) PutUserInlinePolicy(ctx context.Context, userName, policyName string, document policy_engine.PolicyDocument) error {
|
||||
if store, ok := s.CredentialStore.(InlinePolicyStore); ok {
|
||||
return store.PutUserInlinePolicy(ctx, userName, policyName, document)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *PropagatingCredentialStore) GetUserInlinePolicy(ctx context.Context, userName, policyName string) (*policy_engine.PolicyDocument, error) {
|
||||
if store, ok := s.CredentialStore.(InlinePolicyStore); ok {
|
||||
return store.GetUserInlinePolicy(ctx, userName, policyName)
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
func (s *PropagatingCredentialStore) DeleteUserInlinePolicy(ctx context.Context, userName, policyName string) error {
|
||||
if store, ok := s.CredentialStore.(InlinePolicyStore); ok {
|
||||
return store.DeleteUserInlinePolicy(ctx, userName, policyName)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *PropagatingCredentialStore) ListUserInlinePolicies(ctx context.Context, userName string) ([]string, error) {
|
||||
if store, ok := s.CredentialStore.(InlinePolicyStore); ok {
|
||||
return store.ListUserInlinePolicies(ctx, userName)
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
func (s *PropagatingCredentialStore) CreatePolicy(ctx context.Context, name string, document policy_engine.PolicyDocument) error {
|
||||
if pm, ok := s.CredentialStore.(PolicyManager); ok {
|
||||
if err := pm.CreatePolicy(ctx, name, document); err != nil {
|
||||
|
||||
@@ -12,10 +12,10 @@ import (
|
||||
)
|
||||
|
||||
type mockFilerOps struct {
|
||||
countFn func(path util.FullPath) (int, error)
|
||||
deleteFn func(path util.FullPath) error
|
||||
attrsFn func(path util.FullPath) (map[string][]byte, error)
|
||||
isDirKeyObjFn func(path util.FullPath) (bool, error)
|
||||
countFn func(path util.FullPath) (int, error)
|
||||
deleteFn func(path util.FullPath) error
|
||||
attrsFn func(path util.FullPath) (map[string][]byte, error)
|
||||
isDirKeyObjFn func(path util.FullPath) (bool, error)
|
||||
}
|
||||
|
||||
func (m *mockFilerOps) CountDirectoryEntries(_ context.Context, dirPath util.FullPath, _ int) (int, error) {
|
||||
|
||||
+29
-23
@@ -40,28 +40,28 @@ var (
|
||||
)
|
||||
|
||||
type Filer struct {
|
||||
UniqueFilerId int32
|
||||
UniqueFilerEpoch int32
|
||||
Store VirtualFilerStore
|
||||
MasterClient *wdclient.MasterClient
|
||||
fileIdDeletionQueue *util.UnboundedQueue
|
||||
GrpcDialOption grpc.DialOption
|
||||
DirBucketsPath string
|
||||
Cipher bool
|
||||
LocalMetaLogBuffer *log_buffer.LogBuffer
|
||||
metaLogCollection string
|
||||
metaLogReplication string
|
||||
MetaAggregator *MetaAggregator
|
||||
Signature int32
|
||||
FilerConf *FilerConf
|
||||
RemoteStorage *FilerRemoteStorage
|
||||
lazyFetchGroup singleflight.Group
|
||||
lazyListGroup singleflight.Group
|
||||
Dlm *lock_manager.DistributedLockManager
|
||||
MaxFilenameLength uint32
|
||||
deletionQuit chan struct{}
|
||||
DeletionRetryQueue *DeletionRetryQueue
|
||||
EmptyFolderCleaner *empty_folder_cleanup.EmptyFolderCleaner
|
||||
UniqueFilerId int32
|
||||
UniqueFilerEpoch int32
|
||||
Store VirtualFilerStore
|
||||
MasterClient *wdclient.MasterClient
|
||||
fileIdDeletionQueue *util.UnboundedQueue
|
||||
GrpcDialOption grpc.DialOption
|
||||
DirBucketsPath string
|
||||
Cipher bool
|
||||
LocalMetaLogBuffer *log_buffer.LogBuffer
|
||||
metaLogCollection string
|
||||
metaLogReplication string
|
||||
MetaAggregator *MetaAggregator
|
||||
Signature int32
|
||||
FilerConf *FilerConf
|
||||
RemoteStorage *FilerRemoteStorage
|
||||
lazyFetchGroup singleflight.Group
|
||||
lazyListGroup singleflight.Group
|
||||
Dlm *lock_manager.DistributedLockManager
|
||||
MaxFilenameLength uint32
|
||||
deletionQuit chan struct{}
|
||||
DeletionRetryQueue *DeletionRetryQueue
|
||||
EmptyFolderCleaner *empty_folder_cleanup.EmptyFolderCleaner
|
||||
EmptyFolderCleanupDelay time.Duration
|
||||
}
|
||||
|
||||
@@ -82,7 +82,13 @@ func NewFiler(masters pb.ServerDiscovery, grpcDialOption grpc.DialOption, filerH
|
||||
f.UniqueFilerId = -f.UniqueFilerId
|
||||
}
|
||||
|
||||
f.LocalMetaLogBuffer = log_buffer.NewLogBuffer("local", LogFlushInterval, f.logFlushFunc, f.readPersistedLogBufferPosition, notifyFn)
|
||||
// ReadFromDiskFn is intentionally nil here. SubscribeLocalMetadata already
|
||||
// manages disk reads explicitly with shouldReadFromDisk / lastCheckedFlushTsNs
|
||||
// tracking. Setting ReadFromDiskFn would cause LoopProcessLogData to issue a
|
||||
// redundant ReadPersistedLogBuffer call (ListDirectoryEntries + readahead
|
||||
// goroutine) on every 250ms health-check tick when a subscriber encounters
|
||||
// ResumeFromDiskError, adding significant CPU and GC pressure even when idle.
|
||||
f.LocalMetaLogBuffer = log_buffer.NewLogBuffer("local", LogFlushInterval, f.logFlushFunc, nil, notifyFn)
|
||||
f.metaLogCollection = collection
|
||||
f.metaLogReplication = replication
|
||||
|
||||
|
||||
@@ -262,16 +262,3 @@ func (f *Filer) ReadPersistedLogBuffer(startPosition log_buffer.MessagePosition,
|
||||
return
|
||||
}
|
||||
|
||||
func (f *Filer) readPersistedLogBufferPosition(startPosition log_buffer.MessagePosition, stopTsNs int64, eachLogEntryFn log_buffer.EachLogEntryFuncType) (lastReadPosition log_buffer.MessagePosition, isDone bool, err error) {
|
||||
lastReadPosition = startPosition
|
||||
|
||||
lastTsNs, isDone, err := f.ReadPersistedLogBuffer(startPosition, stopTsNs, eachLogEntryFn)
|
||||
if err != nil {
|
||||
return lastReadPosition, isDone, err
|
||||
}
|
||||
if lastTsNs != 0 {
|
||||
lastReadPosition = log_buffer.NewMessagePosition(lastTsNs, 1)
|
||||
}
|
||||
|
||||
return lastReadPosition, isDone, nil
|
||||
}
|
||||
|
||||
@@ -265,8 +265,11 @@ func (fsw *FilerStoreWrapper) DeleteEntry(ctx context.Context, fp util.FullPath)
|
||||
op := ctx.Value("OP")
|
||||
if op != "MV" {
|
||||
glog.V(4).InfofCtx(ctx, "DeleteHardLink %s", existingEntry.FullPath)
|
||||
if err = fsw.DeleteHardLink(ctx, existingEntry.HardLinkId); err != nil {
|
||||
return err
|
||||
if hlErr := fsw.DeleteHardLink(ctx, existingEntry.HardLinkId); hlErr != nil {
|
||||
// Log but continue: the directory entry must be removed
|
||||
// even if hard link counter cleanup fails, otherwise the
|
||||
// parent directory cannot be removed (rmdir ENOTEMPTY).
|
||||
glog.Warningf("DeleteHardLink %s (id %x): %v", existingEntry.FullPath, existingEntry.HardLinkId, hlErr)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -292,8 +295,12 @@ func (fsw *FilerStoreWrapper) DeleteOneEntry(ctx context.Context, existingEntry
|
||||
op := ctx.Value("OP")
|
||||
if op != "MV" {
|
||||
glog.V(4).InfofCtx(ctx, "DeleteHardLink %s", existingEntry.FullPath)
|
||||
if err = fsw.DeleteHardLink(ctx, existingEntry.HardLinkId); err != nil {
|
||||
return err
|
||||
if hlErr := fsw.DeleteHardLink(ctx, existingEntry.HardLinkId); hlErr != nil {
|
||||
// Log the hard link cleanup error but continue to delete
|
||||
// the directory entry. If we return early here, the entry
|
||||
// remains in the store and the parent directory cannot be
|
||||
// removed (rmdir returns ENOTEMPTY).
|
||||
glog.Warningf("DeleteHardLink %s (id %x): %v", existingEntry.FullPath, existingEntry.HardLinkId, hlErr)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -228,13 +228,13 @@ func (ma *MetaAggregator) doSubscribeToOneFiler(f *Filer, self pb.ServerAddress,
|
||||
}
|
||||
|
||||
stream, err := client.SubscribeLocalMetadata(ctx, &filer_pb.SubscribeMetadataRequest{
|
||||
ClientName: "filer:" + string(self),
|
||||
PathPrefix: "/",
|
||||
SinceNs: lastTsNs,
|
||||
ClientId: ma.filer.UniqueFilerId,
|
||||
ClientEpoch: atomic.LoadInt32(&ma.filer.UniqueFilerEpoch),
|
||||
ClientSupportsBatching: true,
|
||||
ClientSupportsMetadataChunks: true,
|
||||
ClientName: "filer:" + string(self),
|
||||
PathPrefix: "/",
|
||||
SinceNs: lastTsNs,
|
||||
ClientId: ma.filer.UniqueFilerId,
|
||||
ClientEpoch: atomic.LoadInt32(&ma.filer.UniqueFilerEpoch),
|
||||
ClientSupportsBatching: true,
|
||||
ClientSupportsMetadataChunks: true,
|
||||
})
|
||||
if err != nil {
|
||||
glog.V(0).Infof("SubscribeLocalMetadata %v: %v", peer, err)
|
||||
|
||||
@@ -14,8 +14,8 @@ import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
util_http "github.com/seaweedfs/seaweedfs/weed/util/http"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
util_http "github.com/seaweedfs/seaweedfs/weed/util/http"
|
||||
"github.com/seaweedfs/seaweedfs/weed/wdclient"
|
||||
)
|
||||
|
||||
@@ -92,10 +92,10 @@ func createMockVolumeServer(chunkData map[string][]byte, latency time.Duration)
|
||||
|
||||
// benchmarkConfig holds parameters for a single benchmark scenario
|
||||
type benchmarkConfig struct {
|
||||
numChunks int
|
||||
chunkSize int
|
||||
latency time.Duration
|
||||
prefetch int // 0 = sequential
|
||||
numChunks int
|
||||
chunkSize int
|
||||
latency time.Duration
|
||||
prefetch int // 0 = sequential
|
||||
}
|
||||
|
||||
func (c benchmarkConfig) name() string {
|
||||
|
||||
+1
-1
@@ -863,7 +863,7 @@ func (l *loggingT) exit(err error) {
|
||||
// file rotation. There are conflicting methods, so the file cannot be embedded.
|
||||
// l.mu is held for all its methods.
|
||||
type syncBuffer struct {
|
||||
logger *loggingT
|
||||
logger *loggingT
|
||||
*bufio.Writer
|
||||
file *os.File
|
||||
sev severity
|
||||
|
||||
@@ -0,0 +1,181 @@
|
||||
package mount
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/glog"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/security"
|
||||
)
|
||||
|
||||
// FileIdEntry holds a pre-allocated file ID from the filer/master, ready for
|
||||
// immediate use by an upload worker without an AssignVolume round-trip.
|
||||
type FileIdEntry struct {
|
||||
FileId string
|
||||
Host string // volume server address (already adjusted for access mode)
|
||||
Auth security.EncodedJwt
|
||||
Time time.Time
|
||||
}
|
||||
|
||||
// FileIdPool pre-allocates file IDs in batches so that chunk uploads can grab
|
||||
// one instantly instead of blocking on an AssignVolume RPC per chunk.
|
||||
//
|
||||
// The pool is refilled in the background when it drops below a low-water mark.
|
||||
// All IDs are allocated with the mount's global (replication, collection, ttl,
|
||||
// diskType, dataCenter) parameters. Path-based storage rules (filer.conf) are
|
||||
// NOT applied to pooled IDs since the pool allocates ahead of any specific file
|
||||
// path. This is an intentional tradeoff for writeback cache performance.
|
||||
type FileIdPool struct {
|
||||
wfs *WFS
|
||||
|
||||
mu sync.Mutex
|
||||
cond *sync.Cond
|
||||
entries []FileIdEntry // available pre-allocated IDs
|
||||
filling bool // true when a background refill is in progress
|
||||
|
||||
poolSize int // target pool capacity
|
||||
batchSize int // how many IDs to request per Assign RPC
|
||||
lowWater int // refill trigger threshold
|
||||
maxAge time.Duration
|
||||
}
|
||||
|
||||
func NewFileIdPool(wfs *WFS) *FileIdPool {
|
||||
concurrency := wfs.option.ConcurrentWriters
|
||||
if concurrency <= 0 {
|
||||
concurrency = 128 // match default async flush worker count
|
||||
}
|
||||
pool := &FileIdPool{
|
||||
wfs: wfs,
|
||||
poolSize: concurrency * 2,
|
||||
batchSize: concurrency,
|
||||
lowWater: concurrency,
|
||||
maxAge: 25 * time.Second, // conservative; JWT TTL is typically 30s+
|
||||
}
|
||||
pool.cond = sync.NewCond(&pool.mu)
|
||||
return pool
|
||||
}
|
||||
|
||||
// Get returns a pre-allocated file ID entry. If the pool is empty and a refill
|
||||
// is in progress, callers wait for it to complete rather than failing. Returns
|
||||
// an error only if the Assign RPC fails.
|
||||
func (p *FileIdPool) Get() (FileIdEntry, error) {
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
|
||||
for {
|
||||
p.evictExpired()
|
||||
|
||||
if len(p.entries) > 0 {
|
||||
entry := p.entries[0]
|
||||
p.entries = p.entries[1:]
|
||||
if len(p.entries) < p.lowWater && !p.filling {
|
||||
p.filling = true
|
||||
go p.doRefill()
|
||||
}
|
||||
return entry, nil
|
||||
}
|
||||
|
||||
// Pool empty.
|
||||
if p.filling {
|
||||
// Wait for the in-flight refill to complete.
|
||||
p.cond.Wait()
|
||||
continue
|
||||
}
|
||||
|
||||
// No refill in progress — start one synchronously.
|
||||
p.filling = true
|
||||
p.mu.Unlock()
|
||||
entries, err := p.assignBatch(p.batchSize)
|
||||
p.mu.Lock()
|
||||
p.filling = false
|
||||
p.cond.Broadcast()
|
||||
|
||||
if err != nil {
|
||||
return FileIdEntry{}, fmt.Errorf("fileIdPool: %w", err)
|
||||
}
|
||||
p.entries = append(p.entries, entries...)
|
||||
// Loop back to pop from entries.
|
||||
}
|
||||
}
|
||||
|
||||
func (p *FileIdPool) evictExpired() {
|
||||
cutoff := time.Now().Add(-p.maxAge)
|
||||
i := 0
|
||||
for i < len(p.entries) && p.entries[i].Time.Before(cutoff) {
|
||||
i++
|
||||
}
|
||||
if i > 0 {
|
||||
p.entries = p.entries[i:]
|
||||
}
|
||||
}
|
||||
|
||||
// doRefill runs in a background goroutine to refill the pool.
|
||||
func (p *FileIdPool) doRefill() {
|
||||
entries, err := p.assignBatch(p.batchSize)
|
||||
if err != nil {
|
||||
glog.V(1).Infof("fileIdPool refill: %v", err)
|
||||
}
|
||||
|
||||
p.mu.Lock()
|
||||
if err == nil {
|
||||
p.entries = append(p.entries, entries...)
|
||||
}
|
||||
p.filling = false
|
||||
p.cond.Broadcast()
|
||||
p.mu.Unlock()
|
||||
}
|
||||
|
||||
// assignBatch requests `count` file IDs from the filer using individual
|
||||
// Count=1 RPCs over a single gRPC connection. Each response includes a
|
||||
// per-fid JWT, so uploads work correctly when JWT security is enabled.
|
||||
//
|
||||
// We use individual requests instead of Count=N because the master generates
|
||||
// one JWT for the base file ID only (master_grpc_server_assign.go:158), and
|
||||
// the volume server validates that the JWT's Fid matches the upload's file ID
|
||||
// exactly (volume_server_handlers.go:367). Sequential IDs derived from a
|
||||
// Count=N response would fail this check.
|
||||
//
|
||||
// Note: the AssignVolumeRequest intentionally omits the Path field. Pooled IDs
|
||||
// use the mount's global storage parameters, not per-path rules from filer.conf
|
||||
// (detectStorageOption / MatchStorageRule). This is a writeback cache tradeoff.
|
||||
func (p *FileIdPool) assignBatch(count int) ([]FileIdEntry, error) {
|
||||
var entries []FileIdEntry
|
||||
err := p.wfs.WithFilerClient(false, func(client filer_pb.SeaweedFilerClient) error {
|
||||
now := time.Now()
|
||||
req := &filer_pb.AssignVolumeRequest{
|
||||
Count: 1,
|
||||
Replication: p.wfs.option.Replication,
|
||||
Collection: p.wfs.option.Collection,
|
||||
TtlSec: p.wfs.option.TtlSec,
|
||||
DiskType: string(p.wfs.option.DiskType),
|
||||
DataCenter: p.wfs.option.DataCenter,
|
||||
ExpectedDataSize: uint64(p.wfs.option.ChunkSizeLimit),
|
||||
}
|
||||
for i := 0; i < count; i++ {
|
||||
resp, assignErr := client.AssignVolume(context.Background(), req)
|
||||
if assignErr != nil {
|
||||
if len(entries) > 0 {
|
||||
break // partial batch is fine
|
||||
}
|
||||
return assignErr
|
||||
}
|
||||
if resp.Error != "" {
|
||||
if len(entries) > 0 {
|
||||
break
|
||||
}
|
||||
return fmt.Errorf("assign: %s", resp.Error)
|
||||
}
|
||||
entries = append(entries, FileIdEntry{
|
||||
FileId: resp.FileId,
|
||||
Host: p.wfs.AdjustedUrl(resp.Location),
|
||||
Auth: security.EncodedJwt(resp.Auth),
|
||||
Time: now,
|
||||
})
|
||||
}
|
||||
return nil
|
||||
})
|
||||
return entries, err
|
||||
}
|
||||
@@ -0,0 +1,280 @@
|
||||
package mount
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/security"
|
||||
"github.com/seaweedfs/seaweedfs/weed/storage/needle"
|
||||
"google.golang.org/grpc"
|
||||
)
|
||||
|
||||
// mockFilerClient simulates AssignVolume RPCs with configurable latency.
|
||||
type mockFilerClient struct {
|
||||
latency time.Duration
|
||||
nextVolId uint32
|
||||
nextKey uint64
|
||||
mu sync.Mutex
|
||||
}
|
||||
|
||||
func (m *mockFilerClient) AssignVolume(_ context.Context, req *filer_pb.AssignVolumeRequest, _ ...grpc.CallOption) (*filer_pb.AssignVolumeResponse, error) {
|
||||
if m.latency > 0 {
|
||||
time.Sleep(m.latency)
|
||||
}
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
|
||||
count := int(req.Count)
|
||||
if count <= 0 {
|
||||
count = 1
|
||||
}
|
||||
vid := needle.VolumeId(m.nextVolId + 1)
|
||||
key := m.nextKey + 1000 // start at a non-zero key for valid hex encoding
|
||||
m.nextKey += uint64(count)
|
||||
m.nextVolId++
|
||||
|
||||
fid := needle.NewFileId(vid, key, 0x12345678)
|
||||
return &filer_pb.AssignVolumeResponse{
|
||||
FileId: fid.String(),
|
||||
Count: int32(count),
|
||||
Auth: "test-jwt-token",
|
||||
Location: &filer_pb.Location{
|
||||
Url: "127.0.0.1:8080",
|
||||
PublicUrl: "127.0.0.1:8080",
|
||||
GrpcPort: 18080,
|
||||
},
|
||||
}, nil
|
||||
}
|
||||
|
||||
// TestFileIdPoolSequentialIds verifies that batch assignment generates
|
||||
// correct sequential file IDs from a single AssignVolume response.
|
||||
func TestFileIdPoolSequentialIds(t *testing.T) {
|
||||
mock := &mockFilerClient{}
|
||||
resp, err := mock.AssignVolume(context.Background(), &filer_pb.AssignVolumeRequest{Count: 5})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
baseFid, parseErr := needle.ParseFileIdFromString(resp.FileId)
|
||||
if parseErr != nil {
|
||||
t.Fatal(parseErr)
|
||||
}
|
||||
|
||||
baseKey := uint64(baseFid.Key)
|
||||
for i := 0; i < int(resp.Count); i++ {
|
||||
fid := needle.NewFileId(baseFid.VolumeId, baseKey+uint64(i), uint32(baseFid.Cookie))
|
||||
t.Logf("ID %d: %s", i, fid.String())
|
||||
parsed, err := needle.ParseFileIdFromString(fid.String())
|
||||
if err != nil {
|
||||
t.Fatalf("Failed to parse sequential ID %d: %v", i, err)
|
||||
}
|
||||
if uint64(parsed.Key) != baseKey+uint64(i) {
|
||||
t.Fatalf("ID %d: expected key %d, got %d", i, baseKey+uint64(i), uint64(parsed.Key))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// TestFileIdPoolExpiry verifies that expired entries are evicted.
|
||||
func TestFileIdPoolExpiry(t *testing.T) {
|
||||
pool := &FileIdPool{
|
||||
maxAge: time.Second,
|
||||
}
|
||||
pool.cond = sync.NewCond(&pool.mu)
|
||||
now := time.Now()
|
||||
pool.entries = []FileIdEntry{
|
||||
{FileId: "1,old", Time: now.Add(-2 * time.Second)},
|
||||
{FileId: "2,old", Time: now.Add(-2 * time.Second)},
|
||||
{FileId: "3,fresh", Time: now},
|
||||
}
|
||||
pool.evictExpired()
|
||||
if len(pool.entries) != 1 {
|
||||
t.Fatalf("expected 1 entry after eviction, got %d", len(pool.entries))
|
||||
}
|
||||
if pool.entries[0].FileId != "3,fresh" {
|
||||
t.Fatalf("expected fresh entry, got %s", pool.entries[0].FileId)
|
||||
}
|
||||
}
|
||||
|
||||
// TestFileIdPoolGetWaitsForRefill verifies that concurrent Get() calls wait
|
||||
// for an in-flight refill instead of returning an error.
|
||||
func TestFileIdPoolGetWaitsForRefill(t *testing.T) {
|
||||
pool := &FileIdPool{
|
||||
poolSize: 10,
|
||||
batchSize: 5,
|
||||
lowWater: 3,
|
||||
maxAge: 30 * time.Second,
|
||||
}
|
||||
pool.cond = sync.NewCond(&pool.mu)
|
||||
|
||||
// Simulate a slow refill in progress.
|
||||
pool.filling = true
|
||||
done := make(chan struct{})
|
||||
go func() {
|
||||
// After a short delay, deliver entries and signal.
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
pool.mu.Lock()
|
||||
now := time.Now()
|
||||
for i := 0; i < 5; i++ {
|
||||
pool.entries = append(pool.entries, FileIdEntry{
|
||||
FileId: fmt.Sprintf("1,%x12345678", 1000+i),
|
||||
Host: "127.0.0.1:8080",
|
||||
Auth: "jwt",
|
||||
Time: now,
|
||||
})
|
||||
}
|
||||
pool.filling = false
|
||||
pool.cond.Broadcast()
|
||||
pool.mu.Unlock()
|
||||
close(done)
|
||||
}()
|
||||
|
||||
// Get() should wait for the refill, not return an error.
|
||||
entry, err := pool.Get()
|
||||
if err != nil {
|
||||
t.Fatalf("Get() returned error while refill in progress: %v", err)
|
||||
}
|
||||
if entry.FileId == "" {
|
||||
t.Fatal("Get() returned empty entry")
|
||||
}
|
||||
<-done
|
||||
}
|
||||
|
||||
// BenchmarkPoolGetVsDirectAssign measures the latency difference between
|
||||
// getting a file ID from the pool vs a direct (simulated) AssignVolume RPC.
|
||||
func BenchmarkPoolGetVsDirectAssign(b *testing.B) {
|
||||
assignLatencies := []time.Duration{
|
||||
0,
|
||||
100 * time.Microsecond,
|
||||
1 * time.Millisecond,
|
||||
5 * time.Millisecond,
|
||||
}
|
||||
|
||||
for _, latency := range assignLatencies {
|
||||
name := "latency=" + latency.String()
|
||||
|
||||
b.Run("DirectAssign/"+name, func(b *testing.B) {
|
||||
mock := &mockFilerClient{latency: latency}
|
||||
b.ResetTimer()
|
||||
for i := 0; i < b.N; i++ {
|
||||
_, err := mock.AssignVolume(context.Background(), &filer_pb.AssignVolumeRequest{Count: 1})
|
||||
if err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
}
|
||||
})
|
||||
|
||||
b.Run("PoolGet/"+name, func(b *testing.B) {
|
||||
pool := &FileIdPool{
|
||||
maxAge: 30 * time.Second,
|
||||
lowWater: 1000000, // disable background refill
|
||||
}
|
||||
pool.cond = sync.NewCond(&pool.mu)
|
||||
// Pre-fill with a fixed-size pool; refill under StopTimer when depleted.
|
||||
const preload = 4096
|
||||
refill := func() {
|
||||
now := time.Now()
|
||||
for i := 0; i < preload; i++ {
|
||||
pool.entries = append(pool.entries, FileIdEntry{
|
||||
FileId: fmt.Sprintf("1,%x12345678", i),
|
||||
Host: "127.0.0.1:8080",
|
||||
Auth: security.EncodedJwt("test-jwt"),
|
||||
Time: now,
|
||||
})
|
||||
}
|
||||
}
|
||||
refill()
|
||||
b.ResetTimer()
|
||||
for i := 0; i < b.N; i++ {
|
||||
pool.mu.Lock()
|
||||
if len(pool.entries) == 0 {
|
||||
pool.mu.Unlock()
|
||||
b.StopTimer()
|
||||
refill()
|
||||
b.StartTimer()
|
||||
pool.mu.Lock()
|
||||
}
|
||||
_ = pool.entries[0]
|
||||
pool.entries = pool.entries[1:]
|
||||
pool.mu.Unlock()
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// BenchmarkConcurrentPoolGet measures pool throughput under concurrent access.
|
||||
func BenchmarkConcurrentPoolGet(b *testing.B) {
|
||||
for _, workers := range []int{1, 4, 16, 64} {
|
||||
b.Run(fmt.Sprintf("workers=%d", workers), func(b *testing.B) {
|
||||
pool := &FileIdPool{
|
||||
maxAge: 30 * time.Second,
|
||||
lowWater: 1000000,
|
||||
}
|
||||
pool.cond = sync.NewCond(&pool.mu)
|
||||
// Pre-fill
|
||||
now := time.Now()
|
||||
total := b.N*workers + 1000
|
||||
pool.entries = make([]FileIdEntry, total)
|
||||
for i := range pool.entries {
|
||||
pool.entries[i] = FileIdEntry{
|
||||
FileId: fmt.Sprintf("1,%x12345678", i+1000),
|
||||
Host: "127.0.0.1:8080",
|
||||
Auth: security.EncodedJwt("test-jwt"),
|
||||
Time: now,
|
||||
}
|
||||
}
|
||||
|
||||
var ops atomic.Int64
|
||||
b.ResetTimer()
|
||||
|
||||
var wg sync.WaitGroup
|
||||
perWorker := b.N
|
||||
for w := 0; w < workers; w++ {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
for i := 0; i < perWorker; i++ {
|
||||
pool.mu.Lock()
|
||||
if len(pool.entries) > 0 {
|
||||
_ = pool.entries[0]
|
||||
pool.entries = pool.entries[1:]
|
||||
ops.Add(1)
|
||||
}
|
||||
pool.mu.Unlock()
|
||||
}
|
||||
}()
|
||||
}
|
||||
wg.Wait()
|
||||
b.ReportMetric(float64(ops.Load())/b.Elapsed().Seconds(), "ids/sec")
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// BenchmarkBatchAssign measures the cost of batch vs individual assign RPCs.
|
||||
func BenchmarkBatchAssign(b *testing.B) {
|
||||
for _, batchSize := range []int{1, 8, 16, 32} {
|
||||
b.Run(fmt.Sprintf("batch=%d", batchSize), func(b *testing.B) {
|
||||
mock := &mockFilerClient{latency: 1 * time.Millisecond}
|
||||
totalIds := 0
|
||||
b.ResetTimer()
|
||||
for i := 0; i < b.N; i++ {
|
||||
resp, err := mock.AssignVolume(context.Background(), &filer_pb.AssignVolumeRequest{Count: int32(batchSize)})
|
||||
if err != nil {
|
||||
b.Fatal(err)
|
||||
}
|
||||
|
||||
baseFid, _ := needle.ParseFileIdFromString(resp.FileId)
|
||||
baseKey := uint64(baseFid.Key)
|
||||
for j := 0; j < int(resp.Count); j++ {
|
||||
_ = needle.NewFileId(baseFid.VolumeId, baseKey+uint64(j), uint32(baseFid.Cookie)).String()
|
||||
totalIds++
|
||||
}
|
||||
}
|
||||
b.ReportMetric(float64(totalIds)/b.Elapsed().Seconds(), "ids/sec")
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -27,6 +27,7 @@ type InodeEntry struct {
|
||||
lastRefresh time.Time
|
||||
updateWindowStart time.Time
|
||||
updateCount int
|
||||
subdirCount int32 // tracked in-memory for POSIX directory nlink
|
||||
}
|
||||
|
||||
func (ie *InodeEntry) resetCacheState() {
|
||||
@@ -250,6 +251,55 @@ func (i *InodeToPath) InvalidateChildrenCache(fullpath util.FullPath) {
|
||||
entry.resetCacheState()
|
||||
}
|
||||
|
||||
// AdjustSubdirCount adjusts the subdirectory count for a directory inode.
|
||||
// delta is typically +1 (mkdir) or -1 (rmdir).
|
||||
func (i *InodeToPath) AdjustSubdirCount(dirPath util.FullPath, delta int32) {
|
||||
i.Lock()
|
||||
defer i.Unlock()
|
||||
inode, found := i.path2inode[dirPath]
|
||||
if !found {
|
||||
return
|
||||
}
|
||||
entry, found := i.inode2path[inode]
|
||||
if !found || !entry.isDirectory {
|
||||
return
|
||||
}
|
||||
entry.subdirCount += delta
|
||||
if entry.subdirCount < 0 {
|
||||
entry.subdirCount = 0
|
||||
}
|
||||
}
|
||||
|
||||
// GetSubdirCount returns the tracked subdirectory count for a directory.
|
||||
func (i *InodeToPath) GetSubdirCount(dirPath util.FullPath) int32 {
|
||||
i.RLock()
|
||||
defer i.RUnlock()
|
||||
inode, found := i.path2inode[dirPath]
|
||||
if !found {
|
||||
return 0
|
||||
}
|
||||
entry, found := i.inode2path[inode]
|
||||
if !found || !entry.isDirectory {
|
||||
return 0
|
||||
}
|
||||
return entry.subdirCount
|
||||
}
|
||||
|
||||
// SetSubdirCount sets the subdirectory count for a directory (used after readdir).
|
||||
func (i *InodeToPath) SetSubdirCount(dirPath util.FullPath, count int32) {
|
||||
i.Lock()
|
||||
defer i.Unlock()
|
||||
inode, found := i.path2inode[dirPath]
|
||||
if !found {
|
||||
return
|
||||
}
|
||||
entry, found := i.inode2path[inode]
|
||||
if !found || !entry.isDirectory {
|
||||
return
|
||||
}
|
||||
entry.subdirCount = count
|
||||
}
|
||||
|
||||
func (i *InodeToPath) TouchDirectory(fullpath util.FullPath) {
|
||||
i.Lock()
|
||||
defer i.Unlock()
|
||||
|
||||
@@ -278,6 +278,25 @@ func (mc *MetaCache) UpdateEntry(ctx context.Context, entry *filer.Entry) error
|
||||
return mc.localStore.UpdateEntry(ctx, entry)
|
||||
}
|
||||
|
||||
// TouchDirMtimeCtime updates the mtime and ctime of a directory entry
|
||||
// directly in the local metadata cache store. This avoids a filer RPC
|
||||
// round-trip and the associated metadata event that would invalidate
|
||||
// recently cached child entries.
|
||||
func (mc *MetaCache) TouchDirMtimeCtime(ctx context.Context, dirPath util.FullPath, now time.Time) error {
|
||||
mc.Lock()
|
||||
defer mc.Unlock()
|
||||
entry, err := mc.localStore.FindEntry(ctx, dirPath)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if entry == nil {
|
||||
return nil
|
||||
}
|
||||
entry.Attr.Mtime = now
|
||||
entry.Attr.Ctime = now
|
||||
return mc.localStore.UpdateEntry(ctx, entry)
|
||||
}
|
||||
|
||||
func (mc *MetaCache) FindEntry(ctx context.Context, fp util.FullPath) (entry *filer.Entry, err error) {
|
||||
mc.RLock()
|
||||
defer mc.RUnlock()
|
||||
|
||||
@@ -42,7 +42,7 @@ func mergeProcessors(mainProcessor func(resp *filer_pb.SubscribeMetadataResponse
|
||||
}
|
||||
}
|
||||
|
||||
func SubscribeMetaEvents(mc *MetaCache, selfSignature int32, client filer_pb.FilerClient, dir string, lastTsNs int64, onRetry func(lastTsNs int64, err error), followers ...*MetadataFollower) error {
|
||||
func SubscribeMetaEvents(mc *MetaCache, selfSignature int32, client filer_pb.FilerClient, dir string, lastTsNs int64, skipSelfEvents bool, onRetry func(lastTsNs int64, err error), followers ...*MetadataFollower) error {
|
||||
|
||||
var prefixes []string
|
||||
for _, follower := range followers {
|
||||
@@ -50,11 +50,14 @@ func SubscribeMetaEvents(mc *MetaCache, selfSignature int32, client filer_pb.Fil
|
||||
}
|
||||
|
||||
processEventFn := func(resp *filer_pb.SubscribeMetadataResponse) error {
|
||||
// Let all events (including self-originated ones) flow through the
|
||||
// applier so that the directory-build buffering and dedup logic
|
||||
// can handle them consistently. The dedupRing in
|
||||
// applyMetadataResponseNow catches duplicates that were already
|
||||
// applied locally via applyLocalMetadataEvent.
|
||||
if skipSelfEvents && resp.EventNotification != nil {
|
||||
for _, sig := range resp.EventNotification.Signatures {
|
||||
if sig == selfSignature {
|
||||
glog.V(4).Infof("skip self-originated event %s", resp.Directory)
|
||||
return nil
|
||||
}
|
||||
}
|
||||
}
|
||||
return mc.ApplyMetadataResponse(context.Background(), resp, SubscriberMetadataResponseApplyOptions)
|
||||
}
|
||||
|
||||
|
||||
+34
-5
@@ -63,9 +63,9 @@ type Option struct {
|
||||
MountMtime time.Time
|
||||
MountParentInode uint64
|
||||
|
||||
VolumeServerAccess string // how to access volume servers
|
||||
Cipher bool // whether encrypt data on volume server
|
||||
UidGidMapper *meta_cache.UidGidMapper
|
||||
VolumeServerAccess string // how to access volume servers
|
||||
Cipher bool // whether encrypt data on volume server
|
||||
UidGidMapper *meta_cache.UidGidMapper
|
||||
IncludeSystemEntries bool
|
||||
|
||||
// Periodic metadata flush interval in seconds (0 to disable)
|
||||
@@ -93,6 +93,12 @@ type Option struct {
|
||||
// When true, Flush() returns immediately and data upload + metadata flush happen in background.
|
||||
WritebackCache bool
|
||||
|
||||
// PosixDirNlink enables POSIX-compliant directory nlink counting
|
||||
// (nlink = 2 + number_of_subdirectories). This requires listing
|
||||
// cached directory entries on every stat, which has a performance cost.
|
||||
// When false (default), directories report nlink=2.
|
||||
PosixDirNlink bool
|
||||
|
||||
uniqueCacheDirForRead string
|
||||
uniqueCacheDirForWrite string
|
||||
}
|
||||
@@ -125,9 +131,14 @@ type WFS struct {
|
||||
refreshingDirs map[util.FullPath]struct{}
|
||||
atimeMu sync.Mutex
|
||||
atimeMap map[uint64]time.Time // inode -> atime, in-memory only, bounded
|
||||
dirMtimeMu sync.Mutex
|
||||
dirMtimeMap map[uint64]time.Time // inode -> mtime/ctime, in-memory overlay for dirs
|
||||
entryValidSec uint64 // kernel FUSE entry cache TTL in seconds
|
||||
attrValidSec uint64 // kernel FUSE attr cache TTL in seconds
|
||||
dirHotWindow time.Duration
|
||||
dirHotThreshold int
|
||||
dirIdleEvict time.Duration
|
||||
fileIdPool *FileIdPool
|
||||
|
||||
// asyncFlushWg tracks pending background flush work items for writebackCache mode.
|
||||
// Must be waited on before unmount cleanup to prevent data loss.
|
||||
@@ -205,14 +216,27 @@ func NewSeaweedFileSystem(option *Option) *WFS {
|
||||
posixLocks: NewPosixLockTable(),
|
||||
refreshingDirs: make(map[util.FullPath]struct{}),
|
||||
atimeMap: make(map[uint64]time.Time, 8192),
|
||||
dirMtimeMap: make(map[uint64]time.Time, 1024),
|
||||
entryValidSec: 1,
|
||||
attrValidSec: 1,
|
||||
dirHotWindow: dirHotWindow,
|
||||
dirHotThreshold: dirHotThreshold,
|
||||
dirIdleEvict: dirIdleEvict,
|
||||
}
|
||||
|
||||
if option.EnableDistributedLock && len(option.FilerAddresses) > 0 {
|
||||
// With writeback caching, this mount is the single writer. Increase kernel
|
||||
// FUSE cache TTLs so the kernel doesn't re-issue Lookup/GetAttr for every
|
||||
// path component and stat — the local meta cache is authoritative.
|
||||
if option.WritebackCache {
|
||||
wfs.entryValidSec = 10
|
||||
wfs.attrValidSec = 10
|
||||
}
|
||||
|
||||
if option.EnableDistributedLock && !option.WritebackCache && len(option.FilerAddresses) > 0 {
|
||||
wfs.lockClient = cluster.NewLockClient(option.GrpcDialOption, option.FilerAddresses[0])
|
||||
glog.V(0).Infof("distributed lock manager enabled for mount")
|
||||
} else if option.EnableDistributedLock && option.WritebackCache {
|
||||
glog.V(0).Infof("distributed lock manager disabled: writeback cache implies single-writer mode")
|
||||
}
|
||||
|
||||
wfs.option.filerIndex = int32(rand.IntN(len(option.FilerAddresses)))
|
||||
@@ -333,7 +357,7 @@ func (wfs *WFS) StartBackgroundTasks() error {
|
||||
}
|
||||
|
||||
startTime := time.Now()
|
||||
go meta_cache.SubscribeMetaEvents(wfs.metaCache, wfs.signature, wfs, wfs.option.FilerMountRootPath, startTime.UnixNano(), func(lastTsNs int64, err error) {
|
||||
go meta_cache.SubscribeMetaEvents(wfs.metaCache, wfs.signature, wfs, wfs.option.FilerMountRootPath, startTime.UnixNano(), wfs.option.WritebackCache, func(lastTsNs int64, err error) {
|
||||
glog.Warningf("meta events follow retry from %v: %v", time.Unix(0, lastTsNs), err)
|
||||
if deleteErr := wfs.metaCache.DeleteFolderChildren(context.Background(), util.FullPath(wfs.option.FilerMountRootPath)); deleteErr != nil {
|
||||
glog.Warningf("meta cache cleanup failed: %v", deleteErr)
|
||||
@@ -344,6 +368,11 @@ func (wfs *WFS) StartBackgroundTasks() error {
|
||||
go wfs.loopFlushDirtyMetadata()
|
||||
go wfs.loopEvictIdleDirCache()
|
||||
|
||||
if wfs.option.WritebackCache {
|
||||
wfs.fileIdPool = NewFileIdPool(wfs)
|
||||
glog.V(0).Infof("file ID pool enabled for writeback cache (batch=%d)", wfs.fileIdPool.batchSize)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
|
||||
+85
-11
@@ -16,19 +16,28 @@ func (wfs *WFS) GetAttr(cancel <-chan struct{}, input *fuse.GetAttrIn, out *fuse
|
||||
glog.V(4).Infof("GetAttr %v", input.NodeId)
|
||||
if input.NodeId == 1 {
|
||||
wfs.setRootAttr(out)
|
||||
if wfs.option.PosixDirNlink {
|
||||
wfs.applyDirNlink(&out.Attr, util.FullPath(wfs.option.FilerMountRootPath))
|
||||
}
|
||||
return fuse.OK
|
||||
}
|
||||
|
||||
inode := input.NodeId
|
||||
_, _, entry, status := wfs.maybeReadEntry(inode)
|
||||
path, _, entry, status := wfs.maybeReadEntry(inode)
|
||||
if status == fuse.OK {
|
||||
out.AttrValid = 1
|
||||
out.AttrValid = wfs.attrValidSec
|
||||
wfs.setAttrByPbEntry(&out.Attr, inode, entry, true)
|
||||
wfs.applyInMemoryAtime(&out.Attr, inode)
|
||||
if entry.IsDirectory {
|
||||
wfs.applyInMemoryDirMtime(&out.Attr, inode)
|
||||
if wfs.option.PosixDirNlink {
|
||||
wfs.applyDirNlink(&out.Attr, path)
|
||||
}
|
||||
}
|
||||
return status
|
||||
} else {
|
||||
if fh, found := wfs.fhMap.FindFileHandle(inode); found {
|
||||
out.AttrValid = 1
|
||||
out.AttrValid = wfs.attrValidSec
|
||||
// Use shared lock to prevent race with Write operations
|
||||
fhActiveLock := wfs.fhLockTable.AcquireLock("GetAttr", fh.fh, util.SharedLock)
|
||||
wfs.setAttrByPbEntry(&out.Attr, inode, fh.entry.GetEntry(), true)
|
||||
@@ -148,7 +157,7 @@ func (wfs *WFS) SetAttr(cancel <-chan struct{}, input *fuse.SetAttrIn, out *fuse
|
||||
entry.Attributes.Ctime = now.Unix()
|
||||
entry.Attributes.CtimeNs = int32(now.Nanosecond())
|
||||
|
||||
out.AttrValid = 1
|
||||
out.AttrValid = wfs.attrValidSec
|
||||
size, includeSize := input.GetSize()
|
||||
if includeSize {
|
||||
out.Attr.Size = size
|
||||
@@ -176,7 +185,7 @@ func (wfs *WFS) setRootAttr(out *fuse.AttrOut) {
|
||||
out.Ctime = now
|
||||
out.Atime = now
|
||||
out.Mode = toSyscallType(os.ModeDir) | uint32(wfs.option.MountMode)
|
||||
out.Nlink = 1
|
||||
out.Nlink = 2
|
||||
}
|
||||
|
||||
func (wfs *WFS) setAttrByPbEntry(out *fuse.Attr, inode uint64, entry *filer_pb.Entry, calculateSize bool) {
|
||||
@@ -208,7 +217,9 @@ func (wfs *WFS) setAttrByPbEntry(out *fuse.Attr, inode uint64, entry *filer_pb.E
|
||||
out.Atimensec = uint32(entry.Attributes.MtimeNs)
|
||||
// In-memory atime overlay is applied by the caller via applyInMemoryAtime.
|
||||
out.Mode = toSyscallMode(os.FileMode(entry.Attributes.FileMode))
|
||||
if entry.HardLinkCounter > 0 {
|
||||
if entry.IsDirectory {
|
||||
out.Nlink = 2
|
||||
} else if entry.HardLinkCounter > 0 {
|
||||
out.Nlink = uint32(entry.HardLinkCounter)
|
||||
} else {
|
||||
out.Nlink = 1
|
||||
@@ -238,7 +249,9 @@ func (wfs *WFS) setAttrByFilerEntry(out *fuse.Attr, inode uint64, entry *filer.E
|
||||
out.Ctimensec = uint32(entry.Attr.Mtime.Nanosecond())
|
||||
}
|
||||
out.Mode = toSyscallMode(entry.Attr.Mode)
|
||||
if entry.HardLinkCounter > 0 {
|
||||
if entry.IsDirectory() {
|
||||
out.Nlink = 2
|
||||
} else if entry.HardLinkCounter > 0 {
|
||||
out.Nlink = uint32(entry.HardLinkCounter)
|
||||
} else {
|
||||
out.Nlink = 1
|
||||
@@ -251,19 +264,31 @@ func (wfs *WFS) setAttrByFilerEntry(out *fuse.Attr, inode uint64, entry *filer.E
|
||||
func (wfs *WFS) outputPbEntry(out *fuse.EntryOut, inode uint64, entry *filer_pb.Entry) {
|
||||
out.NodeId = inode
|
||||
out.Generation = 1
|
||||
out.EntryValid = 1
|
||||
out.AttrValid = 1
|
||||
out.EntryValid = wfs.entryValidSec
|
||||
out.AttrValid = wfs.attrValidSec
|
||||
wfs.setAttrByPbEntry(&out.Attr, inode, entry, true)
|
||||
}
|
||||
|
||||
func (wfs *WFS) outputFilerEntry(out *fuse.EntryOut, inode uint64, entry *filer.Entry) {
|
||||
out.NodeId = inode
|
||||
out.Generation = 1
|
||||
out.EntryValid = 1
|
||||
out.AttrValid = 1
|
||||
out.EntryValid = wfs.entryValidSec
|
||||
out.AttrValid = wfs.attrValidSec
|
||||
wfs.setAttrByFilerEntry(&out.Attr, inode, entry)
|
||||
}
|
||||
|
||||
// touchDirMtimeCtimeBest updates a directory's mtime and ctime using the
|
||||
// best strategy for the current mode:
|
||||
// - WritebackCache: local meta cache only (no filer RPC)
|
||||
// - Normal mode: filer UpdateEntry RPC for POSIX correctness
|
||||
func (wfs *WFS) touchDirMtimeCtimeBest(dirPath util.FullPath) {
|
||||
if wfs.option.WritebackCache {
|
||||
wfs.touchDirMtimeCtimeLocal(dirPath)
|
||||
} else {
|
||||
wfs.touchDirMtimeCtime(dirPath)
|
||||
}
|
||||
}
|
||||
|
||||
// touchDirMtimeCtime updates a directory's mtime and ctime on the filer.
|
||||
// POSIX requires this when entries are created or removed in the directory.
|
||||
func (wfs *WFS) touchDirMtimeCtime(dirPath util.FullPath) {
|
||||
@@ -279,6 +304,46 @@ func (wfs *WFS) touchDirMtimeCtime(dirPath util.FullPath) {
|
||||
wfs.saveEntry(dirPath, dirEntry)
|
||||
}
|
||||
|
||||
// touchDirMtimeCtimeLocal updates a directory's mtime and ctime in an in-memory
|
||||
// overlay, avoiding LevelDB reads and writes entirely. The overlay is applied
|
||||
// by applyInMemoryDirMtime when GetAttr/Lookup reads the directory's attributes.
|
||||
func (wfs *WFS) touchDirMtimeCtimeLocal(dirPath util.FullPath) {
|
||||
if inode, found := wfs.inodeToPath.GetInode(dirPath); found {
|
||||
wfs.setDirMtime(inode, time.Now())
|
||||
}
|
||||
}
|
||||
|
||||
const dirMtimeMapMaxSize = 8192
|
||||
|
||||
func (wfs *WFS) setDirMtime(inode uint64, t time.Time) {
|
||||
wfs.dirMtimeMu.Lock()
|
||||
defer wfs.dirMtimeMu.Unlock()
|
||||
if len(wfs.dirMtimeMap) >= dirMtimeMapMaxSize {
|
||||
for k := range wfs.dirMtimeMap {
|
||||
delete(wfs.dirMtimeMap, k)
|
||||
break
|
||||
}
|
||||
}
|
||||
wfs.dirMtimeMap[inode] = t
|
||||
}
|
||||
|
||||
// applyInMemoryDirMtime overlays the in-memory mtime/ctime onto fuse.Attr
|
||||
// for directories that had recent child mutations.
|
||||
func (wfs *WFS) applyInMemoryDirMtime(out *fuse.Attr, inode uint64) {
|
||||
wfs.dirMtimeMu.Lock()
|
||||
if t, ok := wfs.dirMtimeMap[inode]; ok {
|
||||
sec := uint64(t.Unix())
|
||||
nsec := uint32(t.Nanosecond())
|
||||
if sec > out.Mtime || (sec == out.Mtime && nsec > out.Mtimensec) {
|
||||
out.Mtime = sec
|
||||
out.Mtimensec = nsec
|
||||
out.Ctime = sec
|
||||
out.Ctimensec = nsec
|
||||
}
|
||||
}
|
||||
wfs.dirMtimeMu.Unlock()
|
||||
}
|
||||
|
||||
const atimeMapMaxSize = 8192
|
||||
|
||||
// setAtime stores an in-memory atime for an inode. The map is bounded;
|
||||
@@ -306,6 +371,15 @@ func (wfs *WFS) applyInMemoryAtime(out *fuse.Attr, inode uint64) {
|
||||
wfs.atimeMu.Unlock()
|
||||
}
|
||||
|
||||
// applyDirNlink sets nlink = 2 + number_of_subdirectories for a directory.
|
||||
// Uses the in-memory subdirectory count tracked by mkdir/rmdir/rename.
|
||||
func (wfs *WFS) applyDirNlink(out *fuse.Attr, dirPath util.FullPath) {
|
||||
count := wfs.inodeToPath.GetSubdirCount(dirPath)
|
||||
if count > 0 {
|
||||
out.Nlink = 2 + uint32(count)
|
||||
}
|
||||
}
|
||||
|
||||
func chmod(existing uint32, mode uint32) uint32 {
|
||||
return existing&^07777 | mode&07777
|
||||
}
|
||||
|
||||
@@ -44,6 +44,13 @@ func (wfs *WFS) Lookup(cancel <-chan struct{}, header *fuse.InHeader, name strin
|
||||
|
||||
wfs.outputFilerEntry(out, inode, localEntry)
|
||||
|
||||
if localEntry.IsDirectory() {
|
||||
wfs.applyInMemoryDirMtime(&out.Attr, inode)
|
||||
if wfs.option.PosixDirNlink {
|
||||
wfs.applyDirNlink(&out.Attr, fullFilePath)
|
||||
}
|
||||
}
|
||||
|
||||
return fuse.OK
|
||||
|
||||
}
|
||||
|
||||
@@ -56,7 +56,6 @@ func (wfs *WFS) Mkdir(cancel <-chan struct{}, in *fuse.MkdirIn, name string, out
|
||||
// but BEFORE outputPbEntry writes attributes to the kernel. We restore
|
||||
// explicitly below instead of using defer so the kernel gets local values.
|
||||
|
||||
|
||||
request := &filer_pb.CreateEntryRequest{
|
||||
Directory: string(dirFullPath),
|
||||
Entry: newEntry,
|
||||
@@ -78,7 +77,8 @@ func (wfs *WFS) Mkdir(cancel <-chan struct{}, in *fuse.MkdirIn, name string, out
|
||||
wfs.inodeToPath.InvalidateChildrenCache(dirFullPath)
|
||||
}
|
||||
wfs.inodeToPath.TouchDirectory(dirFullPath)
|
||||
wfs.touchDirMtimeCtime(dirFullPath)
|
||||
wfs.touchDirMtimeCtimeBest(dirFullPath)
|
||||
wfs.inodeToPath.AdjustSubdirCount(dirFullPath, 1)
|
||||
}
|
||||
|
||||
glog.V(3).Infof("mkdir %s: %v", entryFullPath, err)
|
||||
@@ -95,6 +95,11 @@ func (wfs *WFS) Mkdir(cancel <-chan struct{}, in *fuse.MkdirIn, name string, out
|
||||
|
||||
inode := wfs.inodeToPath.Lookup(entryFullPath, newEntry.Attributes.Crtime, true, false, 0, true)
|
||||
|
||||
// The newly created directory is guaranteed to be empty, so mark it as
|
||||
// cached immediately to avoid a needless filer round-trip on the first
|
||||
// Lookup or ReadDir inside this directory.
|
||||
wfs.inodeToPath.MarkChildrenCached(entryFullPath)
|
||||
|
||||
wfs.outputPbEntry(out, inode, newEntry)
|
||||
|
||||
return fuse.OK
|
||||
@@ -155,7 +160,8 @@ func (wfs *WFS) Rmdir(cancel <-chan struct{}, header *fuse.InHeader, name string
|
||||
}
|
||||
wfs.inodeToPath.RemovePath(entryFullPath)
|
||||
wfs.inodeToPath.TouchDirectory(dirFullPath)
|
||||
wfs.touchDirMtimeCtime(dirFullPath)
|
||||
wfs.touchDirMtimeCtimeBest(dirFullPath)
|
||||
wfs.inodeToPath.AdjustSubdirCount(dirFullPath, -1)
|
||||
|
||||
return fuse.OK
|
||||
|
||||
|
||||
@@ -7,6 +7,8 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/seaweedfs/go-fuse/v2/fuse"
|
||||
"google.golang.org/protobuf/proto"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/cluster/lock_manager"
|
||||
"github.com/seaweedfs/seaweedfs/weed/filer"
|
||||
"github.com/seaweedfs/seaweedfs/weed/glog"
|
||||
@@ -69,7 +71,7 @@ func (wfs *WFS) Create(cancel <-chan struct{}, in *fuse.CreateIn, name string, o
|
||||
return code
|
||||
}
|
||||
|
||||
inode, newEntry, code = wfs.createRegularFile(dirFullPath, name, in.Mode, in.Uid, in.Gid, 0, true)
|
||||
inode, newEntry, code = wfs.createRegularFile(dirFullPath, name, in.Mode, in.Uid, in.Gid, 0, true, true)
|
||||
if code == fuse.Status(syscall.EEXIST) && in.Flags&syscall.O_EXCL == 0 {
|
||||
// Race: another process created the file between our check and create.
|
||||
// Reopen the winner's entry.
|
||||
@@ -147,7 +149,7 @@ func (wfs *WFS) Mknod(cancel <-chan struct{}, in *fuse.MknodIn, name string, out
|
||||
return
|
||||
}
|
||||
|
||||
inode, newEntry, code := wfs.createRegularFile(dirFullPath, name, in.Mode, in.Uid, in.Gid, in.Rdev, false)
|
||||
inode, newEntry, code := wfs.createRegularFile(dirFullPath, name, in.Mode, in.Uid, in.Gid, in.Rdev, false, false)
|
||||
if code != fuse.OK {
|
||||
return code
|
||||
}
|
||||
@@ -244,7 +246,7 @@ func (wfs *WFS) Unlink(cancel <-chan struct{}, header *fuse.InHeader, name strin
|
||||
wfs.inodeToPath.InvalidateChildrenCache(dirFullPath)
|
||||
}
|
||||
wfs.inodeToPath.TouchDirectory(dirFullPath)
|
||||
wfs.touchDirMtimeCtime(dirFullPath)
|
||||
wfs.touchDirMtimeCtimeBest(dirFullPath)
|
||||
|
||||
wfs.inodeToPath.RemovePath(entryFullPath)
|
||||
|
||||
@@ -252,7 +254,7 @@ func (wfs *WFS) Unlink(cancel <-chan struct{}, header *fuse.InHeader, name strin
|
||||
|
||||
}
|
||||
|
||||
func (wfs *WFS) createRegularFile(dirFullPath util.FullPath, name string, mode uint32, uid, gid, rdev uint32, deferFilerCreate bool) (inode uint64, newEntry *filer_pb.Entry, code fuse.Status) {
|
||||
func (wfs *WFS) createRegularFile(dirFullPath util.FullPath, name string, mode uint32, uid, gid, rdev uint32, deferFilerCreate bool, skipExistenceCheck bool) (inode uint64, newEntry *filer_pb.Entry, code fuse.Status) {
|
||||
if wfs.IsOverQuotaWithUncommitted() {
|
||||
return 0, nil, fuse.Status(syscall.ENOSPC)
|
||||
}
|
||||
@@ -276,10 +278,12 @@ func (wfs *WFS) createRegularFile(dirFullPath util.FullPath, name string, mode u
|
||||
}
|
||||
|
||||
entryFullPath := dirFullPath.Child(name)
|
||||
if _, status := wfs.maybeLoadEntry(entryFullPath); status == fuse.OK {
|
||||
return 0, nil, fuse.Status(syscall.EEXIST)
|
||||
} else if status != fuse.ENOENT {
|
||||
return 0, nil, status
|
||||
if !skipExistenceCheck {
|
||||
if _, status := wfs.maybeLoadEntry(entryFullPath); status == fuse.OK {
|
||||
return 0, nil, fuse.Status(syscall.EEXIST)
|
||||
} else if status != fuse.ENOENT {
|
||||
return 0, nil, status
|
||||
}
|
||||
}
|
||||
fileMode := toOsFileMode(mode)
|
||||
now := time.Now().Unix()
|
||||
@@ -301,18 +305,29 @@ func (wfs *WFS) createRegularFile(dirFullPath util.FullPath, name string, mode u
|
||||
},
|
||||
}
|
||||
|
||||
if deferFilerCreate {
|
||||
// Defer the filer gRPC call to flush time. The caller (Create) will
|
||||
// build a file handle directly from newEntry, bypassing AcquireHandle.
|
||||
if deferFilerCreate || wfs.option.WritebackCache {
|
||||
// Insert a local placeholder into the metadata cache so that
|
||||
// maybeLoadEntry() can find the file (e.g., duplicate-create checks,
|
||||
// stat, readdir). The actual filer entry is created by flushMetadataToFiler.
|
||||
// maybeLoadEntry() can find the file immediately (e.g., duplicate-
|
||||
// create checks, stat, readdir).
|
||||
// We use InsertEntry directly instead of applyLocalMetadataEvent to avoid
|
||||
// triggering directory hot-threshold eviction that would wipe the entry.
|
||||
if insertErr := wfs.metaCache.InsertEntry(context.Background(), filer.FromPbEntry(string(dirFullPath), newEntry)); insertErr != nil {
|
||||
glog.Warningf("createFile %s: insert local entry: %v", entryFullPath, insertErr)
|
||||
}
|
||||
glog.V(3).Infof("createFile %s: deferred to flush", entryFullPath)
|
||||
wfs.inodeToPath.TouchDirectory(dirFullPath)
|
||||
wfs.touchDirMtimeCtimeBest(dirFullPath)
|
||||
|
||||
if deferFilerCreate {
|
||||
// Fully deferred: the caller (Create) will build a file handle
|
||||
// directly from newEntry. The actual filer entry is created by
|
||||
// flushMetadataToFiler on close.
|
||||
glog.V(3).Infof("createFile %s: deferred to flush", entryFullPath)
|
||||
} else {
|
||||
// Async create: Mknod with writeback caching. The node is
|
||||
// visible locally; fire the filer RPC in the background.
|
||||
wfs.asyncCreateEntry(dirFullPath, newEntry)
|
||||
glog.V(3).Infof("createFile %s: async create", entryFullPath)
|
||||
}
|
||||
return inode, newEntry, fuse.OK
|
||||
}
|
||||
|
||||
@@ -340,7 +355,7 @@ func (wfs *WFS) createRegularFile(dirFullPath util.FullPath, name string, mode u
|
||||
wfs.inodeToPath.InvalidateChildrenCache(dirFullPath)
|
||||
}
|
||||
wfs.inodeToPath.TouchDirectory(dirFullPath)
|
||||
wfs.touchDirMtimeCtime(dirFullPath)
|
||||
wfs.touchDirMtimeCtimeBest(dirFullPath)
|
||||
}
|
||||
|
||||
glog.V(3).Infof("createFile %s: %v", entryFullPath, err)
|
||||
@@ -352,6 +367,51 @@ func (wfs *WFS) createRegularFile(dirFullPath util.FullPath, name string, mode u
|
||||
return inode, newEntry, fuse.OK
|
||||
}
|
||||
|
||||
// asyncCreateEntry sends a CreateEntry RPC to the filer in the background.
|
||||
// The entry is already in the local meta cache; this persists it to the filer.
|
||||
// Used by Mknod with writeback caching — the node is visible locally right away.
|
||||
//
|
||||
// If the filer RPC fails after retries, the local cache entry is removed so the
|
||||
// phantom file does not persist across cache invalidation or mount restart.
|
||||
func (wfs *WFS) asyncCreateEntry(dirFullPath util.FullPath, entry *filer_pb.Entry) {
|
||||
// Clone so the goroutine has its own copy for uid/gid mapping.
|
||||
requestEntry := proto.Clone(entry).(*filer_pb.Entry)
|
||||
dir := string(dirFullPath)
|
||||
entryPath := dirFullPath.Child(entry.Name)
|
||||
go func() {
|
||||
wfs.mapPbIdFromLocalToFiler(requestEntry)
|
||||
request := &filer_pb.CreateEntryRequest{
|
||||
Directory: dir,
|
||||
Entry: requestEntry,
|
||||
Signatures: []int32{wfs.signature},
|
||||
SkipCheckParentDirectory: true,
|
||||
}
|
||||
err := retryMetadataFlush(func() error {
|
||||
resp, createErr := wfs.streamCreateEntry(context.Background(), request)
|
||||
if createErr != nil {
|
||||
return createErr
|
||||
}
|
||||
event := resp.GetMetadataEvent()
|
||||
if event == nil {
|
||||
event = metadataCreateEvent(dir, requestEntry)
|
||||
}
|
||||
if applyErr := wfs.applyLocalMetadataEvent(context.Background(), event); applyErr != nil {
|
||||
glog.Warningf("async createFile %s: metadata apply: %v", entryPath, applyErr)
|
||||
wfs.inodeToPath.InvalidateChildrenCache(dirFullPath)
|
||||
}
|
||||
return nil
|
||||
}, func(nextAttempt, totalAttempts int, backoff time.Duration, err error) {
|
||||
glog.Warningf("async createFile %s: retrying (attempt %d/%d) after %v: %v",
|
||||
entryPath, nextAttempt, totalAttempts, backoff, err)
|
||||
})
|
||||
if err != nil {
|
||||
glog.Errorf("async createFile %s: failed after retries: %v — removing local entry", entryPath, err)
|
||||
wfs.metaCache.DeleteEntry(context.Background(), entryPath)
|
||||
wfs.inodeToPath.InvalidateChildrenCache(dirFullPath)
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
func (wfs *WFS) truncateEntry(entryFullPath util.FullPath, entry *filer_pb.Entry) fuse.Status {
|
||||
if entry == nil {
|
||||
return fuse.EIO
|
||||
|
||||
@@ -144,6 +144,13 @@ func (wfs *WFS) doFlush(fh *FileHandle, uid, gid uint32, allowAsync bool) fuse.S
|
||||
return fuse.OK
|
||||
}
|
||||
|
||||
// Skip metadata flush if the file was unlinked while open.
|
||||
// The filer entry is already gone; flushing would recreate it.
|
||||
if fh.isDeleted {
|
||||
glog.V(3).Infof("doFlush %s fh %d: file was unlinked, skipping metadata flush", fileFullPath, fh.fh)
|
||||
return fuse.OK
|
||||
}
|
||||
|
||||
if isOverQuota {
|
||||
return fuse.Status(syscall.ENOSPC)
|
||||
}
|
||||
|
||||
Some files were not shown because too many files have changed in this diff Show More
Reference in New Issue
Block a user