mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-11 00:55:51 +00:00
Compare commits
14
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
797c2366c5 | ||
|
|
d2723b75ca | ||
|
|
39ab2ef6ad | ||
|
|
5c5d377277 | ||
|
|
955a011f39 | ||
|
|
d074830016 | ||
|
|
e3359badfc | ||
|
|
44fea49816 | ||
|
|
e5ad5e8d4a | ||
|
|
0761be58d3 | ||
|
|
937a168d34 | ||
|
|
479e72b5ab | ||
|
|
46567ac06c | ||
|
|
00fcd5b828 |
@@ -23,6 +23,11 @@ on:
|
||||
- all
|
||||
- standard
|
||||
- large_disk
|
||||
publish:
|
||||
description: 'Publish images and manifests'
|
||||
required: true
|
||||
type: boolean
|
||||
default: false
|
||||
|
||||
permissions:
|
||||
contents: read
|
||||
@@ -33,6 +38,7 @@ jobs:
|
||||
runs-on: ubuntu-latest
|
||||
outputs:
|
||||
variants: ${{ steps.set-variants.outputs.variants }}
|
||||
publish: ${{ steps.set-publish.outputs.publish }}
|
||||
steps:
|
||||
- name: Select variants for this run
|
||||
id: set-variants
|
||||
@@ -43,6 +49,14 @@ jobs:
|
||||
variants='["standard","large_disk"]'
|
||||
fi
|
||||
echo "variants=$variants" >> "$GITHUB_OUTPUT"
|
||||
- name: Select publish mode
|
||||
id: set-publish
|
||||
run: |
|
||||
if [ "${{ github.event_name }}" = "workflow_dispatch" ]; then
|
||||
echo "publish=${{ github.event.inputs.publish }}" >> "$GITHUB_OUTPUT"
|
||||
else
|
||||
echo "publish=true" >> "$GITHUB_OUTPUT"
|
||||
fi
|
||||
|
||||
build:
|
||||
needs: [setup]
|
||||
@@ -112,13 +126,13 @@ jobs:
|
||||
buildkitd-flags: "--debug"
|
||||
buildkitd-config: /tmp/buildkitd.toml
|
||||
- name: Login to Docker Hub
|
||||
if: github.event_name != 'pull_request'
|
||||
if: needs.setup.outputs.publish == 'true'
|
||||
uses: docker/login-action@v4
|
||||
with:
|
||||
username: ${{ secrets.DOCKER_USERNAME }}
|
||||
password: ${{ secrets.DOCKER_PASSWORD }}
|
||||
- name: Login to GHCR
|
||||
if: github.event_name != 'pull_request'
|
||||
if: needs.setup.outputs.publish == 'true'
|
||||
uses: docker/login-action@v4
|
||||
with:
|
||||
registry: ghcr.io
|
||||
@@ -130,7 +144,7 @@ jobs:
|
||||
DOCKER_BUILDKIT: 1
|
||||
with:
|
||||
context: ./docker
|
||||
push: ${{ github.event_name != 'pull_request' }}
|
||||
push: ${{ needs.setup.outputs.publish == 'true' }}
|
||||
file: ./docker/Dockerfile.go_build
|
||||
platforms: linux/${{ matrix.platform }}
|
||||
# Push to GHCR only during build to avoid Docker Hub rate limits
|
||||
@@ -166,33 +180,112 @@ jobs:
|
||||
echo "tag_suffix=" >> $GITHUB_OUTPUT
|
||||
fi
|
||||
- name: Login to GHCR
|
||||
if: needs.setup.outputs.publish == 'true'
|
||||
uses: docker/login-action@v4
|
||||
with:
|
||||
registry: ghcr.io
|
||||
username: ${{ secrets.GHCR_USERNAME }}
|
||||
password: ${{ secrets.GHCR_TOKEN }}
|
||||
- name: Run Trivy vulnerability scanner
|
||||
# Pin to SHA — mutable tags were compromised (GHSA-69fq-xp46-6x23)
|
||||
- name: Checkout for local scan build
|
||||
if: needs.setup.outputs.publish != 'true'
|
||||
uses: actions/checkout@v6
|
||||
with:
|
||||
ref: ${{ github.event_name == 'workflow_dispatch' && github.event.inputs.source_ref || github.ref }}
|
||||
- name: Create BuildKit config for local scan build
|
||||
if: needs.setup.outputs.publish != 'true'
|
||||
run: |
|
||||
cat > /tmp/buildkitd.toml <<EOF
|
||||
[registry."docker.io"]
|
||||
mirrors = ["https://mirror.gcr.io"]
|
||||
EOF
|
||||
- name: Set up Docker Buildx for local scan build
|
||||
if: needs.setup.outputs.publish != 'true'
|
||||
uses: docker/setup-buildx-action@v4
|
||||
with:
|
||||
buildkitd-flags: "--debug"
|
||||
buildkitd-config: /tmp/buildkitd.toml
|
||||
- name: Build local scan image tarball
|
||||
if: needs.setup.outputs.publish != 'true'
|
||||
uses: docker/build-push-action@v7
|
||||
env:
|
||||
DOCKER_BUILDKIT: 1
|
||||
with:
|
||||
context: ./docker
|
||||
file: ./docker/Dockerfile.go_build
|
||||
platforms: linux/amd64
|
||||
outputs: type=docker,dest=/tmp/seaweedfs${{ steps.config.outputs.tag_suffix }}-amd64.tar
|
||||
build-args: |
|
||||
BUILDKIT_INLINE_CACHE=1
|
||||
BRANCH=${{ github.event_name == 'workflow_dispatch' && github.event.inputs.source_ref || github.sha }}
|
||||
${{ matrix.variant == 'large_disk' && 'TAGS=5BytesOffset' || '' }}
|
||||
- name: Trivy report (published image)
|
||||
if: needs.setup.outputs.publish == 'true'
|
||||
# Pin to SHA - mutable tags were compromised (GHSA-69fq-xp46-6x23)
|
||||
uses: aquasecurity/trivy-action@57a97c7e7821a5776cebc9bb87c984fa69cba8f1 # v0.35.0
|
||||
with:
|
||||
# Scan amd64 only — OS packages are identical across architectures
|
||||
scan-type: image
|
||||
# Scan amd64 only - OS packages are identical across architectures
|
||||
# since they all use the same alpine base, so a single-arch scan
|
||||
# provides sufficient coverage without multiplying CI time.
|
||||
image-ref: ghcr.io/chrislusf/seaweedfs:${{ github.event_name == 'workflow_dispatch' && github.event.inputs.image_tag || 'latest' }}${{ steps.config.outputs.tag_suffix }}-amd64
|
||||
scanners: vuln
|
||||
vuln-type: os,library
|
||||
severity: HIGH,CRITICAL
|
||||
ignore-unfixed: true
|
||||
limit-severities-for-sarif: true
|
||||
format: sarif
|
||||
output: trivy-results.sarif
|
||||
exit-code: '0'
|
||||
- name: Trivy report (local tarball)
|
||||
if: needs.setup.outputs.publish != 'true'
|
||||
uses: aquasecurity/trivy-action@57a97c7e7821a5776cebc9bb87c984fa69cba8f1 # v0.35.0
|
||||
with:
|
||||
input: /tmp/seaweedfs${{ steps.config.outputs.tag_suffix }}-amd64.tar
|
||||
scanners: vuln
|
||||
vuln-type: os,library
|
||||
severity: HIGH,CRITICAL
|
||||
exit-code: '1'
|
||||
ignore-unfixed: true
|
||||
limit-severities-for-sarif: true
|
||||
format: sarif
|
||||
output: trivy-results.sarif
|
||||
exit-code: '0'
|
||||
- name: Upload Trivy scan results to GitHub Security
|
||||
uses: github/codeql-action/upload-sarif@v3
|
||||
uses: github/codeql-action/upload-sarif@v4
|
||||
if: always()
|
||||
with:
|
||||
sarif_file: trivy-results.sarif
|
||||
- name: Trivy gate (published image)
|
||||
if: needs.setup.outputs.publish == 'true'
|
||||
# Gate only on fixable high/critical vulnerabilities. Non-fixable
|
||||
# findings are still visible in the SARIF upload above.
|
||||
uses: aquasecurity/trivy-action@57a97c7e7821a5776cebc9bb87c984fa69cba8f1 # v0.35.0
|
||||
with:
|
||||
scan-type: image
|
||||
image-ref: ghcr.io/chrislusf/seaweedfs:${{ github.event_name == 'workflow_dispatch' && github.event.inputs.image_tag || 'latest' }}${{ steps.config.outputs.tag_suffix }}-amd64
|
||||
scanners: vuln
|
||||
vuln-type: os,library
|
||||
severity: HIGH,CRITICAL
|
||||
ignore-unfixed: true
|
||||
format: table
|
||||
exit-code: '1'
|
||||
skip-setup-trivy: true
|
||||
- name: Trivy gate (local tarball)
|
||||
if: needs.setup.outputs.publish != 'true'
|
||||
uses: aquasecurity/trivy-action@57a97c7e7821a5776cebc9bb87c984fa69cba8f1 # v0.35.0
|
||||
with:
|
||||
input: /tmp/seaweedfs${{ steps.config.outputs.tag_suffix }}-amd64.tar
|
||||
scanners: vuln
|
||||
vuln-type: os,library
|
||||
severity: HIGH,CRITICAL
|
||||
ignore-unfixed: true
|
||||
format: table
|
||||
exit-code: '1'
|
||||
skip-setup-trivy: true
|
||||
|
||||
create-manifest:
|
||||
runs-on: ubuntu-latest
|
||||
needs: [setup, build, trivy-scan]
|
||||
if: github.event_name != 'pull_request'
|
||||
if: needs.setup.outputs.publish == 'true' && github.event_name != 'pull_request'
|
||||
strategy:
|
||||
matrix:
|
||||
variant: ${{ fromJSON(needs.setup.outputs.variants) }}
|
||||
|
||||
@@ -16,15 +16,16 @@ RUN cd /go/src/github.com/seaweedfs/seaweedfs/weed \
|
||||
&& export LDFLAGS="-X github.com/seaweedfs/seaweedfs/weed/util/version.COMMIT=$(git rev-parse --short HEAD)" \
|
||||
&& CGO_ENABLED=0 go install -tags "$TAGS" -ldflags "-extldflags -static ${LDFLAGS}"
|
||||
|
||||
# Rust volume server builder (amd64/arm64 only)
|
||||
FROM rust:1-alpine as rust_builder
|
||||
# Rust volume server builder. Alpine packages avoid depending on the
|
||||
# upstream rust:alpine manifest list, which no longer includes linux/386.
|
||||
FROM alpine:3.23 as rust_builder
|
||||
ARG TARGETARCH
|
||||
RUN apk add musl-dev protobuf-dev git
|
||||
COPY --from=builder /go/src/github.com/seaweedfs/seaweedfs/seaweed-volume /build/seaweed-volume
|
||||
COPY --from=builder /go/src/github.com/seaweedfs/seaweedfs/proto /build/proto
|
||||
COPY --from=builder /go/src/github.com/seaweedfs/seaweedfs/weed /build/weed
|
||||
WORKDIR /build/seaweed-volume
|
||||
ARG TAGS
|
||||
RUN if [ "$TARGETARCH" = "amd64" ] || [ "$TARGETARCH" = "arm64" ]; then \
|
||||
apk add --no-cache musl-dev openssl-dev protobuf-dev git rust cargo; \
|
||||
if [ "$TAGS" = "5BytesOffset" ]; then \
|
||||
cargo build --release; \
|
||||
else \
|
||||
@@ -79,9 +80,5 @@ RUN mkdir -p /data/filerldb2 && \
|
||||
VOLUME /data
|
||||
WORKDIR /data
|
||||
|
||||
# Run as non-root by default (satisfies security scanners).
|
||||
# Use `docker run --user root` if you need the entrypoint to fix
|
||||
# /data volume ownership before dropping privileges.
|
||||
USER seaweed
|
||||
|
||||
# Entrypoint will handle permission fixes and user switching
|
||||
ENTRYPOINT ["/entrypoint.sh"]
|
||||
|
||||
@@ -37,9 +37,5 @@ RUN mkdir -p /data/filerldb2 && \
|
||||
VOLUME /data
|
||||
WORKDIR /data
|
||||
|
||||
# Run as non-root by default (satisfies security scanners).
|
||||
# Use `docker run --user root` if you need the entrypoint to fix
|
||||
# /data volume ownership before dropping privileges.
|
||||
USER seaweed
|
||||
|
||||
# Entrypoint will handle permission fixes and user switching
|
||||
ENTRYPOINT ["/entrypoint.sh"]
|
||||
|
||||
@@ -83,7 +83,7 @@ require (
|
||||
github.com/valyala/bytebufferpool v1.0.0
|
||||
github.com/viant/ptrie v1.0.1
|
||||
github.com/xdg-go/pbkdf2 v1.0.0 // indirect
|
||||
github.com/xdg-go/scram v1.1.2 // indirect
|
||||
github.com/xdg-go/scram v1.1.2
|
||||
github.com/xdg-go/stringprep v1.0.4 // indirect
|
||||
github.com/youmark/pkcs8 v0.0.0-20240726163527-a2c0da244d78 // indirect
|
||||
go.etcd.io/etcd/client/v3 v3.6.7
|
||||
@@ -94,7 +94,7 @@ require (
|
||||
gocloud.dev/pubsub/rabbitpubsub v0.45.0
|
||||
golang.org/x/crypto v0.49.0
|
||||
golang.org/x/exp v0.0.0-20260218203240-3dfff04db8fa
|
||||
golang.org/x/image v0.36.0
|
||||
golang.org/x/image v0.38.0
|
||||
golang.org/x/net v0.51.0
|
||||
golang.org/x/oauth2 v0.36.0
|
||||
golang.org/x/sys v0.42.0
|
||||
|
||||
@@ -1838,8 +1838,6 @@ github.com/schollz/progressbar/v3 v3.19.0 h1:Ea18xuIRQXLAUidVDox3AbwfUhD0/1Ivohy
|
||||
github.com/schollz/progressbar/v3 v3.19.0/go.mod h1:IsO3lpbaGuzh8zIMzgY3+J8l4C8GjO0Y9S69eFvNsec=
|
||||
github.com/seaweedfs/cockroachdb-parser v0.0.0-20260225204133-2f342c5ea564 h1:TgxPraf1NmF6XTcUG53ULpLQrKvhtUJxQ3hyekxSDNQ=
|
||||
github.com/seaweedfs/cockroachdb-parser v0.0.0-20260225204133-2f342c5ea564/go.mod h1:JSKCh6uCHBz91lQYFYHCyTrSVIPge4SUFVn28iwMNB0=
|
||||
github.com/seaweedfs/go-fuse/v2 v2.9.1 h1:gnKmfrKreCRGJmekGz5WMnNZqXEf9s9+V2hdWQdvx88=
|
||||
github.com/seaweedfs/go-fuse/v2 v2.9.1/go.mod h1:zABdmWEa6A0bwaBeEOBUeUkGIZlxUhcdv+V1Dcc/U/I=
|
||||
github.com/seaweedfs/go-fuse/v2 v2.9.2 h1:IfP/yFjLGO4rALcJY2Gb39PlebHxLnj7dkIiQAjFres=
|
||||
github.com/seaweedfs/go-fuse/v2 v2.9.2/go.mod h1:zABdmWEa6A0bwaBeEOBUeUkGIZlxUhcdv+V1Dcc/U/I=
|
||||
github.com/seaweedfs/goexif v1.0.3 h1:ve/OjI7dxPW8X9YQsv3JuVMaxEyF9Rvfd04ouL+Bz30=
|
||||
@@ -2242,8 +2240,8 @@ golang.org/x/image v0.0.0-20210607152325-775e3b0c77b9/go.mod h1:023OzeP/+EPmXeap
|
||||
golang.org/x/image v0.0.0-20210628002857-a66eb6448b8d/go.mod h1:023OzeP/+EPmXeapQh35lcL3II3LrY8Ic+EFFKVhULM=
|
||||
golang.org/x/image v0.0.0-20211028202545-6944b10bf410/go.mod h1:023OzeP/+EPmXeapQh35lcL3II3LrY8Ic+EFFKVhULM=
|
||||
golang.org/x/image v0.0.0-20220302094943-723b81ca9867/go.mod h1:023OzeP/+EPmXeapQh35lcL3II3LrY8Ic+EFFKVhULM=
|
||||
golang.org/x/image v0.36.0 h1:Iknbfm1afbgtwPTmHnS2gTM/6PPZfH+z2EFuOkSbqwc=
|
||||
golang.org/x/image v0.36.0/go.mod h1:YsWD2TyyGKiIX1kZlu9QfKIsQ4nAAK9bdgdrIsE7xy4=
|
||||
golang.org/x/image v0.38.0 h1:5l+q+Y9JDC7mBOMjo4/aPhMDcxEptsX+Tt3GgRQRPuE=
|
||||
golang.org/x/image v0.38.0/go.mod h1:/3f6vaXC+6CEanU4KJxbcUZyEePbyKbaLoDOe4ehFYY=
|
||||
golang.org/x/lint v0.0.0-20181026193005-c67002cb31c3/go.mod h1:UVdnD1Gm6xHRNCYTkRU2/jEulfH38KcIWyp/GAMgvoE=
|
||||
golang.org/x/lint v0.0.0-20190227174305-5b3e6a55c961/go.mod h1:wehouNa3lNwaWXcvxsM5YxQ5yQlVC4a0KAMCusXpPoU=
|
||||
golang.org/x/lint v0.0.0-20190301231843-5614ed5bae6f/go.mod h1:UVdnD1Gm6xHRNCYTkRU2/jEulfH38KcIWyp/GAMgvoE=
|
||||
|
||||
+20
-20
@@ -14,7 +14,7 @@ require (
|
||||
replace github.com/seaweedfs/seaweedfs => ../../
|
||||
|
||||
require (
|
||||
cloud.google.com/go/auth v0.17.0 // indirect
|
||||
cloud.google.com/go/auth v0.18.1 // indirect
|
||||
cloud.google.com/go/auth/oauth2adapt v0.2.8 // indirect
|
||||
cloud.google.com/go/compute/metadata v0.9.0 // indirect
|
||||
github.com/Azure/azure-sdk-for-go/sdk/azcore v1.21.0 // indirect
|
||||
@@ -45,19 +45,19 @@ require (
|
||||
github.com/aws/aws-sdk-go v1.55.8 // indirect
|
||||
github.com/aws/aws-sdk-go-v2 v1.41.4 // indirect
|
||||
github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.4 // indirect
|
||||
github.com/aws/aws-sdk-go-v2/config v1.32.7 // indirect
|
||||
github.com/aws/aws-sdk-go-v2/config v1.32.9 // indirect
|
||||
github.com/aws/aws-sdk-go-v2/credentials v1.19.12 // indirect
|
||||
github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.20 // indirect
|
||||
github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.20.12 // indirect
|
||||
github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.20 // indirect
|
||||
github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.20 // indirect
|
||||
github.com/aws/aws-sdk-go-v2/internal/ini v1.8.4 // indirect
|
||||
github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.16 // indirect
|
||||
github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.17 // indirect
|
||||
github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.7 // indirect
|
||||
github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.7 // indirect
|
||||
github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.8 // indirect
|
||||
github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.20 // indirect
|
||||
github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.16 // indirect
|
||||
github.com/aws/aws-sdk-go-v2/service/s3 v1.95.0 // indirect
|
||||
github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.17 // indirect
|
||||
github.com/aws/aws-sdk-go-v2/service/s3 v1.96.0 // indirect
|
||||
github.com/aws/aws-sdk-go-v2/service/signin v1.0.8 // indirect
|
||||
github.com/aws/aws-sdk-go-v2/service/sso v1.30.13 // indirect
|
||||
github.com/aws/aws-sdk-go-v2/service/ssooidc v1.35.17 // indirect
|
||||
@@ -124,8 +124,8 @@ require (
|
||||
github.com/google/btree v1.1.3 // indirect
|
||||
github.com/google/s2a-go v0.1.9 // indirect
|
||||
github.com/google/uuid v1.6.0 // indirect
|
||||
github.com/googleapis/enterprise-certificate-proxy v0.3.7 // indirect
|
||||
github.com/googleapis/gax-go/v2 v2.15.0 // indirect
|
||||
github.com/googleapis/enterprise-certificate-proxy v0.3.11 // indirect
|
||||
github.com/googleapis/gax-go/v2 v2.17.0 // indirect
|
||||
github.com/gorilla/mux v1.8.1 // indirect
|
||||
github.com/gorilla/schema v1.4.1 // indirect
|
||||
github.com/hashicorp/errwrap v1.1.0 // indirect
|
||||
@@ -147,9 +147,9 @@ require (
|
||||
github.com/jtolio/noiseconn v0.0.0-20231127013910-f6d9ecbf1de7 // indirect
|
||||
github.com/jzelinskie/whirlpool v0.0.0-20201016144138-0675e54bb004 // indirect
|
||||
github.com/karlseguin/ccache/v2 v2.0.8 // indirect
|
||||
github.com/klauspost/compress v1.18.4 // indirect
|
||||
github.com/klauspost/compress v1.18.5 // indirect
|
||||
github.com/klauspost/cpuid/v2 v2.3.0 // indirect
|
||||
github.com/klauspost/reedsolomon v1.13.0 // indirect
|
||||
github.com/klauspost/reedsolomon v1.13.3 // indirect
|
||||
github.com/koofr/go-httpclient v0.0.0-20240520111329-e20f8f203988 // indirect
|
||||
github.com/koofr/go-koofrclient v0.0.0-20221207135200-cbd7fc9ad6a6 // indirect
|
||||
github.com/kr/fs v0.1.0 // indirect
|
||||
@@ -237,7 +237,7 @@ require (
|
||||
github.com/yusufpapurcu/wmi v1.2.4 // indirect
|
||||
github.com/zeebo/blake3 v0.2.4 // indirect
|
||||
github.com/zeebo/errs v1.4.0 // indirect
|
||||
github.com/zeebo/xxh3 v1.0.2 // indirect
|
||||
github.com/zeebo/xxh3 v1.1.0 // indirect
|
||||
go.etcd.io/bbolt v1.4.3 // indirect
|
||||
go.mongodb.org/mongo-driver v1.17.9 // indirect
|
||||
go.opentelemetry.io/auto/sdk v1.2.1 // indirect
|
||||
@@ -247,18 +247,18 @@ require (
|
||||
go.opentelemetry.io/otel/trace v1.40.0 // indirect
|
||||
go.yaml.in/yaml/v2 v2.4.3 // indirect
|
||||
go.yaml.in/yaml/v3 v3.0.4 // indirect
|
||||
golang.org/x/crypto v0.48.0 // indirect
|
||||
golang.org/x/exp v0.0.0-20251023183803-a4bb9ffd2546 // indirect
|
||||
golang.org/x/image v0.36.0 // indirect
|
||||
golang.org/x/net v0.49.0 // indirect
|
||||
golang.org/x/crypto v0.49.0 // indirect
|
||||
golang.org/x/exp v0.0.0-20260218203240-3dfff04db8fa // indirect
|
||||
golang.org/x/image v0.38.0 // indirect
|
||||
golang.org/x/net v0.51.0 // indirect
|
||||
golang.org/x/oauth2 v0.36.0 // indirect
|
||||
golang.org/x/sync v0.19.0 // indirect
|
||||
golang.org/x/sync v0.20.0 // indirect
|
||||
golang.org/x/sys v0.42.0 // indirect
|
||||
golang.org/x/term v0.40.0 // indirect
|
||||
golang.org/x/text v0.34.0 // indirect
|
||||
golang.org/x/term v0.41.0 // indirect
|
||||
golang.org/x/text v0.35.0 // indirect
|
||||
golang.org/x/time v0.14.0 // indirect
|
||||
google.golang.org/api v0.258.0 // indirect
|
||||
google.golang.org/genproto/googleapis/rpc v0.0.0-20251213004720-97cd9d5aeac2 // indirect
|
||||
google.golang.org/api v0.267.0 // indirect
|
||||
google.golang.org/genproto/googleapis/rpc v0.0.0-20260203192932-546029d2fa20 // indirect
|
||||
google.golang.org/grpc/security/advancedtls v1.0.0 // indirect
|
||||
google.golang.org/protobuf v1.36.11 // indirect
|
||||
gopkg.in/natefinch/lumberjack.v2 v2.2.1 // indirect
|
||||
|
||||
+48
-48
@@ -13,8 +13,8 @@ cloud.google.com/go v0.56.0/go.mod h1:jr7tqZxxKOVYizybht9+26Z/gUq7tiRzu+ACVAMbKV
|
||||
cloud.google.com/go v0.57.0/go.mod h1:oXiQ6Rzq3RAkkY7N6t3TcE6jE+CIBBbA36lwQ1JyzZs=
|
||||
cloud.google.com/go v0.62.0/go.mod h1:jmCYTdRCQuc1PHIIJ/maLInMho30T/Y0M4hTdTShOYc=
|
||||
cloud.google.com/go v0.65.0/go.mod h1:O5N8zS7uWy9vkA9vayVHs65eM1ubvY4h553ofrNHObY=
|
||||
cloud.google.com/go/auth v0.17.0 h1:74yCm7hCj2rUyyAocqnFzsAYXgJhrG26XCFimrc/Kz4=
|
||||
cloud.google.com/go/auth v0.17.0/go.mod h1:6wv/t5/6rOPAX4fJiRjKkJCvswLwdet7G8+UGXt7nCQ=
|
||||
cloud.google.com/go/auth v0.18.1 h1:IwTEx92GFUo2pJ6Qea0EU3zYvKnTAeRCODxfA/G5UWs=
|
||||
cloud.google.com/go/auth v0.18.1/go.mod h1:GfTYoS9G3CWpRA3Va9doKN9mjPGRS+v41jmZAhBzbrA=
|
||||
cloud.google.com/go/auth/oauth2adapt v0.2.8 h1:keo8NaayQZ6wimpNSmW5OPc283g65QNIiLpZnkHRbnc=
|
||||
cloud.google.com/go/auth/oauth2adapt v0.2.8/go.mod h1:XQ9y31RkqZCcwJWNSx2Xvric3RrU88hAYYbjDWYDL+c=
|
||||
cloud.google.com/go/bigquery v1.0.1/go.mod h1:i/xbL2UlR5RvWAURpBYZTtm/cXjCha9lbfbpx4poX+o=
|
||||
@@ -116,8 +116,8 @@ github.com/aws/aws-sdk-go-v2 v1.41.4 h1:10f50G7WyU02T56ox1wWXq+zTX9I1zxG46HYuG1h
|
||||
github.com/aws/aws-sdk-go-v2 v1.41.4/go.mod h1:mwsPRE8ceUUpiTgF7QmQIJ7lgsKUPQOUl3o72QBrE1o=
|
||||
github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.4 h1:489krEF9xIGkOaaX3CE/Be2uWjiXrkCH6gUX+bZA/BU=
|
||||
github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.4/go.mod h1:IOAPF6oT9KCsceNTvvYMNHy0+kMF8akOjeDvPENWxp4=
|
||||
github.com/aws/aws-sdk-go-v2/config v1.32.7 h1:vxUyWGUwmkQ2g19n7JY/9YL8MfAIl7bTesIUykECXmY=
|
||||
github.com/aws/aws-sdk-go-v2/config v1.32.7/go.mod h1:2/Qm5vKUU/r7Y+zUk/Ptt2MDAEKAfUtKc1+3U1Mo3oY=
|
||||
github.com/aws/aws-sdk-go-v2/config v1.32.9 h1:ktda/mtAydeObvJXlHzyGpK1xcsLaP16zfUPDGoW90A=
|
||||
github.com/aws/aws-sdk-go-v2/config v1.32.9/go.mod h1:U+fCQ+9QKsLW786BCfEjYRj34VVTbPdsLP3CHSYXMOI=
|
||||
github.com/aws/aws-sdk-go-v2/credentials v1.19.12 h1:oqtA6v+y5fZg//tcTWahyN9PEn5eDU/Wpvc2+kJ4aY8=
|
||||
github.com/aws/aws-sdk-go-v2/credentials v1.19.12/go.mod h1:U3R1RtSHx6NB0DvEQFGyf/0sbrpJrluENHdPy1j/3TE=
|
||||
github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.20 h1:zOgq3uezl5nznfoK3ODuqbhVg1JzAGDUhXOsU0IDCAo=
|
||||
@@ -130,18 +130,18 @@ github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.20 h1:tN6W/hg+pkM+tf9XDk
|
||||
github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.20/go.mod h1:YJ898MhD067hSHA6xYCx5ts/jEd8BSOLtQDL3iZsvbc=
|
||||
github.com/aws/aws-sdk-go-v2/internal/ini v1.8.4 h1:WKuaxf++XKWlHWu9ECbMlha8WOEGm0OUEZqm4K/Gcfk=
|
||||
github.com/aws/aws-sdk-go-v2/internal/ini v1.8.4/go.mod h1:ZWy7j6v1vWGmPReu0iSGvRiise4YI5SkR3OHKTZ6Wuc=
|
||||
github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.16 h1:CjMzUs78RDDv4ROu3JnJn/Ig1r6ZD7/T2DXLLRpejic=
|
||||
github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.16/go.mod h1:uVW4OLBqbJXSHJYA9svT9BluSvvwbzLQ2Crf6UPzR3c=
|
||||
github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.17 h1:JqcdRG//czea7Ppjb+g/n4o8i/R50aTBHkA7vu0lK+k=
|
||||
github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.17/go.mod h1:CO+WeGmIdj/MlPel2KwID9Gt7CNq4M65HUfBW97liM0=
|
||||
github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.7 h1:5EniKhLZe4xzL7a+fU3C2tfUN4nWIqlLesfrjkuPFTY=
|
||||
github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.7/go.mod h1:x0nZssQ3qZSnIcePWLvcoFisRXJzcTVvYpAAdYX8+GI=
|
||||
github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.7 h1:DIBqIrJ7hv+e4CmIk2z3pyKT+3B6qVMgRsawHiR3qso=
|
||||
github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.7/go.mod h1:vLm00xmBke75UmpNvOcZQ/Q30ZFjbczeLFqGx5urmGo=
|
||||
github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.8 h1:Z5EiPIzXKewUQK0QTMkutjiaPVeVYXX7KIqhXu/0fXs=
|
||||
github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.8/go.mod h1:FsTpJtvC4U1fyDXk7c71XoDv3HlRm8V3NiYLeYLh5YE=
|
||||
github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.20 h1:2HvVAIq+YqgGotK6EkMf+KIEqTISmTYh5zLpYyeTo1Y=
|
||||
github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.20/go.mod h1:V4X406Y666khGa8ghKmphma/7C0DAtEQYhkq9z4vpbk=
|
||||
github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.16 h1:NSbvS17MlI2lurYgXnCOLvCFX38sBW4eiVER7+kkgsU=
|
||||
github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.16/go.mod h1:SwT8Tmqd4sA6G1qaGdzWCJN99bUmPGHfRwwq3G5Qb+A=
|
||||
github.com/aws/aws-sdk-go-v2/service/s3 v1.95.0 h1:MIWra+MSq53CFaXXAywB2qg9YvVZifkk6vEGl/1Qor0=
|
||||
github.com/aws/aws-sdk-go-v2/service/s3 v1.95.0/go.mod h1:79S2BdqCJpScXZA2y+cpZuocWsjGjJINyXnOsf5DTz8=
|
||||
github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.17 h1:bGeHBsGZx0Dvu/eJC0Lh9adJa3M1xREcndxLNZlve2U=
|
||||
github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.17/go.mod h1:dcW24lbU0CzHusTE8LLHhRLI42ejmINN8Lcr22bwh/g=
|
||||
github.com/aws/aws-sdk-go-v2/service/s3 v1.96.0 h1:oeu8VPlOre74lBA/PMhxa5vewaMIMmILM+RraSyB8KA=
|
||||
github.com/aws/aws-sdk-go-v2/service/s3 v1.96.0/go.mod h1:5jggDlZ2CLQhwJBiZJb4vfk4f0GxWdEDruWKEJ1xOdo=
|
||||
github.com/aws/aws-sdk-go-v2/service/signin v1.0.8 h1:0GFOLzEbOyZABS3PhYfBIx2rNBACYcKty+XGkTgw1ow=
|
||||
github.com/aws/aws-sdk-go-v2/service/signin v1.0.8/go.mod h1:LXypKvk85AROkKhOG6/YEcHFPoX+prKTowKnVdcaIxE=
|
||||
github.com/aws/aws-sdk-go-v2/service/sso v1.30.13 h1:kiIDLZ005EcKomYYITtfsjn7dtOwHDOFy7IbPXKek2o=
|
||||
@@ -385,12 +385,12 @@ github.com/google/s2a-go v0.1.9 h1:LGD7gtMgezd8a/Xak7mEWL0PjoTQFvpRudN895yqKW0=
|
||||
github.com/google/s2a-go v0.1.9/go.mod h1:YA0Ei2ZQL3acow2O62kdp9UlnvMmU7kA6Eutn0dXayM=
|
||||
github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0=
|
||||
github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
|
||||
github.com/googleapis/enterprise-certificate-proxy v0.3.7 h1:zrn2Ee/nWmHulBx5sAVrGgAa0f2/R35S4DJwfFaUPFQ=
|
||||
github.com/googleapis/enterprise-certificate-proxy v0.3.7/go.mod h1:MkHOF77EYAE7qfSuSS9PU6g4Nt4e11cnsDUowfwewLA=
|
||||
github.com/googleapis/enterprise-certificate-proxy v0.3.11 h1:vAe81Msw+8tKUxi2Dqh/NZMz7475yUvmRIkXr4oN2ao=
|
||||
github.com/googleapis/enterprise-certificate-proxy v0.3.11/go.mod h1:RFV7MUdlb7AgEq2v7FmMCfeSMCllAzWxFgRdusoGks8=
|
||||
github.com/googleapis/gax-go/v2 v2.0.4/go.mod h1:0Wqv26UfaUD9n4G6kQubkQ+KchISgw+vpHVxEJEs9eg=
|
||||
github.com/googleapis/gax-go/v2 v2.0.5/go.mod h1:DWXyrwAJ9X0FpwwEdw+IPEYBICEFu5mhpdKc/us6bOk=
|
||||
github.com/googleapis/gax-go/v2 v2.15.0 h1:SyjDc1mGgZU5LncH8gimWo9lW1DtIfPibOG81vgd/bo=
|
||||
github.com/googleapis/gax-go/v2 v2.15.0/go.mod h1:zVVkkxAQHa1RQpg9z2AUCMnKhi0Qld9rcmyfL1OZhoc=
|
||||
github.com/googleapis/gax-go/v2 v2.17.0 h1:RksgfBpxqff0EZkDWYuz9q/uWsTVz+kf43LsZ1J6SMc=
|
||||
github.com/googleapis/gax-go/v2 v2.17.0/go.mod h1:mzaqghpQp4JDh3HvADwrat+6M3MOIDp5YKHhb9PAgDY=
|
||||
github.com/gopherjs/gopherjs v1.17.2 h1:fQnZVsXk8uxXIStYb0N4bGk7jeyTalG/wsZjQ25dO0g=
|
||||
github.com/gopherjs/gopherjs v1.17.2/go.mod h1:pRRIvn/QzFLrKfvEz3qUuEhtE/zLCWfreZ6J5gM2i+k=
|
||||
github.com/gorilla/mux v1.8.1 h1:TuBL49tXwgrFYWhqrNgrUNEY92u81SPhu7sTdzQEiWY=
|
||||
@@ -467,12 +467,12 @@ github.com/keybase/go-keychain v0.0.1 h1:way+bWYa6lDppZoZcgMbYsvC7GxljxrskdNInRt
|
||||
github.com/keybase/go-keychain v0.0.1/go.mod h1:PdEILRW3i9D8JcdM+FmY6RwkHGnhHxXwkPPMeUgOK1k=
|
||||
github.com/kisielk/errcheck v1.5.0/go.mod h1:pFxgyoBC7bSaBwPgfKdkLd5X25qrDl4LWUI2bnpBCr8=
|
||||
github.com/kisielk/gotool v1.0.0/go.mod h1:XhKaO+MFFWcvkIS/tQcRk01m1F5IRFswLeQ+oQHNcck=
|
||||
github.com/klauspost/compress v1.18.4 h1:RPhnKRAQ4Fh8zU2FY/6ZFDwTVTxgJ/EMydqSTzE9a2c=
|
||||
github.com/klauspost/compress v1.18.4/go.mod h1:R0h/fSBs8DE4ENlcrlib3PsXS61voFxhIs2DeRhCvJ4=
|
||||
github.com/klauspost/compress v1.18.5 h1:/h1gH5Ce+VWNLSWqPzOVn6XBO+vJbCNGvjoaGBFW2IE=
|
||||
github.com/klauspost/compress v1.18.5/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ=
|
||||
github.com/klauspost/cpuid/v2 v2.3.0 h1:S4CRMLnYUhGeDFDqkGriYKdfoFlDnMtqTiI/sFzhA9Y=
|
||||
github.com/klauspost/cpuid/v2 v2.3.0/go.mod h1:hqwkgyIinND0mEev00jJYCxPNVRVXFQeu1XKlok6oO0=
|
||||
github.com/klauspost/reedsolomon v1.13.0 h1:E0Cmgf2kMuhZTj6eefnvpKC4/Q4jhCi9YIjcZjK4arc=
|
||||
github.com/klauspost/reedsolomon v1.13.0/go.mod h1:ggJT9lc71Vu+cSOPBlxGvBN6TfAS77qB4fp8vJ05NSA=
|
||||
github.com/klauspost/reedsolomon v1.13.3 h1:01GwnO2xoCSaM0ShP4qwl+FsHg3csFShC6Tu/RS1ji0=
|
||||
github.com/klauspost/reedsolomon v1.13.3/go.mod h1:yjqqjgMTQkBUHSG97/rm4zipffCNbCiZcB3kTqr++sQ=
|
||||
github.com/koofr/go-httpclient v0.0.0-20240520111329-e20f8f203988 h1:CjEMN21Xkr9+zwPmZPaJJw+apzVbjGL5uK/6g9Q2jGU=
|
||||
github.com/koofr/go-httpclient v0.0.0-20240520111329-e20f8f203988/go.mod h1:/agobYum3uo/8V6yPVnq+R82pyVGCeuWW5arT4Txn8A=
|
||||
github.com/koofr/go-koofrclient v0.0.0-20221207135200-cbd7fc9ad6a6 h1:FHVoZMOVRA+6/y4yRlbiR3WvsrOcKBd/f64H7YiWR2U=
|
||||
@@ -739,8 +739,8 @@ github.com/zeebo/errs v1.4.0 h1:XNdoD/RRMKP7HD0UhJnIzUy74ISdGGxURlYG8HSWSfM=
|
||||
github.com/zeebo/errs v1.4.0/go.mod h1:sgbWHsvVuTPHcqJJGQ1WhI5KbWlHYz+2+2C/LSEtCw4=
|
||||
github.com/zeebo/pcg v1.0.1 h1:lyqfGeWiv4ahac6ttHs+I5hwtH/+1mrhlCtVNQM2kHo=
|
||||
github.com/zeebo/pcg v1.0.1/go.mod h1:09F0S9iiKrwn9rlI5yjLkmrug154/YRW6KnnXVDM/l4=
|
||||
github.com/zeebo/xxh3 v1.0.2 h1:xZmwmqxHZA8AI603jOQ0tMqmBr9lPeFwGg6d+xy9DC0=
|
||||
github.com/zeebo/xxh3 v1.0.2/go.mod h1:5NWz9Sef7zIDm2JHfFlcQvNekmcEl9ekUZQQKCYaDcA=
|
||||
github.com/zeebo/xxh3 v1.1.0 h1:s7DLGDK45Dyfg7++yxI0khrfwq9661w9EN78eP/UZVs=
|
||||
github.com/zeebo/xxh3 v1.1.0/go.mod h1:IisAie1LELR4xhVinxWS5+zf1lA4p0MW4T+w+W07F5s=
|
||||
go.etcd.io/bbolt v1.4.3 h1:dEadXpI6G79deX5prL3QRNP6JB8UxVkqo4UPnHaNXJo=
|
||||
go.etcd.io/bbolt v1.4.3/go.mod h1:tKQlpPaYCVFctUIgFKFnAlvbmB3tpy1vkTnDWohtc0E=
|
||||
go.mongodb.org/mongo-driver v1.17.9 h1:IexDdCuuNJ3BHrELgBlyaH9p60JXAvdzWR128q+U5tU=
|
||||
@@ -787,8 +787,8 @@ golang.org/x/crypto v0.13.0/go.mod h1:y6Z2r+Rw4iayiXXAIxJIDAJ1zMW4yaTpebo8fPOliY
|
||||
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.48.0 h1:/VRzVqiRSggnhY7gNRxPauEQ5Drw9haKdM0jqfcCFts=
|
||||
golang.org/x/crypto v0.48.0/go.mod h1:r0kV5h3qnFPlQnBSrULhlsRfryS2pmewsg+XfMgkVos=
|
||||
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/exp v0.0.0-20190121172915-509febef88a4/go.mod h1:CJ0aWSM057203Lf6IL+f9T1iT9GByDxfZKAQTCR3kQA=
|
||||
golang.org/x/exp v0.0.0-20190306152737-a1d7652674e8/go.mod h1:CJ0aWSM057203Lf6IL+f9T1iT9GByDxfZKAQTCR3kQA=
|
||||
golang.org/x/exp v0.0.0-20190510132918-efd6b22b2522/go.mod h1:ZjyILWgesfNpC6sMxTJOJm9Kp84zZh5NQWvqDGG3Qr8=
|
||||
@@ -799,12 +799,12 @@ golang.org/x/exp v0.0.0-20191227195350-da58074b4299/go.mod h1:2RIsYlXP63K8oxa1u0
|
||||
golang.org/x/exp v0.0.0-20200119233911-0405dc783f0a/go.mod h1:2RIsYlXP63K8oxa1u096TMicItID8zy7Y6sNkU49FU4=
|
||||
golang.org/x/exp v0.0.0-20200207192155-f17229e696bd/go.mod h1:J/WKrq2StrnmMY6+EHIKF9dgMWnmCNThgcyBT1FY9mM=
|
||||
golang.org/x/exp v0.0.0-20200224162631-6cc2880d07d6/go.mod h1:3jZMyOhIsHpP37uCMkUooju7aAi5cS1Q23tOzKc+0MU=
|
||||
golang.org/x/exp v0.0.0-20251023183803-a4bb9ffd2546 h1:mgKeJMpvi0yx/sU5GsxQ7p6s2wtOnGAHZWCHUM4KGzY=
|
||||
golang.org/x/exp v0.0.0-20251023183803-a4bb9ffd2546/go.mod h1:j/pmGrbnkbPtQfxEe5D0VQhZC6qKbfKifgD0oM7sR70=
|
||||
golang.org/x/exp v0.0.0-20260218203240-3dfff04db8fa h1:Zt3DZoOFFYkKhDT3v7Lm9FDMEV06GpzjG2jrqW+QTE0=
|
||||
golang.org/x/exp v0.0.0-20260218203240-3dfff04db8fa/go.mod h1:K79w1Vqn7PoiZn+TkNpx3BUWUQksGO3JcVX6qIjytmA=
|
||||
golang.org/x/image v0.0.0-20190227222117-0694c2d4d067/go.mod h1:kZ7UVZpmo3dzQBMxlp+ypCbDeSB+sBbTgSJuh5dn5js=
|
||||
golang.org/x/image v0.0.0-20190802002840-cff245a6509b/go.mod h1:FeLwcggjj3mMvU+oOTbSwawSJRM1uh48EjtB4UJZlP0=
|
||||
golang.org/x/image v0.36.0 h1:Iknbfm1afbgtwPTmHnS2gTM/6PPZfH+z2EFuOkSbqwc=
|
||||
golang.org/x/image v0.36.0/go.mod h1:YsWD2TyyGKiIX1kZlu9QfKIsQ4nAAK9bdgdrIsE7xy4=
|
||||
golang.org/x/image v0.38.0 h1:5l+q+Y9JDC7mBOMjo4/aPhMDcxEptsX+Tt3GgRQRPuE=
|
||||
golang.org/x/image v0.38.0/go.mod h1:/3f6vaXC+6CEanU4KJxbcUZyEePbyKbaLoDOe4ehFYY=
|
||||
golang.org/x/lint v0.0.0-20181026193005-c67002cb31c3/go.mod h1:UVdnD1Gm6xHRNCYTkRU2/jEulfH38KcIWyp/GAMgvoE=
|
||||
golang.org/x/lint v0.0.0-20190227174305-5b3e6a55c961/go.mod h1:wehouNa3lNwaWXcvxsM5YxQ5yQlVC4a0KAMCusXpPoU=
|
||||
golang.org/x/lint v0.0.0-20190301231843-5614ed5bae6f/go.mod h1:UVdnD1Gm6xHRNCYTkRU2/jEulfH38KcIWyp/GAMgvoE=
|
||||
@@ -869,8 +869,8 @@ golang.org/x/net v0.15.0/go.mod h1:idbUs1IY1+zTqbi8yxTbhexhEEk5ur9LInksu6HrEpk=
|
||||
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.49.0 h1:eeHFmOGUTtaaPSGNmjBKpbng9MulQsJURQUAfUwY++o=
|
||||
golang.org/x/net v0.49.0/go.mod h1:/ysNB2EvaqvesRkuLAyjI1ycPZlQHM3q01F02UY/MV8=
|
||||
golang.org/x/net v0.51.0 h1:94R/GTO7mt3/4wIKpcR5gkGmRLOuE/2hNGeWq/GBIFo=
|
||||
golang.org/x/net v0.51.0/go.mod h1:aamm+2QF5ogm02fjy5Bb7CQ0WMt1/WVM7FtyaTLlA9Y=
|
||||
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=
|
||||
@@ -894,8 +894,8 @@ golang.org/x/sync v0.3.0/go.mod h1:FU7BRWz2tNW+3quACPkgCx/L+uEAv1htQ0V83Z9Rj+Y=
|
||||
golang.org/x/sync v0.6.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk=
|
||||
golang.org/x/sync v0.7.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk=
|
||||
golang.org/x/sync v0.10.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk=
|
||||
golang.org/x/sync v0.19.0 h1:vV+1eWNmZ5geRlYjzm2adRgW2/mcpevXNg50YZtPCE4=
|
||||
golang.org/x/sync v0.19.0/go.mod h1:9KTHXmSnoGruLpwFjVSX0lNNA75CykiMECbovNTZqGI=
|
||||
golang.org/x/sync v0.20.0 h1:e0PTpb7pjO8GAtTs2dQ6jYa5BWYlMuX047Dco/pItO4=
|
||||
golang.org/x/sync v0.20.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0=
|
||||
golang.org/x/sys v0.0.0-20180830151530-49385e6e1522/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
|
||||
golang.org/x/sys v0.0.0-20180909124046-d0be0721c37e/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
|
||||
golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
|
||||
@@ -959,8 +959,8 @@ golang.org/x/term v0.12.0/go.mod h1:owVbMEjm3cBLCHdkQu9b1opXd4ETQWc3BhuQGKgXgvU=
|
||||
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.40.0 h1:36e4zGLqU4yhjlmxEaagx2KuYbJq3EwY8K943ZsHcvg=
|
||||
golang.org/x/term v0.40.0/go.mod h1:w2P8uVp06p2iyKKuvXIm7N/y0UCRt3UfJTfZ7oOpglM=
|
||||
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/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=
|
||||
@@ -976,8 +976,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.34.0 h1:oL/Qq0Kdaqxa1KbNeMKwQq0reLCCaFtqu2eNuSeNHbk=
|
||||
golang.org/x/text v0.34.0/go.mod h1:homfLqTYRFyVYemLBFl5GgL/DWEiH5wcsQ5gSh1yziA=
|
||||
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/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=
|
||||
@@ -1030,14 +1030,14 @@ golang.org/x/tools v0.1.12/go.mod h1:hNGJHUnrk76NpqgfD5Aqm5Crs+Hm0VOH/i9J2+nxYbc
|
||||
golang.org/x/tools v0.6.0/go.mod h1:Xwgl3UAJ/d3gWutnCtw505GrjyAbvKui8lOU390QaIU=
|
||||
golang.org/x/tools v0.13.0/go.mod h1:HvlwmtVNQAhOuCjW7xxvovg8wbNq7LwfXh/k7wXUl58=
|
||||
golang.org/x/tools v0.21.1-0.20240508182429-e35e4ccd0d2d/go.mod h1:aiJjzUbINMkxbQROHiO6hDPo2LHcIPhhQsa9DLh0yGk=
|
||||
golang.org/x/tools v0.41.0 h1:a9b8iMweWG+S0OBnlU36rzLp20z1Rp10w+IY2czHTQc=
|
||||
golang.org/x/tools v0.41.0/go.mod h1:XSY6eDqxVNiYgezAVqqCeihT4j1U2CCsqvH3WhQpnlg=
|
||||
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/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
|
||||
golang.org/x/xerrors v0.0.0-20191011141410-1b5146add898/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
|
||||
golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
|
||||
golang.org/x/xerrors v0.0.0-20200804184101-5ec99f83aff1/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
|
||||
gonum.org/v1/gonum v0.16.0 h1:5+ul4Swaf3ESvrOnidPp4GZbzf0mxVQpDCYUQE7OJfk=
|
||||
gonum.org/v1/gonum v0.16.0/go.mod h1:fef3am4MQ93R2HHpKnLk4/Tbh/s0+wqD5nfa6Pnwy4E=
|
||||
gonum.org/v1/gonum v0.17.0 h1:VbpOemQlsSMrYmn7T2OUvQ4dqxQXU+ouZFQsZOx50z4=
|
||||
gonum.org/v1/gonum v0.17.0/go.mod h1:El3tOrEuMpv2UdMrbNlKEh9vd86bmQ6vqIcDwxEOc1E=
|
||||
google.golang.org/api v0.4.0/go.mod h1:8k5glujaEP+g9n7WNsDg8QP6cUVNI86fCNMcbazEtwE=
|
||||
google.golang.org/api v0.7.0/go.mod h1:WtwebWUNSVBH/HAw79HIFXZNqEvBhG+Ra+ax0hx3E3M=
|
||||
google.golang.org/api v0.8.0/go.mod h1:o4eAsZoiT+ibD93RtjEohWalFOjRDx6CVaqeizhEnKg=
|
||||
@@ -1054,8 +1054,8 @@ google.golang.org/api v0.24.0/go.mod h1:lIXQywCXRcnZPGlsd8NbLnOjtAoL6em04bJ9+z0M
|
||||
google.golang.org/api v0.28.0/go.mod h1:lIXQywCXRcnZPGlsd8NbLnOjtAoL6em04bJ9+z0MncE=
|
||||
google.golang.org/api v0.29.0/go.mod h1:Lcubydp8VUV7KeIHD9z2Bys/sm/vGKnG1UHuDBSrHWM=
|
||||
google.golang.org/api v0.30.0/go.mod h1:QGmEvQ87FHZNiUVJkT14jQNYJ4ZJjdRF23ZXz5138Fc=
|
||||
google.golang.org/api v0.258.0 h1:IKo1j5FBlN74fe5isA2PVozN3Y5pwNKriEgAXPOkDAc=
|
||||
google.golang.org/api v0.258.0/go.mod h1:qhOMTQEZ6lUps63ZNq9jhODswwjkjYYguA7fA3TBFww=
|
||||
google.golang.org/api v0.267.0 h1:w+vfWPMPYeRs8qH1aYYsFX68jMls5acWl/jocfLomwE=
|
||||
google.golang.org/api v0.267.0/go.mod h1:Jzc0+ZfLnyvXma3UtaTl023TdhZu6OMBP9tJ+0EmFD0=
|
||||
google.golang.org/appengine v1.1.0/go.mod h1:EbEs0AVv82hx2wNQdGPgUI5lhzA/G0D9YwlJXL52JkM=
|
||||
google.golang.org/appengine v1.4.0/go.mod h1:xpcJRLb0r/rnEns0DIKYYv+WjYCduHsrkT7/EB5XEv4=
|
||||
google.golang.org/appengine v1.5.0/go.mod h1:xpcJRLb0r/rnEns0DIKYYv+WjYCduHsrkT7/EB5XEv4=
|
||||
@@ -1091,12 +1091,12 @@ google.golang.org/genproto v0.0.0-20200618031413-b414f8b61790/go.mod h1:jDfRM7Fc
|
||||
google.golang.org/genproto v0.0.0-20200729003335-053ba62fc06f/go.mod h1:FWY/as6DDZQgahTzZj3fqbO1CbirC29ZNUFHwi0/+no=
|
||||
google.golang.org/genproto v0.0.0-20200804131852-c06518451d9c/go.mod h1:FWY/as6DDZQgahTzZj3fqbO1CbirC29ZNUFHwi0/+no=
|
||||
google.golang.org/genproto v0.0.0-20200825200019-8632dd797987/go.mod h1:FWY/as6DDZQgahTzZj3fqbO1CbirC29ZNUFHwi0/+no=
|
||||
google.golang.org/genproto v0.0.0-20251124214823-79d6a2a48846 h1:dDbsTLIK7EzwUq36kCSAsk0slouq/S0tWHeeGi97cD8=
|
||||
google.golang.org/genproto v0.0.0-20251124214823-79d6a2a48846/go.mod h1:PP0g88Dz3C7hRAfbQCQggeWAXjuqGsNPLE4s7jh0RGU=
|
||||
google.golang.org/genproto/googleapis/api v0.0.0-20251202230838-ff82c1b0f217 h1:fCvbg86sFXwdrl5LgVcTEvNC+2txB5mgROGmRL5mrls=
|
||||
google.golang.org/genproto/googleapis/api v0.0.0-20251202230838-ff82c1b0f217/go.mod h1:+rXWjjaukWZun3mLfjmVnQi18E1AsFbDN9QdJ5YXLto=
|
||||
google.golang.org/genproto/googleapis/rpc v0.0.0-20251213004720-97cd9d5aeac2 h1:2I6GHUeJ/4shcDpoUlLs/2WPnhg7yJwvXtqcMJt9liA=
|
||||
google.golang.org/genproto/googleapis/rpc v0.0.0-20251213004720-97cd9d5aeac2/go.mod h1:7i2o+ce6H/6BluujYR+kqX3GKH+dChPTQU19wjRPiGk=
|
||||
google.golang.org/genproto v0.0.0-20260128011058-8636f8732409 h1:VQZ/yAbAtjkHgH80teYd2em3xtIkkHd7ZhqfH2N9CsM=
|
||||
google.golang.org/genproto v0.0.0-20260128011058-8636f8732409/go.mod h1:rxKD3IEILWEu3P44seeNOAwZN4SaoKaQ/2eTg4mM6EM=
|
||||
google.golang.org/genproto/googleapis/api v0.0.0-20260203192932-546029d2fa20 h1:7ei4lp52gK1uSejlA8AZl5AJjeLUOHBQscRQZUgAcu0=
|
||||
google.golang.org/genproto/googleapis/api v0.0.0-20260203192932-546029d2fa20/go.mod h1:ZdbssH/1SOVnjnDlXzxDHK2MCidiqXtbYccJNzNYPEE=
|
||||
google.golang.org/genproto/googleapis/rpc v0.0.0-20260203192932-546029d2fa20 h1:Jr5R2J6F6qWyzINc+4AM8t5pfUz6beZpHp678GNrMbE=
|
||||
google.golang.org/genproto/googleapis/rpc v0.0.0-20260203192932-546029d2fa20/go.mod h1:j9x/tPzZkyxcgEFkiKEEGxfvyumM01BEtsW8xzOahRQ=
|
||||
google.golang.org/grpc v1.19.0/go.mod h1:mqu4LbDTu4XGKhr4mRzUsmM4RtVoemTSY81AxZiDr8c=
|
||||
google.golang.org/grpc v1.20.1/go.mod h1:10oTOabMzJvdu6/UiuZezV6QK5dSlG84ov/aaiqXj38=
|
||||
google.golang.org/grpc v1.21.1/go.mod h1:oYelfM1adQP15Ek0mdvEgi9Df8B9CZIaU1084ijfRaM=
|
||||
|
||||
@@ -11,9 +11,11 @@ import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/operation"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/volume_server_pb"
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/credentials/insecure"
|
||||
)
|
||||
|
||||
// VolumeServer provides a minimal volume server for erasure coding tests.
|
||||
@@ -196,12 +198,25 @@ func (v *VolumeServer) CopyFile(req *volume_server_pb.CopyFileRequest, stream vo
|
||||
defer file.Close()
|
||||
|
||||
buf := make([]byte, 64*1024)
|
||||
remaining := int64(req.GetStopOffset())
|
||||
for {
|
||||
n, readErr := file.Read(buf)
|
||||
if remaining == 0 {
|
||||
break
|
||||
}
|
||||
|
||||
readBuf := buf
|
||||
if remaining > 0 && remaining < int64(len(buf)) {
|
||||
readBuf = buf[:remaining]
|
||||
}
|
||||
|
||||
n, readErr := file.Read(readBuf)
|
||||
if n > 0 {
|
||||
if err := stream.Send(&volume_server_pb.CopyFileResponse{FileContent: buf[:n]}); err != nil {
|
||||
if err := stream.Send(&volume_server_pb.CopyFileResponse{FileContent: readBuf[:n]}); err != nil {
|
||||
return err
|
||||
}
|
||||
if remaining > 0 {
|
||||
remaining -= int64(n)
|
||||
}
|
||||
}
|
||||
if readErr == io.EOF {
|
||||
break
|
||||
@@ -307,10 +322,21 @@ func (v *VolumeServer) ReadVolumeFileStatus(ctx context.Context, req *volume_ser
|
||||
v.mu.Lock()
|
||||
v.readFileStatusCalls++
|
||||
v.mu.Unlock()
|
||||
|
||||
datInfo, err := os.Stat(v.filePath(req.VolumeId, ".dat"))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
idxInfo, err := os.Stat(v.filePath(req.VolumeId, ".idx"))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &volume_server_pb.ReadVolumeFileStatusResponse{
|
||||
VolumeId: req.VolumeId,
|
||||
DatFileSize: 1024,
|
||||
IdxFileSize: 16,
|
||||
DatFileSize: uint64(datInfo.Size()),
|
||||
IdxFileSize: uint64(idxInfo.Size()),
|
||||
FileCount: 1,
|
||||
}, nil
|
||||
}
|
||||
@@ -349,7 +375,27 @@ func (v *VolumeServer) VolumeCopy(req *volume_server_pb.VolumeCopyRequest, strea
|
||||
v.volumeCopyCalls++
|
||||
v.mu.Unlock()
|
||||
|
||||
if err := stream.Send(&volume_server_pb.VolumeCopyResponse{ProcessedBytes: 1024}); err != nil {
|
||||
dialOption := grpc.WithTransportCredentials(insecure.NewCredentials())
|
||||
var statusResp *volume_server_pb.ReadVolumeFileStatusResponse
|
||||
if err := operation.WithVolumeServerClient(false, pb.ServerAddress(req.SourceDataNode), dialOption,
|
||||
func(client volume_server_pb.VolumeServerClient) error {
|
||||
var readErr error
|
||||
statusResp, readErr = client.ReadVolumeFileStatus(stream.Context(), &volume_server_pb.ReadVolumeFileStatusRequest{
|
||||
VolumeId: req.VolumeId,
|
||||
})
|
||||
return readErr
|
||||
}); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := v.copyRemoteFile(stream.Context(), req.SourceDataNode, req.VolumeId, ".dat", statusResp.DatFileSize, dialOption); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := v.copyRemoteFile(stream.Context(), req.SourceDataNode, req.VolumeId, ".idx", statusResp.IdxFileSize, dialOption); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := stream.Send(&volume_server_pb.VolumeCopyResponse{ProcessedBytes: int64(statusResp.DatFileSize + statusResp.IdxFileSize)}); err != nil {
|
||||
return err
|
||||
}
|
||||
return stream.Send(&volume_server_pb.VolumeCopyResponse{LastAppendAtNs: uint64(time.Now().UnixNano())})
|
||||
@@ -368,3 +414,44 @@ func (v *VolumeServer) VolumeTailReceiver(ctx context.Context, req *volume_serve
|
||||
v.mu.Unlock()
|
||||
return &volume_server_pb.VolumeTailReceiverResponse{}, nil
|
||||
}
|
||||
|
||||
func (v *VolumeServer) copyRemoteFile(ctx context.Context, sourceDataNode string, volumeID uint32, ext string, fileSize uint64, dialOption grpc.DialOption) error {
|
||||
path := v.filePath(volumeID, ext)
|
||||
if err := os.MkdirAll(filepath.Dir(path), 0755); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
file, err := os.Create(path)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer file.Close()
|
||||
|
||||
return operation.WithVolumeServerClient(true, pb.ServerAddress(sourceDataNode), dialOption,
|
||||
func(client volume_server_pb.VolumeServerClient) error {
|
||||
stream, err := client.CopyFile(ctx, &volume_server_pb.CopyFileRequest{
|
||||
VolumeId: volumeID,
|
||||
Ext: ext,
|
||||
StopOffset: fileSize,
|
||||
})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
for {
|
||||
resp, recvErr := stream.Recv()
|
||||
if recvErr == io.EOF {
|
||||
return nil
|
||||
}
|
||||
if recvErr != nil {
|
||||
return recvErr
|
||||
}
|
||||
if len(resp.FileContent) == 0 {
|
||||
continue
|
||||
}
|
||||
if _, err := file.Write(resp.FileContent); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
@@ -16,6 +16,8 @@ import (
|
||||
"google.golang.org/protobuf/proto"
|
||||
)
|
||||
|
||||
const testVolumeDatSize = 1 * 1024 * 1024
|
||||
|
||||
func TestVolumeBalanceExecutionIntegration(t *testing.T) {
|
||||
volumeID := uint32(303)
|
||||
|
||||
@@ -31,6 +33,7 @@ func TestVolumeBalanceExecutionIntegration(t *testing.T) {
|
||||
|
||||
source := pluginworkers.NewVolumeServer(t, "")
|
||||
target := pluginworkers.NewVolumeServer(t, "")
|
||||
pluginworkers.WriteTestVolumeFiles(t, source.BaseDir(), volumeID, testVolumeDatSize)
|
||||
|
||||
job := &plugin_pb.JobSpec{
|
||||
JobId: fmt.Sprintf("balance-job-%d", volumeID),
|
||||
@@ -84,6 +87,9 @@ func TestVolumeBalanceBatchExecutionIntegration(t *testing.T) {
|
||||
|
||||
// Build a batch job with 3 volume moves from source → target.
|
||||
volumeIDs := []uint32{401, 402, 403}
|
||||
for _, vid := range volumeIDs {
|
||||
pluginworkers.WriteTestVolumeFiles(t, source.BaseDir(), vid, testVolumeDatSize)
|
||||
}
|
||||
moves := make([]*worker_pb.BalanceMoveSpec, len(volumeIDs))
|
||||
for i, vid := range volumeIDs {
|
||||
moves[i] = &worker_pb.BalanceMoveSpec{
|
||||
@@ -139,10 +145,11 @@ func TestVolumeBalanceBatchExecutionIntegration(t *testing.T) {
|
||||
require.True(t, deletedVols[vid], "volume %d should have been deleted from source", vid)
|
||||
}
|
||||
|
||||
// Pre-delete verification should have called ReadVolumeFileStatus on both
|
||||
// source and target for each volume.
|
||||
require.Equal(t, len(volumeIDs), source.ReadFileStatusCount(),
|
||||
"each move should read source volume status before delete")
|
||||
// Each move reads source status once before copy and once inside the
|
||||
// target's fake VolumeCopy implementation, then reads target status once
|
||||
// before deleting the source.
|
||||
require.Equal(t, len(volumeIDs)*2, source.ReadFileStatusCount(),
|
||||
"each move should read source volume status before copy and during target copy")
|
||||
require.Equal(t, len(volumeIDs), target.ReadFileStatusCount(),
|
||||
"each move should read target volume status before delete")
|
||||
|
||||
|
||||
@@ -31,6 +31,7 @@ type MountOptions struct {
|
||||
uidMap *string
|
||||
gidMap *string
|
||||
readOnly *bool
|
||||
includeSystemEntries *bool
|
||||
debug *bool
|
||||
debugPort *int
|
||||
localSocket *string
|
||||
@@ -99,6 +100,7 @@ func init() {
|
||||
mountOptions.uidMap = cmdMount.Flag.String("map.uid", "", "map local uid to uid on filer, comma-separated <local_uid>:<filer_uid>")
|
||||
mountOptions.gidMap = cmdMount.Flag.String("map.gid", "", "map local gid to gid on filer, comma-separated <local_gid>:<filer_gid>")
|
||||
mountOptions.readOnly = cmdMount.Flag.Bool("readOnly", false, "read only")
|
||||
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.localSocket = cmdMount.Flag.String("localSocket", "", "default to /tmp/seaweedfs-mount-<mount_dir_hash>.sock")
|
||||
|
||||
@@ -340,6 +340,7 @@ func RunMount(option *MountOptions, umask os.FileMode) bool {
|
||||
VolumeServerAccess: *mountOptions.volumeServerAccess,
|
||||
Cipher: cipher,
|
||||
UidGidMapper: uidGidMapper,
|
||||
IncludeSystemEntries: *option.includeSystemEntries,
|
||||
DisableXAttr: *option.disableXAttr,
|
||||
IsMacOs: runtime.GOOS == "darwin",
|
||||
MetadataFlushSeconds: *option.metadataFlushSeconds,
|
||||
|
||||
@@ -22,6 +22,17 @@ hosts = [
|
||||
topic = "seaweedfs_filer"
|
||||
offsetFile = "./last.offset"
|
||||
offsetSaveIntervalSeconds = 10
|
||||
# SASL Authentication
|
||||
sasl_enabled = false
|
||||
sasl_mechanism = "PLAIN" # PLAIN, SCRAM-SHA-256, SCRAM-SHA-512
|
||||
sasl_username = ""
|
||||
sasl_password = ""
|
||||
# TLS/SSL
|
||||
tls_enabled = false
|
||||
tls_ca_cert = "" # path to CA certificate PEM file
|
||||
tls_client_cert = "" # path to client certificate PEM file (for mTLS)
|
||||
tls_client_key = "" # path to client private key PEM file (for mTLS)
|
||||
tls_insecure_skip_verify = false
|
||||
|
||||
|
||||
[notification.aws_sqs]
|
||||
|
||||
@@ -27,6 +27,7 @@ type FilerOperations interface {
|
||||
CountDirectoryEntries(ctx context.Context, dirPath util.FullPath, limit int) (count int, err error)
|
||||
DeleteEntryMetaAndData(ctx context.Context, p util.FullPath, isRecursive, ignoreRecursiveError, shouldDeleteChunks, isFromOtherCluster bool, signatures []int32, ifNotModifiedAfter int64) error
|
||||
GetEntryAttributes(ctx context.Context, p util.FullPath) (attributes map[string][]byte, err error)
|
||||
IsDirectoryKeyObject(ctx context.Context, p util.FullPath) (bool, error)
|
||||
}
|
||||
|
||||
// folderState tracks the state of a folder for empty folder cleanup
|
||||
@@ -312,6 +313,16 @@ func (efc *EmptyFolderCleaner) executeCleanup(folder string, triggeredBy string)
|
||||
return
|
||||
}
|
||||
|
||||
// Skip explicitly created directory markers (e.g., PUT /bucket/folder/)
|
||||
// These have a MIME type set and should be preserved even when empty
|
||||
if isKeyObj, err := efc.filer.IsDirectoryKeyObject(ctx, util.FullPath(folder)); err != nil {
|
||||
glog.V(2).Infof("EmptyFolderCleaner: error checking directory key object %s: %v", folder, err)
|
||||
return
|
||||
} else if isKeyObj {
|
||||
glog.V(3).Infof("EmptyFolderCleaner: skipping %s (triggered by %s), explicit directory marker", folder, triggeredBy)
|
||||
return
|
||||
}
|
||||
|
||||
// Delete the empty folder
|
||||
glog.Infof("EmptyFolderCleaner: deleting empty folder %s (triggered by %s)", folder, triggeredBy)
|
||||
if err := efc.deleteFolder(ctx, folder); err != nil {
|
||||
|
||||
@@ -12,9 +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)
|
||||
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) {
|
||||
@@ -38,6 +39,13 @@ func (m *mockFilerOps) GetEntryAttributes(_ context.Context, p util.FullPath) (m
|
||||
return m.attrsFn(p)
|
||||
}
|
||||
|
||||
func (m *mockFilerOps) IsDirectoryKeyObject(_ context.Context, p util.FullPath) (bool, error) {
|
||||
if m.isDirKeyObjFn == nil {
|
||||
return false, nil
|
||||
}
|
||||
return m.isDirKeyObjFn(p)
|
||||
}
|
||||
|
||||
func Test_isUnderPath(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
@@ -733,3 +741,70 @@ func TestEmptyFolderCleaner_executeCleanup_bucketPolicyDisabledSkips(t *testing.
|
||||
t.Fatalf("expected folder %s to be skipped, got deletions %v", folder, deleted)
|
||||
}
|
||||
}
|
||||
|
||||
func TestEmptyFolderCleaner_executeCleanup_directoryMarker(t *testing.T) {
|
||||
testCases := []struct {
|
||||
name string
|
||||
isDirKeyObj bool
|
||||
expectDeletion bool
|
||||
}{
|
||||
{
|
||||
name: "skips explicit directory marker",
|
||||
isDirKeyObj: true,
|
||||
expectDeletion: false,
|
||||
},
|
||||
{
|
||||
name: "deletes implicit empty folder",
|
||||
isDirKeyObj: false,
|
||||
expectDeletion: true,
|
||||
},
|
||||
}
|
||||
|
||||
for _, tc := range testCases {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
lockRing := lock_manager.NewLockRing(5 * time.Second)
|
||||
lockRing.SetSnapshot([]pb.ServerAddress{"filer1:8888"})
|
||||
|
||||
var deleted []string
|
||||
mock := &mockFilerOps{
|
||||
countFn: func(_ util.FullPath) (int, error) {
|
||||
return 0, nil
|
||||
},
|
||||
deleteFn: func(path util.FullPath) error {
|
||||
deleted = append(deleted, string(path))
|
||||
return nil
|
||||
},
|
||||
isDirKeyObjFn: func(path util.FullPath) (bool, error) {
|
||||
return tc.isDirKeyObj, nil
|
||||
},
|
||||
}
|
||||
|
||||
cleaner := &EmptyFolderCleaner{
|
||||
filer: mock,
|
||||
lockRing: lockRing,
|
||||
host: "filer1:8888",
|
||||
bucketPath: "/buckets",
|
||||
enabled: true,
|
||||
folderCounts: make(map[string]*folderState),
|
||||
cleanupQueue: NewCleanupQueue(1000, time.Minute),
|
||||
maxCountCheck: 1000,
|
||||
cacheExpiry: time.Minute,
|
||||
processorSleep: time.Second,
|
||||
stopCh: make(chan struct{}),
|
||||
}
|
||||
|
||||
folder := "/buckets/test/folder"
|
||||
cleaner.executeCleanup(folder, "triggered_item")
|
||||
|
||||
if tc.expectDeletion {
|
||||
if len(deleted) != 1 || deleted[0] != folder {
|
||||
t.Fatalf("expected implicit empty folder %s to be deleted, got deletions %v", folder, deleted)
|
||||
}
|
||||
} else {
|
||||
if len(deleted) != 0 {
|
||||
t.Fatalf("expected explicit directory marker %s to be preserved, got deletions %v", folder, deleted)
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -559,3 +559,17 @@ func (f *Filer) GetEntryAttributes(ctx context.Context, p util.FullPath) (map[st
|
||||
}
|
||||
return entry.Extended, nil
|
||||
}
|
||||
|
||||
func (f *Filer) IsDirectoryKeyObject(ctx context.Context, p util.FullPath) (bool, error) {
|
||||
entry, err := f.FindEntry(ctx, p)
|
||||
if err != nil {
|
||||
if errors.Is(err, filer_pb.ErrNotFound) {
|
||||
return false, nil
|
||||
}
|
||||
return false, err
|
||||
}
|
||||
if entry == nil {
|
||||
return false, nil
|
||||
}
|
||||
return entry.IsDirectory() && entry.Mime != "", nil
|
||||
}
|
||||
|
||||
@@ -179,6 +179,10 @@ func PrepareStreamContentWithThrottler(ctx context.Context, masterClient wdclien
|
||||
jwt := jwtFunc(chunkView.FileId)
|
||||
written, err := retriedStreamFetchChunkData(ctx, writer, urlStrings, jwt, chunkView.CipherKey, chunkView.IsGzipped, chunkView.IsFullChunk(), chunkView.OffsetInChunk, int(chunkView.ViewSize))
|
||||
|
||||
if err != nil && ctx.Err() != nil {
|
||||
return ctx.Err()
|
||||
}
|
||||
|
||||
// If read failed, try to invalidate cache and re-lookup
|
||||
if err != nil && written == 0 {
|
||||
if invalidator, ok := masterClient.(CacheInvalidator); ok {
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package filer
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"testing"
|
||||
|
||||
@@ -173,3 +174,37 @@ func TestRetryLogicSkipsSameUrls(t *testing.T) {
|
||||
t.Error("Expected different URLs to not be equal")
|
||||
}
|
||||
}
|
||||
|
||||
func TestCanceledStreamSkipsCacheInvalidation(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
fileId := "3,canceled"
|
||||
|
||||
mock := &mockMasterClient{
|
||||
lookupFunc: func(ctx context.Context, fid string) ([]string, error) {
|
||||
return []string{"http://server:8080"}, nil
|
||||
},
|
||||
}
|
||||
|
||||
chunks := []*filer_pb.FileChunk{
|
||||
{
|
||||
FileId: fileId,
|
||||
Offset: 0,
|
||||
Size: 10,
|
||||
},
|
||||
}
|
||||
|
||||
streamFn, err := PrepareStreamContentWithThrottler(ctx, mock, noJwtFunc, chunks, 0, 10, 0)
|
||||
if err != nil {
|
||||
t.Fatalf("PrepareStreamContentWithThrottler failed: %v", err)
|
||||
}
|
||||
|
||||
cancel()
|
||||
|
||||
err = streamFn(&bytes.Buffer{})
|
||||
if err != context.Canceled {
|
||||
t.Fatalf("expected context.Canceled, got %v", err)
|
||||
}
|
||||
if len(mock.invalidatedFileIds) != 0 {
|
||||
t.Fatalf("expected no cache invalidation on cancellation, got %v", mock.invalidatedFileIds)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -27,18 +27,19 @@ type MetaCache struct {
|
||||
localStore filer.VirtualFilerStore
|
||||
leveldbStore *leveldb.LevelDBStore // direct reference for batch operations
|
||||
sync.RWMutex
|
||||
uidGidMapper *UidGidMapper
|
||||
markCachedFn func(fullpath util.FullPath)
|
||||
isCachedFn func(fullpath util.FullPath) bool
|
||||
invalidateFunc func(fullpath util.FullPath, entry *filer_pb.Entry)
|
||||
onDirectoryUpdate func(dir util.FullPath)
|
||||
visitGroup singleflight.Group // deduplicates concurrent EnsureVisited calls for the same path
|
||||
applyCh chan metadataApplyRequest
|
||||
applyDone chan struct{}
|
||||
applyStateMu sync.Mutex
|
||||
applyClosed bool
|
||||
buildingDirs map[util.FullPath]*directoryBuildState
|
||||
dedupRing dedupRingBuffer
|
||||
uidGidMapper *UidGidMapper
|
||||
markCachedFn func(fullpath util.FullPath)
|
||||
isCachedFn func(fullpath util.FullPath) bool
|
||||
invalidateFunc func(fullpath util.FullPath, entry *filer_pb.Entry)
|
||||
onDirectoryUpdate func(dir util.FullPath)
|
||||
visitGroup singleflight.Group // deduplicates concurrent EnsureVisited calls for the same path
|
||||
applyCh chan metadataApplyRequest
|
||||
applyDone chan struct{}
|
||||
applyStateMu sync.Mutex
|
||||
applyClosed bool
|
||||
buildingDirs map[util.FullPath]*directoryBuildState
|
||||
dedupRing dedupRingBuffer
|
||||
includeSystemEntries bool
|
||||
}
|
||||
|
||||
var errMetaCacheClosed = errors.New("metadata cache is shut down")
|
||||
@@ -84,17 +85,18 @@ type metadataApplyRequest struct {
|
||||
done chan error
|
||||
}
|
||||
|
||||
func NewMetaCache(dbFolder string, uidGidMapper *UidGidMapper, root util.FullPath,
|
||||
func NewMetaCache(dbFolder string, uidGidMapper *UidGidMapper, root util.FullPath, includeSystemEntries bool,
|
||||
markCachedFn func(path util.FullPath), isCachedFn func(path util.FullPath) bool, invalidateFunc func(util.FullPath, *filer_pb.Entry), onDirectoryUpdate func(dir util.FullPath)) *MetaCache {
|
||||
leveldbStore, virtualStore := openMetaStore(dbFolder)
|
||||
mc := &MetaCache{
|
||||
root: root,
|
||||
localStore: virtualStore,
|
||||
leveldbStore: leveldbStore,
|
||||
markCachedFn: markCachedFn,
|
||||
isCachedFn: isCachedFn,
|
||||
uidGidMapper: uidGidMapper,
|
||||
onDirectoryUpdate: onDirectoryUpdate,
|
||||
root: root,
|
||||
localStore: virtualStore,
|
||||
leveldbStore: leveldbStore,
|
||||
markCachedFn: markCachedFn,
|
||||
isCachedFn: isCachedFn,
|
||||
uidGidMapper: uidGidMapper,
|
||||
onDirectoryUpdate: onDirectoryUpdate,
|
||||
includeSystemEntries: includeSystemEntries,
|
||||
invalidateFunc: func(fullpath util.FullPath, entry *filer_pb.Entry) {
|
||||
invalidateFunc(fullpath, entry)
|
||||
},
|
||||
@@ -182,6 +184,29 @@ func (mc *MetaCache) atomicUpdateEntryFromFilerLocked(ctx context.Context, oldPa
|
||||
return nil
|
||||
}
|
||||
|
||||
func (mc *MetaCache) shouldHideEntry(fullpath util.FullPath) bool {
|
||||
if mc.includeSystemEntries {
|
||||
return false
|
||||
}
|
||||
dir, name := fullpath.DirAndName()
|
||||
return IsHiddenSystemEntry(dir, name)
|
||||
}
|
||||
|
||||
func (mc *MetaCache) purgeEntryLocked(ctx context.Context, fullpath util.FullPath, isDirectory bool) error {
|
||||
if fullpath == "" {
|
||||
return nil
|
||||
}
|
||||
if err := mc.localStore.DeleteEntry(ctx, fullpath); err != nil {
|
||||
return err
|
||||
}
|
||||
if isDirectory {
|
||||
if err := mc.localStore.DeleteFolderChildren(ctx, fullpath); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (mc *MetaCache) ApplyMetadataResponse(ctx context.Context, resp *filer_pb.SubscribeMetadataResponse, options MetadataResponseApplyOptions) error {
|
||||
if resp == nil || resp.EventNotification == nil {
|
||||
return nil
|
||||
@@ -520,7 +545,9 @@ func (mc *MetaCache) applyMetadataResponseLocked(ctx context.Context, resp *file
|
||||
}
|
||||
|
||||
var oldPath util.FullPath
|
||||
var newPath util.FullPath
|
||||
var newEntry *filer.Entry
|
||||
hideNewPath := false
|
||||
if message.OldEntry != nil {
|
||||
oldPath = util.NewFullPath(resp.Directory, message.OldEntry.Name)
|
||||
}
|
||||
@@ -530,11 +557,20 @@ func (mc *MetaCache) applyMetadataResponseLocked(ctx context.Context, resp *file
|
||||
if message.NewParentPath != "" {
|
||||
dir = message.NewParentPath
|
||||
}
|
||||
newEntry = filer.FromPbEntry(dir, message.NewEntry)
|
||||
newPath = util.NewFullPath(dir, message.NewEntry.Name)
|
||||
hideNewPath = mc.shouldHideEntry(newPath)
|
||||
if !hideNewPath {
|
||||
newEntry = filer.FromPbEntry(dir, message.NewEntry)
|
||||
}
|
||||
}
|
||||
|
||||
mc.Lock()
|
||||
err := mc.atomicUpdateEntryFromFilerLocked(ctx, oldPath, newEntry, allowUncachedInsert)
|
||||
if err == nil && hideNewPath {
|
||||
if purgeErr := mc.purgeEntryLocked(ctx, newPath, message.NewEntry.IsDirectory); purgeErr != nil {
|
||||
err = purgeErr
|
||||
}
|
||||
}
|
||||
// When a directory is deleted or moved, remove its cached descendants
|
||||
// so stale children cannot be served from the local cache.
|
||||
if err == nil && oldPath != "" && message.OldEntry != nil && message.OldEntry.IsDirectory {
|
||||
|
||||
@@ -2,6 +2,7 @@ package meta_cache
|
||||
|
||||
import (
|
||||
"context"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sync"
|
||||
"testing"
|
||||
@@ -296,6 +297,114 @@ func TestApplyMetadataResponseDeduplicatesRepeatedFilerEvent(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestApplyMetadataResponseSkipsHiddenSystemEntryWhenDisabled(t *testing.T) {
|
||||
mc, _, _, _ := newTestMetaCache(t, map[util.FullPath]bool{
|
||||
"/": true,
|
||||
})
|
||||
defer mc.Shutdown()
|
||||
|
||||
createResp := &filer_pb.SubscribeMetadataResponse{
|
||||
Directory: "/",
|
||||
EventNotification: &filer_pb.EventNotification{
|
||||
NewEntry: &filer_pb.Entry{
|
||||
Name: "topics",
|
||||
Attributes: &filer_pb.FuseAttributes{
|
||||
Crtime: 1,
|
||||
Mtime: 1,
|
||||
FileMode: uint32(os.ModeDir | 0o755),
|
||||
},
|
||||
IsDirectory: true,
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
if err := mc.ApplyMetadataResponse(context.Background(), createResp, SubscriberMetadataResponseApplyOptions); err != nil {
|
||||
t.Fatalf("apply create: %v", err)
|
||||
}
|
||||
|
||||
entry, err := mc.FindEntry(context.Background(), util.FullPath("/topics"))
|
||||
if err != filer_pb.ErrNotFound {
|
||||
t.Fatalf("find hidden entry error = %v, want %v", err, filer_pb.ErrNotFound)
|
||||
}
|
||||
if entry != nil {
|
||||
t.Fatalf("hidden entry still cached: %+v", entry)
|
||||
}
|
||||
}
|
||||
|
||||
func TestApplyMetadataResponsePurgesHiddenDestinationPath(t *testing.T) {
|
||||
mc, _, _, _ := newTestMetaCache(t, map[util.FullPath]bool{
|
||||
"/": true,
|
||||
"/src": true,
|
||||
})
|
||||
defer mc.Shutdown()
|
||||
|
||||
if err := mc.InsertEntry(context.Background(), &filer.Entry{
|
||||
FullPath: "/topics",
|
||||
Attr: filer.Attr{
|
||||
Crtime: time.Unix(1, 0),
|
||||
Mtime: time.Unix(1, 0),
|
||||
Mode: os.ModeDir | 0o755,
|
||||
},
|
||||
}); err != nil {
|
||||
t.Fatalf("insert stale hidden dir: %v", err)
|
||||
}
|
||||
if err := mc.InsertEntry(context.Background(), &filer.Entry{
|
||||
FullPath: "/topics/leaked.txt",
|
||||
Attr: filer.Attr{
|
||||
Crtime: time.Unix(1, 0),
|
||||
Mtime: time.Unix(1, 0),
|
||||
Mode: 0o644,
|
||||
FileSize: 7,
|
||||
},
|
||||
}); err != nil {
|
||||
t.Fatalf("insert leaked hidden child: %v", err)
|
||||
}
|
||||
if err := mc.InsertEntry(context.Background(), &filer.Entry{
|
||||
FullPath: "/src/visible",
|
||||
Attr: filer.Attr{
|
||||
Crtime: time.Unix(1, 0),
|
||||
Mtime: time.Unix(1, 0),
|
||||
Mode: os.ModeDir | 0o755,
|
||||
},
|
||||
}); err != nil {
|
||||
t.Fatalf("insert source dir: %v", err)
|
||||
}
|
||||
|
||||
renameResp := &filer_pb.SubscribeMetadataResponse{
|
||||
Directory: "/src",
|
||||
EventNotification: &filer_pb.EventNotification{
|
||||
OldEntry: &filer_pb.Entry{
|
||||
Name: "visible",
|
||||
IsDirectory: true,
|
||||
},
|
||||
NewEntry: &filer_pb.Entry{
|
||||
Name: "topics",
|
||||
Attributes: &filer_pb.FuseAttributes{
|
||||
Crtime: 2,
|
||||
Mtime: 2,
|
||||
FileMode: uint32(os.ModeDir | 0o755),
|
||||
},
|
||||
IsDirectory: true,
|
||||
},
|
||||
NewParentPath: "/",
|
||||
},
|
||||
}
|
||||
|
||||
if err := mc.ApplyMetadataResponse(context.Background(), renameResp, SubscriberMetadataResponseApplyOptions); err != nil {
|
||||
t.Fatalf("apply rename: %v", err)
|
||||
}
|
||||
|
||||
if entry, err := mc.FindEntry(context.Background(), util.FullPath("/src/visible")); err != filer_pb.ErrNotFound || entry != nil {
|
||||
t.Fatalf("source dir after rename = %+v, %v; want nil, %v", entry, err, filer_pb.ErrNotFound)
|
||||
}
|
||||
if entry, err := mc.FindEntry(context.Background(), util.FullPath("/topics")); err != filer_pb.ErrNotFound || entry != nil {
|
||||
t.Fatalf("hidden destination after rename = %+v, %v; want nil, %v", entry, err, filer_pb.ErrNotFound)
|
||||
}
|
||||
if entry, err := mc.FindEntry(context.Background(), util.FullPath("/topics/leaked.txt")); err != filer_pb.ErrNotFound || entry != nil {
|
||||
t.Fatalf("hidden child after rename = %+v, %v; want nil, %v", entry, err, filer_pb.ErrNotFound)
|
||||
}
|
||||
}
|
||||
|
||||
func newTestMetaCache(t *testing.T, cached map[util.FullPath]bool) (*MetaCache, map[util.FullPath]bool, *recordedPaths, *recordedPaths) {
|
||||
t.Helper()
|
||||
|
||||
@@ -312,6 +421,7 @@ func newTestMetaCache(t *testing.T, cached map[util.FullPath]bool) (*MetaCache,
|
||||
filepath.Join(t.TempDir(), "meta"),
|
||||
mapper,
|
||||
util.FullPath("/"),
|
||||
false,
|
||||
func(path util.FullPath) {
|
||||
cachedMu.Lock()
|
||||
defer cachedMu.Unlock()
|
||||
|
||||
@@ -107,7 +107,7 @@ func doEnsureVisited(ctx context.Context, mc *MetaCache, client filer_pb.FilerCl
|
||||
var err error
|
||||
snapshotTsNs, err = filer_pb.ReadDirAllEntriesWithSnapshot(ctx, client, path, "", func(pbEntry *filer_pb.Entry, isLast bool) error {
|
||||
entry := filer.FromPbEntry(string(path), pbEntry)
|
||||
if IsHiddenSystemEntry(string(path), entry.Name()) {
|
||||
if !mc.includeSystemEntries && IsHiddenSystemEntry(string(path), entry.Name()) {
|
||||
return nil
|
||||
}
|
||||
|
||||
|
||||
@@ -65,6 +65,7 @@ type Option struct {
|
||||
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)
|
||||
// This protects chunks from being purged by volume.fsck for long-running writes
|
||||
@@ -203,6 +204,7 @@ func NewSeaweedFileSystem(option *Option) *WFS {
|
||||
|
||||
wfs.metaCache = meta_cache.NewMetaCache(path.Join(option.getUniqueCacheDirForRead(), "meta"), option.UidGidMapper,
|
||||
util.FullPath(option.FilerMountRootPath),
|
||||
option.IncludeSystemEntries,
|
||||
func(path util.FullPath) {
|
||||
wfs.inodeToPath.MarkChildrenCached(path)
|
||||
}, func(path util.FullPath) bool {
|
||||
|
||||
@@ -302,7 +302,7 @@ func (wfs *WFS) readDirectoryDirect(input *fuse.ReadIn, out *fuse.DirEntryList,
|
||||
if input.Offset >= dh.entryStreamOffset {
|
||||
if len(dh.entryStream) == 0 && input.Offset > dh.entryStreamOffset {
|
||||
skipCount := uint32(input.Offset-dh.entryStreamOffset) + batchSize
|
||||
entries, snapshotTs, err := loadDirectoryEntriesDirect(context.Background(), wfs, wfs.option.UidGidMapper, dirPath, "", false, skipCount, dh.snapshotTsNs)
|
||||
entries, snapshotTs, err := loadDirectoryEntriesDirect(context.Background(), wfs, wfs.option.UidGidMapper, dirPath, "", false, skipCount, dh.snapshotTsNs, wfs.option.IncludeSystemEntries)
|
||||
if err != nil {
|
||||
glog.Errorf("list filer directory: %v", err)
|
||||
return fuse.EIO
|
||||
@@ -331,7 +331,7 @@ func (wfs *WFS) readDirectoryDirect(input *fuse.ReadIn, out *fuse.DirEntryList,
|
||||
}
|
||||
}
|
||||
|
||||
entries, snapshotTs, err := loadDirectoryEntriesDirect(context.Background(), wfs, wfs.option.UidGidMapper, dirPath, lastEntryName, false, batchSize, dh.snapshotTsNs)
|
||||
entries, snapshotTs, err := loadDirectoryEntriesDirect(context.Background(), wfs, wfs.option.UidGidMapper, dirPath, lastEntryName, false, batchSize, dh.snapshotTsNs, wfs.option.IncludeSystemEntries)
|
||||
if err != nil {
|
||||
glog.Errorf("list filer directory: %v", err)
|
||||
return fuse.EIO
|
||||
@@ -360,13 +360,13 @@ func (wfs *WFS) readDirectoryDirect(input *fuse.ReadIn, out *fuse.DirEntryList,
|
||||
return fuse.OK
|
||||
}
|
||||
|
||||
func loadDirectoryEntriesDirect(ctx context.Context, client filer_pb.FilerClient, uidGidMapper *meta_cache.UidGidMapper, dirPath util.FullPath, startFileName string, includeStart bool, limit uint32, snapshotTsNs int64) ([]*filer.Entry, int64, error) {
|
||||
func loadDirectoryEntriesDirect(ctx context.Context, client filer_pb.FilerClient, uidGidMapper *meta_cache.UidGidMapper, dirPath util.FullPath, startFileName string, includeStart bool, limit uint32, snapshotTsNs int64, includeSystemEntries bool) ([]*filer.Entry, int64, error) {
|
||||
entries := make([]*filer.Entry, 0, limit)
|
||||
var actualSnapshotTsNs int64
|
||||
err := client.WithFilerClient(false, func(sc filer_pb.SeaweedFilerClient) error {
|
||||
var innerErr error
|
||||
actualSnapshotTsNs, innerErr = filer_pb.DoSeaweedListWithSnapshot(ctx, sc, dirPath, "", func(entry *filer_pb.Entry, isLast bool) error {
|
||||
if meta_cache.IsHiddenSystemEntry(string(dirPath), entry.Name) {
|
||||
if !includeSystemEntries && meta_cache.IsHiddenSystemEntry(string(dirPath), entry.Name) {
|
||||
return nil
|
||||
}
|
||||
if uidGidMapper != nil && entry.Attributes != nil {
|
||||
|
||||
@@ -84,7 +84,7 @@ func TestLoadDirectoryEntriesDirectFiltersHiddenEntriesAndMapsIds(t *testing.T)
|
||||
},
|
||||
}
|
||||
|
||||
entries, _, err := loadDirectoryEntriesDirect(context.Background(), client, mapper, util.FullPath("/"), "", false, 10, 0)
|
||||
entries, _, err := loadDirectoryEntriesDirect(context.Background(), client, mapper, util.FullPath("/"), "", false, 10, 0, false)
|
||||
if err != nil {
|
||||
t.Fatalf("loadDirectoryEntriesDirect: %v", err)
|
||||
}
|
||||
@@ -98,3 +98,25 @@ func TestLoadDirectoryEntriesDirectFiltersHiddenEntriesAndMapsIds(t *testing.T)
|
||||
t.Fatalf("mapped uid/gid = %d/%d, want 10/20", entries[0].Attr.Uid, entries[0].Attr.Gid)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadDirectoryEntriesDirectShowsSystemEntriesWhenEnabled(t *testing.T) {
|
||||
client := &directoryFilerAccessor{
|
||||
client: &directoryListClient{
|
||||
responses: []*filer_pb.ListEntriesResponse{
|
||||
{Entry: &filer_pb.Entry{Name: "topics"}},
|
||||
{Entry: &filer_pb.Entry{Name: "visible"}},
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
entries, _, err := loadDirectoryEntriesDirect(context.Background(), client, nil, util.FullPath("/"), "", false, 10, 0, true)
|
||||
if err != nil {
|
||||
t.Fatalf("loadDirectoryEntriesDirect: %v", err)
|
||||
}
|
||||
if got := len(entries); got != 2 {
|
||||
t.Fatalf("entry count = %d, want 2", got)
|
||||
}
|
||||
if entries[0].Name() != "topics" || entries[1].Name() != "visible" {
|
||||
t.Fatalf("entry names = %q, %q, want topics, visible", entries[0].Name(), entries[1].Name())
|
||||
}
|
||||
}
|
||||
|
||||
@@ -338,6 +338,7 @@ func newCopyRangeTestWFSWithMetaCache(t *testing.T) *WFS {
|
||||
filepath.Join(t.TempDir(), "meta"),
|
||||
uidGidMapper,
|
||||
root,
|
||||
false,
|
||||
func(path util.FullPath) {
|
||||
wfs.inodeToPath.MarkChildrenCached(path)
|
||||
},
|
||||
|
||||
@@ -120,6 +120,7 @@ func newCreateTestWFS(t *testing.T) (*WFS, *createEntryTestServer) {
|
||||
filepath.Join(t.TempDir(), "meta"),
|
||||
uidGidMapper,
|
||||
root,
|
||||
false,
|
||||
func(path util.FullPath) {
|
||||
wfs.inodeToPath.MarkChildrenCached(path)
|
||||
},
|
||||
|
||||
@@ -23,6 +23,7 @@ func TestHandleRenameResponseLeavesUncachedTargetOutOfCache(t *testing.T) {
|
||||
filepath.Join(t.TempDir(), "meta"),
|
||||
uidGidMapper,
|
||||
root,
|
||||
false,
|
||||
func(path util.FullPath) {
|
||||
inodeToPath.MarkChildrenCached(path)
|
||||
},
|
||||
|
||||
@@ -1,6 +1,8 @@
|
||||
package kafka
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"github.com/Shopify/sarama"
|
||||
"github.com/seaweedfs/seaweedfs/weed/glog"
|
||||
"github.com/seaweedfs/seaweedfs/weed/notification"
|
||||
@@ -27,15 +29,29 @@ func (k *KafkaQueue) Initialize(configuration util.Configuration, prefix string)
|
||||
return k.initialize(
|
||||
configuration.GetStringSlice(prefix+"hosts"),
|
||||
configuration.GetString(prefix+"topic"),
|
||||
SASLTLSConfig{
|
||||
SASLEnabled: configuration.GetBool(prefix + "sasl_enabled"),
|
||||
SASLMechanism: configuration.GetString(prefix + "sasl_mechanism"),
|
||||
SASLUsername: configuration.GetString(prefix + "sasl_username"),
|
||||
SASLPassword: configuration.GetString(prefix + "sasl_password"),
|
||||
TLSEnabled: configuration.GetBool(prefix + "tls_enabled"),
|
||||
TLSCACert: configuration.GetString(prefix + "tls_ca_cert"),
|
||||
TLSClientCert: configuration.GetString(prefix + "tls_client_cert"),
|
||||
TLSClientKey: configuration.GetString(prefix + "tls_client_key"),
|
||||
TLSInsecureSkipVerify: configuration.GetBool(prefix + "tls_insecure_skip_verify"),
|
||||
},
|
||||
)
|
||||
}
|
||||
|
||||
func (k *KafkaQueue) initialize(hosts []string, topic string) (err error) {
|
||||
func (k *KafkaQueue) initialize(hosts []string, topic string, saslTLS SASLTLSConfig) (err error) {
|
||||
config := sarama.NewConfig()
|
||||
config.Producer.RequiredAcks = sarama.WaitForLocal
|
||||
config.Producer.Partitioner = sarama.NewHashPartitioner
|
||||
config.Producer.Return.Successes = true
|
||||
config.Producer.Return.Errors = true
|
||||
if err = ConfigureSASLTLS(config, saslTLS); err != nil {
|
||||
return fmt.Errorf("kafka producer security configuration: %w", err)
|
||||
}
|
||||
k.producer, err = sarama.NewAsyncProducer(hosts, config)
|
||||
if err != nil {
|
||||
return err
|
||||
|
||||
@@ -0,0 +1,112 @@
|
||||
package kafka
|
||||
|
||||
import (
|
||||
"crypto/tls"
|
||||
"crypto/x509"
|
||||
"fmt"
|
||||
"os"
|
||||
"strings"
|
||||
|
||||
"github.com/Shopify/sarama"
|
||||
"github.com/xdg-go/scram"
|
||||
)
|
||||
|
||||
// SASLTLSConfig holds SASL and TLS configuration for Kafka connections.
|
||||
type SASLTLSConfig struct {
|
||||
SASLEnabled bool
|
||||
SASLMechanism string
|
||||
SASLUsername string
|
||||
SASLPassword string
|
||||
|
||||
TLSEnabled bool
|
||||
TLSCACert string
|
||||
TLSClientCert string
|
||||
TLSClientKey string
|
||||
TLSInsecureSkipVerify bool
|
||||
}
|
||||
|
||||
// ConfigureSASLTLS applies SASL and TLS settings to a sarama config.
|
||||
func ConfigureSASLTLS(config *sarama.Config, st SASLTLSConfig) error {
|
||||
if st.SASLEnabled {
|
||||
config.Net.SASL.Enable = true
|
||||
config.Net.SASL.User = st.SASLUsername
|
||||
config.Net.SASL.Password = st.SASLPassword
|
||||
|
||||
mechanism := strings.ToUpper(st.SASLMechanism)
|
||||
switch mechanism {
|
||||
case "PLAIN", "":
|
||||
config.Net.SASL.Mechanism = sarama.SASLTypePlaintext
|
||||
case "SCRAM-SHA-256":
|
||||
config.Net.SASL.Mechanism = sarama.SASLTypeSCRAMSHA256
|
||||
config.Net.SASL.SCRAMClientGeneratorFunc = func() sarama.SCRAMClient {
|
||||
return &scramClient{HashGeneratorFcn: scram.SHA256}
|
||||
}
|
||||
case "SCRAM-SHA-512":
|
||||
config.Net.SASL.Mechanism = sarama.SASLTypeSCRAMSHA512
|
||||
config.Net.SASL.SCRAMClientGeneratorFunc = func() sarama.SCRAMClient {
|
||||
return &scramClient{HashGeneratorFcn: scram.SHA512}
|
||||
}
|
||||
default:
|
||||
return fmt.Errorf("unsupported SASL mechanism: %s", mechanism)
|
||||
}
|
||||
}
|
||||
|
||||
if st.TLSEnabled {
|
||||
if (st.TLSClientCert == "") != (st.TLSClientKey == "") {
|
||||
return fmt.Errorf("both tls_client_cert and tls_client_key must be provided for mTLS, or neither")
|
||||
}
|
||||
|
||||
tlsConfig := &tls.Config{
|
||||
MinVersion: tls.VersionTLS12,
|
||||
InsecureSkipVerify: st.TLSInsecureSkipVerify,
|
||||
}
|
||||
|
||||
if st.TLSCACert != "" {
|
||||
caCert, err := os.ReadFile(st.TLSCACert)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to read CA certificate: %w", err)
|
||||
}
|
||||
caCertPool := x509.NewCertPool()
|
||||
if !caCertPool.AppendCertsFromPEM(caCert) {
|
||||
return fmt.Errorf("failed to parse CA certificate")
|
||||
}
|
||||
tlsConfig.RootCAs = caCertPool
|
||||
}
|
||||
|
||||
if st.TLSClientCert != "" && st.TLSClientKey != "" {
|
||||
cert, err := tls.LoadX509KeyPair(st.TLSClientCert, st.TLSClientKey)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to load client certificate/key: %w", err)
|
||||
}
|
||||
tlsConfig.Certificates = []tls.Certificate{cert}
|
||||
}
|
||||
|
||||
config.Net.TLS.Enable = true
|
||||
config.Net.TLS.Config = tlsConfig
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// scramClient implements the sarama.SCRAMClient interface.
|
||||
type scramClient struct {
|
||||
*scram.ClientConversation
|
||||
scram.HashGeneratorFcn
|
||||
}
|
||||
|
||||
func (c *scramClient) Begin(userName, password, authzID string) (err error) {
|
||||
client, err := c.HashGeneratorFcn.NewClient(userName, password, authzID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
c.ClientConversation = client.NewConversation()
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c *scramClient) Step(challenge string) (string, error) {
|
||||
return c.ClientConversation.Step(challenge)
|
||||
}
|
||||
|
||||
func (c *scramClient) Done() bool {
|
||||
return c.ClientConversation.Done()
|
||||
}
|
||||
@@ -0,0 +1,65 @@
|
||||
package kafka
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"github.com/Shopify/sarama"
|
||||
)
|
||||
|
||||
func TestConfigureSASLTLSRejectsPartialMTLSConfig(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
cfg SASLTLSConfig
|
||||
}{
|
||||
{
|
||||
name: "missing key",
|
||||
cfg: SASLTLSConfig{
|
||||
TLSEnabled: true,
|
||||
TLSClientCert: "/tmp/client.crt",
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "missing cert",
|
||||
cfg: SASLTLSConfig{
|
||||
TLSEnabled: true,
|
||||
TLSClientKey: "/tmp/client.key",
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
err := ConfigureSASLTLS(sarama.NewConfig(), tt.cfg)
|
||||
if err == nil {
|
||||
t.Fatal("expected error")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "both tls_client_cert and tls_client_key must be provided") {
|
||||
t.Fatalf("unexpected error: %v", err)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestConfigureSASLTLSConfiguresSCRAMSHA256(t *testing.T) {
|
||||
config := sarama.NewConfig()
|
||||
err := ConfigureSASLTLS(config, SASLTLSConfig{
|
||||
SASLEnabled: true,
|
||||
SASLMechanism: "SCRAM-SHA-256",
|
||||
SASLUsername: "alice",
|
||||
SASLPassword: "secret",
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("ConfigureSASLTLS returned error: %v", err)
|
||||
}
|
||||
|
||||
if !config.Net.SASL.Enable {
|
||||
t.Fatal("expected SASL to be enabled")
|
||||
}
|
||||
if config.Net.SASL.Mechanism != sarama.SASLTypeSCRAMSHA256 {
|
||||
t.Fatalf("unexpected mechanism: %v", config.Net.SASL.Mechanism)
|
||||
}
|
||||
if config.Net.SASL.SCRAMClientGeneratorFunc == nil {
|
||||
t.Fatal("expected SCRAM client generator")
|
||||
}
|
||||
}
|
||||
@@ -9,6 +9,7 @@ import (
|
||||
|
||||
"github.com/Shopify/sarama"
|
||||
"github.com/seaweedfs/seaweedfs/weed/glog"
|
||||
kafkanotif "github.com/seaweedfs/seaweedfs/weed/notification/kafka"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
"google.golang.org/protobuf/proto"
|
||||
@@ -36,25 +37,38 @@ func (k *KafkaInput) Initialize(configuration util.Configuration, prefix string)
|
||||
configuration.GetString(prefix+"topic"),
|
||||
configuration.GetString(prefix+"offsetFile"),
|
||||
configuration.GetInt(prefix+"offsetSaveIntervalSeconds"),
|
||||
kafkanotif.SASLTLSConfig{
|
||||
SASLEnabled: configuration.GetBool(prefix + "sasl_enabled"),
|
||||
SASLMechanism: configuration.GetString(prefix + "sasl_mechanism"),
|
||||
SASLUsername: configuration.GetString(prefix + "sasl_username"),
|
||||
SASLPassword: configuration.GetString(prefix + "sasl_password"),
|
||||
TLSEnabled: configuration.GetBool(prefix + "tls_enabled"),
|
||||
TLSCACert: configuration.GetString(prefix + "tls_ca_cert"),
|
||||
TLSClientCert: configuration.GetString(prefix + "tls_client_cert"),
|
||||
TLSClientKey: configuration.GetString(prefix + "tls_client_key"),
|
||||
TLSInsecureSkipVerify: configuration.GetBool(prefix + "tls_insecure_skip_verify"),
|
||||
},
|
||||
)
|
||||
}
|
||||
|
||||
func (k *KafkaInput) initialize(hosts []string, topic string, offsetFile string, offsetSaveIntervalSeconds int) (err error) {
|
||||
func (k *KafkaInput) initialize(hosts []string, topic string, offsetFile string, offsetSaveIntervalSeconds int, saslTLS kafkanotif.SASLTLSConfig) (err error) {
|
||||
config := sarama.NewConfig()
|
||||
config.Consumer.Return.Errors = true
|
||||
if err = kafkanotif.ConfigureSASLTLS(config, saslTLS); err != nil {
|
||||
return fmt.Errorf("kafka consumer security configuration: %w", err)
|
||||
}
|
||||
k.consumer, err = sarama.NewConsumer(hosts, config)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
} else {
|
||||
glog.V(0).Infof("connected to %v", hosts)
|
||||
return fmt.Errorf("create kafka consumer: %w", err)
|
||||
}
|
||||
glog.V(0).Infof("connected to %v", hosts)
|
||||
|
||||
k.topic = topic
|
||||
k.messageChan = make(chan *sarama.ConsumerMessage, 1)
|
||||
|
||||
partitions, err := k.consumer.Partitions(topic)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
return fmt.Errorf("get kafka partitions for topic %q: %w", topic, err)
|
||||
}
|
||||
|
||||
progress := loadProgress(offsetFile)
|
||||
@@ -77,7 +91,7 @@ func (k *KafkaInput) initialize(hosts []string, topic string, offsetFile string,
|
||||
}
|
||||
partitionConsumer, err := k.consumer.ConsumePartition(topic, partition, offset)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
return fmt.Errorf("consume kafka topic %q partition %d: %w", topic, partition, err)
|
||||
}
|
||||
go func() {
|
||||
for {
|
||||
|
||||
@@ -1039,13 +1039,6 @@ func getMD5HashBase64(data []byte) string {
|
||||
return base64.StdEncoding.EncodeToString(getMD5Sum(data))
|
||||
}
|
||||
|
||||
// getSHA256Sum returns SHA-256 sum of given data.
|
||||
func getSHA256Sum(data []byte) []byte {
|
||||
hash := sha256.New()
|
||||
hash.Write(data)
|
||||
return hash.Sum(nil)
|
||||
}
|
||||
|
||||
// getMD5Sum returns MD5 sum of given data.
|
||||
func getMD5Sum(data []byte) []byte {
|
||||
hash := md5.New()
|
||||
@@ -1053,11 +1046,6 @@ func getMD5Sum(data []byte) []byte {
|
||||
return hash.Sum(nil)
|
||||
}
|
||||
|
||||
// getMD5Hash returns MD5 hash in hex encoding of given data.
|
||||
func getMD5Hash(data []byte) string {
|
||||
return hex.EncodeToString(getMD5Sum(data))
|
||||
}
|
||||
|
||||
var ignoredHeaders = map[string]bool{
|
||||
"Authorization": true,
|
||||
"Content-Type": true,
|
||||
|
||||
@@ -539,10 +539,6 @@ func testJWTAuthentication(t *testing.T, iam *IdentityAccessManagement, token st
|
||||
return iam.authenticateJWTWithIAM(req)
|
||||
}
|
||||
|
||||
func testJWTAuthorization(t *testing.T, iam *IdentityAccessManagement, identity *Identity, action Action, bucket, object, token string) bool {
|
||||
return testJWTAuthorizationWithRole(t, iam, identity, action, bucket, object, token, "TestRole")
|
||||
}
|
||||
|
||||
func testJWTAuthorizationWithRole(t *testing.T, iam *IdentityAccessManagement, identity *Identity, action Action, bucket, object, token, roleName string) bool {
|
||||
// Create test request
|
||||
req := httptest.NewRequest("GET", "/"+bucket+"/"+object, http.NoBody)
|
||||
|
||||
@@ -1,67 +0,0 @@
|
||||
package s3api
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"crypto/md5"
|
||||
"encoding/base64"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"testing"
|
||||
|
||||
"github.com/gorilla/mux"
|
||||
"github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants"
|
||||
)
|
||||
|
||||
// ResponseRecorder that also implements http.Flusher
|
||||
type recorderFlusher struct{ *httptest.ResponseRecorder }
|
||||
|
||||
func (r recorderFlusher) Flush() {}
|
||||
|
||||
// TestSSECRangeRequestsSupported verifies that HTTP Range requests are now supported
|
||||
// for SSE-C encrypted objects since the IV is stored in metadata and CTR mode allows seeking
|
||||
func TestSSECRangeRequestsSupported(t *testing.T) {
|
||||
// Create a request with Range header and valid SSE-C headers
|
||||
req := httptest.NewRequest(http.MethodGet, "/b/o", nil)
|
||||
req.Header.Set("Range", "bytes=10-20")
|
||||
req.Header.Set(s3_constants.AmzServerSideEncryptionCustomerAlgorithm, "AES256")
|
||||
|
||||
key := make([]byte, 32)
|
||||
for i := range key {
|
||||
key[i] = byte(i)
|
||||
}
|
||||
s := md5.Sum(key)
|
||||
keyMD5 := base64.StdEncoding.EncodeToString(s[:])
|
||||
|
||||
req.Header.Set(s3_constants.AmzServerSideEncryptionCustomerKey, base64.StdEncoding.EncodeToString(key))
|
||||
req.Header.Set(s3_constants.AmzServerSideEncryptionCustomerKeyMD5, keyMD5)
|
||||
|
||||
// Attach mux vars to avoid panic in error writer
|
||||
req = mux.SetURLVars(req, map[string]string{"bucket": "b", "object": "o"})
|
||||
|
||||
// Create a mock HTTP response that simulates SSE-C encrypted object metadata
|
||||
proxyResponse := &http.Response{
|
||||
StatusCode: 200,
|
||||
Header: make(http.Header),
|
||||
Body: io.NopCloser(bytes.NewReader([]byte("mock encrypted data"))),
|
||||
}
|
||||
proxyResponse.Header.Set(s3_constants.AmzServerSideEncryptionCustomerAlgorithm, "AES256")
|
||||
proxyResponse.Header.Set(s3_constants.AmzServerSideEncryptionCustomerKeyMD5, keyMD5)
|
||||
|
||||
// Call the function under test - should no longer reject range requests
|
||||
s3a := &S3ApiServer{
|
||||
option: &S3ApiServerOption{
|
||||
BucketsPath: "/buckets",
|
||||
},
|
||||
}
|
||||
rec := httptest.NewRecorder()
|
||||
w := recorderFlusher{rec}
|
||||
// Pass nil for entry since this test focuses on Range request handling
|
||||
statusCode, _ := s3a.handleSSECResponse(req, proxyResponse, w, nil)
|
||||
|
||||
// Range requests should now be allowed to proceed (will be handled by filer layer)
|
||||
// The exact status code depends on the object existence and filer response
|
||||
if statusCode == http.StatusRequestedRangeNotSatisfiable {
|
||||
t.Fatalf("Range requests should no longer be rejected for SSE-C objects, got status %d", statusCode)
|
||||
}
|
||||
}
|
||||
@@ -331,47 +331,6 @@ func TestDetectPrimarySSETypeS3(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestAddSSES3HeadersToResponse tests that SSE-S3 headers are added to responses
|
||||
func TestAddSSES3HeadersToResponse(t *testing.T) {
|
||||
s3a := &S3ApiServer{}
|
||||
|
||||
entry := &filer_pb.Entry{
|
||||
Extended: map[string][]byte{
|
||||
s3_constants.AmzServerSideEncryption: []byte("AES256"),
|
||||
},
|
||||
Attributes: &filer_pb.FuseAttributes{},
|
||||
Chunks: []*filer_pb.FileChunk{
|
||||
{
|
||||
FileId: "1,123",
|
||||
Offset: 0,
|
||||
Size: 1024,
|
||||
SseType: filer_pb.SSEType_SSE_S3,
|
||||
SseMetadata: []byte("metadata"),
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
proxyResponse := &http.Response{
|
||||
Header: make(http.Header),
|
||||
}
|
||||
|
||||
s3a.addSSEHeadersToResponse(proxyResponse, entry)
|
||||
|
||||
algorithm := proxyResponse.Header.Get(s3_constants.AmzServerSideEncryption)
|
||||
if algorithm != "AES256" {
|
||||
t.Errorf("Expected SSE algorithm AES256, got %s", algorithm)
|
||||
}
|
||||
|
||||
// Should NOT have SSE-C or SSE-KMS specific headers
|
||||
if proxyResponse.Header.Get(s3_constants.AmzServerSideEncryptionCustomerAlgorithm) != "" {
|
||||
t.Error("Should not have SSE-C customer algorithm header")
|
||||
}
|
||||
|
||||
if proxyResponse.Header.Get(s3_constants.AmzServerSideEncryptionAwsKmsKeyId) != "" {
|
||||
t.Error("Should not have SSE-KMS key ID header")
|
||||
}
|
||||
}
|
||||
|
||||
// TestSSES3EncryptionWithBaseIV tests multipart encryption with base IV
|
||||
func TestSSES3EncryptionWithBaseIV(t *testing.T) {
|
||||
// Generate SSE-S3 key
|
||||
|
||||
@@ -24,7 +24,6 @@ import (
|
||||
"github.com/seaweedfs/seaweedfs/weed/s3api/s3_constants"
|
||||
"github.com/seaweedfs/seaweedfs/weed/s3api/s3err"
|
||||
util_http "github.com/seaweedfs/seaweedfs/weed/util/http"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util/mem"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/glog"
|
||||
)
|
||||
@@ -249,6 +248,10 @@ func newStreamErrorWithResponse(err error) *StreamError {
|
||||
return &StreamError{Err: err, ResponseWritten: true}
|
||||
}
|
||||
|
||||
func shouldWriteStreamingErrorResponse(err error) bool {
|
||||
return err != nil && !errors.Is(err, context.Canceled)
|
||||
}
|
||||
|
||||
func mimeDetect(r *http.Request, dataReader io.Reader) io.ReadCloser {
|
||||
mimeBuffer := make([]byte, 512)
|
||||
size, _ := dataReader.Read(mimeBuffer)
|
||||
@@ -880,7 +883,15 @@ func (s3a *S3ApiServer) GetObjectHandler(w http.ResponseWriter, r *http.Request)
|
||||
err = s3a.streamFromVolumeServersWithSSE(w, r, objectEntryForSSE, primarySSEType, bucket, object, versionId)
|
||||
streamTime = time.Since(tStream)
|
||||
if err != nil {
|
||||
glog.Errorf("GetObjectHandler: failed to stream %s/%s from volume servers: %v", bucket, object, err)
|
||||
switch {
|
||||
case errors.Is(err, context.Canceled):
|
||||
glog.V(3).Infof("GetObjectHandler: client disconnected while streaming %s/%s: %v", bucket, object, err)
|
||||
return
|
||||
case errors.Is(err, context.DeadlineExceeded):
|
||||
glog.Warningf("GetObjectHandler: deadline exceeded while streaming %s/%s: %v", bucket, object, err)
|
||||
default:
|
||||
glog.Errorf("GetObjectHandler: failed to stream %s/%s from volume servers: %v", bucket, object, err)
|
||||
}
|
||||
// Check if the streaming function already wrote an HTTP response
|
||||
var streamErr *StreamError
|
||||
if errors.As(err, &streamErr) && streamErr.ResponseWritten {
|
||||
@@ -892,7 +903,7 @@ func (s3a *S3ApiServer) GetObjectHandler(w http.ResponseWriter, r *http.Request)
|
||||
// Check if error is due to volume server rate limiting (HTTP 429)
|
||||
if errors.Is(err, util_http.ErrTooManyRequests) {
|
||||
s3err.WriteErrorResponse(w, r, s3err.ErrRequestBytesExceed)
|
||||
} else {
|
||||
} else if shouldWriteStreamingErrorResponse(err) {
|
||||
s3err.WriteErrorResponse(w, r, s3err.ErrInternalError)
|
||||
}
|
||||
return
|
||||
@@ -1028,7 +1039,15 @@ func (s3a *S3ApiServer) streamFromVolumeServers(w http.ResponseWriter, r *http.R
|
||||
resolvedChunks, _, err := filer.ResolveChunkManifest(ctx, lookupFileIdFn, chunks, offset, offset+size)
|
||||
chunkResolveTime = time.Since(tChunkResolve)
|
||||
if err != nil {
|
||||
glog.Errorf("streamFromVolumeServers: failed to resolve chunks: %v", err)
|
||||
if errors.Is(err, context.Canceled) {
|
||||
glog.V(3).Infof("streamFromVolumeServers: request canceled while resolving chunks: %v", err)
|
||||
return err
|
||||
}
|
||||
if errors.Is(err, context.DeadlineExceeded) {
|
||||
glog.Warningf("streamFromVolumeServers: request deadline exceeded while resolving chunks: %v", err)
|
||||
} else {
|
||||
glog.Errorf("streamFromVolumeServers: failed to resolve chunks: %v", err)
|
||||
}
|
||||
// Write S3-compliant XML error response
|
||||
s3err.WriteErrorResponse(w, r, s3err.ErrInternalError)
|
||||
return newStreamErrorWithResponse(fmt.Errorf("failed to resolve chunks: %v", err))
|
||||
@@ -1048,7 +1067,15 @@ func (s3a *S3ApiServer) streamFromVolumeServers(w http.ResponseWriter, r *http.R
|
||||
)
|
||||
streamPrepTime = time.Since(tStreamPrep)
|
||||
if err != nil {
|
||||
glog.Errorf("streamFromVolumeServers: failed to prepare stream: %v", err)
|
||||
if errors.Is(err, context.Canceled) {
|
||||
glog.V(3).Infof("streamFromVolumeServers: request canceled while preparing stream: %v", err)
|
||||
return err
|
||||
}
|
||||
if errors.Is(err, context.DeadlineExceeded) {
|
||||
glog.Warningf("streamFromVolumeServers: request deadline exceeded while preparing stream: %v", err)
|
||||
} else {
|
||||
glog.Errorf("streamFromVolumeServers: failed to prepare stream: %v", err)
|
||||
}
|
||||
// Write S3-compliant XML error response
|
||||
s3err.WriteErrorResponse(w, r, s3err.ErrInternalError)
|
||||
return newStreamErrorWithResponse(fmt.Errorf("failed to prepare stream: %v", err))
|
||||
@@ -2394,45 +2421,6 @@ func (s3a *S3ApiServer) HeadObjectHandler(w http.ResponseWriter, r *http.Request
|
||||
w.WriteHeader(http.StatusOK)
|
||||
}
|
||||
|
||||
func captureCORSHeaders(w http.ResponseWriter, headersToCapture []string) map[string]string {
|
||||
captured := make(map[string]string)
|
||||
for _, corsHeader := range headersToCapture {
|
||||
if value := w.Header().Get(corsHeader); value != "" {
|
||||
captured[corsHeader] = value
|
||||
}
|
||||
}
|
||||
return captured
|
||||
}
|
||||
|
||||
func restoreCORSHeaders(w http.ResponseWriter, capturedCORSHeaders map[string]string) {
|
||||
for corsHeader, value := range capturedCORSHeaders {
|
||||
w.Header().Set(corsHeader, value)
|
||||
}
|
||||
}
|
||||
|
||||
// writeFinalResponse handles the common response writing logic shared between
|
||||
// passThroughResponse and handleSSECResponse
|
||||
func writeFinalResponse(w http.ResponseWriter, proxyResponse *http.Response, bodyReader io.Reader, capturedCORSHeaders map[string]string) (statusCode int, bytesTransferred int64) {
|
||||
// Restore CORS headers that were set by middleware
|
||||
restoreCORSHeaders(w, capturedCORSHeaders)
|
||||
|
||||
if proxyResponse.Header.Get("Content-Range") != "" && proxyResponse.StatusCode == 200 {
|
||||
statusCode = http.StatusPartialContent
|
||||
} else {
|
||||
statusCode = proxyResponse.StatusCode
|
||||
}
|
||||
w.WriteHeader(statusCode)
|
||||
|
||||
// Stream response data
|
||||
buf := mem.Allocate(128 * 1024)
|
||||
defer mem.Free(buf)
|
||||
bytesTransferred, err := io.CopyBuffer(w, bodyReader, buf)
|
||||
if err != nil {
|
||||
glog.V(1).Infof("response read %d bytes: %v", bytesTransferred, err)
|
||||
}
|
||||
return statusCode, bytesTransferred
|
||||
}
|
||||
|
||||
// fetchObjectEntry fetches the filer entry for an object
|
||||
// Returns nil if not found (not an error), or propagates other errors
|
||||
func (s3a *S3ApiServer) fetchObjectEntry(bucket, object string) (*filer_pb.Entry, error) {
|
||||
@@ -2458,187 +2446,6 @@ func (s3a *S3ApiServer) fetchObjectEntryRequired(bucket, object string) (*filer_
|
||||
return fetchedEntry, nil
|
||||
}
|
||||
|
||||
// copyResponseHeaders copies headers from proxy response to the response writer,
|
||||
// excluding internal SeaweedFS headers and optionally excluding body-related headers
|
||||
func copyResponseHeaders(w http.ResponseWriter, proxyResponse *http.Response, excludeBodyHeaders bool) {
|
||||
for k, v := range proxyResponse.Header {
|
||||
// Always exclude internal SeaweedFS headers
|
||||
if s3_constants.IsSeaweedFSInternalHeader(k) {
|
||||
continue
|
||||
}
|
||||
// Optionally exclude body-related headers that might change after decryption
|
||||
if excludeBodyHeaders && (k == "Content-Length" || k == "Content-Encoding") {
|
||||
continue
|
||||
}
|
||||
w.Header()[k] = v
|
||||
}
|
||||
}
|
||||
|
||||
func passThroughResponse(proxyResponse *http.Response, w http.ResponseWriter) (statusCode int, bytesTransferred int64) {
|
||||
// Capture existing CORS headers that may have been set by middleware
|
||||
capturedCORSHeaders := captureCORSHeaders(w, corsHeaders)
|
||||
|
||||
// Copy headers from proxy response (excluding internal SeaweedFS headers)
|
||||
copyResponseHeaders(w, proxyResponse, false)
|
||||
|
||||
return writeFinalResponse(w, proxyResponse, proxyResponse.Body, capturedCORSHeaders)
|
||||
}
|
||||
|
||||
// handleSSECResponse handles SSE-C decryption and response processing
|
||||
func (s3a *S3ApiServer) handleSSECResponse(r *http.Request, proxyResponse *http.Response, w http.ResponseWriter, entry *filer_pb.Entry) (statusCode int, bytesTransferred int64) {
|
||||
// Check if the object has SSE-C metadata
|
||||
sseAlgorithm := proxyResponse.Header.Get(s3_constants.AmzServerSideEncryptionCustomerAlgorithm)
|
||||
sseKeyMD5 := proxyResponse.Header.Get(s3_constants.AmzServerSideEncryptionCustomerKeyMD5)
|
||||
isObjectEncrypted := sseAlgorithm != "" && sseKeyMD5 != ""
|
||||
|
||||
// Parse SSE-C headers from request once (avoid duplication)
|
||||
customerKey, err := ParseSSECHeaders(r)
|
||||
if err != nil {
|
||||
errCode := MapSSECErrorToS3Error(err)
|
||||
s3err.WriteErrorResponse(w, r, errCode)
|
||||
return http.StatusBadRequest, 0
|
||||
}
|
||||
|
||||
if isObjectEncrypted {
|
||||
// This object was encrypted with SSE-C, validate customer key
|
||||
if customerKey == nil {
|
||||
s3err.WriteErrorResponse(w, r, s3err.ErrSSECustomerKeyMissing)
|
||||
return http.StatusBadRequest, 0
|
||||
}
|
||||
|
||||
// SSE-C MD5 is base64 and case-sensitive
|
||||
if customerKey.KeyMD5 != sseKeyMD5 {
|
||||
// For GET/HEAD requests, AWS S3 returns 403 Forbidden for a key mismatch.
|
||||
s3err.WriteErrorResponse(w, r, s3err.ErrAccessDenied)
|
||||
return http.StatusForbidden, 0
|
||||
}
|
||||
|
||||
// SSE-C encrypted objects support HTTP Range requests
|
||||
// The IV is stored in metadata and CTR mode allows seeking to any offset
|
||||
// Range requests will be handled by the filer layer with proper offset-based decryption
|
||||
|
||||
// Check if this is a chunked or small content SSE-C object
|
||||
// Use the entry parameter passed from the caller (avoids redundant lookup)
|
||||
if entry != nil {
|
||||
// Check for SSE-C chunks
|
||||
sseCChunks := 0
|
||||
for _, chunk := range entry.GetChunks() {
|
||||
if chunk.GetSseType() == filer_pb.SSEType_SSE_C {
|
||||
sseCChunks++
|
||||
}
|
||||
}
|
||||
|
||||
if sseCChunks >= 1 {
|
||||
|
||||
// Handle chunked SSE-C objects - each chunk needs independent decryption
|
||||
multipartReader, decErr := s3a.createMultipartSSECDecryptedReader(r, proxyResponse, entry)
|
||||
if decErr != nil {
|
||||
glog.Errorf("Failed to create multipart SSE-C decrypted reader: %v", decErr)
|
||||
s3err.WriteErrorResponse(w, r, s3err.ErrInternalError)
|
||||
return http.StatusInternalServerError, 0
|
||||
}
|
||||
|
||||
// Capture existing CORS headers
|
||||
capturedCORSHeaders := captureCORSHeaders(w, corsHeaders)
|
||||
|
||||
// Copy headers from proxy response (excluding internal SeaweedFS headers)
|
||||
copyResponseHeaders(w, proxyResponse, false)
|
||||
|
||||
// Set proper headers for range requests
|
||||
rangeHeader := r.Header.Get("Range")
|
||||
if rangeHeader != "" {
|
||||
|
||||
// Parse range header (e.g., "bytes=0-99")
|
||||
if len(rangeHeader) > 6 && rangeHeader[:6] == "bytes=" {
|
||||
rangeSpec := rangeHeader[6:]
|
||||
parts := strings.Split(rangeSpec, "-")
|
||||
if len(parts) == 2 {
|
||||
startOffset, endOffset := int64(0), int64(-1)
|
||||
if parts[0] != "" {
|
||||
startOffset, _ = strconv.ParseInt(parts[0], 10, 64)
|
||||
}
|
||||
if parts[1] != "" {
|
||||
endOffset, _ = strconv.ParseInt(parts[1], 10, 64)
|
||||
}
|
||||
|
||||
if endOffset >= startOffset {
|
||||
// Specific range - set proper Content-Length and Content-Range headers
|
||||
rangeLength := endOffset - startOffset + 1
|
||||
totalSize := proxyResponse.Header.Get("Content-Length")
|
||||
|
||||
w.Header().Set("Content-Length", strconv.FormatInt(rangeLength, 10))
|
||||
w.Header().Set("Content-Range", fmt.Sprintf("bytes %d-%d/%s", startOffset, endOffset, totalSize))
|
||||
// writeFinalResponse will set status to 206 if Content-Range is present
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return writeFinalResponse(w, proxyResponse, multipartReader, capturedCORSHeaders)
|
||||
} else if len(entry.GetChunks()) == 0 && len(entry.Content) > 0 {
|
||||
// Small content SSE-C object stored directly in entry.Content
|
||||
|
||||
// Fall through to traditional single-object SSE-C handling below
|
||||
}
|
||||
}
|
||||
|
||||
// Single-part SSE-C object: Get IV from proxy response headers (stored during upload)
|
||||
ivBase64 := proxyResponse.Header.Get(s3_constants.SeaweedFSSSEIVHeader)
|
||||
if ivBase64 == "" {
|
||||
glog.Errorf("SSE-C encrypted single-part object missing IV in metadata")
|
||||
s3err.WriteErrorResponse(w, r, s3err.ErrInternalError)
|
||||
return http.StatusInternalServerError, 0
|
||||
}
|
||||
|
||||
iv, err := base64.StdEncoding.DecodeString(ivBase64)
|
||||
if err != nil {
|
||||
glog.Errorf("Failed to decode IV from metadata: %v", err)
|
||||
s3err.WriteErrorResponse(w, r, s3err.ErrInternalError)
|
||||
return http.StatusInternalServerError, 0
|
||||
}
|
||||
|
||||
// Create decrypted reader with IV from metadata
|
||||
decryptedReader, decErr := CreateSSECDecryptedReader(proxyResponse.Body, customerKey, iv)
|
||||
if decErr != nil {
|
||||
glog.Errorf("Failed to create SSE-C decrypted reader: %v", decErr)
|
||||
s3err.WriteErrorResponse(w, r, s3err.ErrInternalError)
|
||||
return http.StatusInternalServerError, 0
|
||||
}
|
||||
|
||||
// Capture existing CORS headers that may have been set by middleware
|
||||
capturedCORSHeaders := captureCORSHeaders(w, corsHeaders)
|
||||
|
||||
// Copy headers from proxy response (excluding body-related headers that might change and internal SeaweedFS headers)
|
||||
copyResponseHeaders(w, proxyResponse, true)
|
||||
|
||||
// Set correct Content-Length for SSE-C (only for full object requests)
|
||||
// With IV stored in metadata, the encrypted length equals the original length
|
||||
if proxyResponse.Header.Get("Content-Range") == "" {
|
||||
// Full object request: encrypted length equals original length (IV not in stream)
|
||||
if contentLengthStr := proxyResponse.Header.Get("Content-Length"); contentLengthStr != "" {
|
||||
// Content-Length is already correct since IV is stored in metadata, not in data stream
|
||||
w.Header().Set("Content-Length", contentLengthStr)
|
||||
}
|
||||
}
|
||||
// For range requests, let the actual bytes transferred determine the response length
|
||||
|
||||
// Add SSE-C response headers
|
||||
w.Header().Set(s3_constants.AmzServerSideEncryptionCustomerAlgorithm, sseAlgorithm)
|
||||
w.Header().Set(s3_constants.AmzServerSideEncryptionCustomerKeyMD5, sseKeyMD5)
|
||||
|
||||
return writeFinalResponse(w, proxyResponse, decryptedReader, capturedCORSHeaders)
|
||||
} else {
|
||||
// Object is not encrypted, but check if customer provided SSE-C headers unnecessarily
|
||||
if customerKey != nil {
|
||||
s3err.WriteErrorResponse(w, r, s3err.ErrSSECustomerKeyNotNeeded)
|
||||
return http.StatusBadRequest, 0
|
||||
}
|
||||
|
||||
// Normal pass-through response
|
||||
return passThroughResponse(proxyResponse, w)
|
||||
}
|
||||
}
|
||||
|
||||
// addObjectLockHeadersToResponse extracts object lock metadata from entry Extended attributes
|
||||
// and adds the appropriate S3 headers to the response
|
||||
func (s3a *S3ApiServer) addObjectLockHeadersToResponse(w http.ResponseWriter, entry *filer_pb.Entry) {
|
||||
@@ -2680,54 +2487,6 @@ func (s3a *S3ApiServer) addObjectLockHeadersToResponse(w http.ResponseWriter, en
|
||||
}
|
||||
}
|
||||
|
||||
// addSSEHeadersToResponse converts stored SSE metadata from entry.Extended to HTTP response headers
|
||||
// Uses intelligent prioritization: only set headers for the PRIMARY encryption type to avoid conflicts
|
||||
func (s3a *S3ApiServer) addSSEHeadersToResponse(proxyResponse *http.Response, entry *filer_pb.Entry) {
|
||||
if entry == nil || entry.Extended == nil {
|
||||
return
|
||||
}
|
||||
|
||||
// Determine the primary encryption type by examining chunks (most reliable)
|
||||
primarySSEType := s3a.detectPrimarySSEType(entry)
|
||||
|
||||
// Only set headers for the PRIMARY encryption type
|
||||
switch primarySSEType {
|
||||
case s3_constants.SSETypeC:
|
||||
// Add only SSE-C headers
|
||||
if algorithmBytes, exists := entry.Extended[s3_constants.AmzServerSideEncryptionCustomerAlgorithm]; exists && len(algorithmBytes) > 0 {
|
||||
proxyResponse.Header.Set(s3_constants.AmzServerSideEncryptionCustomerAlgorithm, string(algorithmBytes))
|
||||
}
|
||||
|
||||
if keyMD5Bytes, exists := entry.Extended[s3_constants.AmzServerSideEncryptionCustomerKeyMD5]; exists && len(keyMD5Bytes) > 0 {
|
||||
proxyResponse.Header.Set(s3_constants.AmzServerSideEncryptionCustomerKeyMD5, string(keyMD5Bytes))
|
||||
}
|
||||
|
||||
if ivBytes, exists := entry.Extended[s3_constants.SeaweedFSSSEIV]; exists && len(ivBytes) > 0 {
|
||||
ivBase64 := base64.StdEncoding.EncodeToString(ivBytes)
|
||||
proxyResponse.Header.Set(s3_constants.SeaweedFSSSEIVHeader, ivBase64)
|
||||
}
|
||||
|
||||
case s3_constants.SSETypeKMS:
|
||||
// Add only SSE-KMS headers
|
||||
if sseAlgorithm, exists := entry.Extended[s3_constants.AmzServerSideEncryption]; exists && len(sseAlgorithm) > 0 {
|
||||
proxyResponse.Header.Set(s3_constants.AmzServerSideEncryption, string(sseAlgorithm))
|
||||
}
|
||||
|
||||
if kmsKeyID, exists := entry.Extended[s3_constants.AmzServerSideEncryptionAwsKmsKeyId]; exists && len(kmsKeyID) > 0 {
|
||||
proxyResponse.Header.Set(s3_constants.AmzServerSideEncryptionAwsKmsKeyId, string(kmsKeyID))
|
||||
}
|
||||
|
||||
case s3_constants.SSETypeS3:
|
||||
// Add only SSE-S3 headers
|
||||
proxyResponse.Header.Set(s3_constants.AmzServerSideEncryption, SSES3Algorithm)
|
||||
|
||||
default:
|
||||
// Unencrypted or unknown - don't set any SSE headers
|
||||
}
|
||||
|
||||
glog.V(3).Infof("addSSEHeadersToResponse: processed %d extended metadata entries", len(entry.Extended))
|
||||
}
|
||||
|
||||
// detectPrimarySSEType determines the primary SSE type by examining chunk metadata
|
||||
func (s3a *S3ApiServer) detectPrimarySSEType(entry *filer_pb.Entry) string {
|
||||
// Safety check: handle nil entry
|
||||
@@ -3183,140 +2942,6 @@ func (r *SSERangeReader) Read(p []byte) (n int, err error) {
|
||||
return n, err
|
||||
}
|
||||
|
||||
// createMultipartSSECDecryptedReader creates a decrypted reader for multipart SSE-C objects
|
||||
// Each chunk has its own IV and encryption key from the original multipart parts
|
||||
func (s3a *S3ApiServer) createMultipartSSECDecryptedReader(r *http.Request, proxyResponse *http.Response, entry *filer_pb.Entry) (io.Reader, error) {
|
||||
ctx := r.Context()
|
||||
|
||||
// Parse SSE-C headers from the request for decryption key
|
||||
customerKey, err := ParseSSECHeaders(r)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("invalid SSE-C headers for multipart decryption: %v", err)
|
||||
}
|
||||
|
||||
// Entry is passed from caller to avoid redundant filer lookup
|
||||
|
||||
// Sort chunks by offset to ensure correct order
|
||||
chunks := entry.GetChunks()
|
||||
sort.Slice(chunks, func(i, j int) bool {
|
||||
return chunks[i].GetOffset() < chunks[j].GetOffset()
|
||||
})
|
||||
|
||||
// Check for Range header to optimize chunk processing
|
||||
var startOffset, endOffset int64 = 0, -1
|
||||
rangeHeader := r.Header.Get("Range")
|
||||
if rangeHeader != "" {
|
||||
// Parse range header (e.g., "bytes=0-99")
|
||||
if len(rangeHeader) > 6 && rangeHeader[:6] == "bytes=" {
|
||||
rangeSpec := rangeHeader[6:]
|
||||
parts := strings.Split(rangeSpec, "-")
|
||||
if len(parts) == 2 {
|
||||
if parts[0] != "" {
|
||||
startOffset, _ = strconv.ParseInt(parts[0], 10, 64)
|
||||
}
|
||||
if parts[1] != "" {
|
||||
endOffset, _ = strconv.ParseInt(parts[1], 10, 64)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Filter chunks to only those needed for the range request
|
||||
var neededChunks []*filer_pb.FileChunk
|
||||
for _, chunk := range chunks {
|
||||
chunkStart := chunk.GetOffset()
|
||||
chunkEnd := chunkStart + int64(chunk.GetSize()) - 1
|
||||
|
||||
// Check if this chunk overlaps with the requested range
|
||||
if endOffset == -1 {
|
||||
// No end specified, take all chunks from startOffset
|
||||
if chunkEnd >= startOffset {
|
||||
neededChunks = append(neededChunks, chunk)
|
||||
}
|
||||
} else {
|
||||
// Specific range: check for overlap
|
||||
if chunkStart <= endOffset && chunkEnd >= startOffset {
|
||||
neededChunks = append(neededChunks, chunk)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Create readers for only the needed chunks
|
||||
var readers []io.Reader
|
||||
|
||||
for _, chunk := range neededChunks {
|
||||
|
||||
// Get this chunk's encrypted data
|
||||
chunkReader, err := s3a.createEncryptedChunkReader(ctx, chunk)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to create chunk reader: %v", err)
|
||||
}
|
||||
|
||||
if chunk.GetSseType() == filer_pb.SSEType_SSE_C {
|
||||
// For SSE-C chunks, extract the IV from the stored per-chunk metadata (unified approach)
|
||||
if len(chunk.GetSseMetadata()) > 0 {
|
||||
// Deserialize the SSE-C metadata stored in the unified metadata field
|
||||
ssecMetadata, decErr := DeserializeSSECMetadata(chunk.GetSseMetadata())
|
||||
if decErr != nil {
|
||||
chunkReader.Close()
|
||||
return nil, fmt.Errorf("failed to deserialize SSE-C metadata for chunk %s: %v", chunk.GetFileIdString(), decErr)
|
||||
}
|
||||
|
||||
// Decode the IV from the metadata
|
||||
iv, ivErr := base64.StdEncoding.DecodeString(ssecMetadata.IV)
|
||||
if ivErr != nil {
|
||||
chunkReader.Close()
|
||||
return nil, fmt.Errorf("failed to decode IV for SSE-C chunk %s: %v", chunk.GetFileIdString(), ivErr)
|
||||
}
|
||||
|
||||
partOffset := ssecMetadata.PartOffset
|
||||
if partOffset < 0 {
|
||||
chunkReader.Close()
|
||||
return nil, fmt.Errorf("invalid SSE-C part offset %d for chunk %s", partOffset, chunk.GetFileIdString())
|
||||
}
|
||||
|
||||
// Use stored IV and advance CTR stream by PartOffset within the encrypted stream
|
||||
decryptedReader, decErr := CreateSSECDecryptedReaderWithOffset(chunkReader, customerKey, iv, uint64(partOffset))
|
||||
if decErr != nil {
|
||||
chunkReader.Close()
|
||||
return nil, fmt.Errorf("failed to create SSE-C decrypted reader for chunk %s: %v", chunk.GetFileIdString(), decErr)
|
||||
}
|
||||
readers = append(readers, decryptedReader)
|
||||
} else {
|
||||
chunkReader.Close()
|
||||
return nil, fmt.Errorf("SSE-C chunk %s missing required metadata", chunk.GetFileIdString())
|
||||
}
|
||||
} else {
|
||||
// Non-SSE-C chunk, use as-is
|
||||
readers = append(readers, chunkReader)
|
||||
}
|
||||
}
|
||||
|
||||
multiReader := NewMultipartSSEReader(readers)
|
||||
|
||||
// Apply range logic if a range was requested
|
||||
if rangeHeader != "" && startOffset >= 0 {
|
||||
if endOffset == -1 {
|
||||
// Open-ended range (e.g., "bytes=100-")
|
||||
return &SSERangeReader{
|
||||
reader: multiReader,
|
||||
offset: startOffset,
|
||||
remaining: -1, // Read until EOF
|
||||
}, nil
|
||||
} else {
|
||||
// Specific range (e.g., "bytes=0-99")
|
||||
rangeLength := endOffset - startOffset + 1
|
||||
return &SSERangeReader{
|
||||
reader: multiReader,
|
||||
offset: startOffset,
|
||||
remaining: rangeLength,
|
||||
}, nil
|
||||
}
|
||||
}
|
||||
|
||||
return multiReader, nil
|
||||
}
|
||||
|
||||
// PartBoundaryInfo holds information about a part's chunk boundaries
|
||||
type PartBoundaryInfo struct {
|
||||
PartNumber int `json:"part"`
|
||||
|
||||
@@ -9,22 +9,6 @@ import (
|
||||
"github.com/stretchr/testify/assert"
|
||||
)
|
||||
|
||||
// mockAccountManager implements AccountManager for testing
|
||||
type mockAccountManager struct {
|
||||
accounts map[string]string
|
||||
}
|
||||
|
||||
func (m *mockAccountManager) GetAccountNameById(id string) string {
|
||||
if name, exists := m.accounts[id]; exists {
|
||||
return name
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
func (m *mockAccountManager) GetAccountIdByEmail(email string) string {
|
||||
return ""
|
||||
}
|
||||
|
||||
func TestNewListEntryOwnerDisplayName(t *testing.T) {
|
||||
// Create S3ApiServer with a properly initialized IAM
|
||||
s3a := &S3ApiServer{
|
||||
|
||||
@@ -0,0 +1,48 @@
|
||||
package s3api
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestShouldWriteStreamingErrorResponse(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
err error
|
||||
expected bool
|
||||
}{
|
||||
{
|
||||
name: "nil error",
|
||||
err: nil,
|
||||
expected: false,
|
||||
},
|
||||
{
|
||||
name: "context canceled",
|
||||
err: context.Canceled,
|
||||
expected: false,
|
||||
},
|
||||
{
|
||||
name: "wrapped context canceled",
|
||||
err: &StreamError{Err: context.Canceled},
|
||||
expected: false,
|
||||
},
|
||||
{
|
||||
name: "deadline exceeded",
|
||||
err: context.DeadlineExceeded,
|
||||
expected: true,
|
||||
},
|
||||
{
|
||||
name: "wrapped deadline exceeded",
|
||||
err: &StreamError{Err: context.DeadlineExceeded},
|
||||
expected: true,
|
||||
},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
if got := shouldWriteStreamingErrorResponse(tt.err); got != tt.expected {
|
||||
t.Fatalf("shouldWriteStreamingErrorResponse(%v) = %v, want %v", tt.err, got, tt.expected)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -388,11 +388,7 @@ func (fs *FilerServer) DeleteEntry(ctx context.Context, req *filer_pb.DeleteEntr
|
||||
|
||||
func (fs *FilerServer) AssignVolume(ctx context.Context, req *filer_pb.AssignVolumeRequest) (resp *filer_pb.AssignVolumeResponse, err error) {
|
||||
|
||||
if req.DiskType == "" {
|
||||
req.DiskType = fs.option.DiskType
|
||||
}
|
||||
|
||||
so, err := fs.detectStorageOption(ctx, req.Path, req.Collection, req.Replication, req.TtlSec, req.DiskType, req.DataCenter, req.Rack, req.DataNode)
|
||||
so, err := fs.resolveAssignStorageOption(ctx, req)
|
||||
if err != nil {
|
||||
glog.V(3).InfofCtx(ctx, "AssignVolume: %v", err)
|
||||
return &filer_pb.AssignVolumeResponse{Error: fmt.Sprintf("assign volume: %v", err)}, nil
|
||||
@@ -424,6 +420,21 @@ func (fs *FilerServer) AssignVolume(ctx context.Context, req *filer_pb.AssignVol
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (fs *FilerServer) resolveAssignStorageOption(ctx context.Context, req *filer_pb.AssignVolumeRequest) (*operation.StorageOption, error) {
|
||||
so, err := fs.detectStorageOption(ctx, req.Path, req.Collection, req.Replication, req.TtlSec, req.DiskType, req.DataCenter, req.Rack, req.DataNode)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// Mirror the HTTP write path: only apply the filer's default disk when the
|
||||
// matched locationPrefix rule did not already select one.
|
||||
if so.DiskType == "" {
|
||||
so.DiskType = fs.option.DiskType
|
||||
}
|
||||
|
||||
return so, nil
|
||||
}
|
||||
|
||||
func (fs *FilerServer) CollectionList(ctx context.Context, req *filer_pb.CollectionListRequest) (resp *filer_pb.CollectionListResponse, err error) {
|
||||
|
||||
glog.V(4).InfofCtx(ctx, "CollectionList %v", req)
|
||||
|
||||
@@ -0,0 +1,68 @@
|
||||
package weed_server
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/filer"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
)
|
||||
|
||||
func TestResolveAssignStorageOptionUsesBucketRuleBeforeFilerDiskDefault(t *testing.T) {
|
||||
fc := filer.NewFilerConf()
|
||||
if err := fc.SetLocationConf(&filer_pb.FilerConf_PathConf{
|
||||
LocationPrefix: "/buckets/zot",
|
||||
DiskType: "disk",
|
||||
}); err != nil {
|
||||
t.Fatalf("set location conf: %v", err)
|
||||
}
|
||||
|
||||
fs := &FilerServer{
|
||||
option: &FilerOption{
|
||||
DiskType: "hdd",
|
||||
},
|
||||
filer: &filer.Filer{
|
||||
DirBucketsPath: "/buckets",
|
||||
FilerConf: fc,
|
||||
MaxFilenameLength: 255,
|
||||
},
|
||||
}
|
||||
|
||||
so, err := fs.resolveAssignStorageOption(context.Background(), &filer_pb.AssignVolumeRequest{
|
||||
Path: "/buckets/zot/.uploads/upload-id/0001_part.part",
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("resolve assign storage option: %v", err)
|
||||
}
|
||||
|
||||
if got, want := so.Collection, "zot"; got != want {
|
||||
t.Fatalf("collection = %q, want %q", got, want)
|
||||
}
|
||||
if got, want := so.DiskType, "disk"; got != want {
|
||||
t.Fatalf("disk type = %q, want %q", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestResolveAssignStorageOptionFallsBackToFilerDiskDefault(t *testing.T) {
|
||||
fs := &FilerServer{
|
||||
option: &FilerOption{
|
||||
DiskType: "hdd",
|
||||
},
|
||||
filer: &filer.Filer{
|
||||
DirBucketsPath: "/buckets",
|
||||
FilerConf: filer.NewFilerConf(),
|
||||
MaxFilenameLength: 255,
|
||||
},
|
||||
}
|
||||
|
||||
so, err := fs.resolveAssignStorageOption(context.Background(), &filer_pb.AssignVolumeRequest{
|
||||
Path: "/tmp/unmatched/file.bin",
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("resolve assign storage option: %v", err)
|
||||
}
|
||||
|
||||
if got, want := so.DiskType, "hdd"; got != want {
|
||||
t.Fatalf("disk type = %q, want %q", got, want)
|
||||
}
|
||||
}
|
||||
@@ -321,6 +321,11 @@ func (fs *FilerServer) tusWriteData(ctx context.Context, session *TusSession, of
|
||||
return 0, fmt.Errorf("detect storage option: %w", err)
|
||||
}
|
||||
|
||||
// When DiskType is empty, use filer's -disk
|
||||
if so.DiskType == "" {
|
||||
so.DiskType = fs.option.DiskType
|
||||
}
|
||||
|
||||
// Read first bytes for MIME type detection
|
||||
sniffSize := int64(512)
|
||||
if contentLength < sniffSize {
|
||||
|
||||
@@ -4,7 +4,6 @@ import (
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"math"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
@@ -146,14 +145,14 @@ func (t *ErasureCodingTask) Execute(ctx context.Context, params *worker_pb.TaskP
|
||||
// Step 1: Mark volume readonly
|
||||
t.ReportProgressWithStage(10.0, "Marking volume readonly")
|
||||
t.GetLogger().Info("Marking volume readonly")
|
||||
if err := t.markVolumeReadonly(); err != nil {
|
||||
if err := t.markVolumeReadonly(ctx); err != nil {
|
||||
return fmt.Errorf("failed to mark volume readonly: %v", err)
|
||||
}
|
||||
|
||||
// Step 2: Copy volume files to worker
|
||||
t.ReportProgressWithStage(25.0, "Copying volume files to worker")
|
||||
t.GetLogger().Info("Copying volume files to worker")
|
||||
localFiles, err := t.copyVolumeFilesToWorker(taskWorkDir)
|
||||
localFiles, err := t.copyVolumeFilesToWorker(ctx, taskWorkDir)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to copy volume files: %v", err)
|
||||
}
|
||||
@@ -183,7 +182,7 @@ func (t *ErasureCodingTask) Execute(ctx context.Context, params *worker_pb.TaskP
|
||||
// Step 6: Delete original volume
|
||||
t.ReportProgressWithStage(90.0, "Deleting original volume")
|
||||
t.GetLogger().Info("Deleting original volume")
|
||||
if err := t.deleteOriginalVolume(); err != nil {
|
||||
if err := t.deleteOriginalVolume(ctx); err != nil {
|
||||
return fmt.Errorf("failed to delete original volume: %v", err)
|
||||
}
|
||||
|
||||
@@ -250,10 +249,10 @@ func (t *ErasureCodingTask) GetProgress() float64 {
|
||||
// Helper methods for actual EC operations
|
||||
|
||||
// markVolumeReadonly marks the volume as readonly on the source server
|
||||
func (t *ErasureCodingTask) markVolumeReadonly() error {
|
||||
func (t *ErasureCodingTask) markVolumeReadonly(ctx context.Context) error {
|
||||
return operation.WithVolumeServerClient(false, pb.ServerAddress(t.server), t.grpcDialOption,
|
||||
func(client volume_server_pb.VolumeServerClient) error {
|
||||
_, err := client.VolumeMarkReadonly(context.Background(), &volume_server_pb.VolumeMarkReadonlyRequest{
|
||||
_, err := client.VolumeMarkReadonly(ctx, &volume_server_pb.VolumeMarkReadonlyRequest{
|
||||
VolumeId: t.volumeID,
|
||||
})
|
||||
return err
|
||||
@@ -261,18 +260,26 @@ func (t *ErasureCodingTask) markVolumeReadonly() error {
|
||||
}
|
||||
|
||||
// copyVolumeFilesToWorker copies .dat and .idx files from source server to local worker
|
||||
func (t *ErasureCodingTask) copyVolumeFilesToWorker(workDir string) (map[string]string, error) {
|
||||
func (t *ErasureCodingTask) copyVolumeFilesToWorker(ctx context.Context, workDir string) (map[string]string, error) {
|
||||
localFiles := make(map[string]string)
|
||||
|
||||
fileStatus, err := t.readSourceVolumeFileStatus(ctx)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to read source volume file status: %v", err)
|
||||
}
|
||||
|
||||
t.GetLogger().WithFields(map[string]interface{}{
|
||||
"volume_id": t.volumeID,
|
||||
"source": t.server,
|
||||
"working_dir": workDir,
|
||||
"volume_id": t.volumeID,
|
||||
"source": t.server,
|
||||
"working_dir": workDir,
|
||||
"compaction_revision": fileStatus.GetCompactionRevision(),
|
||||
"dat_file_size_bytes": fileStatus.GetDatFileSize(),
|
||||
"idx_file_size_bytes": fileStatus.GetIdxFileSize(),
|
||||
}).Info("Starting volume file copy from source server")
|
||||
|
||||
// Copy .dat file
|
||||
datFile := filepath.Join(workDir, fmt.Sprintf("%d.dat", t.volumeID))
|
||||
if err := t.copyFileFromSource(".dat", datFile); err != nil {
|
||||
if err := t.copyFileFromSource(ctx, ".dat", datFile, fileStatus.GetCompactionRevision(), fileStatus.GetDatFileSize()); err != nil {
|
||||
return nil, fmt.Errorf("failed to copy .dat file: %v", err)
|
||||
}
|
||||
localFiles["dat"] = datFile
|
||||
@@ -289,7 +296,7 @@ func (t *ErasureCodingTask) copyVolumeFilesToWorker(workDir string) (map[string]
|
||||
|
||||
// Copy .idx file
|
||||
idxFile := filepath.Join(workDir, fmt.Sprintf("%d.idx", t.volumeID))
|
||||
if err := t.copyFileFromSource(".idx", idxFile); err != nil {
|
||||
if err := t.copyFileFromSource(ctx, ".idx", idxFile, fileStatus.GetCompactionRevision(), fileStatus.GetIdxFileSize()); err != nil {
|
||||
return nil, fmt.Errorf("failed to copy .idx file: %v", err)
|
||||
}
|
||||
localFiles["idx"] = idxFile
|
||||
@@ -307,15 +314,38 @@ func (t *ErasureCodingTask) copyVolumeFilesToWorker(workDir string) (map[string]
|
||||
return localFiles, nil
|
||||
}
|
||||
|
||||
func (t *ErasureCodingTask) readSourceVolumeFileStatus(ctx context.Context) (*volume_server_pb.ReadVolumeFileStatusResponse, error) {
|
||||
var statusResp *volume_server_pb.ReadVolumeFileStatusResponse
|
||||
err := operation.WithVolumeServerClient(false, pb.ServerAddress(t.server), t.grpcDialOption,
|
||||
func(client volume_server_pb.VolumeServerClient) error {
|
||||
var readErr error
|
||||
statusResp, readErr = client.ReadVolumeFileStatus(ctx, &volume_server_pb.ReadVolumeFileStatusRequest{
|
||||
VolumeId: t.volumeID,
|
||||
})
|
||||
return readErr
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if statusResp.GetDatFileSize() == 0 {
|
||||
return nil, fmt.Errorf("volume %d on %s reports zero dat file size", t.volumeID, t.server)
|
||||
}
|
||||
if statusResp.GetIdxFileSize() == 0 {
|
||||
return nil, fmt.Errorf("volume %d on %s reports zero idx file size with non-empty dat", t.volumeID, t.server)
|
||||
}
|
||||
return statusResp, nil
|
||||
}
|
||||
|
||||
// copyFileFromSource copies a file from source server to local path using gRPC streaming
|
||||
func (t *ErasureCodingTask) copyFileFromSource(ext, localPath string) error {
|
||||
func (t *ErasureCodingTask) copyFileFromSource(ctx context.Context, ext, localPath string, compactionRevision uint32, stopOffset uint64) error {
|
||||
return operation.WithVolumeServerClient(false, pb.ServerAddress(t.server), t.grpcDialOption,
|
||||
func(client volume_server_pb.VolumeServerClient) error {
|
||||
stream, err := client.CopyFile(context.Background(), &volume_server_pb.CopyFileRequest{
|
||||
VolumeId: t.volumeID,
|
||||
Collection: t.collection,
|
||||
Ext: ext,
|
||||
StopOffset: uint64(math.MaxInt64),
|
||||
stream, err := client.CopyFile(ctx, &volume_server_pb.CopyFileRequest{
|
||||
VolumeId: t.volumeID,
|
||||
Collection: t.collection,
|
||||
Ext: ext,
|
||||
CompactionRevision: compactionRevision,
|
||||
StopOffset: stopOffset,
|
||||
})
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to initiate file copy: %v", err)
|
||||
@@ -348,6 +378,9 @@ func (t *ErasureCodingTask) copyFileFromSource(ext, localPath string) error {
|
||||
}
|
||||
}
|
||||
|
||||
if totalBytes != int64(stopOffset) {
|
||||
return fmt.Errorf("short copy of %s: got %d bytes, expected %d", ext, totalBytes, stopOffset)
|
||||
}
|
||||
glog.V(1).Infof("Successfully copied %s (%d bytes) from %s to %s", ext, totalBytes, t.server, localPath)
|
||||
return nil
|
||||
})
|
||||
@@ -468,7 +501,7 @@ func (t *ErasureCodingTask) mountEcShards() error {
|
||||
}
|
||||
|
||||
// deleteOriginalVolume deletes the original volume and all its replicas from all servers
|
||||
func (t *ErasureCodingTask) deleteOriginalVolume() error {
|
||||
func (t *ErasureCodingTask) deleteOriginalVolume(ctx context.Context) error {
|
||||
// Get replicas from task parameters (set during detection)
|
||||
replicas := t.getReplicas()
|
||||
|
||||
@@ -497,7 +530,7 @@ func (t *ErasureCodingTask) deleteOriginalVolume() error {
|
||||
|
||||
err := operation.WithVolumeServerClient(false, pb.ServerAddress(replicaServer), t.grpcDialOption,
|
||||
func(client volume_server_pb.VolumeServerClient) error {
|
||||
_, err := client.VolumeDelete(context.Background(), &volume_server_pb.VolumeDeleteRequest{
|
||||
_, err := client.VolumeDelete(ctx, &volume_server_pb.VolumeDeleteRequest{
|
||||
VolumeId: t.volumeID,
|
||||
OnlyEmpty: false, // Force delete since we've created EC shards
|
||||
})
|
||||
|
||||
@@ -0,0 +1,108 @@
|
||||
package erasure_coding
|
||||
|
||||
import (
|
||||
"context"
|
||||
"io"
|
||||
"net/http"
|
||||
"os"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/test/volume_server/framework"
|
||||
"github.com/seaweedfs/seaweedfs/test/volume_server/matrix"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/volume_server_pb"
|
||||
"github.com/stretchr/testify/require"
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/credentials/insecure"
|
||||
)
|
||||
|
||||
func TestCopyVolumeFilesToWorkerUsesCurrentCompactionRevision(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("skipping integration test in short mode")
|
||||
}
|
||||
|
||||
clusterHarness := framework.StartVolumeCluster(t, matrix.P1())
|
||||
conn, grpcClient := framework.DialVolumeServer(t, clusterHarness.VolumeGRPCAddress())
|
||||
defer conn.Close()
|
||||
|
||||
const volumeID = uint32(951)
|
||||
framework.AllocateVolume(t, grpcClient, volumeID, "")
|
||||
|
||||
httpClient := framework.NewHTTPClient()
|
||||
|
||||
liveFID := framework.NewFileID(volumeID, 1001, 0x1111AAAA)
|
||||
liveUploadResp := framework.UploadBytes(t, httpClient, clusterHarness.VolumeAdminURL(), liveFID, []byte("live-payload-for-ec-copy"))
|
||||
_ = framework.ReadAllAndClose(t, liveUploadResp)
|
||||
require.Equal(t, http.StatusCreated, liveUploadResp.StatusCode)
|
||||
|
||||
deletedFID := framework.NewFileID(volumeID, 1002, 0x2222BBBB)
|
||||
deletedUploadResp := framework.UploadBytes(t, httpClient, clusterHarness.VolumeAdminURL(), deletedFID, []byte("deleted-payload-for-vacuum"))
|
||||
_ = framework.ReadAllAndClose(t, deletedUploadResp)
|
||||
require.Equal(t, http.StatusCreated, deletedUploadResp.StatusCode)
|
||||
|
||||
deleteReq, err := http.NewRequest(http.MethodDelete, clusterHarness.VolumeAdminURL()+"/"+deletedFID, nil)
|
||||
require.NoError(t, err)
|
||||
deleteResp := framework.DoRequest(t, httpClient, deleteReq)
|
||||
_ = framework.ReadAllAndClose(t, deleteResp)
|
||||
require.Equal(t, http.StatusAccepted, deleteResp.StatusCode)
|
||||
|
||||
compactVolumeOnce(t, grpcClient, volumeID)
|
||||
|
||||
task := NewErasureCodingTask(
|
||||
"copy-after-compaction",
|
||||
clusterHarness.VolumeServerAddress(),
|
||||
volumeID,
|
||||
"",
|
||||
grpc.WithTransportCredentials(insecure.NewCredentials()),
|
||||
)
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
||||
defer cancel()
|
||||
|
||||
require.NoError(t, task.markVolumeReadonly(ctx))
|
||||
|
||||
fileStatus, err := task.readSourceVolumeFileStatus(ctx)
|
||||
require.NoError(t, err)
|
||||
require.Greater(t, fileStatus.GetCompactionRevision(), uint32(0))
|
||||
|
||||
localFiles, err := task.copyVolumeFilesToWorker(ctx, t.TempDir())
|
||||
require.NoError(t, err)
|
||||
|
||||
datInfo, err := os.Stat(localFiles["dat"])
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, int64(fileStatus.GetDatFileSize()), datInfo.Size())
|
||||
|
||||
idxInfo, err := os.Stat(localFiles["idx"])
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, int64(fileStatus.GetIdxFileSize()), idxInfo.Size())
|
||||
}
|
||||
|
||||
func compactVolumeOnce(t *testing.T, grpcClient volume_server_pb.VolumeServerClient, volumeID uint32) {
|
||||
t.Helper()
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
|
||||
defer cancel()
|
||||
|
||||
compactStream, err := grpcClient.VacuumVolumeCompact(ctx, &volume_server_pb.VacuumVolumeCompactRequest{
|
||||
VolumeId: volumeID,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
|
||||
for {
|
||||
_, err = compactStream.Recv()
|
||||
if err == io.EOF {
|
||||
break
|
||||
}
|
||||
require.NoError(t, err)
|
||||
}
|
||||
|
||||
_, err = grpcClient.VacuumVolumeCommit(ctx, &volume_server_pb.VacuumVolumeCommitRequest{
|
||||
VolumeId: volumeID,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
|
||||
_, err = grpcClient.VacuumVolumeCleanup(ctx, &volume_server_pb.VacuumVolumeCleanupRequest{
|
||||
VolumeId: volumeID,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
}
|
||||
Reference in New Issue
Block a user