Compare commits

...

7 Commits

Author SHA1 Message Date
Feng Ruohang 9101fe78db ci: refresh README repository cards daily at 00:00 UTC
Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 15:27:48 +08:00
Feng Ruohang 82948d6306 Merge pull request #198 from mrjavadseydi/fix/issue-79-list-multipart-uploads
Retain the original contributor commit and integrate the reviewed durable listing, cancellation confirmation and migration fixes. Capacity acceptance and delayed creation-write fencing remain open in issue #79.

Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 13:40:10 +08:00
Feng Ruohang dfb4b2a1d2 fix: return HTTP 503 when multipart scan capacity is exhausted
Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 13:24:52 +08:00
Feng Ruohang 143f6970d8 fix: make multipart discovery and cancellation explicit and bounded
Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 13:16:37 +08:00
Feng Ruohang ea5dac5a99 Merge current main into multipart listing contribution
Signed-off-by: Feng Ruohang <rh@vonng.com>
2026-09-16 12:51:20 +08:00
Feng Ruohang 3c26a8b0b5 Merge pull request #209 from pgsty/codex/credit-console-share-reporter
fix: integrate the bounded Console sharing proxy and credit its reporter
2026-09-16 12:19:29 +08:00
mr javad seydi 4cbb074ccd fix: make multipart upload listing S3-compatible
Signed-off-by: mr javad seydi <seydi.birjand@gmail.com>
2026-09-15 23:17:55 +03:30
29 changed files with 2723 additions and 164 deletions
+3
View File
@@ -93,6 +93,9 @@ jobs:
- name: Run conditional PUT tests under race detector
run: go test -race ./cmd -run '^Test(PoolsConditionalPut|SinglePoolConditionalPutHTTP)' -count=1 -timeout=5m
- name: Run multipart listing and cancellation tests under race detector
run: go test -race ./cmd -run '^Test(MultipartListing|MultipartAbort|PaginateMultipartUploads|ListMultipartUploads)' -count=1 -timeout=5m
crosscompile:
name: Cross Compile
runs-on: ubuntu-latest
+100
View File
@@ -0,0 +1,100 @@
name: Repository Cards
on:
schedule:
# 00:00 UTC = 08:00 Asia/Shanghai. GitHub may queue scheduled runs.
- cron: "0 0 * * *"
workflow_dispatch:
push:
branches: [main]
paths:
- .github/workflows/repository-cards.yml
- .github/silo.svg
- buildscripts/repository-cards/**
pull_request:
branches: [main]
paths:
- .github/workflows/repository-cards.yml
- .github/silo.svg
- buildscripts/repository-cards/**
permissions:
contents: read
concurrency:
group: repository-cards-${{ github.event.pull_request.number || 'publish' }}
cancel-in-progress: false
jobs:
check:
name: Validate repository cards
runs-on: ubuntu-latest
timeout-minutes: 5
steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7
with:
persist-credentials: false
- uses: actions/setup-python@ece7cb06caefa5fff74198d8649806c4678c61a1 # v6
with:
python-version: "3.13"
- run: python -m pip install -r buildscripts/repository-cards/requirements.txt
- run: python -m unittest discover -s buildscripts/repository-cards -p 'test_*.py' -v
publish:
name: Update README images
needs: check
if: github.repository == 'pgsty/silo' && github.ref == 'refs/heads/main' && github.event_name != 'pull_request'
runs-on: ubuntu-latest
timeout-minutes: 15
permissions:
contents: write
issues: read
pull-requests: read
env:
OUTPUT_BRANCH: codex/repository-cards
steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7
- uses: actions/setup-python@ece7cb06caefa5fff74198d8649806c4678c61a1 # v6
with:
python-version: "3.13"
- run: python -m pip install -r buildscripts/repository-cards/requirements.txt
- name: Load the generated-assets branch
shell: bash
run: |
set -euo pipefail
artifacts="$RUNNER_TEMP/repository-cards"
if git ls-remote --exit-code --heads origin "$OUTPUT_BRANCH"; then
git fetch --depth=1 origin "$OUTPUT_BRANCH"
git worktree add --detach "$artifacts" FETCH_HEAD
else
status=$?
# Exit 2 means no matching ref; transport/auth failures must stop.
if [ "$status" -ne 2 ]; then exit "$status"; fi
git worktree add --detach "$artifacts" HEAD
git -C "$artifacts" checkout --orphan "$OUTPUT_BRANCH"
git -C "$artifacts" rm -rf .
fi
- name: Refresh contributor and star cards
env:
GH_TOKEN: ${{ github.token }}
run: python buildscripts/repository-cards/update.py --output "$RUNNER_TEMP/repository-cards"
- name: Publish changed assets
shell: bash
run: |
set -euo pipefail
cd "$RUNNER_TEMP/repository-cards"
git add README.md history.json curated.json contributors.json \
contributors-light.svg contributors-dark.svg \
star-history-light.svg star-history-dark.svg
if git diff --cached --quiet; then
echo "Repository cards are already current."
exit 0
fi
git config user.name 'github-actions[bot]'
git config user.email '41898282+github-actions[bot]@users.noreply.github.com'
git commit -s -m "chore: update repository cards $(date -u +%F)"
# A normal push preserves history and refuses concurrent overwrites.
git push origin "HEAD:refs/heads/$OUTPUT_BRANCH"
+20
View File
@@ -49,6 +49,26 @@ and [complete commit range](https://github.com/pgsty/silo/compare/RELEASE.2026-0
### Object storage and replication
- Make `ListMultipartUploads` discover quorum-valid uploads from durable state
across pools, erasure sets and drives, then apply S3 prefix, delimiter,
marker, ordering and 1,000-entry pagination semantics globally (#198). New uploads
store their canonical bucket and key as reserved fields in the existing
quorum-written `xl.meta`; completion removes those upload-only fields. Native
markers remain usable after their upload is completed or canceled. Strict
listing returns a diagnostic 503 for legacy uploads or uncertain coverage;
`api multipart_listing=legacy` is an explicit temporary migration mode.
Upgrade every writer, drain old uploads and check the read-only admin
`multipart-preflight` report before relying on strict listing. Per-process
admission, directory-entry, worker and time budgets bound scan scheduling;
each page still scans durable state. See [issue #79](https://github.com/pgsty/silo/issues/79)
and its [design record](https://silo.pgsty.com/blog/design/list-multipart-uploads/).
Thanks to mr javad seydi (@mrjavadseydi) for the original implementation.
- Confirm multipart cancellation on a strict majority of each relevant set,
and allow retries after partial deletion. Uncertain pools or insufficient
confirmations return 503 rather than acknowledging a cancellation whose
static remnants can later become readable. **Known boundary:** creation
writes that finish after a storage timeout can still restore an upload after
successful cancellation; this change does not add a durable creation fence.
- Preserve object tags during multi-pool metadata reconciliation by reading the
resolved tag field together with its revision (#189). Previously, reconciliation
could replace existing tags with an empty value.
+17 -46
View File
@@ -123,54 +123,25 @@ Report vulnerabilities privately as described in [`SECURITY.md`](SECURITY.md); e
## Contributors
**42 community contributors** build SILO, Console, mcli, shared packages, and related projects. The list includes maintainers and every human Issue or PR author, ordered by merged PRs, other PRs, then issue reports. Gold rings highlight significant contributions.
<!-- Generated by silo.pgsty.com/bin/contributors.py. -->
<p align="center">
<a href="https://github.com/Vonng"><img src="https://silo.pgsty.com/images/contributors/Vonng.svg" width="60" height="60" alt="@Vonng" title="@Vonng — Maintains SILO, Console, mcli, shared packages, releases, and documentation"></a>
<a href="https://github.com/h5vx"><img src="https://silo.pgsty.com/images/contributors/h5vx.svg" width="60" height="60" alt="@h5vx" title="@h5vx — Implemented per-bucket CORS configuration and enforcement"></a>
<a href="https://github.com/mrjavadseydi"><img src="https://silo.pgsty.com/images/contributors/mrjavadseydi.svg" width="60" height="60" alt="@mrjavadseydi" title="@mrjavadseydi — Fixed effective bucket quota metrics; proposed access-frequency ILM"></a>
<a href="https://github.com/Dansyuqri"><img src="https://silo.pgsty.com/images/contributors/Dansyuqri.svg" width="60" height="60" alt="@Dansyuqri" title="@Dansyuqri — Added ChecksumType to multipart completion responses"></a>
<a href="https://github.com/ycjlin"><img src="https://silo.pgsty.com/images/contributors/ycjlin.svg" width="60" height="60" alt="@ycjlin" title="@ycjlin — Fixed missing-bucket ListObjects semantics"></a>
<a href="https://github.com/pinginfo"><img src="https://silo.pgsty.com/images/contributors/pinginfo.svg" width="60" height="60" alt="@pinginfo" title="@pinginfo — Repaired bucket notification streaming"></a>
<a href="https://github.com/ZouhairCharef"><img src="https://silo.pgsty.com/images/contributors/ZouhairCharef.svg" width="60" height="60" alt="@ZouhairCharef" title="@ZouhairCharef — Patched CVE-2026-34986 in go-jose"></a>
<a href="https://github.com/mfredenhagen"><img src="https://silo.pgsty.com/images/contributors/mfredenhagen.svg" width="60" height="60" alt="@mfredenhagen" title="@mfredenhagen — Patched CVE-2026-39883 in OpenTelemetry"></a>
<a href="https://github.com/waterkip"><img src="https://silo.pgsty.com/images/contributors/waterkip.svg" width="60" height="60" alt="@waterkip" title="@waterkip — Repointed documentation links to the SILO portal"></a>
<a href="https://github.com/mikemikimike"><img src="https://silo.pgsty.com/images/contributors/mikemikimike.svg" width="60" height="60" alt="@mikemikimike" title="@mikemikimike — Contributed the replicated SSE-C plaintext part-size fix"></a>
<a href="https://github.com/metaneutrons"><img src="https://silo.pgsty.com/images/contributors/metaneutrons.svg" width="60" height="60" alt="@metaneutrons" title="@metaneutrons — Reported and proposed explicit-version delete authorization"></a>
<a href="https://github.com/magicxor"><img src="https://silo.pgsty.com/images/contributors/magicxor.svg" width="60" height="60" alt="@magicxor" title="@magicxor — Reported and proposed conditional DELETE support for If-Match"></a>
<a href="https://github.com/davinkevin"><img src="https://silo.pgsty.com/images/contributors/davinkevin.svg" width="60" height="60" alt="@davinkevin" title="@davinkevin — Proposed the distroless container image and dependency automation"></a>
<a href="https://github.com/lem21h"><img src="https://silo.pgsty.com/images/contributors/lem21h.svg" width="48" height="48" alt="@lem21h" title="@lem21h — Proposed robustness and goroutine improvements"></a>
<a href="https://github.com/sulin37392"><img src="https://silo.pgsty.com/images/contributors/sulin37392.svg" width="48" height="48" alt="@sulin37392" title="@sulin37392 — Proposed dependency updates"></a>
<a href="https://github.com/cbornet"><img src="https://silo.pgsty.com/images/contributors/cbornet.svg" width="60" height="60" alt="@cbornet" title="@cbornet — Reported multipart and streaming checksum defects and missing-bucket semantics"></a>
<a href="https://github.com/vampywiz17"><img src="https://silo.pgsty.com/images/contributors/vampywiz17.svg" width="60" height="60" alt="@vampywiz17" title="@vampywiz17 — Reported LDAP TLS and Console login regressions"></a>
<a href="https://github.com/orenyomtov"><img src="https://silo.pgsty.com/images/contributors/orenyomtov.svg" width="60" height="60" alt="@orenyomtov" title="@orenyomtov — Reported the unsigned-header CopyObject cross-object read (SN-2026-011)"></a>
<a href="https://github.com/jiri-pejchal"><img src="https://silo.pgsty.com/images/contributors/jiri-pejchal.svg" width="60" height="60" alt="@jiri-pejchal" title="@jiri-pejchal — Reported the Console share proxy exposure of internal metrics"></a>
<a href="https://github.com/mumu-lab"><img src="https://silo.pgsty.com/images/contributors/mumu-lab.svg" width="48" height="48" alt="@mumu-lab" title="@mumu-lab — Reported bucket quota metrics reading a deprecated field"></a>
<a href="https://github.com/jvasile"><img src="https://silo.pgsty.com/images/contributors/jvasile.svg" width="48" height="48" alt="@jvasile" title="@jvasile — Reported missing user, group, and defaults in Debian packages"></a>
<a href="https://github.com/pmezhuev"><img src="https://silo.pgsty.com/images/contributors/pmezhuev.svg" width="48" height="48" alt="@pmezhuev" title="@pmezhuev — Reported missing RPM package signatures"></a>
<a href="https://github.com/TLINDEN"><img src="https://silo.pgsty.com/images/contributors/TLINDEN.svg" width="48" height="48" alt="@TLINDEN" title="@TLINDEN — Reported the missing client in release tarballs"></a>
<a href="https://github.com/makinikm"><img src="https://silo.pgsty.com/images/contributors/makinikm.svg" width="48" height="48" alt="@makinikm" title="@makinikm — Reported the missing client in the container image"></a>
<a href="https://github.com/meesudzu"><img src="https://silo.pgsty.com/images/contributors/meesudzu.svg" width="48" height="48" alt="@meesudzu" title="@meesudzu — Requested the migration guide from upstream MinIO"></a>
<a href="https://github.com/kuldeep-link11"><img src="https://silo.pgsty.com/images/contributors/kuldeep-link11.svg" width="48" height="48" alt="@kuldeep-link11" title="@kuldeep-link11 — Reported NATS JWT credentials and target reload issues"></a>
<a href="https://github.com/sargarass"><img src="https://silo.pgsty.com/images/contributors/sargarass.svg" width="48" height="48" alt="@sargarass" title="@sargarass — Reported ListMultipartUploads prefix and pagination semantics"></a>
<a href="https://github.com/liuhaodongliu990-cmyk"><img src="https://silo.pgsty.com/images/contributors/liuhaodongliu990-cmyk.svg" width="48" height="48" alt="@liuhaodongliu990-cmyk" title="@liuhaodongliu990-cmyk — Reported indeterminate progress for prefix downloads"></a>
<a href="https://github.com/Xavier-777"><img src="https://silo.pgsty.com/images/contributors/Xavier-777.svg" width="48" height="48" alt="@Xavier-777" title="@Xavier-777 — Reported Console lifecycle management and file preview gaps"></a>
<a href="https://github.com/spaceg00se-r"><img src="https://silo.pgsty.com/images/contributors/spaceg00se-r.svg" width="48" height="48" alt="@spaceg00se-r" title="@spaceg00se-r — Requested cpuv1 support and reported a workflow token failure"></a>
<a href="https://github.com/kh0mka"><img src="https://silo.pgsty.com/images/contributors/kh0mka.svg" width="48" height="48" alt="@kh0mka" title="@kh0mka — Reported inter-node I/O timeouts in ReadFileStreamHandler"></a>
<a href="https://github.com/bagutzu"><img src="https://silo.pgsty.com/images/contributors/bagutzu.svg" width="48" height="48" alt="@bagutzu" title="@bagutzu — Requested KES-compatible external KMS and OpenBao support"></a>
<a href="https://github.com/DestroyLee"><img src="https://silo.pgsty.com/images/contributors/DestroyLee.svg" width="48" height="48" alt="@DestroyLee" title="@DestroyLee — Reported the missing documentation navigation"></a>
<a href="https://github.com/mosesdd"><img src="https://silo.pgsty.com/images/contributors/mosesdd.svg" width="48" height="48" alt="@mosesdd" title="@mosesdd — Requested a maintained Helm chart"></a>
<a href="https://github.com/zylpsrs"><img src="https://silo.pgsty.com/images/contributors/zylpsrs.svg" width="48" height="48" alt="@zylpsrs" title="@zylpsrs — Reported missing Console tiering and site replication"></a>
<a href="https://github.com/heroes1412"><img src="https://silo.pgsty.com/images/contributors/heroes1412.svg" width="48" height="48" alt="@heroes1412" title="@heroes1412 — Reported the unusable profiling option"></a>
<a href="https://github.com/redfoxfox"><img src="https://silo.pgsty.com/images/contributors/redfoxfox.svg" width="48" height="48" alt="@redfoxfox" title="@redfoxfox — Reported Chinese documentation availability"></a>
<a href="https://github.com/jiadzh"><img src="https://silo.pgsty.com/images/contributors/jiadzh.svg" width="48" height="48" alt="@jiadzh" title="@jiadzh — Requested Windows build guidance"></a>
<a href="https://github.com/AntonOfTheWoods"><img src="https://silo.pgsty.com/images/contributors/AntonOfTheWoods.svg" width="48" height="48" alt="@AntonOfTheWoods" title="@AntonOfTheWoods — Asked for clarity on Helm chart and operator options"></a>
<a href="https://github.com/chalukyaj"><img src="https://silo.pgsty.com/images/contributors/chalukyaj.svg" width="48" height="48" alt="@chalukyaj" title="@chalukyaj — Proposed making the SILO Operator easier to discover"></a>
<a href="https://github.com/nsanitate"><img src="https://silo.pgsty.com/images/contributors/nsanitate.svg" width="48" height="48" alt="@nsanitate" title="@nsanitate — Proposed CNCF Sandbox governance"></a>
<a href="https://github.com/Kesavaambati"><img src="https://silo.pgsty.com/images/contributors/Kesavaambati.svg" width="48" height="48" alt="@Kesavaambati" title="@Kesavaambati — Asked about community support and image maintenance"></a>
</p>
Every human issue or pull-request author is part of the SILO community, including open and unmerged work. Merged fixes, adopted proposals, and actionable reports receive priority, with first participation guiding the remaining order. Gold rings highlight reviewed significant contributions.
[View the full contribution record](CONTRIBUTORS.md) for each person's proposals, fixes, and reports.
<a href="CONTRIBUTORS.md">
<picture>
<source media="(prefers-color-scheme: dark)" srcset="https://raw.githubusercontent.com/pgsty/silo/codex/repository-cards/contributors-dark.svg">
<img src="https://raw.githubusercontent.com/pgsty/silo/codex/repository-cards/contributors-light.svg" alt="SILO community contributors">
</picture>
</a>
[View contribution notes and actual PR status](CONTRIBUTORS.md).
## Star History
<picture>
<source media="(prefers-color-scheme: dark)" srcset="https://raw.githubusercontent.com/pgsty/silo/codex/repository-cards/star-history-dark.svg">
<img src="https://raw.githubusercontent.com/pgsty/silo/codex/repository-cards/star-history-light.svg" alt="SILO GitHub star history">
</picture>
## Background
+17 -46
View File
@@ -102,54 +102,25 @@ S3 API、`MINIO_*` 环境变量、`minio_*` 指标、`x-minio-*` 头、`/minio/*
## 贡献者
**42 位社区贡献者**共同建设 SILO、Console、mcli、公共包与相关项目。名单包含维护者,以及所有提出 Issue 或 PR 的真人作者;按已合并 PR、其他 PR、Issue 报告排序,黄圈标记显著贡献。
<!-- Generated by silo.pgsty.com/bin/contributors.py. -->
<p align="center">
<a href="https://github.com/Vonng"><img src="https://silo.pgsty.com/images/contributors/Vonng.svg" width="60" height="60" alt="@Vonng" title="@Vonng — 维护 SILO、Console、mcli、公共包、发行与文档"></a>
<a href="https://github.com/h5vx"><img src="https://silo.pgsty.com/images/contributors/h5vx.svg" width="60" height="60" alt="@h5vx" title="@h5vx — 实现单桶 CORS 配置与请求执行"></a>
<a href="https://github.com/mrjavadseydi"><img src="https://silo.pgsty.com/images/contributors/mrjavadseydi.svg" width="60" height="60" alt="@mrjavadseydi" title="@mrjavadseydi — 修复有效桶配额指标,并提交按访问频率分层的 ILM 方案"></a>
<a href="https://github.com/Dansyuqri"><img src="https://silo.pgsty.com/images/contributors/Dansyuqri.svg" width="60" height="60" alt="@Dansyuqri" title="@Dansyuqri — 为分片上传完成响应补充 ChecksumType"></a>
<a href="https://github.com/ycjlin"><img src="https://silo.pgsty.com/images/contributors/ycjlin.svg" width="60" height="60" alt="@ycjlin" title="@ycjlin — 修复缺失桶的 ListObjects 语义"></a>
<a href="https://github.com/pinginfo"><img src="https://silo.pgsty.com/images/contributors/pinginfo.svg" width="60" height="60" alt="@pinginfo" title="@pinginfo — 修复桶通知的流式输出"></a>
<a href="https://github.com/ZouhairCharef"><img src="https://silo.pgsty.com/images/contributors/ZouhairCharef.svg" width="60" height="60" alt="@ZouhairCharef" title="@ZouhairCharef — 修复 go-jose 中的 CVE-2026-34986"></a>
<a href="https://github.com/mfredenhagen"><img src="https://silo.pgsty.com/images/contributors/mfredenhagen.svg" width="60" height="60" alt="@mfredenhagen" title="@mfredenhagen — 修复 OpenTelemetry 中的 CVE-2026-39883"></a>
<a href="https://github.com/waterkip"><img src="https://silo.pgsty.com/images/contributors/waterkip.svg" width="60" height="60" alt="@waterkip" title="@waterkip — 将文档链接指向 SILO 门户"></a>
<a href="https://github.com/mikemikimike"><img src="https://silo.pgsty.com/images/contributors/mikemikimike.svg" width="60" height="60" alt="@mikemikimike" title="@mikemikimike — 提交 SSE-C 复制分片明文尺寸修复"></a>
<a href="https://github.com/metaneutrons"><img src="https://silo.pgsty.com/images/contributors/metaneutrons.svg" width="60" height="60" alt="@metaneutrons" title="@metaneutrons — 报告并提交显式版本删除鉴权方案"></a>
<a href="https://github.com/magicxor"><img src="https://silo.pgsty.com/images/contributors/magicxor.svg" width="60" height="60" alt="@magicxor" title="@magicxor — 报告并提交 DELETE If-Match 条件请求支持方案"></a>
<a href="https://github.com/davinkevin"><img src="https://silo.pgsty.com/images/contributors/davinkevin.svg" width="60" height="60" alt="@davinkevin" title="@davinkevin — 提交 distroless 容器镜像与依赖自动更新方案"></a>
<a href="https://github.com/lem21h"><img src="https://silo.pgsty.com/images/contributors/lem21h.svg" width="48" height="48" alt="@lem21h" title="@lem21h — 提交健壮性与 goroutine 改进"></a>
<a href="https://github.com/sulin37392"><img src="https://silo.pgsty.com/images/contributors/sulin37392.svg" width="48" height="48" alt="@sulin37392" title="@sulin37392 — 提交依赖更新"></a>
<a href="https://github.com/cbornet"><img src="https://silo.pgsty.com/images/contributors/cbornet.svg" width="60" height="60" alt="@cbornet" title="@cbornet — 报告分片与流式校验和缺陷及缺失桶语义问题"></a>
<a href="https://github.com/vampywiz17"><img src="https://silo.pgsty.com/images/contributors/vampywiz17.svg" width="60" height="60" alt="@vampywiz17" title="@vampywiz17 — 报告 LDAP TLS 与 Console 登录回归"></a>
<a href="https://github.com/orenyomtov"><img src="https://silo.pgsty.com/images/contributors/orenyomtov.svg" width="60" height="60" alt="@orenyomtov" title="@orenyomtov — 报告未签名头导致的 CopyObject 跨对象读取(SN-2026-011)"></a>
<a href="https://github.com/jiri-pejchal"><img src="https://silo.pgsty.com/images/contributors/jiri-pejchal.svg" width="60" height="60" alt="@jiri-pejchal" title="@jiri-pejchal — 报告 Console 分享代理暴露内部指标的问题"></a>
<a href="https://github.com/mumu-lab"><img src="https://silo.pgsty.com/images/contributors/mumu-lab.svg" width="48" height="48" alt="@mumu-lab" title="@mumu-lab — 报告桶配额指标读取已弃用字段的问题"></a>
<a href="https://github.com/jvasile"><img src="https://silo.pgsty.com/images/contributors/jvasile.svg" width="48" height="48" alt="@jvasile" title="@jvasile — 报告 Debian 包缺少用户、用户组与默认配置"></a>
<a href="https://github.com/pmezhuev"><img src="https://silo.pgsty.com/images/contributors/pmezhuev.svg" width="48" height="48" alt="@pmezhuev" title="@pmezhuev — 报告 RPM 包缺少 GPG 签名"></a>
<a href="https://github.com/TLINDEN"><img src="https://silo.pgsty.com/images/contributors/TLINDEN.svg" width="48" height="48" alt="@TLINDEN" title="@TLINDEN — 报告发布压缩包缺少客户端"></a>
<a href="https://github.com/makinikm"><img src="https://silo.pgsty.com/images/contributors/makinikm.svg" width="48" height="48" alt="@makinikm" title="@makinikm — 报告容器镜像缺少客户端"></a>
<a href="https://github.com/meesudzu"><img src="https://silo.pgsty.com/images/contributors/meesudzu.svg" width="48" height="48" alt="@meesudzu" title="@meesudzu — 提出从上游 MinIO 迁移的指南需求"></a>
<a href="https://github.com/kuldeep-link11"><img src="https://silo.pgsty.com/images/contributors/kuldeep-link11.svg" width="48" height="48" alt="@kuldeep-link11" title="@kuldeep-link11 — 报告 NATS JWT 凭据与通知目标重载问题"></a>
<a href="https://github.com/sargarass"><img src="https://silo.pgsty.com/images/contributors/sargarass.svg" width="48" height="48" alt="@sargarass" title="@sargarass — 报告 ListMultipartUploads 前缀与分页语义问题"></a>
<a href="https://github.com/liuhaodongliu990-cmyk"><img src="https://silo.pgsty.com/images/contributors/liuhaodongliu990-cmyk.svg" width="48" height="48" alt="@liuhaodongliu990-cmyk" title="@liuhaodongliu990-cmyk — 报告前缀下载进度显示异常"></a>
<a href="https://github.com/Xavier-777"><img src="https://silo.pgsty.com/images/contributors/Xavier-777.svg" width="48" height="48" alt="@Xavier-777" title="@Xavier-777 — 报告 Console 生命周期管理与文件预览缺失"></a>
<a href="https://github.com/spaceg00se-r"><img src="https://silo.pgsty.com/images/contributors/spaceg00se-r.svg" width="48" height="48" alt="@spaceg00se-r" title="@spaceg00se-r — 提出 cpuv1 支持需求并报告工作流令牌错误"></a>
<a href="https://github.com/kh0mka"><img src="https://silo.pgsty.com/images/contributors/kh0mka.svg" width="48" height="48" alt="@kh0mka" title="@kh0mka — 报告 ReadFileStreamHandler 节点间 I/O 超时"></a>
<a href="https://github.com/bagutzu"><img src="https://silo.pgsty.com/images/contributors/bagutzu.svg" width="48" height="48" alt="@bagutzu" title="@bagutzu — 提出兼容 KES 的外部 KMS 与 OpenBao 支持需求"></a>
<a href="https://github.com/DestroyLee"><img src="https://silo.pgsty.com/images/contributors/DestroyLee.svg" width="48" height="48" alt="@DestroyLee" title="@DestroyLee — 报告文档目录导航缺失"></a>
<a href="https://github.com/mosesdd"><img src="https://silo.pgsty.com/images/contributors/mosesdd.svg" width="48" height="48" alt="@mosesdd" title="@mosesdd — 提出维护 Helm Chart 的需求"></a>
<a href="https://github.com/zylpsrs"><img src="https://silo.pgsty.com/images/contributors/zylpsrs.svg" width="48" height="48" alt="@zylpsrs" title="@zylpsrs — 报告 Console 缺少分层与站点复制"></a>
<a href="https://github.com/heroes1412"><img src="https://silo.pgsty.com/images/contributors/heroes1412.svg" width="48" height="48" alt="@heroes1412" title="@heroes1412 — 报告性能分析选项不可用"></a>
<a href="https://github.com/redfoxfox"><img src="https://silo.pgsty.com/images/contributors/redfoxfox.svg" width="48" height="48" alt="@redfoxfox" title="@redfoxfox — 报告中文文档站点不可用"></a>
<a href="https://github.com/jiadzh"><img src="https://silo.pgsty.com/images/contributors/jiadzh.svg" width="48" height="48" alt="@jiadzh" title="@jiadzh — 提出 Windows 构建指导需求"></a>
<a href="https://github.com/AntonOfTheWoods"><img src="https://silo.pgsty.com/images/contributors/AntonOfTheWoods.svg" width="48" height="48" alt="@AntonOfTheWoods" title="@AntonOfTheWoods — 提出明确 Helm Chart 与 Operator 选项的需求"></a>
<a href="https://github.com/chalukyaj"><img src="https://silo.pgsty.com/images/contributors/chalukyaj.svg" width="48" height="48" alt="@chalukyaj" title="@chalukyaj — 提出改善 SILO Operator 可发现性的建议"></a>
<a href="https://github.com/nsanitate"><img src="https://silo.pgsty.com/images/contributors/nsanitate.svg" width="48" height="48" alt="@nsanitate" title="@nsanitate — 提出加入 CNCF Sandbox 的治理建议"></a>
<a href="https://github.com/Kesavaambati"><img src="https://silo.pgsty.com/images/contributors/Kesavaambati.svg" width="48" height="48" alt="@Kesavaambati" title="@Kesavaambati — 提出社区支持与容器镜像维护问题"></a>
</p>
每位 issue 或 PR 的作者都是 SILO 社区的一员,包括尚未合并的工作。已合并的修复、被采纳的方案和有效报告优先展示,其余参考首次参与时间;金色圆环突出经过审核的显著贡献。
[查看完整贡献记录](CONTRIBUTORS.md),了解每位贡献者的提案、修复与问题报告。
<a href="CONTRIBUTORS.md">
<picture>
<source media="(prefers-color-scheme: dark)" srcset="https://raw.githubusercontent.com/pgsty/silo/codex/repository-cards/contributors-dark.svg">
<img src="https://raw.githubusercontent.com/pgsty/silo/codex/repository-cards/contributors-light.svg" alt="SILO 社区贡献者">
</picture>
</a>
[查看贡献记录与实际 PR 状态](CONTRIBUTORS.md)。
## Star History
<picture>
<source media="(prefers-color-scheme: dark)" srcset="https://raw.githubusercontent.com/pgsty/silo/codex/repository-cards/star-history-dark.svg">
<img src="https://raw.githubusercontent.com/pgsty/silo/codex/repository-cards/star-history-light.svg" alt="SILO GitHub 星标历史">
</picture>
## 背景
@@ -134,6 +134,7 @@
"MINIO_API_GZIP_OBJECTS",
"MINIO_API_LEGACY_BUCKET_RESOURCE_MATCH",
"MINIO_API_LIST_QUORUM",
"MINIO_API_MULTIPART_LISTING",
"MINIO_API_OBJECT_MAX_VERSIONS",
"MINIO_API_ODIRECT",
"MINIO_API_REMOTE_TRANSPORT_DEADLINE",
@@ -766,6 +767,7 @@
"/metrics/v3",
"/minio/grid/",
"/minio/grid/lock/",
"/multipart-preflight",
"/netperf",
"/notification",
"/oauth2/callback",
+1
View File
@@ -0,0 +1 @@
__pycache__/
+163
View File
@@ -0,0 +1,163 @@
# Copyright (c) 2026 Feng Ruohang
#
# This program is free software: you can redistribute it and/or modify
# it under the terms of the GNU Affero General Public License as published by
# the Free Software Foundation, either version 3 of the License, or
# (at your option) any later version.
#
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU Affero General Public License for more details.
#
# You should have received a copy of the GNU Affero General Public License
# along with this program. If not, see <http://www.gnu.org/licenses/>.
"""Pure, self-contained SVG rendering for SILO's README cards."""
from datetime import date, timedelta
from html import escape
import math
import xml.etree.ElementTree as ET
themes = {
'light': dict(bg='#ffffff', wash='#f2f7fc', edge='#d9e3ee', ink='#16222e',
muted='#62758a', blue='#1d588c', copper='#b4762e', grid='#e5edf5',
line='#2b6ca3', ring='#dce5ef', field='#f7f9fc', label='#3d4e61'),
'dark': dict(bg='#101923', wash='#152738', edge='#2b3c50', ink='#e8eef6',
muted='#93a3b8', blue='#7fb8e8', copper='#e0a35c', grid='#263749',
line='#5da2dd', ring='#3a4e63', field='#0b1119', label='#b6c2d2'),
}
def read_emblem(path):
emblem = ET.parse(path).getroot()
body = ''.join(ET.tostring(child, encoding='unicode') for child in emblem
if child.tag.rsplit('}', 1)[-1] in ('defs', 'g'))
return '\n'.join(line.rstrip() for line in body.splitlines()).strip()
def txt(x, y, value, size=14, color=None, weight=400, anchor='start', mono=False, spacing=None):
family = 'Menlo,Consolas,monospace' if mono else 'Arial,Helvetica,sans-serif'
extra = f' letter-spacing="{spacing}"' if spacing is not None else ''
return (f'<text x="{x}" y="{y}" font-family="{family}" font-size="{size}" '
f'font-weight="{weight}" fill="{color}" text-anchor="{anchor}"{extra}>'
f'{escape(str(value))}</text>')
def start(height, theme, title, description, emblem_body):
t = themes[theme]
return [f'<svg xmlns="http://www.w3.org/2000/svg" width="1000" height="{height}" '
f'viewBox="0 0 1000 {height}" role="img" aria-labelledby="title desc">',
f'<title id="title">{escape(title)}</title><desc id="desc">{escape(description)}</desc>',
'<defs><linearGradient id="surface" x1="0" y1="1" x2="1" y2="0">'
f'<stop offset="0" stop-color="{t["bg"]}"/>'
f'<stop offset="1" stop-color="{t["wash"]}"/></linearGradient>'
'<linearGradient id="accent" x1="0" y1="0" x2="1" y2="0">'
f'<stop offset="0" stop-color="{t["blue"]}"/>'
f'<stop offset="1" stop-color="{t["copper"]}"/></linearGradient>'
'<linearGradient id="area" x1="0" y1="0" x2="0" y2="1">'
f'<stop offset="0" stop-color="{t["line"]}" stop-opacity=".22"/>'
f'<stop offset="1" stop-color="{t["line"]}" stop-opacity=".015"/>'
'</linearGradient></defs>',
f'<rect x=".75" y=".75" width="998.5" height="{height-1.5}" rx="22" '
f'fill="url(#surface)" stroke="{t["edge"]}" stroke-width="1.5"/>',
f'<svg x="40" y="25" width="25" height="25" viewBox="230 213 570 570">{emblem_body}</svg>']
def heading(parts, t, eyebrow, title, subtitle, value, value_label):
parts.extend([
txt(76, 43, eyebrow, 11, t['muted'], 600, mono=True, spacing=1.7),
txt(40, 94, title, 32, t['ink'], 700),
txt(41, 123, subtitle, 14, t['muted']),
txt(958, 89, f'{value:,}', 45, t['ink'], 700, anchor='end'),
txt(957, 114, value_label, 10, t['muted'], 600, anchor='end', mono=True, spacing=1.5),
f'<path d="M40 146 H960" stroke="{t["edge"]}"/>',
])
def contributors(theme, people, snapshot, emblem_body):
height = 206 + 76 * math.ceil(len(people) / 10)
t = themes[theme]
parts = start(height, theme, f'SILO community — {len(people)} contributors',
f'The existing SILO community roll, including code, proposals and reports across related projects. '
f'Gold rings retain the existing significant-contribution designation. Snapshot {snapshot}.', emblem_body)
heading(parts, t, 'SILO / COMMUNITY', 'Contributors',
'Code, proposals & reports across SILO and related projects', len(people), 'COMMUNITY CONTRIBUTORS')
for row in range(math.ceil(len(people) / 10)):
group = people[row * 10:(row + 1) * 10]
row_width = len(group) * 91
for col, person in enumerate(group):
x = (1000 - row_width) / 2 + col * 91 + 45.5
y = 199 + row * 76
identifier = f'avatar-{row}-{col}'
featured = bool(person.get('featured'))
parts.append(f'<g><title>@{escape(person["handle"])} — {escape(person["what"])}</title>')
parts.append(f'<defs><clipPath id="{identifier}"><circle cx="{x}" cy="{y}" r="29"/></clipPath></defs>')
if featured:
parts.append(f'<circle cx="{x}" cy="{y}" r="34" fill="{t["copper"]}" opacity=".09"/>')
if person.get('avatarDataUrl'):
parts.append(f'<image x="{x-29}" y="{y-29}" width="58" height="58" '
f'clip-path="url(#{identifier})" href="{escape(person["avatarDataUrl"])}"/>')
else:
parts.append(f'<circle cx="{x}" cy="{y}" r="29" fill="{t["ring"]}"/>')
parts.append(txt(x, y + 9, person['handle'][0].upper(), 26, t['ink'], 700, 'middle'))
parts.append(f'<circle cx="{x}" cy="{y}" r="30.5" fill="none" '
f'stroke="{t["copper"] if featured else t["ring"]}" stroke-width="{2 if featured else 1.25}"/></g>')
parts.extend([
f'<path d="M40 {height-42} H960" stroke="{t["edge"]}"/>',
f'<circle cx="47" cy="{height-21}" r="4" fill="none" stroke="{t["copper"]}" stroke-width="1.5"/>',
txt(61, height-17, 'Gold rings mark significant contributions', 12, t['muted']),
txt(959, height-17, f'AS OF {snapshot}', 10, t['muted'], 500, 'end', mono=True, spacing=.6),
'</svg>',
])
return ''.join(parts)
def stars(theme, history, snapshot, emblem_body):
points = history['points']
star_count = points[-1]['stars']
t = themes[theme]
provenance = ('Initial history reconstructed · Daily totals since ' + history['bootstrap']['through'] + ' · UTC'
if history['bootstrap']['reconstructed'] else 'Observed daily star totals · UTC')
parts = start(558, theme, f'SILO star history — {star_count:,} stars',
f'GitHub repository pgsty/silo. {star_count:,} stars as of {snapshot}. ' +
provenance, emblem_body)
heading(parts, t, 'SILO / GITHUB', 'Star History', 'pgsty/silo', star_count, 'GITHUB STARS')
left, right, top, bottom = 76, 958, 177, 440
begin = date.fromisoformat(points[0]['date'])
end = date.fromisoformat(points[-1]['date'])
days = max(1, (end - begin).days)
maximum = max(500, math.ceil(max(p['stars'] for p in points) / 500) * 500)
tick_step = 10 ** max(0, int(math.log10(maximum)))
xy = lambda day, n: (left + (right-left)*(date.fromisoformat(day)-begin).days / days,
bottom-(bottom-top)*n/maximum)
for value in range(0, maximum+1, tick_step):
y = xy(points[0]['date'], value)[1]
parts.append(f'<path d="M{left} {y:.2f} H{right}" stroke="{t["grid"]}" stroke-dasharray="4 6"/>')
parts.append(txt(left-16, round(y+4, 2), f'{value / 1000:g}k' if value >= 1000 else str(value), 12, t['muted'], anchor='end'))
dates = sorted({begin + timedelta(days=round((end-begin).days*i/5)) for i in range(6)})
ticks = [(d.isoformat(), d.strftime('%b %Y') if days > 90 else d.strftime('%b %d')) for d in dates]
for day, label in ticks:
x = xy(day, 0)[0]
parts.append(f'<path d="M{x:.2f} {top} V{bottom}" stroke="{t["grid"]}" stroke-opacity=".65"/>')
anchor = 'start' if day == points[0]['date'] else 'end' if day == snapshot else 'middle'
parts.append(txt(round(x, 2), 466, label, 12, t['muted'], anchor=anchor))
coords = [xy(p['date'], p['stars']) for p in points]
line = 'M' + ' L'.join(f'{x:.2f} {y:.2f}' for x, y in coords)
area = line + f' L{coords[-1][0]:.2f} {bottom} L{left} {bottom} Z'
parts.extend([
f'<path d="{area}" fill="url(#area)"/>',
f'<path d="{line}" fill="none" stroke="url(#accent)" stroke-width="3" '
'stroke-linecap="round" stroke-linejoin="round"/>',
f'<path d="M{left} {bottom} H{right}" stroke="{t["edge"]}"/>',
])
x, y = coords[-1]
parts.extend([
f'<circle cx="{x}" cy="{y}" r="10" fill="{t["copper"]}" opacity=".12"/>',
f'<circle cx="{x}" cy="{y}" r="5" fill="{t["copper"]}" stroke="{t["bg"]}" stroke-width="2"/>',
f'<path d="M40 491 H960" stroke="{t["edge"]}"/>',
txt(40, 516, f'{begin:%b %Y} — {end:%b %Y}'.upper(), 10, t['muted'], 500, mono=True, spacing=.7),
txt(959, 516, f'SNAPSHOT {snapshot}', 10, t['muted'], 500, 'end', mono=True, spacing=.6),
txt(40, 539, provenance, 11, t['muted']),
'</svg>',
])
return ''.join(parts)
@@ -0,0 +1 @@
PyYAML==6.0.3
+178
View File
@@ -0,0 +1,178 @@
# Copyright (c) 2026 Feng Ruohang
#
# This program is free software: you can redistribute it and/or modify
# it under the terms of the GNU Affero General Public License as published by
# the Free Software Foundation, either version 3 of the License, or
# (at your option) any later version.
#
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU Affero General Public License for more details.
#
# You should have received a copy of the GNU Affero General Public License
# along with this program. If not, see <http://www.gnu.org/licenses/>.
"""Regression checks for historical accuracy, contributor scope and SVG safety."""
import base64
import json
from pathlib import Path
import tempfile
import unittest
from unittest.mock import patch
from urllib.error import URLError
import xml.etree.ElementTree as ET
import render
import update
NS = {'s': 'http://www.w3.org/2000/svg'}
PNG = base64.b64decode('iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAQAAAC1HAwCAAAAC0lEQVR42mP8/x8AAwMCAO+a7mgAAAAASUVORK5CYII=')
def person(handle='Alice', group='reports', featured=False):
return {'handle': handle, 'group': group, 'featured': featured,
'what': 'A reviewed contribution', 'firstContribution': '2026-09-01'}
def history():
return {'repository': 'pgsty/silo',
'bootstrap': {'through': '2026-09-15', 'reconstructed': True},
'points': [{'date': '2026-09-14', 'stars': 100}, {'date': '2026-09-15', 'stars': 105}]}
class HistoryTests(unittest.TestCase):
def test_new_day_preserves_old_counts_and_unstars(self):
before = history()
after = update.update_history(before, '2026-09-16', 103)
self.assertEqual(after['points'][:-1], before['points'])
self.assertEqual(after['points'][-1], {'date': '2026-09-16', 'stars': 103})
self.assertEqual(before, history())
def test_same_day_rerun_replaces_instead_of_appending(self):
first = update.update_history(history(), '2026-09-15', 107)
self.assertEqual(len(first['points']), 2)
self.assertEqual(first, update.update_history(first, '2026-09-15', 107))
def test_missing_days_are_not_invented(self):
result = update.update_history(history(), '2026-09-18', 106)
self.assertEqual([p['date'] for p in result['points']], ['2026-09-14', '2026-09-15', '2026-09-18'])
def test_rejects_wrong_repository_and_corrupt_history(self):
cases = []
wrong = history(); wrong['repository'] = 'someone/else'; cases.append(wrong)
duplicate = history(); duplicate['points'].append(duplicate['points'][-1]); cases.append(duplicate)
unordered = history(); unordered['points'].reverse(); cases.append(unordered)
negative = history(); negative['points'][0]['stars'] = -1; cases.append(negative)
future = history(); future['points'][-1]['date'] = '2026-09-20'; cases.append(future)
for case in cases:
with self.subTest(case=case), self.assertRaises(ValueError):
update.update_history(case, '2026-09-16', 100)
def test_first_run_has_no_fabricated_history(self):
result = update.update_history(None, '2026-09-16', 10)
self.assertFalse(result['bootstrap']['reconstructed'])
self.assertEqual(result['points'], [{'date': '2026-09-16', 'stars': 10}])
class ContributorTests(unittest.TestCase):
def test_retries_truncated_json_before_using_it(self):
with patch('update.request', side_effect=[b'{"partial":', b'{"ok":true}']), patch('update.time.sleep'):
self.assertEqual(update.GitHub('').get('repos/pgsty/silo'), {'ok': True})
def test_paginates_past_one_full_page(self):
class API(update.GitHub):
def __init__(self): self.calls = []
def get(self, path):
self.calls.append(path)
return list(range(100)) if 'page=1&' in path else [100]
api = API()
self.assertEqual(len(list(api.issues('pgsty/silo'))), 101)
self.assertIn('state=all', api.calls[0])
self.assertIn('page=2&', api.calls[1])
def test_bots_deduplication_unmerged_work_and_reviewed_credit(self):
def issue(login, kind='issue', user_type='User'):
item = {'user': {'login': login, 'type': user_type, 'avatar_url': ''}, 'created_at': '2026-09-02T00:00:00Z'}
if kind != 'issue': item['pull_request'] = {'merged_at': None if kind == 'open' else '2026-09-03T00:00:00Z'}
return item
class API:
def issues(self, _repo):
return [issue('alice'), issue('Bob', 'open'), issue('Bob', 'merged'),
issue('Carol', 'open'), issue('Copilot'), issue('robot', user_type='Bot')]
curated = {'repositories': ['pgsty/silo', 'pgsty/mc'], 'bots': ['Copilot'],
'people': [person('Alice', featured=True), person('Reporter'), person('Copilot')]}
result = update.collect_people(API(), curated)
self.assertEqual({p['handle'] for p in result}, {'Alice', 'Bob', 'Carol', 'Reporter'})
self.assertEqual(result[0]['handle'], 'Bob')
self.assertEqual(result[1]['handle'], 'Carol')
self.assertTrue(next(p for p in result if p['handle'] == 'Alice')['featured'])
self.assertEqual(next(p for p in result if p['handle'] == 'Bob')['group'], 'code')
self.assertFalse(next(p for p in result if p['handle'] == 'Carol')['featured'])
def test_newer_reviewed_preview_survives_until_site_catches_up(self):
remote = {'updated': '2026-09-16T03:00:00+00:00'}
cached = {'updated': '2026-09-16T04:00:00+00:00'}
self.assertIs(update.select_curated(remote, cached), cached)
newer = {'updated': '2026-09-17T03:00:00+00:00'}
self.assertIs(update.select_curated(newer, cached), newer)
def test_avatar_failure_reuses_raster_cache(self):
previous = {**person(), 'avatarDataUrl': update.raster_data_url(PNG)}
with patch('update.request', side_effect=URLError('unavailable')):
result = update.add_avatars(None, [{**person(), 'avatarUrl': 'https://avatars.githubusercontent.com/u/1'}], [previous])
self.assertEqual(result[0]['avatarDataUrl'], previous['avatarDataUrl'])
with self.assertRaises(ValueError): update.raster_data_url(b'<svg onload="bad()"/>')
with self.assertRaises(ValueError): update.cached_avatar({'avatarDataUrl': 'data:image/svg+xml;base64,PHN2Zy8+'})
def test_fetch_failure_leaves_published_assets_untouched(self):
class API:
def get(self, path):
if path == 'repos/pgsty/silo': return {'full_name': 'pgsty/silo', 'stargazers_count': 106}
raise URLError('roster unavailable')
with tempfile.TemporaryDirectory() as directory:
out = Path(directory)
original = json.dumps(history())
(out / 'history.json').write_text(original)
(out / 'contributors-light.svg').write_text('previous image')
with self.assertRaises(URLError): update.refresh(out, API(), Path('.'))
self.assertEqual((out / 'history.json').read_text(), original)
self.assertEqual((out / 'contributors-light.svg').read_text(), 'previous image')
class RenderTests(unittest.TestCase):
def test_real_emblem_generates_clean_xml(self):
emblem = render.read_emblem(Path(__file__).resolve().parents[2] / '.github/silo.svg')
svg = render.contributors('light', [person()], '2026-09-16', emblem)
ET.fromstring(svg)
self.assertTrue(all(line == line.rstrip() for line in svg.splitlines()))
def test_all_avatars_fit_when_the_roster_grows(self):
people = [{**person(f'person-{i}'), 'avatarDataUrl': update.raster_data_url(PNG)} for i in range(151)]
for theme in ('light', 'dark'):
root = ET.fromstring(render.contributors(theme, people, '2026-09-16', ''))
images = root.findall('.//s:image', NS)
self.assertEqual(len(images), 151)
footer = float(root.attrib['height']) - 42
self.assertTrue(all(float(i.attrib['y']) + float(i.attrib['height']) < footer for i in images))
self.assertTrue(all(i.attrib['href'].startswith('data:image/png;base64,') for i in images))
def test_untrusted_text_is_escaped(self):
data = [{**person(), 'what': '<script>alert("x")</script> & contributions'}]
svg = render.contributors('light', data, '2026-09-16', '')
root = ET.fromstring(svg)
self.assertEqual(root.findall('.//s:script', NS), [])
self.assertIn('&lt;script&gt;', svg)
def test_single_point_and_decreasing_star_history_render(self):
for data in (update.update_history(None, '2026-09-16', 0), update.update_history(history(), '2026-09-16', 90)):
for theme in ('light', 'dark'):
svg = render.stars(theme, data, '2026-09-16', '')
ET.fromstring(svg)
self.assertNotIn('nan', svg.lower())
self.assertNotIn('inf', svg.lower())
if __name__ == '__main__':
unittest.main()
+295
View File
@@ -0,0 +1,295 @@
#!/usr/bin/env python3
# Copyright (c) 2026 Feng Ruohang
#
# This program is free software: you can redistribute it and/or modify
# it under the terms of the GNU Affero General Public License as published by
# the Free Software Foundation, either version 3 of the License, or
# (at your option) any later version.
#
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
# GNU Affero General Public License for more details.
#
# You should have received a copy of the GNU Affero General Public License
# along with this program. If not, see <http://www.gnu.org/licenses/>.
"""Refresh the generated-asset checkout; publishing is handled by the workflow."""
import argparse
import base64
from concurrent.futures import ThreadPoolExecutor
from datetime import date, datetime, timezone
import json
from http.client import IncompleteRead
import os
from pathlib import Path
import re
import sys
import time
from urllib.error import HTTPError, URLError
from urllib.parse import urlparse
from urllib.request import Request, urlopen
import xml.etree.ElementTree as ET
import yaml
import render
REPOSITORY = 'pgsty/silo'
SOURCE = 'repos/pgsty/silo.pgsty.com/contents/data/home/contributors.yaml?ref=main'
GROUPS = ('code', 'proposed', 'reports')
HANDLE = re.compile(r'[A-Za-z0-9][A-Za-z0-9-]{0,38}\Z')
def request(url, token='', limit=8 * 1024 * 1024):
headers = {'User-Agent': 'silo-repository-cards', 'Accept': 'application/vnd.github+json'}
if urlparse(url).netloc == 'api.github.com':
headers['X-GitHub-Api-Version'] = '2022-11-28'
if token:
headers['Authorization'] = f'Bearer {token}'
for attempt in range(3):
try:
with urlopen(Request(url, headers=headers), timeout=25) as response:
data = response.read(limit + 1)
if len(data) > limit:
raise ValueError('Response exceeds the size limit')
expected = response.headers.get('Content-Length')
if expected is not None and len(data) != int(expected):
raise URLError('Incomplete response body')
return data
except HTTPError as exc:
if exc.code < 500 or attempt == 2:
raise
except (URLError, TimeoutError, IncompleteRead):
if attempt == 2:
raise
time.sleep(attempt + 1)
class GitHub:
def __init__(self, token):
self.token = token
def get(self, path):
for attempt in range(3):
try:
return json.loads(request('https://api.github.com/' + path, self.token))
except (json.JSONDecodeError, UnicodeDecodeError) as exc:
if attempt == 2:
raise ValueError(f'Incomplete or invalid GitHub JSON: {path}') from exc
time.sleep(attempt + 1)
def issues(self, repository):
page = 1
while True:
batch = self.get(f'repos/{repository}/issues?state=all&per_page=100&page={page}&sort=created&direction=asc')
if not isinstance(batch, list):
raise ValueError(f'Invalid issues response for {repository}')
yield from batch
if len(batch) < 100:
return
page += 1
def curated_snapshot(data, revision):
updated = str(data['updated'])
datetime.fromisoformat(updated)
repositories = [item['repo'] for item in data['repositories']]
if not repositories or any(not re.fullmatch(r'pgsty/[A-Za-z0-9_.-]+', repo) for repo in repositories):
raise ValueError('Invalid contributor repository scope')
people = []
for group in GROUPS:
for entry in data[group]:
if not HANDLE.fullmatch(entry['handle']):
raise ValueError('Invalid GitHub contributor handle')
people.append({
'handle': entry['handle'], 'group': group,
'featured': bool(entry.get('featured')), 'what': entry['what'],
'firstContribution': str(entry.get('firstContribution', '9999-12-31')),
})
if not people or len({p['handle'].lower() for p in people}) != len(people):
raise ValueError('Empty or duplicate contributor roster')
return {'updated': updated, 'revision': revision, 'repositories': repositories,
'bots': data.get('bots', ['Copilot', 'dependabot[bot]']), 'people': people}
def select_curated(remote, cached):
# The initial, approved preview can contain reviewed credit not published by
# the companion site yet. Keep that newer snapshot until the site catches up.
if cached and datetime.fromisoformat(cached['updated']) > datetime.fromisoformat(remote['updated']):
return cached
return remote
def collect_people(api, curated):
bots = {name.lower() for name in curated['bots']}
people = {p['handle'].lower(): dict(p) for p in curated['people']
if p['handle'].lower() not in bots and not p['handle'].lower().endswith('[bot]')}
order = {p['handle'].lower(): index for index, p in enumerate(curated['people'])}
for repository in curated['repositories']:
print(f'Reading issue and PR authors: {repository}', flush=True)
for issue in api.issues(repository):
user = issue.get('user') or {}
handle = user.get('login', '')
key = handle.lower()
if user.get('type') != 'User' or key in bots or key.endswith('[bot]'):
continue
if not HANDLE.fullmatch(handle):
raise ValueError('Invalid issue author')
pr = issue.get('pull_request')
group = 'code' if pr and pr.get('merged_at') else 'proposed' if pr else 'reports'
first = issue['created_at'][:10]
date.fromisoformat(first)
person = people.setdefault(key, {
'handle': handle, 'group': group, 'featured': False,
'what': 'Contributed an issue or pull request to SILO and related projects',
'firstContribution': first,
})
person['avatarUrl'] = user.get('avatar_url', '')
person['firstContribution'] = min(person['firstContribution'], first)
if GROUPS.index(group) < GROUPS.index(person['group']):
person['group'] = group
if not people:
raise ValueError('No human contributors were collected')
return sorted(people.values(), key=lambda p: (
GROUPS.index(p['group']), not p['featured'],
order.get(p['handle'].lower(), len(order)), p['firstContribution'], p['handle'].lower()))
def raster_data_url(data):
if data.startswith(b'\x89PNG\r\n\x1a\n'):
mime = 'image/png'
elif data.startswith(b'\xff\xd8\xff'):
mime = 'image/jpeg'
elif data.startswith((b'GIF87a', b'GIF89a')):
mime = 'image/gif'
elif data[:4] == b'RIFF' and data[8:12] == b'WEBP':
mime = 'image/webp'
else:
raise ValueError('Avatar is not a raster image')
return f'data:{mime};base64,' + base64.b64encode(data).decode('ascii')
def cached_avatar(person):
value = person.get('avatarDataUrl', '')
if not value:
return ''
prefix, encoded = value.split(',', 1)
if prefix not in ('data:image/png;base64', 'data:image/jpeg;base64', 'data:image/gif;base64', 'data:image/webp;base64'):
raise ValueError('Invalid cached avatar format')
raw = base64.b64decode(encoded, validate=True)
if len(raw) > 512 * 1024 or raster_data_url(raw) != value:
raise ValueError('Invalid cached avatar')
return value
def add_avatars(api, people, previous):
cached = {p['handle'].lower(): cached_avatar(p) for p in previous}
def update(person):
person = dict(person)
try:
url = person.pop('avatarUrl', '') or api.get('users/' + person['handle'])['avatar_url']
parsed = urlparse(url)
if parsed.scheme != 'https' or parsed.netloc != 'avatars.githubusercontent.com':
raise ValueError('Unexpected avatar host')
data = request(url + ('&' if '?' in url else '?') + 's=96', limit=512 * 1024)
person['avatarDataUrl'] = raster_data_url(data)
except (HTTPError, URLError, TimeoutError, IncompleteRead, ValueError, KeyError) as exc:
person.pop('avatarUrl', None)
person['avatarDataUrl'] = cached.get(person['handle'].lower(), '')
print(f'Avatar fallback for @{person["handle"]}: {type(exc).__name__}', file=sys.stderr)
return person
with ThreadPoolExecutor(max_workers=6) as pool:
return list(pool.map(update, people))
def update_history(history, day, stars):
date.fromisoformat(day)
if type(stars) is not int or stars < 0:
raise ValueError('Invalid repository star count')
if history is None:
history = {'repository': REPOSITORY, 'bootstrap': {'through': day, 'reconstructed': False}, 'points': []}
if history['repository'] != REPOSITORY:
raise ValueError('Star history belongs to a different repository')
date.fromisoformat(history['bootstrap']['through'])
dates = []
for point in history['points']:
date.fromisoformat(point['date'])
if type(point['stars']) is not int or point['stars'] < 0:
raise ValueError('Invalid historical star count')
dates.append(point['date'])
if dates != sorted(set(dates)) or any(d > day for d in dates):
raise ValueError('History contains duplicate, unordered, or future dates')
# Replace today's observation, preserve previous days, and allow unstars.
points = [dict(p) for p in history['points'] if p['date'] != day]
points.append({'date': day, 'stars': stars})
return {**history, 'points': points}
def read_json(path, default=None):
return json.loads(path.read_text()) if path.exists() else default
def refresh(output, api, source_root):
day = datetime.now(timezone.utc).date().isoformat()
metadata = api.get('repos/' + REPOSITORY)
if metadata['full_name'].lower() != REPOSITORY:
raise ValueError('Unexpected repository metadata')
history = update_history(read_json(output / 'history.json'), day, metadata['stargazers_count'])
source = api.get(SOURCE)
reviewed = yaml.safe_load(base64.b64decode(source['content'], validate=False))
curated = select_curated(curated_snapshot(reviewed, source['sha']), read_json(output / 'curated.json'))
people = collect_people(api, curated)
previous = read_json(output / 'contributors.json', {}).get('people', [])
people = add_avatars(api, people, previous)
emblem = render.read_emblem(source_root / '.github/silo.svg')
payloads = {}
for theme in ('light', 'dark'):
payloads[f'contributors-{theme}.svg'] = render.contributors(theme, people, day, emblem) + '\n'
payloads[f'star-history-{theme}.svg'] = render.stars(theme, history, day, emblem) + '\n'
for svg in payloads.values():
ET.fromstring(svg)
for name, data in {
'history.json': history,
'curated.json': curated,
'contributors.json': {'repository': REPOSITORY, 'updated': day, 'people': people},
}.items():
payloads[name] = json.dumps(data, indent=2, ensure_ascii=False) + '\n'
payloads['README.md'] = f'''# SILO repository cards
Generated by [Repository Cards](https://github.com/pgsty/silo/actions/workflows/repository-cards.yml)
at 00:00 UTC daily (08:00 Asia/Shanghai). GitHub may queue scheduled runs.
Snapshot: {day}. {metadata['stargazers_count']:,} stars; {len(people)} community contributors.
- `contributors-light.svg` / `contributors-dark.svg`: human issue and PR authors across the SILO project scope, plus reviewed acknowledgements. Bots are excluded. Gold rings follow the reviewed companion-site roster; new authors are collected automatically.
- `star-history-light.svg` / `star-history-dark.svg`: initial history reconstructed from the then-current stargazers; later points are daily observed totals, including decreases. Missing days are not fabricated.
- `curated.json`: a cache of reviewed contributor credit from `pgsty/silo.pgsty.com/data/home/contributors.yaml`. The approved initial preview may be newer than the published site; a newer reviewed snapshot is retained until the site catches up.
- `contributors.json`: generated contributor data and embedded raster avatars. Failed avatar refreshes use the previous image, or an initial when no image is available.
- `history.json`: persistent daily totals. Keep this file when regenerating images.
The SVGs are self-contained. Source and instructions live on the default branch;
this branch contains generated assets only. Do not merge it into `main`.
'''
# Collect and validate everything before touching the publication checkout.
output.mkdir(parents=True, exist_ok=True)
for filename, text in payloads.items():
(output / filename).write_text(text)
print(f'{day}: {len(people)} contributors; {metadata["stargazers_count"]:,} stars; {len(history["points"])} history points')
def main():
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument('--output', type=Path, required=True)
args = parser.parse_args()
configured = os.environ.get('GITHUB_REPOSITORY', REPOSITORY)
if configured.lower() != REPOSITORY:
raise SystemExit('This workflow is scoped to pgsty/silo')
refresh(args.output, GitHub(os.environ.get('GH_TOKEN', '')), Path(__file__).resolve().parents[2])
if __name__ == '__main__':
main()
+1
View File
@@ -163,6 +163,7 @@ func registerAdminRouter(router *mux.Router, enableConfigOps bool) {
// StorageInfo operations
adminRouter.Methods(http.MethodGet).Path(adminVersion + "/storageinfo").HandlerFunc(adminMiddleware(adminAPI.StorageInfoHandler, traceAllFlag))
adminRouter.Methods(http.MethodGet).Path(adminVersion + "/multipart-preflight").HandlerFunc(adminMiddleware(adminAPI.MultipartPreflightHandler, traceAllFlag))
// DataUsageInfo operations
adminRouter.Methods(http.MethodGet).Path(adminVersion + "/datausageinfo").HandlerFunc(adminMiddleware(adminAPI.DataUsageInfoHandler, traceAllFlag))
// Metrics operation
+24
View File
@@ -450,6 +450,9 @@ const (
ErrAdminNoSecretKey
ErrIAMNotInitialized
ErrMultipartListingLegacy
ErrMultipartListingIdentity
ErrSlowDown
apiErrCodeEnd // This is used only for the testing code
)
@@ -1336,6 +1339,21 @@ var errorCodes = errorCodeMap{
Description: "IAM sub-system not initialized yet, please try again.",
HTTPStatusCode: http.StatusServiceUnavailable,
},
ErrMultipartListingLegacy: {
Code: "MultipartListingNotReady",
Description: "Legacy multipart uploads prevent a complete listing. Upgrade all writers, drain old uploads and run the multipart preflight check.",
HTTPStatusCode: http.StatusServiceUnavailable,
},
ErrMultipartListingIdentity: {
Code: "MultipartListingMetadataInvalid",
Description: "Multipart upload metadata is inconsistent. Run the multipart preflight check to locate the affected storage set.",
HTTPStatusCode: http.StatusServiceUnavailable,
},
ErrSlowDown: {
Code: "SlowDown",
Description: "Please reduce your request rate",
HTTPStatusCode: http.StatusServiceUnavailable,
},
ErrBucketMetadataNotInitialized: {
Code: "XMinioBucketMetadataNotInitialized",
Description: "Bucket metadata not initialized yet, please try again.",
@@ -2173,6 +2191,10 @@ func toAPIErrorCode(ctx context.Context, err error) (apiErr APIErrorCode) {
err = unwrapAll(err)
switch err {
case errMultipartListingLegacy:
apiErr = ErrMultipartListingLegacy
case errMultipartListingIdentity:
apiErr = ErrMultipartListingIdentity
case errCompleteMultipartChecksumMismatch, errCompleteMultipartChecksumTypeMismatch:
apiErr = ErrBadDigest
case errMissingPartChecksum:
@@ -2303,6 +2325,8 @@ func toAPIErrorCode(ctx context.Context, err error) (apiErr APIErrorCode) {
}
switch err.(type) {
case SlowDown:
apiErr = ErrSlowDown
case StorageFull:
apiErr = ErrStorageFull
case hash.BadDigest:
+1 -1
View File
@@ -41,7 +41,7 @@ import (
const (
maxObjectList = 1000 // Limit number of objects in a listObjectsResponse/listObjectsVersionsResponse.
maxDeleteList = 1000 // Limit number of objects deleted in a delete call.
maxUploadsList = 10000 // Limit number of uploads in a listUploadsResponse.
maxUploadsList = 1000 // Limit number of uploads in a listUploadsResponse.
maxPartsList = 10000 // Limit number of parts in a listPartsResponse.
)
File diff suppressed because one or more lines are too long
-8
View File
@@ -277,14 +277,6 @@ func (api objectAPIHandlers) ListMultipartUploadsHandler(w http.ResponseWriter,
return
}
if keyMarker != "" {
// Marker not common with prefix is not implemented.
if !HasPrefix(keyMarker, prefix) {
writeErrorResponse(ctx, w, errorCodes.ToAPIErr(ErrNotImplemented), r.URL)
return
}
}
listMultipartsInfo, err := objectAPI.ListMultipartUploads(ctx, bucket, prefix, keyMarker, uploadIDMarker, delimiter, maxUploads)
if err != nil {
writeErrorResponse(ctx, w, toAPIError(ctx, err), r.URL)
+3 -3
View File
@@ -425,7 +425,7 @@ func testListMultipartUploadsHandler(obj ObjectLayer, instanceType, bucketName s
shouldPass: true,
},
// Test case - 4.
// Setting Invalid prefix and marker combination.
// A key marker outside the prefix is valid and produces an empty page.
{
bucket: bucketName,
prefix: "asia",
@@ -435,8 +435,8 @@ func testListMultipartUploadsHandler(obj ObjectLayer, instanceType, bucketName s
maxUploads: "0",
accessKey: credentials.AccessKey,
secretKey: credentials.SecretKey,
expectedRespStatus: http.StatusNotImplemented,
shouldPass: false,
expectedRespStatus: http.StatusOK,
shouldPass: true,
},
// Test case - 5.
// Invalid upload id and marker combination.
+385
View File
@@ -0,0 +1,385 @@
// Copyright (c) 2026 mr javad seydi and Ruohang Feng
//
// This file is part of Silo Object Storage stack.
// SPDX-License-Identifier: AGPL-3.0-or-later
package cmd
import (
"context"
"encoding/json"
"errors"
"fmt"
"net/http"
"strings"
"sync"
"sync/atomic"
"time"
"github.com/pgsty/silo-pkg/v3/policy"
)
const multipartScanEntryLimit = 100_000
var (
multipartScanSlots = make(chan struct{}, 2)
errMultipartListingLegacy = errors.New("legacy multipart uploads require a coordinated upgrade and drain")
errMultipartListingIdentity = errors.New("multipart upload identity is invalid")
)
// One budget and admission slot cover the entire request, including all pools.
// A slot is released only after the scan workers have actually stopped.
type multipartScan struct {
ctx context.Context
cancel context.CancelFunc
remaining atomic.Int64
metadataSlots chan struct{}
preflight bool
sets []multipartScanSet
}
type multipartScanSet struct {
Pool int `json:"pool"`
Set int `json:"set"`
Drives int `json:"drives"`
ScannedDrives int `json:"scannedDrives"`
UncoveredDrives []int `json:"uncoveredDrives,omitempty"`
Candidates int `json:"candidates"`
LegacyUploads int `json:"legacyUploads"`
OldestLegacy time.Time `json:"oldestLegacy,omitempty"`
Error string `json:"error,omitempty"`
}
type multipartPreflightReport struct {
Ready bool `json:"ready"`
Complete bool `json:"complete"`
Mode string `json:"mode"`
ScannedEntries int64 `json:"scannedEntries"`
LegacyUploads int `json:"legacyUploads"`
Sets []multipartScanSet `json:"sets"`
}
func startMultipartScan(ctx context.Context, preflight bool) (*multipartScan, error) {
select {
case multipartScanSlots <- struct{}{}:
case <-ctx.Done():
return nil, ctx.Err()
default:
return nil, SlowDown{}
}
ctx, cancel := context.WithTimeout(ctx, 30*time.Second)
s := &multipartScan{ctx: ctx, cancel: cancel, preflight: preflight, metadataSlots: make(chan struct{}, multipartMetadataScanConcurrency)}
s.remaining.Store(multipartScanEntryLimit)
return s, nil
}
func (s *multipartScan) close() {
s.cancel()
<-multipartScanSlots
}
func (s *multipartScan) listDir(disk StorageAPI, bucket, dir string) ([]string, error) {
if err := s.ctx.Err(); err != nil {
return nil, SlowDown{}
}
remaining := s.remaining.Load()
if remaining <= 0 {
return nil, SlowDown{}
}
// Request one extra entry to detect overflow, never a silently partial page.
entries, err := disk.ListDir(s.ctx, bucket, minioMetaMultipartBucket, dir, int(remaining)+1)
if errors.Is(err, errFileNotFound) {
return nil, nil
}
if err != nil {
return nil, err
}
if s.remaining.Add(-int64(len(entries))) < 0 {
return nil, SlowDown{}
}
return entries, nil
}
func (s *multipartScan) listUploadDirs(disk StorageAPI, bucket string) ([]string, error) {
hashDirs, err := s.listDir(disk, bucket, "")
if err != nil {
return nil, err
}
var candidates []string
for _, hashDir := range hashDirs {
if !strings.HasSuffix(hashDir, SlashSeparator) {
continue
}
hashDir = strings.TrimSuffix(hashDir, SlashSeparator)
uploadDirs, err := s.listDir(disk, bucket, hashDir)
if err != nil {
return nil, err
}
for _, uploadDir := range uploadDirs {
if strings.HasSuffix(uploadDir, SlashSeparator) {
candidates = append(candidates, pathJoin(hashDir, strings.TrimSuffix(uploadDir, SlashSeparator)))
}
}
}
return candidates, nil
}
// Identity is immutable at hash/uploadUUID. Check the hash before using even
// a single source drive to exclude another bucket. Missing/bad identities must
// fall back to the full metadata read; they cannot prove absence.
func (er erasureObjects) multipartIdentity(fi FileInfo, shaDir string) (string, string, bool) {
bucket, object := fi.Metadata[multipartMetaBucket], fi.Metadata[multipartMetaObject]
return bucket, object, bucket != "" && object != "" && IsValidBucketName(bucket) &&
IsValidObjectPrefix(object) && er.getMultipartSHADir(bucket, object) == shaDir
}
func (er erasureObjects) readMultipartUploadCandidate(s *multipartScan, bucket, candidate string, source StorageAPI) (MultipartInfo, bool, bool, error) {
if s.ctx.Err() != nil {
return MultipartInfo{}, false, false, SlowDown{}
}
shaDir, uploadUUID, ok := strings.Cut(candidate, SlashSeparator)
if !ok || shaDir == "" || uploadUUID == "" || strings.Contains(uploadUUID, SlashSeparator) {
return MultipartInfo{}, false, false, errMultipartListingIdentity
}
if !s.preflight && bucket != "" && source != nil {
fi, err := source.ReadVersion(s.ctx, bucket, minioMetaMultipartBucket, candidate, "", ReadOptions{})
if err == nil {
storedBucket, _, valid := er.multipartIdentity(fi, shaDir)
if valid && storedBucket != bucket {
return MultipartInfo{}, false, false, nil
}
}
}
if s.ctx.Err() != nil {
return MultipartInfo{}, false, false, SlowDown{}
}
select {
case s.metadataSlots <- struct{}{}:
defer func() { <-s.metadataSlots }()
case <-s.ctx.Done():
return MultipartInfo{}, false, false, SlowDown{}
}
disks := er.getDisks()
metadata, errs := readAllFileInfo(s.ctx, disks, bucket, minioMetaMultipartBucket, candidate, "", false, false)
if s.preflight {
// Readiness is stronger than listing liveness: even one readable legacy
// copy must be drained, and an unreadable copy cannot certify readiness.
var oldest MultipartInfo
legacy := false
_, nativeID := multipartUploadTime(uploadUUID)
for i, err := range errs {
if errors.Is(err, errFileNotFound) || errors.Is(err, errFileVersionNotFound) {
continue
}
if err != nil {
return MultipartInfo{}, false, false, err
}
fi := metadata[i]
storedBucket, storedObject, valid := er.multipartIdentity(fi, shaDir)
if storedBucket != "" && storedObject != "" && !valid {
return MultipartInfo{}, false, false, errMultipartListingIdentity
}
if !valid || !nativeID {
info := multipartUploadInfo(storedBucket, storedObject, uploadUUID, fi.ModTime)
if !legacy || info.Initiated.Before(oldest.Initiated) {
oldest = info
}
legacy = true
}
}
if legacy {
return oldest, false, true, nil
}
}
readQuorum, _, err := objectQuorumFromMeta(s.ctx, metadata, errs, er.defaultParityCount)
if err != nil {
return MultipartInfo{}, false, false, err
}
_, modTime, etag := listOnlineDisks(disks, metadata, errs, readQuorum)
if err := reduceReadQuorumErrs(s.ctx, errs, objectOpIgnoredErrs, readQuorum); err != nil {
return MultipartInfo{}, false, false, err
}
fi, err := pickValidFileInfo(s.ctx, metadata, modTime, etag, readQuorum)
if err != nil {
return MultipartInfo{}, false, false, err
}
storedBucket, storedObject, valid := er.multipartIdentity(fi, shaDir)
info := multipartUploadInfo(storedBucket, storedObject, uploadUUID, fi.ModTime)
if storedBucket == "" || storedObject == "" {
return info, false, true, nil
}
if !valid {
return MultipartInfo{}, false, false, fmt.Errorf("%w: %s", errMultipartListingIdentity, candidate)
}
if _, ok := multipartUploadTime(uploadUUID); !ok {
return info, false, true, nil
}
if bucket != "" && bucket != storedBucket {
return MultipartInfo{}, false, false, nil
}
return info, true, false, nil
}
func (er erasureObjects) scanMultipartUploads(s *multipartScan, bucket string, poolIdx, setIdx int) ([]MultipartInfo, bool, error) {
disks := er.getDisks()
report := multipartScanSet{Pool: poolIdx, Set: setIdx, Drives: er.setDriveCount}
var resultErr error
defer func() {
if resultErr != nil {
report.Error = resultErr.Error()
}
s.sets = append(s.sets, report)
}()
var candidateMu sync.Mutex
candidates := make(map[string]StorageAPI)
errs := make([]error, len(disks))
var wg sync.WaitGroup
for i, disk := range disks {
wg.Add(1)
go func() {
defer wg.Done()
if disk == nil || !disk.IsOnline() {
errs[i] = errDiskNotFound
return
}
paths, err := s.listUploadDirs(disk, bucket)
errs[i] = err
if err != nil {
return
}
candidateMu.Lock()
defer candidateMu.Unlock()
for _, p := range paths {
candidates[p] = disk
}
}()
}
wg.Wait()
for i, err := range errs {
if err == nil {
report.ScannedDrives++
} else {
report.UncoveredDrives = append(report.UncoveredDrives, i)
if _, limited := err.(SlowDown); limited {
resultErr = err
}
}
}
report.Candidates = len(candidates)
// A successful scan must intersect every metadata read quorum, including
// records left with R copies by a partially failed cancellation.
if report.ScannedDrives < er.setDriveCount/2+1 && resultErr == nil {
resultErr = toObjectErr(errErasureReadQuorum, bucket)
}
if resultErr != nil {
return nil, false, resultErr
}
type candidate struct {
path string
source StorageAPI
}
jobs := make(chan candidate)
var mu sync.Mutex
var uploads []MultipartInfo
for range min(16, len(candidates)) {
wg.Add(1)
go func() {
defer wg.Done()
for c := range jobs {
mu.Lock()
stopped := resultErr != nil
mu.Unlock()
if stopped || s.ctx.Err() != nil {
continue
}
upload, found, legacy, err := er.readMultipartUploadCandidate(s, bucket, c.path, c.source)
if errors.Is(err, errFileNotFound) || errors.Is(err, errFileVersionNotFound) {
continue
}
mu.Lock()
switch {
case err != nil && resultErr == nil:
resultErr = err
case legacy:
report.LegacyUploads++
if report.OldestLegacy.IsZero() || upload.Initiated.Before(report.OldestLegacy) {
report.OldestLegacy = upload.Initiated
}
case found && !s.preflight:
uploads = append(uploads, upload)
}
mu.Unlock()
}
}()
}
dispatch:
for p, source := range candidates {
select {
case jobs <- candidate{p, source}:
case <-s.ctx.Done():
break dispatch
}
}
close(jobs)
wg.Wait()
if s.ctx.Err() != nil {
resultErr = SlowDown{}
}
return uploads, report.LegacyUploads != 0, resultErr
}
// multipartPreflight scans all pools, sets and drives, independent of caches.
// A majority suffices for normal listing; upgrade readiness requires every
// drive to have been inspected. The operator must also upgrade all writers.
func (z *erasureServerPools) multipartPreflight(ctx context.Context) (multipartPreflightReport, error) {
s, err := startMultipartScan(ctx, true)
if err != nil {
return multipartPreflightReport{}, err
}
defer s.close()
report := multipartPreflightReport{Complete: true, Mode: "strict"}
if globalAPIConfig.getMultipartListingLegacy() {
report.Mode = "legacy"
}
for p, pool := range z.serverPools {
for i, set := range pool.sets {
_, _, _ = set.scanMultipartUploads(s, "", p, i)
}
}
report.Sets = s.sets
for _, set := range s.sets {
report.LegacyUploads += set.LegacyUploads
if set.Error != "" || set.ScannedDrives != set.Drives {
report.Complete = false
}
}
report.ScannedEntries = multipartScanEntryLimit - s.remaining.Load()
report.Ready = report.Complete && report.LegacyUploads == 0
return report, nil
}
// MultipartPreflightHandler is a read-only storage-admin diagnostic. It never
// accepts an arbitrary deletion path or changes upload lifetime settings.
func (a adminAPIHandlers) MultipartPreflightHandler(w http.ResponseWriter, r *http.Request) {
ctx := r.Context()
obj, _ := validateAdminReq(ctx, w, r, policy.StorageInfoAdminAction)
if obj == nil {
return
}
z, ok := obj.(*erasureServerPools)
if !ok {
writeErrorResponseJSON(ctx, w, errorCodes.ToAPIErr(ErrNotImplemented), r.URL)
return
}
report, err := z.multipartPreflight(ctx)
if err != nil {
writeErrorResponseJSON(ctx, w, toAdminAPIErr(ctx, err), r.URL)
return
}
b, err := json.Marshal(report)
if err != nil {
writeErrorResponseJSON(ctx, w, toAdminAPIErr(ctx, err), r.URL)
return
}
writeSuccessResponseJSON(w, b)
}
+796
View File
@@ -0,0 +1,796 @@
// Copyright (c) 2026 Ruohang Feng
// SPDX-License-Identifier: AGPL-3.0-or-later
package cmd
import (
"context"
"encoding/base64"
"encoding/json"
"encoding/xml"
"errors"
"fmt"
"net/http"
"net/http/httptest"
"slices"
"strconv"
"sync"
"sync/atomic"
"testing"
"time"
"github.com/minio/minio/internal/config/storageclass"
)
// Check the real handler and storage path, including a successful abort between
// pages. Choose IDs increasing in both initiation time and lexical order so the
// failure does not depend on differing interpretations of S3 marker ordering.
func TestMultipartListingAbortBetweenHTTPPages(t *testing.T) {
z, _ := consistencyPools(t)
bucket, router, err := initAPIHandlerTest(t.Context(), z, []string{"ListMultipartUploads", "AbortMultipart"}, MakeBucketOptions{})
if err != nil {
t.Fatal(err)
}
var firstID, secondID string
for attempt := 0; attempt < 32; attempt++ {
one, err := z.NewMultipartUpload(t.Context(), bucket, "a", ObjectOptions{})
if err != nil {
t.Fatal(err)
}
two, err := z.NewMultipartUpload(t.Context(), bucket, "a", ObjectOptions{})
if err != nil {
t.Fatal(err)
}
if one.UploadID < two.UploadID {
firstID, secondID = one.UploadID, two.UploadID
break
}
for _, id := range []string{one.UploadID, two.UploadID} {
if err := z.AbortMultipartUpload(t.Context(), bucket, "a", id, ObjectOptions{}); err != nil {
t.Fatal(err)
}
}
}
if firstID == "" {
t.Fatal("could not construct increasing upload IDs")
}
if _, err := z.NewMultipartUpload(t.Context(), bucket, "b", ObjectOptions{}); err != nil {
t.Fatal(err)
}
request := func(method, u string) *httptest.ResponseRecorder {
t.Helper()
req, err := newTestSignedRequestV4(method, u, 0, nil, globalActiveCred.AccessKey, globalActiveCred.SecretKey, nil)
if err != nil {
t.Fatal(err)
}
rec := httptest.NewRecorder()
router.ServeHTTP(rec, req)
return rec
}
list := func(keyMarker, uploadMarker string, limit int) ListMultipartUploadsResponse {
t.Helper()
rec := request(http.MethodGet, getListMultipartUploadsURLWithParams("", bucket, "", keyMarker, uploadMarker, "", strconv.Itoa(limit)))
if rec.Code != http.StatusOK {
t.Fatalf("list HTTP %d: %s", rec.Code, rec.Body.String())
}
var result ListMultipartUploadsResponse
if err := xml.Unmarshal(rec.Body.Bytes(), &result); err != nil {
t.Fatal(err)
}
return result
}
first := list("", "", 1)
if len(first.Uploads) != 1 || first.Uploads[0].UploadID != firstID || !first.IsTruncated {
t.Fatalf("unexpected first page: %+v", first)
}
rec := request(http.MethodDelete, getAbortMultipartUploadURL("", bucket, "a", firstID))
if rec.Code != http.StatusNoContent {
t.Fatalf("abort HTTP %d: %s", rec.Code, rec.Body.String())
}
rest := list(first.NextKeyMarker, first.NextUploadIDMarker, 10)
var keys []string
for _, upload := range rest.Uploads {
keys = append(keys, upload.Key)
}
t.Logf("GET first page=200; DELETE marker=204; GET next page=200 keys=%v truncated=%v", keys, rest.IsTruncated)
if _, err := z.GetMultipartInfo(t.Context(), bucket, "a", secondID, ObjectOptions{}); err != nil {
t.Fatalf("remaining upload is not valid: %v", err)
}
if len(rest.Uploads) != 2 || rest.Uploads[0].UploadID != secondID {
t.Fatalf("valid remaining upload for key a is missing from continuation: %+v", rest.Uploads)
}
}
type multipartListingFaultDisk struct {
StorageAPI
read func(context.Context, string, string, string, string, ReadOptions) (FileInfo, error)
delete func(context.Context, string, string, DeleteOptions) error
}
func (d multipartListingFaultDisk) ReadVersion(ctx context.Context, original, volume, object, version string, opts ReadOptions) (FileInfo, error) {
if d.read != nil {
return d.read(ctx, original, volume, object, version, opts)
}
return d.StorageAPI.ReadVersion(ctx, original, volume, object, version, opts)
}
func (d multipartListingFaultDisk) Delete(ctx context.Context, volume, object string, opts DeleteOptions) error {
if d.delete != nil {
return d.delete(ctx, volume, object, opts)
}
return d.StorageAPI.Delete(ctx, volume, object, opts)
}
func TestMultipartListingLegacyPreflight(t *testing.T) {
z, set, bucket := multipartListingFixture(t)
const other = "multipart-legacy-other"
if err := z.MakeBucket(t.Context(), other, MakeBucketOptions{}); err != nil {
t.Fatal(err)
}
old, err := z.NewMultipartUpload(t.Context(), other, "old", ObjectOptions{})
if err != nil {
t.Fatal(err)
}
fi, metadata, err := set.checkUploadIDExists(t.Context(), other, "old", old.UploadID, true)
if err != nil {
t.Fatal(err)
}
for i := range metadata {
delete(metadata[i].Metadata, multipartMetaBucket)
delete(metadata[i].Metadata, multipartMetaObject)
}
if _, err = writeAllMetadata(t.Context(), set.getDisks(), other, minioMetaMultipartBucket,
set.getUploadIDDir(other, "old", old.UploadID), metadata, fi.WriteQuorum(set.defaultWQuorum())); err != nil {
t.Fatal(err)
}
if _, err = z.NewMultipartUpload(t.Context(), bucket, "a", ObjectOptions{}); err != nil {
t.Fatal(err)
}
z.mpCache.Clear()
_, err = z.ListMultipartUploads(t.Context(), bucket, "", "", "", "", 10)
if !errors.Is(err, errMultipartListingLegacy) {
t.Fatalf("old upload in another bucket: %v", err)
}
report, err := z.multipartPreflight(t.Context())
if err != nil || report.Ready || !report.Complete || report.LegacyUploads != 1 {
t.Fatalf("legacy preflight: %+v %v", report, err)
}
if err = z.AbortMultipartUpload(t.Context(), other, "old", old.UploadID, ObjectOptions{}); err != nil {
t.Fatal(err)
}
report, err = z.multipartPreflight(t.Context())
if err != nil || !report.Ready || report.LegacyUploads != 0 {
t.Fatalf("drained preflight: %+v %v", report, err)
}
got, err := z.ListMultipartUploads(t.Context(), bucket, "", "", "", "", 10)
if err != nil {
t.Fatal(err)
}
requireMultipartUploadKeys(t, got, "a")
}
func TestMultipartListingIdentityFallback(t *testing.T) {
for _, kind := range []string{"old", "corrupt", "read-failure", "wrong-bucket"} {
t.Run(kind, func(t *testing.T) {
z, set, bucket := multipartListingFixture(t)
if _, err := z.NewMultipartUpload(t.Context(), bucket, "a", ObjectOptions{}); err != nil {
t.Fatal(err)
}
original := set.getDisks
t.Cleanup(func() { set.getDisks = original })
var first atomic.Bool
set.getDisks = func() []StorageAPI {
disks := append([]StorageAPI(nil), original()...)
for i, d := range disks {
disks[i] = multipartListingFaultDisk{StorageAPI: d, read: func(ctx context.Context, b, v, p, version string, opts ReadOptions) (FileInfo, error) {
fi, err := d.ReadVersion(ctx, b, v, p, version, opts)
if v == minioMetaMultipartBucket && first.CompareAndSwap(false, true) {
if kind == "read-failure" {
return FileInfo{}, errDiskNotFound
}
if kind == "corrupt" {
return FileInfo{}, errFileCorrupt
}
fi.Metadata = cloneMSS(fi.Metadata)
if kind == "old" {
delete(fi.Metadata, multipartMetaBucket)
} else {
fi.Metadata[multipartMetaBucket] = "another-bucket"
}
}
return fi, err
}}
}
return disks
}
got, err := z.ListMultipartUploads(t.Context(), bucket, "", "", "", "", 10)
if err != nil {
t.Fatal(err)
}
requireMultipartUploadKeys(t, got, "a")
})
}
}
func TestMultipartAbortPoolsAndRetry(t *testing.T) {
z, bucket := consistencyPools(t)
mp, err := z.serverPools[1].NewMultipartUpload(t.Context(), bucket, "a", ObjectOptions{})
if err != nil {
t.Fatal(err)
}
emptySet := z.serverPools[0].getHashedSet("a")
original := emptySet.getDisks
t.Cleanup(func() { emptySet.getDisks = original })
emptySet.getDisks = func() []StorageAPI {
disks := append([]StorageAPI(nil), original()...)
for i, d := range disks {
disks[i] = multipartListingFaultDisk{StorageAPI: d, read: func(context.Context, string, string, string, string, ReadOptions) (FileInfo, error) {
return FileInfo{}, errDiskNotFound
}}
}
return disks
}
err = z.AbortMultipartUpload(t.Context(), bucket, "a", mp.UploadID, ObjectOptions{})
if err == nil || toAPIError(t.Context(), err).HTTPStatusCode != 503 {
t.Fatalf("unknown pool must not acknowledge cancellation: %v", err)
}
emptySet.getDisks = original
err = z.AbortMultipartUpload(t.Context(), bucket, "a", mp.UploadID, ObjectOptions{})
if _, ok := err.(InvalidUploadID); !ok {
t.Fatalf("retry after all pools confirm absence: %v", err)
}
mp, err = z.serverPools[1].NewMultipartUpload(t.Context(), bucket, "a", ObjectOptions{})
if err != nil {
t.Fatal(err)
}
if err = z.AbortMultipartUpload(t.Context(), bucket, "a", mp.UploadID, ObjectOptions{}); err != nil {
t.Fatal(err)
}
if _, err = z.serverPools[1].GetMultipartInfo(t.Context(), bucket, "a", mp.UploadID, ObjectOptions{}); err == nil {
t.Fatal("empty first pool hid the actual upload")
}
}
func TestMultipartAbortRetryBelowReadQuorum(t *testing.T) {
z, set, bucket := multipartListingFixture(t)
mp, err := z.NewMultipartUpload(t.Context(), bucket, "a", ObjectOptions{})
if err != nil {
t.Fatal(err)
}
original := set.getDisks
t.Cleanup(func() { set.getDisks = original })
set.getDisks = func() []StorageAPI {
disks := append([]StorageAPI(nil), original()...)
disks[3] = nil
d := disks[2]
disks[2] = multipartListingFaultDisk{StorageAPI: d, delete: func(ctx context.Context, v, p string, opts DeleteOptions) error {
if err := d.Delete(ctx, v, p, opts); err != nil {
return err
}
return context.DeadlineExceeded // operation finished, acknowledgement lost
}}
return disks
}
if err = z.AbortMultipartUpload(t.Context(), bucket, "a", mp.UploadID, ObjectOptions{}); err == nil {
t.Fatal("lost acknowledgement must fail")
}
set.getDisks = original
if err = z.AbortMultipartUpload(t.Context(), bucket, "a", mp.UploadID, ObjectOptions{}); err != nil {
t.Fatalf("one remaining metadata copy prevented retry: %v", err)
}
}
func TestMultipartListingMarkerHTTP(t *testing.T) {
z, _ := consistencyPools(t)
bucket, router, err := initAPIHandlerTest(t.Context(), z, []string{"ListMultipartUploads"}, MakeBucketOptions{})
if err != nil {
t.Fatal(err)
}
for _, tc := range []struct {
key, marker string
status int
}{
{"", "not-base64=", 200},
{"a", "not-base64=", 404},
{"a", base64.RawURLEncoding.EncodeToString([]byte("not-native")), 400},
{"a", multipartListingTestID(time.Unix(100, 0), 1), 200},
} {
u := getListMultipartUploadsURLWithParams("", bucket, "", tc.key, tc.marker, "", "10")
req, err := newTestSignedRequestV4(http.MethodGet, u, 0, nil, globalActiveCred.AccessKey, globalActiveCred.SecretKey, nil)
if err != nil {
t.Fatal(err)
}
rec := httptest.NewRecorder()
router.ServeHTTP(rec, req)
if rec.Code != tc.status {
t.Fatalf("key=%q marker=%q: %d %s", tc.key, tc.marker, rec.Code, rec.Body.String())
}
}
}
func TestMultipartPreflightAdminHTTP(t *testing.T) {
bed, err := prepareAdminErasureTestBed(t.Context())
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { bed.done(); bed.objLayer.Shutdown(context.Background()); removeRoots(bed.erasureDirs) })
const target = "/minio/admin/v3/multipart-preflight"
rec := httptest.NewRecorder()
bed.router.ServeHTTP(rec, httptest.NewRequest(http.MethodGet, target, nil))
if rec.Code != 403 {
t.Fatalf("anonymous preflight: %d", rec.Code)
}
req, err := newTestSignedRequestV4(http.MethodGet, target, 0, nil, globalActiveCred.AccessKey, globalActiveCred.SecretKey, nil)
if err != nil {
t.Fatal(err)
}
rec = httptest.NewRecorder()
bed.router.ServeHTTP(rec, req)
if rec.Code != 200 {
t.Fatalf("admin preflight: %d %s", rec.Code, rec.Body.String())
}
var report multipartPreflightReport
if err = json.Unmarshal(rec.Body.Bytes(), &report); err != nil || !report.Ready || !report.Complete {
t.Fatalf("preflight: %+v %v", report, err)
}
for range cap(multipartScanSlots) {
scan, err := startMultipartScan(t.Context(), false)
if err != nil {
t.Fatal(err)
}
defer scan.close()
}
rec = httptest.NewRecorder()
bed.router.ServeHTTP(rec, req)
if rec.Code != http.StatusServiceUnavailable {
t.Fatalf("busy admin preflight: %d %s", rec.Code, rec.Body.String())
}
}
type multipartLateCreateDisk struct {
StorageAPI
late chan<- func() error
}
func (d multipartLateCreateDisk) WriteMetadata(_ context.Context, original, volume, object string, fi FileInfo) error {
if volume == minioMetaMultipartBucket {
d.late <- func() error { return d.StorageAPI.WriteMetadata(context.Background(), original, volume, object, fi) }
return context.DeadlineExceeded
}
return d.StorageAPI.WriteMetadata(context.Background(), original, volume, object, fi)
}
// This characterization is deliberately NOT an assertion of terminal abort
// correctness. It preserves an executable example of the separately scoped
// late-creation-write limitation. Change the expectation when fencing is added.
func TestMultipartAbortLateCreateBoundary(t *testing.T) {
obj, dirs, err := prepareErasure(t.Context(), 16)
if err != nil {
t.Fatal(err)
}
z := obj.(*erasureServerPools)
t.Cleanup(func() { z.Shutdown(context.Background()); removeRoots(dirs) })
saved := globalStorageClass
globalStorageClass.Update(storageclass.Config{Standard: storageclass.StorageClass{Parity: 8}})
t.Cleanup(func() { globalStorageClass.Update(saved) })
const bucket = "multipart-late-create"
if err = z.MakeBucket(t.Context(), bucket, MakeBucketOptions{}); err != nil {
t.Fatal(err)
}
set := z.serverPools[0].getHashedSet("a")
original := set.getDisks
t.Cleanup(func() { set.getDisks = original })
late := make(chan func() error, 16)
set.getDisks = func() []StorageAPI {
disks := append([]StorageAPI(nil), original()...)
for i := 9; i < 16; i++ {
disks[i] = multipartLateCreateDisk{disks[i], late}
}
return disks
}
mp, err := z.NewMultipartUpload(t.Context(), bucket, "a", ObjectOptions{})
if err != nil {
t.Fatal(err)
}
if len(late) != 7 {
t.Fatalf("expected seven timed-out creation writes, got %d", len(late))
}
set.getDisks = func() []StorageAPI {
disks := append([]StorageAPI(nil), original()...)
for i := 2; i < 9; i++ {
disks[i] = nil
}
return disks
}
if err = z.AbortMultipartUpload(t.Context(), bucket, "a", mp.UploadID, ObjectOptions{}); err != nil {
t.Fatal(err)
}
for range 7 {
if err = (<-late)(); err != nil {
t.Fatal(err)
}
}
set.getDisks = original
_, metadata, err := set.checkUploadIDExists(t.Context(), bucket, "a", mp.UploadID, true)
if err != nil {
t.Fatalf("late creation boundary changed; review and remove the documented limitation: %v", err)
}
count := 0
for _, fi := range metadata {
if fi.IsValid() {
count++
}
}
if count != 14 {
t.Fatalf("expected fourteen resurrected copies, got %d", count)
}
t.Log("KNOWN UNRESOLVED BOUNDARY: seven delayed creation writes plus seven offline copies restore a writable upload after acknowledged cancellation")
}
type multipartListingCountingDisk struct {
StorageAPI
directoryCalls *atomic.Int64
metadataCalls *atomic.Int64
}
func (d multipartListingCountingDisk) ListDir(ctx context.Context, original, volume, dir string, count int) ([]string, error) {
if volume == minioMetaMultipartBucket {
d.directoryCalls.Add(1)
}
return d.StorageAPI.ListDir(ctx, original, volume, dir, count)
}
func (d multipartListingCountingDisk) ReadVersion(ctx context.Context, original, volume, object, version string, opts ReadOptions) (FileInfo, error) {
if volume == minioMetaMultipartBucket {
d.metadataCalls.Add(1)
}
return d.StorageAPI.ReadVersion(ctx, original, volume, object, version, opts)
}
// Counts storage API calls; this is not a deployment throughput benchmark.
func TestMultipartListingScanCosts(t *testing.T) {
obj, dirs, err := prepareErasure(t.Context(), 4)
if err != nil {
t.Fatal(err)
}
z := obj.(*erasureServerPools)
t.Cleanup(func() { z.Shutdown(context.Background()); removeRoots(dirs) })
const bucket, otherBucket = "r9-scan-target", "r9-scan-unrelated"
for _, name := range []string{bucket, otherBucket} {
if err := z.MakeBucket(t.Context(), name, MakeBucketOptions{}); err != nil {
t.Fatal(err)
}
}
if _, err := z.NewMultipartUpload(t.Context(), bucket, "dir/target", ObjectOptions{}); err != nil {
t.Fatal(err)
}
var directoryCalls, metadataCalls atomic.Int64
for _, pool := range z.serverPools {
for _, set := range pool.sets {
original := set.getDisks
set.getDisks = func() []StorageAPI {
disks := original()
wrapped := make([]StorageAPI, len(disks))
for i, disk := range disks {
if disk != nil {
wrapped[i] = multipartListingCountingDisk{disk, &directoryCalls, &metadataCalls}
}
}
return wrapped
}
t.Cleanup(func() { set.getDisks = original })
}
}
measure := func(label string) (int64, int64) {
t.Helper()
directoryCalls.Store(0)
metadataCalls.Store(0)
result, err := z.ListMultipartUploads(t.Context(), bucket, "dir/", "", "", "", 1)
if err != nil {
t.Fatal(err)
}
requireMultipartUploadKeys(t, result, "dir/target")
d, m := directoryCalls.Load(), metadataCalls.Load()
t.Logf("%s: max-uploads=1 returned=%d ListDir=%d ReadVersion=%d", label, len(result.Uploads), d, m)
return d, m
}
_, baseline := measure("no unrelated uploads")
for n := 0; n < 32; n++ {
if _, err := z.NewMultipartUpload(t.Context(), otherBucket, fmt.Sprintf("other/%03d", n), ObjectOptions{}); err != nil {
t.Fatal(err)
}
}
_, loaded := measure("32 uploads in another bucket")
_, repeated := measure("same query repeated")
if loaded != baseline+32 || repeated != loaded {
t.Fatalf("unexpected scan accounting: baseline=%d loaded=%d repeated=%d", baseline, loaded, repeated)
}
}
func multipartListingTestID(created time.Time, n int) string {
return base64.RawURLEncoding.EncodeToString([]byte(fmt.Sprintf("00000000-0000-4000-8000-000000000000.00000000-0000-4000-8000-%012xx%d", n, created.UnixNano())))
}
func multipartListingFixture(t *testing.T) (*erasureServerPools, *erasureObjects, string) {
t.Helper()
obj, dirs, err := prepareErasure(t.Context(), 4)
if err != nil {
t.Fatal(err)
}
z := obj.(*erasureServerPools)
t.Cleanup(func() { z.Shutdown(context.Background()); removeRoots(dirs) })
const bucket = "multipart-listing-test"
if err := z.MakeBucket(t.Context(), bucket, MakeBucketOptions{}); err != nil {
t.Fatal(err)
}
return z, z.serverPools[0].getHashedSet("a"), bucket
}
func TestMultipartListingMissingMarkers(t *testing.T) {
base := time.Unix(100, 0)
sameTime := []string{multipartListingTestID(base.Add(time.Second), 1), multipartListingTestID(base.Add(time.Second), 2), multipartListingTestID(base.Add(time.Second), 3)}
slices.Sort(sameTime)
uploads := []MultipartInfo{
{Bucket: "bucket", Object: "a", UploadID: multipartListingTestID(base, 1), Initiated: base},
{Bucket: "bucket", Object: "a", UploadID: sameTime[1], Initiated: base.Add(time.Second)},
{Bucket: "bucket", Object: "b", UploadID: multipartListingTestID(base, 4), Initiated: base},
}
for _, tc := range []struct {
name string
created time.Time
id int
want []string
}{
{"before", base.Add(-time.Second), 0, []string{"a", "a", "b"}},
{"between", base.Add(time.Second / 2), 2, []string{"a", "b"}},
{"same-time-before", base.Add(time.Second), 2, []string{"a", "b"}},
{"same-time-after", base.Add(time.Second), 4, []string{"b"}},
{"after", base.Add(2 * time.Second), 5, []string{"b"}},
} {
t.Run(tc.name, func(t *testing.T) {
marker := multipartListingTestID(tc.created, tc.id)
if tc.name == "same-time-before" {
marker = sameTime[0]
}
if tc.name == "same-time-after" {
marker = sameTime[2]
}
if err := checkListMultipartArgs(t.Context(), "bucket", "", "a", marker, ""); err != nil {
t.Fatal(err)
}
got := paginateMultipartUploads(uploads, "", "a", marker, "", 10)
if !slices.Equal(multipartUploadKeys(got.Uploads), tc.want) {
t.Fatalf("got %v want %v", multipartUploadKeys(got.Uploads), tc.want)
}
})
}
}
func TestMultipartListingAbortRecovery(t *testing.T) {
for _, offline := range []int{1, 2} {
t.Run(fmt.Sprint(offline), func(t *testing.T) {
z, set, bucket := multipartListingFixture(t)
mp, err := z.NewMultipartUpload(t.Context(), bucket, "a", ObjectOptions{})
if err != nil {
t.Fatal(err)
}
original := set.getDisks
t.Cleanup(func() { set.getDisks = original })
set.getDisks = func() []StorageAPI {
disks := append([]StorageAPI(nil), original()...)
for i := 0; i < offline; i++ {
disks[i] = nil
}
return disks
}
err = z.AbortMultipartUpload(t.Context(), bucket, "a", mp.UploadID, ObjectOptions{})
if offline == 1 && err != nil {
t.Fatal(err)
}
if offline == 2 && (err == nil || toAPIError(t.Context(), err).HTTPStatusCode != 503) {
t.Fatalf("two offline drives must not acknowledge cancellation: %v", err)
}
set.getDisks = original
if offline == 2 {
// Retry a partially completed deletion after the original disks return.
if err = z.AbortMultipartUpload(t.Context(), bucket, "a", mp.UploadID, ObjectOptions{}); err != nil {
t.Fatal(err)
}
}
got, err := z.ListMultipartUploads(t.Context(), bucket, "", "", "", "", 100)
if err != nil || len(got.Uploads) != 0 {
t.Fatalf("canceled upload reappeared: %+v %v", got, err)
}
})
}
}
func TestMultipartListingCoverage(t *testing.T) {
z, set, bucket := multipartListingFixture(t)
mp, err := z.NewMultipartUpload(t.Context(), bucket, "a", ObjectOptions{})
if err != nil {
t.Fatal(err)
}
original := set.getDisks
t.Cleanup(func() { set.getDisks = original })
disks := original()
// Leave a readable two-copy record, simulating a partially failed abort.
for _, d := range disks[:2] {
if err = d.Delete(t.Context(), minioMetaMultipartBucket, set.getUploadIDDir(bucket, "a", mp.UploadID), DeleteOptions{Recursive: true}); err != nil {
t.Fatal(err)
}
}
set.getDisks = func() []StorageAPI { return []StorageAPI{disks[0], disks[1], nil, nil} }
_, err = z.ListMultipartUploads(t.Context(), bucket, "", "", "", "", 10)
if err == nil || toAPIError(t.Context(), err).HTTPStatusCode != 503 {
t.Fatalf("incomplete discovery returned success: %v", err)
}
set.getDisks = func() []StorageAPI { return []StorageAPI{nil, disks[1], disks[2], disks[3]} }
got, err := z.ListMultipartUploads(t.Context(), bucket, "", "", "", "", 10)
if err != nil {
t.Fatal(err)
}
requireMultipartUploadKeys(t, got, "a")
report, err := z.multipartPreflight(t.Context())
if err != nil || report.Ready || report.Complete || len(report.Sets[0].UncoveredDrives) != 1 {
t.Fatalf("offline upgrade preflight: %+v %v", report, err)
}
}
func TestMultipartListingBudgetAndAdmission(t *testing.T) {
_, set, bucket := multipartListingFixture(t)
if _, err := set.NewMultipartUpload(t.Context(), bucket, "a", ObjectOptions{}); err != nil {
t.Fatal(err)
}
scan, err := startMultipartScan(t.Context(), false)
if err != nil {
t.Fatal(err)
}
defer scan.close()
scan.remaining.Store(1)
_, _, err = set.scanMultipartUploads(scan, bucket, 0, 0)
var limited SlowDown
if !errors.As(err, &limited) {
t.Fatalf("budget: %v", err)
}
if apiErr := toAPIError(t.Context(), err); apiErr.HTTPStatusCode != http.StatusServiceUnavailable || apiErr.Code != "SlowDown" {
t.Fatalf("budget error mapping: %+v", apiErr)
}
second, err := startMultipartScan(t.Context(), false)
if err != nil {
t.Fatal(err)
}
defer second.close()
if third, err := startMultipartScan(t.Context(), false); err == nil {
third.close()
t.Fatal("third scan admitted")
}
}
func TestMultipartListingAdmissionHTTP(t *testing.T) {
z, _, _ := multipartListingFixture(t)
bucket, router, err := initAPIHandlerTest(t.Context(), z, []string{"ListMultipartUploads"}, MakeBucketOptions{})
if err != nil {
t.Fatal(err)
}
for range cap(multipartScanSlots) {
scan, err := startMultipartScan(t.Context(), false)
if err != nil {
t.Fatal(err)
}
defer scan.close()
}
req, err := newTestSignedRequestV4(http.MethodGet,
getListMultipartUploadsURLWithParams("", bucket, "", "", "", "", "1"),
0, nil, globalActiveCred.AccessKey, globalActiveCred.SecretKey, nil)
if err != nil {
t.Fatal(err)
}
rec := httptest.NewRecorder()
router.ServeHTTP(rec, req)
var response APIErrorResponse
if err := xml.Unmarshal(rec.Body.Bytes(), &response); err != nil {
t.Fatal(err)
}
if rec.Code != http.StatusServiceUnavailable || response.Code != "SlowDown" {
t.Fatalf("admission returned %d %s", rec.Code, rec.Body.String())
}
}
func TestMultipartListingPreflightMinorityLegacy(t *testing.T) {
z, set, bucket := multipartListingFixture(t)
mp, err := z.NewMultipartUpload(t.Context(), bucket, "old", ObjectOptions{})
if err != nil {
t.Fatal(err)
}
path := set.getUploadIDDir(bucket, "old", mp.UploadID)
disks := set.getDisks()
fi, err := disks[0].ReadVersion(t.Context(), bucket, minioMetaMultipartBucket, path, "", ReadOptions{})
if err != nil {
t.Fatal(err)
}
delete(fi.Metadata, multipartMetaBucket)
delete(fi.Metadata, multipartMetaObject)
if err = disks[0].WriteMetadata(t.Context(), bucket, minioMetaMultipartBucket, path, fi); err != nil {
t.Fatal(err)
}
for _, d := range disks[1:] {
if err = d.Delete(t.Context(), minioMetaMultipartBucket, path, DeleteOptions{Recursive: true}); err != nil {
t.Fatal(err)
}
}
report, err := z.multipartPreflight(t.Context())
if err != nil || report.Ready || !report.Complete || report.LegacyUploads != 1 {
t.Fatalf("minority legacy copy must prevent readiness: %+v %v", report, err)
}
}
func TestMultipartListingCancellationRetainsAdmission(t *testing.T) {
z, set, bucket := multipartListingFixture(t)
for i := range 40 {
if _, err := z.NewMultipartUpload(t.Context(), bucket, fmt.Sprintf("key-%02d", i), ObjectOptions{}); err != nil {
t.Fatal(err)
}
}
original := set.getDisks
t.Cleanup(func() { set.getDisks = original })
entered := make(chan struct{}, 40)
release := make(chan struct{})
var once sync.Once
defer once.Do(func() { close(release) })
var reads atomic.Int32
set.getDisks = func() []StorageAPI {
disks := append([]StorageAPI(nil), original()...)
for i, d := range disks {
disks[i] = multipartListingFaultDisk{StorageAPI: d, read: func(ctx context.Context, b, v, p, version string, opts ReadOptions) (FileInfo, error) {
if v != minioMetaMultipartBucket {
return d.ReadVersion(ctx, b, v, p, version, opts)
}
reads.Add(1)
entered <- struct{}{}
// Model an RPC which does not return immediately on cancellation.
<-release
return FileInfo{}, ctx.Err()
}}
}
return disks
}
second, err := startMultipartScan(t.Context(), false)
if err != nil {
t.Fatal(err)
}
defer second.close()
ctx, cancel := context.WithCancel(t.Context())
defer cancel()
done := make(chan error, 1)
go func() {
_, err := z.ListMultipartUploads(ctx, bucket, "", "", "", "", 10)
done <- err
}()
for range 16 {
select {
case <-entered:
case <-time.After(5 * time.Second):
t.Fatal("identity workers did not start")
}
}
cancel()
if extra, err := startMultipartScan(t.Context(), false); err == nil {
extra.close()
t.Fatal("canceled request released admission while RPCs were still running")
}
start := time.Now()
once.Do(func() { close(release) })
select {
case err := <-done:
if err == nil {
t.Fatal("canceled scan returned a successful partial list")
}
case <-time.After(2 * time.Second):
t.Fatal("scan did not stop after blocked RPCs returned")
}
if got := reads.Load(); got != 16 {
t.Fatalf("scheduled more identity/metadata reads after cancellation: %d", got)
}
t.Logf("16 identity RPCs bounded; admission held until return; cancellation settled in %s", time.Since(start))
}
+246 -35
View File
@@ -31,6 +31,7 @@ import (
"sync"
"time"
"github.com/google/uuid"
"github.com/klauspost/readahead"
"github.com/minio/minio-go/v7/pkg/set"
"github.com/minio/minio/internal/config/storageclass"
@@ -44,6 +45,15 @@ import (
"github.com/pgsty/silo-pkg/v3/sync/errgroup"
)
const (
multipartMetaBucket = ReservedMetadataPrefixLower + "multipart-v1-bucket"
multipartMetaObject = ReservedMetadataPrefixLower + "multipart-v1-object"
// ponytail: keep scan concurrency fixed until the multipart-list benchmark
// establishes a better adaptive limit.
multipartMetadataScanConcurrency = 4
)
func (er erasureObjects) getUploadIDDir(bucket, object, uploadID string) string {
uploadUUID := uploadID
uploadBytes, err := base64.RawURLEncoding.DecodeString(uploadID)
@@ -251,14 +261,152 @@ func (er erasureObjects) cleanupStaleUploadsOnDisk(ctx context.Context, disk Sto
})
}
// ListMultipartUploads - lists all the pending multipart
// uploads for a particular object in a bucket.
//
// Implements minimal S3 compatible ListMultipartUploads API. We do
// not support prefix based listing, this is a deliberate attempt
// towards simplification of multipart APIs.
// The resulting ListMultipartsInfo structure is unmarshalled directly as XML.
func (er erasureObjects) ListMultipartUploads(ctx context.Context, bucket, object, keyMarker, uploadIDMarker, delimiter string, maxUploads int) (result ListMultipartsInfo, err error) {
func multipartUploadInfo(bucket, object, uploadUUID string, fallback time.Time) MultipartInfo {
initiated := fallback
if parsed, ok := multipartUploadTime(uploadUUID); ok {
initiated = parsed
}
return MultipartInfo{
Bucket: bucket,
Object: object,
UploadID: base64.RawURLEncoding.EncodeToString(fmt.Appendf(nil, "%s.%s", globalDeploymentID(), uploadUUID)),
Initiated: initiated,
}
}
// multipartUploadTime is shared by stored records and continuation markers.
// The time is part of the immutable upload ID, so a removed marker still
// identifies the same ordering boundary.
func multipartUploadTime(uploadUUID string) (time.Time, bool) {
if len(uploadUUID) < 38 || uploadUUID[36] != 'x' {
return time.Time{}, false
}
if _, err := uuid.Parse(uploadUUID[:36]); err != nil {
return time.Time{}, false
}
ns, err := strconv.ParseInt(uploadUUID[37:], 10, 64)
if err != nil || ns <= 0 || strconv.FormatInt(ns, 10) != uploadUUID[37:] {
return time.Time{}, false
}
return time.Unix(0, ns), true
}
func multipartMarkerTime(uploadID string) (time.Time, bool) {
b, err := base64.RawURLEncoding.DecodeString(uploadID)
if err != nil {
return time.Time{}, false
}
_, uploadUUID, ok := strings.Cut(string(b), ".")
if !ok {
return time.Time{}, false
}
return multipartUploadTime(uploadUUID)
}
type multipartListEntry struct {
upload *MultipartInfo
commonPrefix string
}
// paginateMultipartUploads applies the S3 ordering, prefix, delimiter, marker,
// and page rules exactly once after all pools and sets have been merged.
func paginateMultipartUploads(uploads []MultipartInfo, prefix, keyMarker, uploadIDMarker, delimiter string, maxUploads int) ListMultipartsInfo {
if maxUploads > maxUploadsList {
maxUploads = maxUploadsList
}
result := ListMultipartsInfo{
MaxUploads: maxUploads,
KeyMarker: keyMarker,
UploadIDMarker: uploadIDMarker,
Prefix: prefix,
Delimiter: delimiter,
}
deduplicated := make([]MultipartInfo, 0, len(uploads))
seenUploads := make(map[string]struct{}, len(uploads))
for _, upload := range uploads {
identity := upload.Bucket + "\x00" + upload.Object + "\x00" + upload.UploadID
if _, ok := seenUploads[identity]; ok {
continue
}
seenUploads[identity] = struct{}{}
deduplicated = append(deduplicated, upload)
}
sort.Slice(deduplicated, func(i, j int) bool {
if deduplicated[i].Object != deduplicated[j].Object {
return deduplicated[i].Object < deduplicated[j].Object
}
if !deduplicated[i].Initiated.Equal(deduplicated[j].Initiated) {
return deduplicated[i].Initiated.Before(deduplicated[j].Initiated)
}
return deduplicated[i].UploadID < deduplicated[j].UploadID
})
markerTime, _ := multipartMarkerTime(uploadIDMarker)
seenPrefixes := make(map[string]struct{})
entries := make([]multipartListEntry, 0, len(deduplicated))
for i := range deduplicated {
upload := &deduplicated[i]
if !strings.HasPrefix(upload.Object, prefix) {
continue
}
if keyMarker != "" {
switch strings.Compare(upload.Object, keyMarker) {
case -1:
continue
case 0:
if uploadIDMarker == "" || upload.Initiated.Before(markerTime) ||
(upload.Initiated.Equal(markerTime) && upload.UploadID <= uploadIDMarker) {
continue
}
}
}
if delimiter != "" {
remainder := strings.TrimPrefix(upload.Object, prefix)
if i := strings.Index(remainder, delimiter); i >= 0 {
commonPrefix := prefix + remainder[:i+len(delimiter)]
if keyMarker != "" && commonPrefix <= keyMarker {
continue
}
if _, ok := seenPrefixes[commonPrefix]; ok {
continue
}
seenPrefixes[commonPrefix] = struct{}{}
entries = append(entries, multipartListEntry{commonPrefix: commonPrefix})
continue
}
}
entries = append(entries, multipartListEntry{upload: upload})
}
if maxUploads <= 0 {
return result
}
pageSize := min(maxUploads, len(entries))
for _, entry := range entries[:pageSize] {
if entry.upload != nil {
result.Uploads = append(result.Uploads, *entry.upload)
continue
}
result.CommonPrefixes = append(result.CommonPrefixes, entry.commonPrefix)
}
result.IsTruncated = pageSize < len(entries)
if result.IsTruncated && pageSize > 0 {
last := entries[pageSize-1]
if last.upload != nil {
result.NextKeyMarker = last.upload.Object
result.NextUploadIDMarker = last.upload.UploadID
} else {
result.NextKeyMarker = last.commonPrefix
}
}
return result
}
// listMultipartUploadsExact preserves the hashed exact-object lookup used by
// multipart write placement and by rolling-upgrade legacy mode.
func (er erasureObjects) listMultipartUploadsExact(ctx context.Context, bucket, object, keyMarker, uploadIDMarker, delimiter string, maxUploads int) (result ListMultipartsInfo, err error) {
auditObjectErasureSet(ctx, "ListMultipartUploads", object, &er)
result.MaxUploads = maxUploads
@@ -311,20 +459,16 @@ func (er erasureObjects) ListMultipartUploads(ctx context.Context, bucket, objec
if populatedUploadIDs.Contains(uploadID) {
continue
}
// If present, use time stored in ID.
startTime := time.Now()
if split := strings.Split(uploadID, "x"); len(split) == 2 {
t, err := strconv.ParseInt(split[1], 10, 64)
if err == nil {
startTime = time.Unix(0, t)
var fallback time.Time
if _, ok := multipartUploadTime(uploadID); !ok {
fi, err := disk.ReadVersion(ctx, bucket, minioMetaMultipartBucket,
pathJoin(er.getMultipartSHADir(bucket, object), uploadID), "", ReadOptions{})
if err != nil {
return result, toObjectErr(err, bucket, object)
}
fallback = fi.ModTime
}
uploads = append(uploads, MultipartInfo{
Bucket: bucket,
Object: object,
UploadID: base64.RawURLEncoding.EncodeToString(fmt.Appendf(nil, "%s.%s", globalDeploymentID(), uploadID)),
Initiated: startTime,
})
uploads = append(uploads, multipartUploadInfo(bucket, object, uploadID, fallback))
populatedUploadIDs.Add(uploadID)
}
@@ -365,6 +509,25 @@ func (er erasureObjects) ListMultipartUploads(ctx context.Context, bucket, objec
return result, nil
}
func (er erasureObjects) ListMultipartUploads(ctx context.Context, bucket, prefix, keyMarker, uploadIDMarker, delimiter string, maxUploads int) (ListMultipartsInfo, error) {
if err := checkListMultipartArgs(ctx, bucket, prefix, keyMarker, uploadIDMarker, delimiter); err != nil {
return ListMultipartsInfo{}, err
}
scan, err := startMultipartScan(ctx, false)
if err != nil {
return ListMultipartsInfo{}, err
}
defer scan.close()
uploads, legacy, err := er.scanMultipartUploads(scan, bucket, 0, 0)
if err != nil {
return ListMultipartsInfo{}, err
}
if legacy {
return ListMultipartsInfo{}, errMultipartListingLegacy
}
return paginateMultipartUploads(uploads, prefix, keyMarker, uploadIDMarker, delimiter, maxUploads), nil
}
// newMultipartUpload - wrapper for initializing a new multipart
// request; returns a unique upload id.
//
@@ -402,6 +565,8 @@ func (er erasureObjects) newMultipartUpload(ctx context.Context, bucket string,
}
userDefined := cloneMSS(opts.UserDefined)
userDefined[multipartMetaBucket] = bucket
userDefined[multipartMetaObject] = object
if opts.PreserveETag != "" {
userDefined["etag"] = opts.PreserveETag
}
@@ -1459,6 +1624,8 @@ func (er erasureObjects) CompleteMultipartUpload(ctx context.Context, bucket str
// Remove superfluous internal headers.
delete(fi.Metadata, hash.MinIOMultipartChecksum)
delete(fi.Metadata, hash.MinIOMultipartChecksumType)
delete(fi.Metadata, multipartMetaBucket)
delete(fi.Metadata, multipartMetaObject)
// Save the final object size and modtime.
fi.Size = objectSize
@@ -1586,22 +1753,66 @@ func (er erasureObjects) CompleteMultipartUpload(ctx context.Context, bucket str
return fi.ToObjectInfo(bucket, object, opts.Versioned || opts.VersionSuspended), nil
}
// AbortMultipartUpload - aborts an ongoing multipart operation
// signified by the input uploadID. This is an atomic operation
// doesn't require clients to initiate multiple such requests.
//
// All parts are purged from all disks and reference to the uploadID
// would be removed from the system, rollback is not possible on this
// operation.
func (er erasureObjects) AbortMultipartUpload(ctx context.Context, bucket, object, uploadID string, opts ObjectOptions) (err error) {
// abortMultipartUpload confirms absence on a strict majority of this set.
// Unlike an existence read followed by best-effort deletion, it also permits
// retrying a partial deletion which no longer has a readable metadata quorum.
// This does not fence creation writes still executing after a storage timeout.
func (er erasureObjects) abortMultipartUpload(ctx context.Context, bucket, object, uploadID string, opts ObjectOptions) (bool, error) {
if !opts.NoAuditLog {
auditObjectErasureSet(ctx, "AbortMultipartUpload", object, &er)
}
// Cleanup all uploaded parts.
defer er.deleteAll(ctx, minioMetaMultipartBucket, er.getUploadIDDir(bucket, object, uploadID))
// Validates if upload ID exists.
_, _, err = er.checkUploadIDExists(ctx, bucket, object, uploadID, false)
return toObjectErr(err, bucket, object, uploadID)
b, err := base64.RawURLEncoding.DecodeString(uploadID)
if err != nil {
return false, MalformedUploadID{UploadID: uploadID}
}
_, internalID, ok := strings.Cut(string(b), ".")
if !ok || internalID == "" || internalID == "." || internalID == ".." || strings.ContainsAny(internalID, "/\\") {
return false, InvalidUploadID{Bucket: bucket, Object: object, UploadID: uploadID}
}
disks := er.getDisks()
uploadPath := er.getUploadIDDir(bucket, object, uploadID)
_, errs := readAllFileInfo(ctx, disks, bucket, minioMetaMultipartBucket, uploadPath, "", false, false)
quorum := er.setDriveCount/2 + 1
found, absent := false, 0
for _, err := range errs {
switch {
case err == nil, errors.Is(err, errFileCorrupt):
found = true
case errors.Is(err, errFileNotFound), errors.Is(err, errFileVersionNotFound):
absent++
}
}
if absent >= quorum {
return found, nil
}
if !found {
return false, toObjectErr(errErasureReadQuorum, bucket, object, uploadID)
}
g := errgroup.WithNErrs(len(disks))
for i, disk := range disks {
g.Go(func() error {
if disk == nil {
return errDiskNotFound
}
err := disk.Delete(ctx, minioMetaMultipartBucket, uploadPath, DeleteOptions{Recursive: true, Immediate: false})
if errors.Is(err, errFileNotFound) || errors.Is(err, errFileVersionNotFound) {
return nil
}
return err
}, i)
}
return true, toObjectErr(reduceWriteQuorumErrs(ctx, g.Wait(), nil, quorum), bucket, object, uploadID)
}
// AbortMultipartUpload confirms logical cancellation. Offline part data may
// still need stale-upload cleanup after its drives return.
func (er erasureObjects) AbortMultipartUpload(ctx context.Context, bucket, object, uploadID string, opts ObjectOptions) (err error) {
found, err := er.abortMultipartUpload(ctx, bucket, object, uploadID, opts)
if err != nil {
return err
}
if !found {
return InvalidUploadID{Bucket: bucket, Object: object, UploadID: uploadID}
}
return nil
}
+54 -17
View File
@@ -1891,7 +1891,41 @@ func (z *erasureServerPools) ListMultipartUploads(ctx context.Context, bucket, p
if err := checkListMultipartArgs(ctx, bucket, prefix, keyMarker, uploadIDMarker, delimiter); err != nil {
return ListMultipartsInfo{}, err
}
if _, err := z.GetBucketInfo(ctx, bucket, BucketOptions{}); err != nil {
return ListMultipartsInfo{}, toObjectErr(err, bucket)
}
if globalAPIConfig.getMultipartListingLegacy() {
return z.listMultipartUploadsLegacy(ctx, bucket, prefix, keyMarker, uploadIDMarker, delimiter, maxUploads)
}
scan, err := startMultipartScan(ctx, false)
if err != nil {
return ListMultipartsInfo{}, err
}
defer scan.close()
var uploads []MultipartInfo
var keyless bool
for idx, pool := range z.serverPools {
if z.IsSuspended(idx) {
continue
}
poolUploads, poolKeyless, err := pool.scanMultipartUploads(scan, bucket, idx)
if err != nil {
return ListMultipartsInfo{}, err
}
uploads = append(uploads, poolUploads...)
keyless = keyless || poolKeyless
}
// The old format cannot be enumerated authoritatively. Migration mode is
// explicit: another bucket must never silently change this API's semantics.
if keyless {
return ListMultipartsInfo{}, errMultipartListingLegacy
}
return paginateMultipartUploads(uploads, prefix, keyMarker, uploadIDMarker, delimiter, maxUploads), nil
}
func (z *erasureServerPools) listMultipartUploadsLegacy(ctx context.Context, bucket, prefix, keyMarker, uploadIDMarker, delimiter string, maxUploads int) (ListMultipartsInfo, error) {
poolResult := ListMultipartsInfo{}
poolResult.MaxUploads = maxUploads
poolResult.KeyMarker = keyMarker
@@ -1917,15 +1951,14 @@ func (z *erasureServerPools) ListMultipartUploads(ctx context.Context, bucket, p
}
if z.SinglePool() {
return z.serverPools[0].ListMultipartUploads(ctx, bucket, prefix, keyMarker, uploadIDMarker, delimiter, maxUploads)
return z.serverPools[0].getHashedSet(prefix).listMultipartUploadsExact(ctx, bucket, prefix, keyMarker, uploadIDMarker, delimiter, maxUploads)
}
for idx, pool := range z.serverPools {
if z.IsSuspended(idx) {
continue
}
result, err := pool.ListMultipartUploads(ctx, bucket, prefix, keyMarker, uploadIDMarker,
delimiter, maxUploads)
result, err := pool.getHashedSet(prefix).listMultipartUploadsExact(ctx, bucket, prefix, keyMarker, uploadIDMarker, delimiter, maxUploads)
if err != nil {
return result, err
}
@@ -1961,7 +1994,7 @@ func (z *erasureServerPools) NewMultipartUpload(ctx context.Context, bucket, obj
continue
}
result, err := pool.ListMultipartUploads(ctx, bucket, object, "", "", "", maxUploadsList)
result, err := pool.listMultipartUploadsExact(ctx, bucket, object)
if err != nil {
return nil, err
}
@@ -2123,9 +2156,13 @@ func (z *erasureServerPools) AbortMultipartUpload(ctx context.Context, bucket, o
if err := checkAbortMultipartArgs(ctx, bucket, object, uploadID); err != nil {
return err
}
if _, err := z.GetBucketInfo(ctx, bucket, BucketOptions{}); err != nil {
return toObjectErr(err, bucket)
}
defer func() {
if err == nil {
_, absent := err.(InvalidUploadID)
if err == nil || absent {
z.mpCache.Delete(uploadID)
globalNotificationSys.DeleteUploadID(ctx, uploadID)
}
@@ -2139,23 +2176,23 @@ func (z *erasureServerPools) AbortMultipartUpload(ctx context.Context, bucket, o
ctx = lkctx.Context()
defer lk.Unlock(lkctx)
if z.SinglePool() {
return z.serverPools[0].AbortMultipartUpload(ctx, bucket, object, uploadID, opts)
}
found := false
var firstErr error
for idx, pool := range z.serverPools {
if z.IsSuspended(idx) {
continue
}
err := pool.AbortMultipartUpload(ctx, bucket, object, uploadID, opts)
if err == nil {
return nil
poolFound, err := pool.getHashedSet(object).abortMultipartUpload(ctx, bucket, object, uploadID, opts)
found = found || poolFound
if err != nil && firstErr == nil {
firstErr = err
}
if _, ok := err.(InvalidUploadID); ok {
// upload id not found move to next pool
continue
}
return err
}
if firstErr != nil {
return firstErr
}
if found {
return nil
}
return InvalidUploadID{
Bucket: bucket,
+34 -4
View File
@@ -880,10 +880,40 @@ func (s *erasureSets) CopyObject(ctx context.Context, srcBucket, srcObject, dstB
}
func (s *erasureSets) ListMultipartUploads(ctx context.Context, bucket, prefix, keyMarker, uploadIDMarker, delimiter string, maxUploads int) (result ListMultipartsInfo, err error) {
// In list multipart uploads we are going to treat input prefix as the object,
// this means that we are not supporting directory navigation.
set := s.getHashedSet(prefix)
return set.ListMultipartUploads(ctx, bucket, prefix, keyMarker, uploadIDMarker, delimiter, maxUploads)
if err := checkListMultipartArgs(ctx, bucket, prefix, keyMarker, uploadIDMarker, delimiter); err != nil {
return ListMultipartsInfo{}, err
}
scan, err := startMultipartScan(ctx, false)
if err != nil {
return ListMultipartsInfo{}, err
}
defer scan.close()
uploads, legacy, err := s.scanMultipartUploads(scan, bucket, 0)
if err != nil {
return ListMultipartsInfo{}, err
}
if legacy {
return ListMultipartsInfo{}, errMultipartListingLegacy
}
return paginateMultipartUploads(uploads, prefix, keyMarker, uploadIDMarker, delimiter, maxUploads), nil
}
func (s *erasureSets) scanMultipartUploads(scan *multipartScan, bucket string, poolIdx int) ([]MultipartInfo, bool, error) {
var uploads []MultipartInfo
var keyless bool
for i, set := range s.sets {
setUploads, setKeyless, err := set.scanMultipartUploads(scan, bucket, poolIdx, i)
if err != nil {
return nil, false, err
}
uploads = append(uploads, setUploads...)
keyless = keyless || setKeyless
}
return uploads, keyless, nil
}
func (s *erasureSets) listMultipartUploadsExact(ctx context.Context, bucket, object string) (ListMultipartsInfo, error) {
return s.getHashedSet(object).listMultipartUploadsExact(ctx, bucket, object, "", "", "", maxUploadsList)
}
// Initiate a new multipart upload on a hashedSet based on object name.
+8
View File
@@ -50,6 +50,7 @@ type apiConfig struct {
transitionWorkers int
staleUploadsExpiry time.Duration
multipartListingLegacy bool
staleUploadsCleanupInterval time.Duration
deleteCleanupInterval time.Duration
enableODirect bool
@@ -181,6 +182,7 @@ func (t *apiConfig) init(cfg api.Config, setDriveCounts []int, legacy bool) {
t.transitionWorkers = cfg.TransitionWorkers
t.staleUploadsExpiry = cfg.StaleUploadsExpiry
t.multipartListingLegacy = cfg.MultipartListing == "legacy"
t.deleteCleanupInterval = cfg.DeleteCleanupInterval
t.enableODirect = cfg.EnableODirect
t.gzipObjects = cfg.GzipObjects
@@ -206,6 +208,12 @@ func (t *apiConfig) odirectEnabled() bool {
return t.enableODirect
}
func (t *apiConfig) getMultipartListingLegacy() bool {
t.mu.RLock()
defer t.mu.RUnlock()
return t.multipartListingLegacy
}
func (t *apiConfig) shouldGzipObjects() bool {
t.mu.RLock()
defer t.mu.RUnlock()
+298
View File
@@ -0,0 +1,298 @@
// Copyright (c) 2026 mr javad seydi
//
// This file is part of Silo Object Storage stack
//
// This program is free software: you can redistribute it and/or modify
// it under the terms of the GNU Affero General Public License as published by
// the Free Software Foundation, either version 3 of the License, or
// (at your option) any later version.
//
// This program is distributed in the hope that it will be useful
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Affero General Public License for more details.
//
// You should have received a copy of the GNU Affero General Public License
// along with this program. If not, see <http://www.gnu.org/licenses/>.
package cmd
import (
"bytes"
"errors"
"fmt"
"slices"
"testing"
"time"
)
func multipartUploadKeys(uploads []MultipartInfo) []string {
keys := make([]string, len(uploads))
for i := range uploads {
keys[i] = uploads[i].Object
}
return keys
}
func requireMultipartUploadKeys(t *testing.T, got ListMultipartsInfo, want ...string) {
t.Helper()
if keys := multipartUploadKeys(got.Uploads); !slices.Equal(keys, want) {
t.Fatalf("uploads = %v, want %v", keys, want)
}
}
func TestListMultipartUploadsS3Compatibility(t *testing.T) {
obj, dirs, err := prepareErasureSets32(t.Context())
if err != nil {
t.Fatal(err)
}
z := obj.(*erasureServerPools)
t.Cleanup(func() {
z.Shutdown(t.Context())
removeRoots(dirs)
})
const bucket = "multipart-list-compat"
if err = z.MakeBucket(t.Context(), bucket, MakeBucketOptions{}); err != nil {
t.Fatal(err)
}
objects := []string{"t/a_b/p1", "t/a_b/p2", "t/c_d/p1", "u/x"}
sets := z.serverPools[0]
firstSet := sets.getHashedSetIndex(objects[0])
if !slices.ContainsFunc(objects[1:], func(object string) bool {
return sets.getHashedSetIndex(object) != firstSet
}) {
for n := 0; ; n++ {
object := fmt.Sprintf("v/cross-set-%d", n)
if sets.getHashedSetIndex(object) != firstSet {
objects = append(objects, object)
break
}
}
}
uploadIDs := make(map[string]string, len(objects))
for _, object := range objects {
mp, err := z.NewMultipartUpload(t.Context(), bucket, object, ObjectOptions{})
if err != nil {
t.Fatalf("NewMultipartUpload(%q): %v", object, err)
}
uploadIDs[object] = mp.UploadID
}
// Durable multipart metadata, rather than this node-local cache, must be
// authoritative after a restart or when another node handles the request.
z.mpCache.Range(func(uploadID string, _ MultipartInfo) bool {
z.mpCache.Delete(uploadID)
return true
})
all, err := z.ListMultipartUploads(t.Context(), bucket, "", "", "", "", 100)
if err != nil {
t.Fatal(err)
}
requireMultipartUploadKeys(t, all, objects...)
if all.IsTruncated || all.NextKeyMarker != "" || all.NextUploadIDMarker != "" {
t.Fatalf("complete listing has truncation state: %+v", all)
}
first, err := z.ListMultipartUploads(t.Context(), bucket, "", "", "", "", 1)
if err != nil {
t.Fatal(err)
}
requireMultipartUploadKeys(t, first, objects[0])
if !first.IsTruncated || first.NextKeyMarker != objects[0] || first.NextUploadIDMarker != uploadIDs[objects[0]] {
t.Fatalf("first page markers = (%q, %q, %t), want (%q, %q, true)",
first.NextKeyMarker, first.NextUploadIDMarker, first.IsTruncated,
objects[0], uploadIDs[objects[0]])
}
rest, err := z.ListMultipartUploads(t.Context(), bucket, "", first.NextKeyMarker, first.NextUploadIDMarker, "", 100)
if err != nil {
t.Fatal(err)
}
requireMultipartUploadKeys(t, rest, objects[1:]...)
afterKey, err := z.ListMultipartUploads(t.Context(), bucket, "", objects[1], "", "", 100)
if err != nil {
t.Fatal(err)
}
requireMultipartUploadKeys(t, afterKey, objects[2:]...)
prefixed, err := z.ListMultipartUploads(t.Context(), bucket, "t/", "", "", "", 100)
if err != nil {
t.Fatal(err)
}
requireMultipartUploadKeys(t, prefixed, objects[:3]...)
nested, err := z.ListMultipartUploads(t.Context(), bucket, "t/a_b/", "", "", "", 100)
if err != nil {
t.Fatal(err)
}
requireMultipartUploadKeys(t, nested, objects[:2]...)
grouped, err := z.ListMultipartUploads(t.Context(), bucket, "t/", "", "", SlashSeparator, 100)
if err != nil {
t.Fatal(err)
}
requireMultipartUploadKeys(t, grouped)
if want := []string{"t/a_b/", "t/c_d/"}; !slices.Equal(grouped.CommonPrefixes, want) {
t.Fatalf("common prefixes = %v, want %v", grouped.CommonPrefixes, want)
}
groupPage, err := z.ListMultipartUploads(t.Context(), bucket, "t/", "", "", SlashSeparator, 1)
if err != nil {
t.Fatal(err)
}
if want := []string{"t/a_b/"}; !slices.Equal(groupPage.CommonPrefixes, want) {
t.Fatalf("first common-prefix page = %v, want %v", groupPage.CommonPrefixes, want)
}
if !groupPage.IsTruncated || groupPage.NextKeyMarker != "t/a_b/" || groupPage.NextUploadIDMarker != "" {
t.Fatalf("common-prefix page markers = (%q, %q, %t)",
groupPage.NextKeyMarker, groupPage.NextUploadIDMarker, groupPage.IsTruncated)
}
groupRest, err := z.ListMultipartUploads(t.Context(), bucket, "t/", groupPage.NextKeyMarker, "", SlashSeparator, 1)
if err != nil {
t.Fatal(err)
}
if want := []string{"t/c_d/"}; !slices.Equal(groupRest.CommonPrefixes, want) {
t.Fatalf("second common-prefix page = %v, want %v", groupRest.CommonPrefixes, want)
}
if groupRest.IsTruncated {
t.Fatalf("last common-prefix page is truncated: %+v", groupRest)
}
part, err := z.PutObjectPart(t.Context(), bucket, objects[0], uploadIDs[objects[0]], 1,
mustGetPutObjReader(t, bytes.NewBufferString("part"), 4, "", ""), ObjectOptions{})
if err != nil {
t.Fatal(err)
}
completed, err := z.CompleteMultipartUpload(t.Context(), bucket, objects[0], uploadIDs[objects[0]],
[]CompletePart{{PartNumber: 1, ETag: part.ETag}}, ObjectOptions{})
if err != nil {
t.Fatal(err)
}
for _, key := range []string{multipartMetaBucket, multipartMetaObject} {
if _, ok := completed.UserDefined[key]; ok {
t.Errorf("completed object retained upload-only metadata %q", key)
}
}
// Simulate an upload written by a pre-upgrade server. Its key cannot be
// recovered by scanning the hashed namespace. Strict listing must fail;
// the old path is available only through an explicit migration setting.
legacyObject := objects[1]
er := sets.getHashedSet(legacyObject)
fi, metadata, err := er.checkUploadIDExists(t.Context(), bucket, legacyObject, uploadIDs[legacyObject], true)
if err != nil {
t.Fatal(err)
}
for i := range metadata {
delete(metadata[i].Metadata, multipartMetaBucket)
delete(metadata[i].Metadata, multipartMetaObject)
}
if _, err = writeAllMetadata(t.Context(), er.getDisks(), bucket, minioMetaMultipartBucket,
er.getUploadIDDir(bucket, legacyObject, uploadIDs[legacyObject]), metadata, fi.WriteQuorum(er.defaultWQuorum())); err != nil {
t.Fatal(err)
}
_, err = z.ListMultipartUploads(t.Context(), bucket, legacyObject, "", "", "", 100)
if !errors.Is(err, errMultipartListingLegacy) {
t.Fatalf("legacy strict listing: %v", err)
}
globalAPIConfig.mu.Lock()
oldLegacy := globalAPIConfig.multipartListingLegacy
globalAPIConfig.multipartListingLegacy = true
globalAPIConfig.mu.Unlock()
t.Cleanup(func() {
globalAPIConfig.mu.Lock()
globalAPIConfig.multipartListingLegacy = oldLegacy
globalAPIConfig.mu.Unlock()
})
legacy, err := z.ListMultipartUploads(t.Context(), bucket, legacyObject, "", "", "", 100)
if err != nil {
t.Fatal(err)
}
requireMultipartUploadKeys(t, legacy, legacyObject)
}
func TestPaginateMultipartUploads(t *testing.T) {
base := time.Unix(100, 0)
id1, id2, id3 := multipartListingTestID(base, 1), multipartListingTestID(base.Add(time.Second), 2), multipartListingTestID(base, 3)
uploads := []MultipartInfo{
{Bucket: "bucket", Object: "b", UploadID: id3, Initiated: base},
{Bucket: "bucket", Object: "a", UploadID: id2, Initiated: base.Add(time.Second)},
{Bucket: "bucket", Object: "a", UploadID: id1, Initiated: base},
{Bucket: "bucket", Object: "a", UploadID: id1, Initiated: base}, // duplicate discovery
}
first := paginateMultipartUploads(uploads, "", "", "", "", 1)
requireMultipartUploadKeys(t, first, "a")
if !first.IsTruncated || first.NextKeyMarker != "a" || first.NextUploadIDMarker != id1 {
t.Fatalf("first page = %+v", first)
}
second := paginateMultipartUploads(uploads, "", first.NextKeyMarker, first.NextUploadIDMarker, "", 1)
if len(second.Uploads) != 1 || second.Uploads[0].Object != "a" || second.Uploads[0].UploadID != id2 {
t.Fatalf("second page uploads = %+v", second.Uploads)
}
if !second.IsTruncated || second.NextKeyMarker != "a" || second.NextUploadIDMarker != id2 {
t.Fatalf("second page = %+v", second)
}
last := paginateMultipartUploads(uploads, "", second.NextKeyMarker, second.NextUploadIDMarker, "", 1)
requireMultipartUploadKeys(t, last, "b")
if last.IsTruncated || last.NextKeyMarker != "" || last.NextUploadIDMarker != "" {
t.Fatalf("last page = %+v", last)
}
missingUploadMarker := paginateMultipartUploads(uploads, "", "a", multipartListingTestID(base.Add(2*time.Second), 4), "", 10)
requireMultipartUploadKeys(t, missingUploadMarker, "b")
if err := checkListMultipartArgs(t.Context(), "bucket", "", "", "not-base64=", ""); err != nil {
t.Fatalf("upload-id-marker without key-marker must be ignored: %v", err)
}
overLimit := make([]MultipartInfo, maxUploadsList+1)
for i := range overLimit {
overLimit[i] = MultipartInfo{Bucket: "bucket", Object: fmt.Sprintf("%04d", i), UploadID: fmt.Sprint(i)}
}
capped := paginateMultipartUploads(overLimit, "", "", "", "", maxUploadsList+1)
if capped.MaxUploads != maxUploadsList || len(capped.Uploads) != maxUploadsList || !capped.IsTruncated {
t.Fatalf("over-limit page = MaxUploads %d, uploads %d, truncated %t",
capped.MaxUploads, len(capped.Uploads), capped.IsTruncated)
}
}
func TestListMultipartUploadsGlobalPageAcrossPools(t *testing.T) {
z, bucket := consistencyPools(t)
objects := []string{"a/one", "b/two", "c/three", "d/four"}
for i, object := range objects {
if _, err := z.serverPools[i%len(z.serverPools)].NewMultipartUpload(t.Context(), bucket, object, ObjectOptions{}); err != nil {
t.Fatalf("NewMultipartUpload(%q): %v", object, err)
}
}
z.mpCache.Range(func(uploadID string, _ MultipartInfo) bool {
z.mpCache.Delete(uploadID)
return true
})
page, err := z.ListMultipartUploads(t.Context(), bucket, "", "", "", "", 2)
if err != nil {
t.Fatal(err)
}
requireMultipartUploadKeys(t, page, objects[:2]...)
if !page.IsTruncated || page.NextKeyMarker != objects[1] {
t.Fatalf("first global page = %+v", page)
}
rest, err := z.ListMultipartUploads(t.Context(), bucket, "", page.NextKeyMarker, page.NextUploadIDMarker, "", 2)
if err != nil {
t.Fatal(err)
}
requireMultipartUploadKeys(t, rest, objects[2:]...)
if rest.IsTruncated {
t.Fatalf("last global page is truncated: %+v", rest)
}
}
+6 -1
View File
@@ -20,6 +20,7 @@ package cmd
import (
"context"
"encoding/base64"
"errors"
"runtime"
"strings"
@@ -83,7 +84,8 @@ func checkListMultipartArgs(ctx context.Context, bucket, prefix, keyMarker, uplo
if err := checkListObjsArgs(ctx, bucket, prefix, keyMarker); err != nil {
return err
}
if uploadIDMarker != "" {
// S3 ignores upload-id-marker when key-marker is absent.
if uploadIDMarker != "" && keyMarker != "" {
if HasSuffix(keyMarker, SlashSeparator) {
return InvalidUploadIDKeyCombination{
UploadIDMarker: uploadIDMarker,
@@ -96,6 +98,9 @@ func checkListMultipartArgs(ctx context.Context, bucket, prefix, keyMarker, uplo
UploadID: uploadIDMarker,
}
}
if _, ok := multipartMarkerTime(uploadIDMarker); !ok {
return InvalidArgument{Bucket: bucket, Object: keyMarker, Err: errors.New("upload-id-marker must contain a native multipart upload ID")}
}
}
return nil
}
+9
View File
@@ -109,6 +109,15 @@ func testStorageAPIListDir(t *testing.T, storage StorageAPI) {
}
}
}
for _, name := range []string{"one", "two", "three"} {
if err := storage.AppendFile(t.Context(), "foo", "bounded/"+name, []byte("x")); err != nil {
t.Fatal(err)
}
}
entries, err := storage.ListDir(t.Context(), "", "foo", "bounded", 2)
if err != nil || len(entries) != 2 {
t.Fatalf("ListDir count was lost in storage/RPC path: %v %v", entries, err)
}
}
func testStorageAPIReadAll(t *testing.T, storage StorageAPI) {
+8
View File
@@ -45,6 +45,7 @@ const (
apiTransitionWorkers = "transition_workers"
apiStaleUploadsCleanupInterval = "stale_uploads_cleanup_interval"
apiStaleUploadsExpiry = "stale_uploads_expiry"
apiMultipartListing = "multipart_listing"
apiDeleteCleanupInterval = "delete_cleanup_interval"
apiDisableODirect = "disable_odirect"
apiODirect = "odirect"
@@ -67,6 +68,7 @@ const (
EnvAPIStaleUploadsCleanupInterval = "MINIO_API_STALE_UPLOADS_CLEANUP_INTERVAL"
EnvAPIStaleUploadsExpiry = "MINIO_API_STALE_UPLOADS_EXPIRY"
EnvAPIMultipartListing = "MINIO_API_MULTIPART_LISTING"
EnvAPIDeleteCleanupInterval = "MINIO_API_DELETE_CLEANUP_INTERVAL"
EnvDeleteCleanupInterval = "MINIO_DELETE_CLEANUP_INTERVAL"
EnvAPIODirect = "MINIO_API_ODIRECT"
@@ -133,6 +135,7 @@ var (
Key: apiStaleUploadsExpiry,
Value: "24h",
},
config.KV{Key: apiMultipartListing, Value: "strict"},
config.KV{
Key: apiDeleteCleanupInterval,
Value: "5m",
@@ -178,6 +181,7 @@ type Config struct {
TransitionWorkers int `json:"transition_workers"`
StaleUploadsCleanupInterval time.Duration `json:"stale_uploads_cleanup_interval"`
StaleUploadsExpiry time.Duration `json:"stale_uploads_expiry"`
MultipartListing string `json:"multipart_listing"`
DeleteCleanupInterval time.Duration `json:"delete_cleanup_interval"`
EnableODirect bool `json:"enable_odirect"`
GzipObjects bool `json:"gzip_objects"`
@@ -320,6 +324,10 @@ func LookupConfig(kvs config.KVS) (cfg Config, err error) {
return cfg, err
}
cfg.StaleUploadsExpiry = staleUploadsExpiry
cfg.MultipartListing = env.Get(EnvAPIMultipartListing, kvs.GetWithDefault(apiMultipartListing, DefaultKVS))
if cfg.MultipartListing != "strict" && cfg.MultipartListing != "legacy" {
return cfg, fmt.Errorf("%s must be strict or legacy", apiMultipartListing)
}
cfg.SyncEvents = env.Get(EnvAPISyncEvents, kvs.Get(apiSyncEvents)) == config.EnableOn
+41
View File
@@ -0,0 +1,41 @@
// Copyright (c) 2026 Ruohang Feng
// SPDX-License-Identifier: AGPL-3.0-or-later
package api
import (
"testing"
"github.com/minio/minio/internal/config"
)
func TestMultipartListingMigrationMode(t *testing.T) {
for _, tc := range []struct {
name, stored, override, want string
invalid bool
}{
{name: "default", want: "strict"},
{name: "explicit-migration", stored: "legacy", want: "legacy"},
{name: "environment-override", stored: "strict", override: "legacy", want: "legacy"},
{name: "invalid-stored", stored: "automatic", invalid: true},
{name: "invalid-environment", stored: "strict", override: "automatic", invalid: true},
} {
t.Run(tc.name, func(t *testing.T) {
t.Setenv(EnvAPIMultipartListing, tc.override)
var kvs config.KVS
if tc.stored != "" {
kvs = config.KVS{{Key: apiMultipartListing, Value: tc.stored}}
}
cfg, err := LookupConfig(kvs)
if tc.invalid {
if err == nil {
t.Fatal("invalid migration mode silently accepted")
}
return
}
if err != nil || cfg.MultipartListing != tc.want {
t.Fatalf("mode=%q err=%v, want %q", cfg.MultipartListing, err, tc.want)
}
})
}
}
+6
View File
@@ -26,6 +26,12 @@ var (
// Help holds configuration keys and their default values for api subsystem.
Help = config.HelpKVS{
config.HelpKV{
Key: apiMultipartListing,
Description: "multipart listing mode: strict, or temporary legacy mode during coordinated upgrade" + defaultHelpPostfix(apiMultipartListing),
Type: "string",
Optional: true,
},
config.HelpKV{
Key: apiRequestsMax,
Description: `set the maximum number of concurrent requests (default: auto)`,