时隔一个季度,对 OtterIO 进行了一些维护工作,应该对一些有需求的朋友有用,那么就记录下来吧。
写在前面
六月整理 OtterIO 时(项目地址 soulteary/otterio),我主要处理了代码基线、功能裁剪和部署迁移。当完成 HTTP 路由调整、依赖升级和一轮历史安全问题检查后,项目也就有了容器镜像和二进制发布。
OtterIO 是一个 S3 兼容对象存储服务,基于 MinIO 最后一个 Apache 2.0 release 附近的代码,起点约为 RELEASE.2021-04-22T15-44-28Z。项目继续以 Apache License 2.0 分发,保留原始版权与来源说明(项目说明)。
当前这轮维护中,一个安全反馈暴露了预签名上传的问题:接收方可能通过附加未签名的请求头,把上传改成服务端复制。分布式存储接口也补上了路径和元数据检查。
** 如果你还在使用六月份发布的软件,记得升级容器镜像的版本。**
本文以北京时间 2026 年 10 月 5 日发布的 RELEASE.2026-10-04T22-12-15Z 为准,并与六月七日的正式版本对照。版本标签使用 UTC,所以日期仍是十月四日。
之前的整理过程,记录在这三篇文章里:
- 重新审视 MinIO:许可证、归档、社区 fork 与我的 Apache 2.0 基线
- 从 MinIO 到 OtterIO:整理一条 Apache 2.0 开源对象存储代码线
- 把 MinIO 示例迁到 OtterIO:使用、部署与迁移验证
修复 SigV4 请求头校验
感谢 Act Security 的 Oren Yomtov 报告这个问题,并提供详细分析和复现用例(问题与修复说明)。
预签名 URL 常用于让用户直接上传对象:应用服务端生成带签名的 URL,指定上传位置和有效期,再把它交给用户。用户凭这个 URL 就能上传,不需要拿到服务端的访问密钥。
受影响的代码只检查签名中列出的请求头,没有拦住额外附加的操作请求头。例如,拿到上传 URL 的人可以加上未签名的 x-amz-copy-source,指定一个源对象,让服务端把上传改成复制。服务端随后会使用生成签名的身份所拥有的读取权限,把源对象复制到 URL 指定的位置。
复制目的地仍由签名 URL 限定;要读取复制后的内容,还需要具备目的地的读取权限。六月七日的正式版已确认受影响,使用这个版本的用户需要升级到包含修复的版本。
修复后,服务端会检查请求中的 x-amz-* 和 x-otterio-* 请求头。决定或修改 S3 操作的请求头必须受到签名保护,如果这些操作请求头未被签名覆盖,服务端会返回 HTTP 403 AccessDenied。这项检查同时用于预签名 URL、通过 Authorization 请求头认证的 SigV4 请求,以及流式上传的初始签名。
这里保留了协议允许的写法。x-amz-content-sha256 不一定要列入 SignedHeaders,因为它的有效值已经作为 payload hash 参与签名。同一参数也可以同时出现在已签名的 URL 查询参数和请求头里,但两处的值必须完全一致,包括多值的数量和顺序。
如果自己封装了 S3 客户端,需要在生成签名前设置好这些操作请求头。先签名、再附加复制来源等请求头的写法,会被新版拒绝,应调整客户端的签名流程。合法复制仍然可用,源对象的 IAM 权限检查也继续保留。
回归测试使用独立的 S3 SDK 生成签名,再通过实际使用的 Fiber 路由,分别在 FS 和纠删码存储上执行请求。测试确认普通上传和合法复制继续可用,附加未签名复制头的请求会被拒绝,分片上传的数据也不会被这种请求改变。只具备写入权限的签名身份,即使正确签名了复制请求,也不能绕过源对象的读取权限检查(回归测试)。
检查内部存储接口的路径和元数据
OtterIO 为分布式纠删码部署注册了一组内部 storage REST 接口,供节点之间读写数据。这些接口使用由集群 root 凭据签名的 JWT 进行认证,匿名 S3 请求或有限权限的 S3 访问密钥不足以调用。
即使请求通过认证,传来的路径和元数据仍然需要检查。有内部调用权限的节点也可能发来异常输入,导致操作越过存储目录,或让接收方因畸形元数据而崩溃。
这次参考 SILO SN-2026-002 公告,核查了 OtterIO 继承的同类缺陷,并独立实现相应校验。两条代码线的接口和传输方式不同,修复范围以 OtterIO 的实际实现为准(内部存储安全说明)。
在拼接路径之前检查输入
路径清理和拼接会消掉原始输入中的部分路径段。如果等到这些处理完成后再检查,可能已经看不出输入里的问题。新版先检查未经规范化的路径,再进行拼接;检查范围也不只包括 URL 参数,还包括元数据请求体中的对象名和数据目录。
重命名时,来源和目标路径都要检查。批量建卷和删除版本时,则先检查整个批次,避免已经修改一部分数据,才发现后面的路径不合法。
删除目录时,还要正确判断目录之间的包含关系。例如,backup 和 backup-old 名称相近,但后者并不是前者的子目录。新版按路径段判断范围,避免仅凭字符串前缀做出错误判断。已有的嵌套系统目录、目录修复路径和根目录列举等合法用法仍然保留(存储校验实现)。
这些检查针对的是路径字符串,不能替代操作系统对数据目录的保护,也不能保证防住符号链接或并发文件系统改动造成的目录逃逸。数据目录仍应由服务控制,节点通信也应保持隔离。
检查元数据的数值和结构
纠删码参数、分片数组和对象长度会用于索引、大小计算和完整性判断。数组长度不一致可能导致越界访问,无效或溢出的长度可能触发异常内存分配,错误的大小信息还可能把截断的分片误判为健康。新版在使用这些值前增加了校验,并拒绝不支持的 bitrot 校验算法。
仅限制请求体大小还不够。MessagePack 可以用很少的字节声明一个很大的数组或映射,解码器如果直接按声明分配内存,小请求也可能造成大量内存占用。因此,新版在调用正式解码器之前,先检查数组、映射的元素数量和嵌套深度。
对内部缓冲读写和 FileInfo 元数据请求体,还增加了单次 64 MiB 的大小限制。这项限制不影响 S3 对象大小,大对象继续使用已有的流式读写路径。自行实现内部客户端时,需要拆分过大的缓冲操作和版本删除批次。
回归测试会用有效的节点 JWT 调用内部 REST 接口,并确认正常写入能够成功。在此基础上,再验证异常路径和元数据会被拒绝,确保这些请求确实经过了输入校验(存储安全回归)。
升级前应备份元数据,升级后再验证已有对象的读取、修复和复制。新版会把不合法的历史元数据视为损坏数据并拒绝处理,不会自动改写或删除;遇到这类问题,需要调查原因,并从已知良好的副本恢复。
分布式部署需要更新所有节点。
新版容器要求显式配置凭据
六月八日那篇文章里,容器启动示例省略了凭据。新版在运行 server 和 gateway 前,会先检查用户名和密码是否完整配置,并拒绝使用默认密码。因此,旧命令需要补上凭据才能继续使用(容器升级说明)。
已经配置自定义凭据的部署,可以继续沿用原值;六月九日文章中的显式凭据示例也仍然适用。升级前需要确认这些值已经传入容器。
隔离的本地演示仍可通过 OTTERIO_ALLOW_DEFAULT_CREDENTIALS=1 使用默认凭据,但这个开关不会接受只提供用户名或密码的不完整配置。
新建本地开发或测试环境,可以先生成一组凭据,再传给固定版本的镜像。下面的例子用 OpenSSL 生成密码,S3 和控制台端口都绑定到 127.0.0.1:
mkdir -p ./data
export OTTERIO_ROOT_USER=otterio-admin
export OTTERIO_ROOT_PASSWORD="$(openssl rand -hex 32)"
docker run -d \
--name otterio \
-p 127.0.0.1:9000:9000 \
-p 127.0.0.1:9001:9001 \
-e OTTERIO_ROOT_USER \
-e OTTERIO_ROOT_PASSWORD \
-v "$PWD/data:/data" \
ghcr.io/soulteary/otterio:RELEASE.2026-10-04T22-12-15Z \
server --address ":9000" --console-address ":9001" /data
启动后,S3 接口地址是 http://127.0.0.1:9000,控制台地址是 http://127.0.0.1:9001。生成的用户名和密码要保存,后续连接或重建容器时继续使用。上面的密码生成步骤用于新建环境,已有部署升级时应保留原来的凭据。如果容器名或端口已经被占用,需要修改示例中的对应值。
凭据文件和目录权限
用户名和密码可以分别从环境变量或文件加载。例如,用户名通过 OTTERIO_ROOT_USER 设置,密码通过 OTTERIO_ROOT_PASSWORD_FILE 指定的容器内文件加载,这种组合可以正常使用。
如果要求启动时必须读到某个凭据文件,应通过 _FILE 指定它在容器内的绝对路径,并确保文件可读、为普通文件且内容非空。同一项配置同时提供非空环境值和已有文件,会报冲突。
新变量 OTTERIO_ROOT_USER / OTTERIO_ROOT_PASSWORD 优先于旧变量 OTTERIO_ACCESS_KEY / OTTERIO_SECRET_KEY,但必须成对使用。不能从新、旧两组变量中各取一项,拼成一组凭据。
凭据检查由容器入口脚本执行,直接运行裸机二进制时没有这项强制要求。镜像默认 UID 也没有改变,已有数据卷的归属不会自动调整。
需要以非 root 用户运行时,可以使用安全 Compose 示例。这份配置指定了非 root UID/GID,使用只读根文件系统,限制临时目录大小,移除 capabilities,并禁止提权。使用前应确保选定的容器用户能够写入数据目录、读取凭据文件。
补充 Windows 启动和并发回归测试
六月八日那篇文章里,Windows CI 主要验证代码能否编译。后来社区反馈了目录枚举和服务启动的问题,并提供分析和补丁建议。相应修复完成后,CI 也补上了回归测试(贡献记录)。
现在,Windows CI 会执行选定的目录回归用例,并实际启动编译出来的服务,验证 S3 读写。完整的 Go 测试仍在 Linux 上运行,其中包含依赖 POSIX 文件系统语义的用例(参考当前 Go workflow)。
这轮还修复了流式响应中的并发读写问题。HTTP 路由继续使用 Fiber v3,处理协程和发送侧需要交接响应头、状态码等信息。新版在首次提交响应时固定一份快照,发送侧只读取这份结果,后续修改不会影响已提交的响应。
响应提交前发生的异常会交回原有错误处理流程,提交后的异常则中断数据流;响应内容仍然边写边传。对应测试覆盖提交前后的异常、部分数据已经写出,以及读取端提前退出等情况(流式响应实现)。
源码构建的 Go 最低版本由 1.26 提升到 1.27.1。服务端依赖和前端锁文件一并更新,生成代码及嵌入服务中的控制台产物也已重新构建。直接运行发布镜像或二进制时,无需安装 Go 或前端构建工具(依赖与兼容记录)。
Linux 上的 Race 检查会先列出需要测试的包,再运行这些包中的全部用例,不额外筛选测试。-count=1 用来避免复用之前成功的测试结果。如果获取包列表失败,任务会直接停止,避免漏测后仍报告通过。
再次核对发布出去的镜像和二进制
源码通过测试之后,我们还要确认发布的是同一份代码。六月已有镜像、二进制和校验和。这次发布会核对待发布的源码提交,实际运行已推送的镜像,并把上传的附件重新下载比对。
本次正式版提供 Linux amd64、arm64、ppc64le,macOS amd64、arm64,以及 Windows amd64 六种系统与架构组合的二进制,附带校验文件和发布清单(版本下载)。
本地 Docker 构建使用当前源码
原来的 Dockerfile 会在构建时克隆远端仓库,再切到 main 编译。当时的策略是,始终以合并到 main 的代码为准。所以,本地修改不会进入这条构建路径,生成的镜像可能与准备测试的源码不同。
现在,本地 Docker 构建直接使用传入的源码上下文,包含本地修改和已提交的控制台资源。来源提交号默认为 unknown,可通过 VCS_REF 显式标注;这个参数只记录来源,不会切换构建使用的代码。正式发布仍从固定提交编译二进制,再由 Dockerfile.ci 打包镜像(本地源码构建说明)。
先验证固定版本再更新稳定标签
现在的发布流程按下面的步骤执行:
- 确认针对该源码提交运行的 Go、Lint 和 Release checks 三个工作流都已完成且成功;这些工作流由
main分支推送触发。 - 从这个固定提交构建二进制和版本镜像,暂不更新稳定标签。
- 按 digest 拉取已经推送到 GHCR 的镜像,在 Linux amd64 上运行启动和 S3 读写冒烟验证。
- 将二进制、校验文件和发布清单上传到 GitHub 草稿 Release,逐个下载并与原产物比对。比对通过后公开这个版本,仍保留现有的
latest指向。 - 在串行流程中核对已发布版本的源码提交、各仓库固定版本标签指向的镜像 digest,并比较版本先后,阻止较旧版本覆盖当前的
latest。通过检查后,用已验证的 digest 更新稳定镜像标签,最后更新 GitHub latest。
固定版本的公开发布和稳定标签更新共用并发锁,这两类任务会串行执行,避免彼此在核对与更新之间改变版本状态。镜像标签更新后,还会再次核对 digest(发布工作流)。
新增的 release-manifest.json 记录了发布标签、源码提交和镜像 digest,可以用来核对运行中的版本。使用镜像时可以按 digest 固定内容,下载二进制后则用附带的 SHA256 文件校验(本次发布清单)。
清单并不提供发布签名,也不证明可以逐字节复现同一次构建。GitHub、GHCR 和 Docker Hub 的更新也分步完成,不能作为一个原子事务提交。某一步失败后,应先核对各平台已经完成的状态,再恢复对应步骤。
本版的 Go 检查和正式发布流程均已成功执行,包含 Race、平台检查、附件比对和稳定标签更新。发布记录中的镜像运行验证覆盖 Linux amd64,其他镜像架构仍需单独验证(Go 检查记录、正式发布记录)。
升级后还要验证自己的应用
六月的迁移实验工具会生成随机对象,记录源端大小和 SHA256,再从目标端下载校验,检查两端内容是否一致。升级后仍可沿用这个方法,并结合应用的实际用法检查:
- 验证应用的上传、下载和删除流程,并抽查已有对象能否读取。
- 确认未签名操作头会被拒绝,合法签名复制仍然可用,缺少源对象读取权限时复制失败。
- 验证分片上传,以及超过 64 MiB 对象的流式读写。
- 检查环境变量和 secret 文件是否正确加载、数据目录权限是否匹配,并在实际使用的平台上启动服务。
- 分布式部署还要确认所有节点的版本,并检查已有数据的修复和复制。
之前的一篇文章里,已经验证过 MinIO 不同衍生版本之间的数据迁移。这次只测试 OtterIO 六月版到十月版的迁移:
2026/10/05 06:00:06 SYNC SUMMARY
2026/10/05 06:00:06 source endpoints: 1
2026/10/05 06:00:06 objects per source: 12
2026/10/05 06:00:06 expected objects: 12
2026/10/05 06:00:06 uploaded source objects: 12
2026/10/05 06:00:06 verified target objects: 12
2026/10/05 06:00:06 target endpoint: http://otterio-261004:9000
2026/10/05 06:00:06 target bucket: migration-test
2026/10/05 06:00:06 target from/ objects: 12
2026/10/05 06:00:06 target from/otterio-260607/ objects: 12
2026/10/05 06:00:06 result: ALL OBJECTS SYNCED AND VERIFIED
测试程序
测试使用的环境配置:
ROOT_USER=admin
ROOT_PASSWORD=otterio123456
容器配置文件(docker-compose.yaml):
services:
otterio-261004:
image: soulteary/otterio:RELEASE.2026-10-04T22-12-15Z
container_name: otterio-261004
command: server /data
environment:
OTTERIO_ROOT_USER: ${ROOT_USER}
OTTERIO_ROOT_PASSWORD: ${ROOT_PASSWORD}
volumes:
- ./data/otterio-261004:/data
ports:
- "9101:9000"
networks:
- object-storage-lab
restart: unless-stopped
otterio-260607:
image: soulteary/otterio:RELEASE.2026-06-07T11-32-46Z
container_name: otterio-260607
command: server /data
environment:
OTTERIO_ROOT_USER: ${ROOT_USER}
OTTERIO_ROOT_PASSWORD: ${ROOT_PASSWORD}
volumes:
- ./data/otterio-260607:/data
ports:
- "9102:9000"
networks:
- object-storage-lab
restart: unless-stopped
networks:
object-storage-lab:
name: object-storage-lab
driver: bridge
测试程序:
package main
import (
"bytes"
"context"
"crypto/rand"
"crypto/sha256"
"encoding/hex"
"errors"
"flag"
"fmt"
"io"
"log"
"math/big"
"mime"
"os"
"path/filepath"
"sort"
"strings"
"time"
"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/config"
awscred "github.com/aws/aws-sdk-go-v2/credentials"
"github.com/aws/aws-sdk-go-v2/service/s3"
s3types "github.com/aws/aws-sdk-go-v2/service/s3/types"
)
type Endpoint struct {
Name string
Endpoint string
AccessKey string
SecretKey string
Region string
Client *s3.Client
}
type ObjectRecord struct {
SourceName string
Bucket string
Key string
Size int64
SHA256 string
}
type Options struct {
Bucket string
TargetName string
ObjectCount int
MinSize int64
MaxSize int64
Cleanup bool
ListObjects bool
}
func env(key, fallback string) string {
if v := strings.TrimSpace(os.Getenv(key)); v != "" {
return v
}
return fallback
}
func parseSize(s string) (int64, error) {
original := s
s = strings.TrimSpace(strings.ToLower(s))
multiplier := int64(1)
switch {
case strings.HasSuffix(s, "kb"):
multiplier = 1024
s = strings.TrimSuffix(s, "kb")
case strings.HasSuffix(s, "mb"):
multiplier = 1024 * 1024
s = strings.TrimSuffix(s, "mb")
case strings.HasSuffix(s, "gb"):
multiplier = 1024 * 1024 * 1024
s = strings.TrimSuffix(s, "gb")
case strings.HasSuffix(s, "b"):
multiplier = 1
s = strings.TrimSuffix(s, "b")
}
var n int64
_, err := fmt.Sscanf(strings.TrimSpace(s), "%d", &n)
if err != nil {
return 0, fmt.Errorf("invalid size %q: %w", original, err)
}
if n < 0 {
return 0, fmt.Errorf("invalid size %q: size must be positive", original)
}
return n * multiplier, nil
}
func randomInt64(min, max int64) (int64, error) {
if max < min {
return 0, fmt.Errorf("max size must be >= min size")
}
if min == max {
return min, nil
}
span := max - min + 1
if span <= 0 {
return 0, fmt.Errorf("invalid random range: min=%d max=%d", min, max)
}
n, err := rand.Int(rand.Reader, big.NewInt(span))
if err != nil {
return 0, err
}
return min + n.Int64(), nil
}
func randomBytes(size int64) ([]byte, string, error) {
if size < 0 {
return nil, "", fmt.Errorf("size must be positive")
}
buf := make([]byte, size)
if _, err := rand.Read(buf); err != nil {
return nil, "", err
}
sum := sha256.Sum256(buf)
return buf, hex.EncodeToString(sum[:]), nil
}
func sha256File(path string) (string, int64, error) {
f, err := os.Open(path)
if err != nil {
return "", 0, err
}
defer f.Close()
h := sha256.New()
n, err := io.Copy(h, f)
if err != nil {
return "", 0, err
}
return hex.EncodeToString(h.Sum(nil)), n, nil
}
func newS3Client(ctx context.Context, ep Endpoint) (*s3.Client, error) {
cfg, err := config.LoadDefaultConfig(
ctx,
config.WithRegion(ep.Region),
config.WithCredentialsProvider(
awscred.NewStaticCredentialsProvider(ep.AccessKey, ep.SecretKey, ""),
),
)
if err != nil {
return nil, err
}
client := s3.NewFromConfig(cfg, func(o *s3.Options) {
o.BaseEndpoint = aws.String(ep.Endpoint)
o.UsePathStyle = true
})
return client, nil
}
func ensureBucket(ctx context.Context, client *s3.Client, bucket string) error {
_, err := client.CreateBucket(ctx, &s3.CreateBucketInput{
Bucket: aws.String(bucket),
})
if err == nil {
return nil
}
var owned *s3types.BucketAlreadyOwnedByYou
if errors.As(err, &owned) {
return nil
}
var exists *s3types.BucketAlreadyExists
if errors.As(err, &exists) {
return nil
}
msg := err.Error()
if strings.Contains(msg, "BucketAlreadyOwnedByYou") ||
strings.Contains(msg, "BucketAlreadyExists") ||
strings.Contains(msg, "Your previous request to create the named bucket succeeded") {
return nil
}
return err
}
func deleteBucketObjects(ctx context.Context, client *s3.Client, bucket string) error {
var token *string
for {
out, err := client.ListObjectsV2(ctx, &s3.ListObjectsV2Input{
Bucket: aws.String(bucket),
ContinuationToken: token,
})
if err != nil {
if strings.Contains(err.Error(), "NoSuchBucket") {
return nil
}
return err
}
for _, obj := range out.Contents {
key := aws.ToString(obj.Key)
if key == "" {
continue
}
_, err := client.DeleteObject(ctx, &s3.DeleteObjectInput{
Bucket: aws.String(bucket),
Key: aws.String(key),
})
if err != nil {
return fmt.Errorf("delete object %s: %w", key, err)
}
log.Printf("[cleanup] deleted: s3://%s/%s", bucket, key)
}
if out.IsTruncated == nil || !*out.IsTruncated {
break
}
token = out.NextContinuationToken
}
return nil
}
func contentTypeForKey(key string) string {
ext := filepath.Ext(key)
if ext == "" {
return "application/octet-stream"
}
if ct := mime.TypeByExtension(ext); ct != "" {
return ct
}
switch ext {
case ".bin":
return "application/octet-stream"
case ".json":
return "application/json"
case ".txt":
return "text/plain; charset=utf-8"
default:
return "application/octet-stream"
}
}
func makeObjectKey(sourceName string, index int) string {
now := time.Now().UTC().Format("20060102T150405Z")
patterns := []string{
"small/%s-%03d-random.txt",
"medium/%s-%03d-random.bin",
"nested/year=2026/month=06/day=09/%s-%03d-file.txt",
"unicode/%s-%03d-中文文件名.txt",
"unicode/%s-%03d-空 格 文件.txt",
"metadata/%s-%03d-content-type.json",
}
pattern := patterns[(index-1)%len(patterns)]
return fmt.Sprintf(pattern, now+"-"+sourceName, index)
}
func uploadRandomObjects(ctx context.Context, ep Endpoint, opt Options) ([]ObjectRecord, error) {
log.Printf("[%s] ensure bucket: %s", ep.Name, opt.Bucket)
if err := ensureBucket(ctx, ep.Client, opt.Bucket); err != nil {
return nil, fmt.Errorf("[%s] create bucket: %w", ep.Name, err)
}
records := make([]ObjectRecord, 0, opt.ObjectCount)
for i := 1; i <= opt.ObjectCount; i++ {
size, err := randomInt64(opt.MinSize, opt.MaxSize)
if err != nil {
return nil, err
}
key := makeObjectKey(ep.Name, i)
data, digest, err := randomBytes(size)
if err != nil {
return nil, err
}
_, err = ep.Client.PutObject(ctx, &s3.PutObjectInput{
Bucket: aws.String(opt.Bucket),
Key: aws.String(key),
Body: bytes.NewReader(data),
ContentLength: aws.Int64(size),
ContentType: aws.String(contentTypeForKey(key)),
Metadata: map[string]string{
"source": ep.Name,
"sha256": digest,
},
})
if err != nil {
return nil, fmt.Errorf("[%s] put object %s: %w", ep.Name, key, err)
}
log.Printf("[%s] uploaded: s3://%s/%s size=%d sha256=%s",
ep.Name,
opt.Bucket,
key,
size,
digest,
)
records = append(records, ObjectRecord{
SourceName: ep.Name,
Bucket: opt.Bucket,
Key: key,
Size: size,
SHA256: digest,
})
}
return records, nil
}
func downloadToTempFile(ctx context.Context, client *s3.Client, bucket, key string) (string, error) {
out, err := client.GetObject(ctx, &s3.GetObjectInput{
Bucket: aws.String(bucket),
Key: aws.String(key),
})
if err != nil {
return "", err
}
defer out.Body.Close()
tmp, err := os.CreateTemp("", "s3-migrate-*")
if err != nil {
return "", err
}
defer tmp.Close()
if _, err := io.Copy(tmp, out.Body); err != nil {
_ = os.Remove(tmp.Name())
return "", err
}
return tmp.Name(), nil
}
func uploadFile(
ctx context.Context,
client *s3.Client,
bucket string,
key string,
path string,
sourceName string,
digest string,
size int64,
) error {
f, err := os.Open(path)
if err != nil {
return err
}
defer f.Close()
_, err = client.PutObject(ctx, &s3.PutObjectInput{
Bucket: aws.String(bucket),
Key: aws.String(key),
Body: f,
ContentLength: aws.Int64(size),
ContentType: aws.String(contentTypeForKey(key)),
Metadata: map[string]string{
"source": sourceName,
"source-sha256": digest,
},
})
return err
}
func syncOneObject(ctx context.Context, source Endpoint, target Endpoint, record ObjectRecord, opt Options) error {
targetKey := fmt.Sprintf("from/%s/%s", source.Name, record.Key)
tmpPath, err := downloadToTempFile(ctx, source.Client, record.Bucket, record.Key)
if err != nil {
return fmt.Errorf("[%s] download %s: %w", source.Name, record.Key, err)
}
defer os.Remove(tmpPath)
digest, size, err := sha256File(tmpPath)
if err != nil {
return fmt.Errorf("[%s] sha256 temp file %s: %w", source.Name, record.Key, err)
}
if size != record.Size || digest != record.SHA256 {
return fmt.Errorf(
"[%s] local verify failed key=%s source_size=%d local_size=%d source_sha=%s local_sha=%s",
source.Name,
record.Key,
record.Size,
size,
record.SHA256,
digest,
)
}
if err := uploadFile(ctx, target.Client, opt.Bucket, targetKey, tmpPath, source.Name, digest, size); err != nil {
return fmt.Errorf("[otterio] upload %s: %w", targetKey, err)
}
log.Printf("[%s -> %s] synced: s3://%s/%s size=%d sha256=%s",
source.Name,
target.Name,
opt.Bucket,
targetKey,
size,
digest,
)
return nil
}
func verifyTargetObject(ctx context.Context, target Endpoint, opt Options, record ObjectRecord) error {
targetKey := fmt.Sprintf("from/%s/%s", record.SourceName, record.Key)
tmpPath, err := downloadToTempFile(ctx, target.Client, opt.Bucket, targetKey)
if err != nil {
return fmt.Errorf("[verify] download target %s: %w", targetKey, err)
}
defer os.Remove(tmpPath)
digest, size, err := sha256File(tmpPath)
if err != nil {
return fmt.Errorf("[verify] sha256 target %s: %w", targetKey, err)
}
if size != record.Size || digest != record.SHA256 {
return fmt.Errorf(
"[verify] mismatch target=%s source_size=%d target_size=%d source_sha=%s target_sha=%s",
targetKey,
record.Size,
size,
record.SHA256,
digest,
)
}
log.Printf("[verify] ok: s3://%s/%s size=%d sha256=%s",
opt.Bucket,
targetKey,
size,
digest,
)
return nil
}
func countObjectsWithPrefix(ctx context.Context, client *s3.Client, bucket, prefix string) (int, error) {
var token *string
count := 0
for {
out, err := client.ListObjectsV2(ctx, &s3.ListObjectsV2Input{
Bucket: aws.String(bucket),
Prefix: aws.String(prefix),
ContinuationToken: token,
})
if err != nil {
return 0, err
}
count += len(out.Contents)
if out.IsTruncated == nil || !*out.IsTruncated {
break
}
token = out.NextContinuationToken
}
return count, nil
}
func listTargetObjects(ctx context.Context, target Endpoint, bucket string) error {
var token *string
var keys []string
for {
out, err := target.Client.ListObjectsV2(ctx, &s3.ListObjectsV2Input{
Bucket: aws.String(bucket),
ContinuationToken: token,
})
if err != nil {
return err
}
for _, obj := range out.Contents {
keys = append(keys, fmt.Sprintf("%s %d", aws.ToString(obj.Key), aws.ToInt64(obj.Size)))
}
if out.IsTruncated == nil || !*out.IsTruncated {
break
}
token = out.NextContinuationToken
}
sort.Strings(keys)
log.Printf("[otterio] object list:")
for _, key := range keys {
log.Printf(" %s", key)
}
return nil
}
func printSyncSummary(
ctx context.Context,
target Endpoint,
sources []Endpoint,
opt Options,
allRecords []ObjectRecord,
verifiedCount int,
) error {
expected := len(sources) * opt.ObjectCount
log.Printf("")
log.Printf("SYNC SUMMARY")
log.Printf(" source endpoints: %d", len(sources))
log.Printf(" objects per source: %d", opt.ObjectCount)
log.Printf(" expected objects: %d", expected)
log.Printf(" uploaded source objects: %d", len(allRecords))
log.Printf(" verified target objects: %d", verifiedCount)
log.Printf(" target endpoint: %s", target.Endpoint)
log.Printf(" target bucket: %s", opt.Bucket)
if len(allRecords) != expected {
return fmt.Errorf("uploaded object count mismatch: expected=%d actual=%d", expected, len(allRecords))
}
if verifiedCount != expected {
return fmt.Errorf("verified object count mismatch: expected=%d actual=%d", expected, verifiedCount)
}
totalFromPrefixCount, err := countObjectsWithPrefix(ctx, target.Client, opt.Bucket, "from/")
if err != nil {
return fmt.Errorf("count target objects: %w", err)
}
log.Printf(" target from/ objects: %d", totalFromPrefixCount)
if totalFromPrefixCount != expected {
return fmt.Errorf("target object count mismatch under from/: expected=%d actual=%d", expected, totalFromPrefixCount)
}
for _, source := range sources {
prefix := fmt.Sprintf("from/%s/", source.Name)
count, err := countObjectsWithPrefix(ctx, target.Client, opt.Bucket, prefix)
if err != nil {
return fmt.Errorf("count target prefix %s: %w", prefix, err)
}
log.Printf(" target %s objects: %d", prefix, count)
if count != opt.ObjectCount {
return fmt.Errorf("target prefix count mismatch prefix=%s expected=%d actual=%d", prefix, opt.ObjectCount, count)
}
}
log.Printf(" result: ALL OBJECTS SYNCED AND VERIFIED")
return nil
}
func main() {
var (
bucket string
objectCount int
minSizeText string
maxSizeText string
cleanup bool
listObjects bool
)
flag.StringVar(&bucket, "bucket", env("S3_TEST_BUCKET", "migration-test"), "bucket name")
flag.IntVar(&objectCount, "count", 12, "random object count per source")
flag.StringVar(&minSizeText, "min-size", "1kb", "min random object size, e.g. 1kb, 1mb")
flag.StringVar(&maxSizeText, "max-size", "10mb", "max random object size, e.g. 10mb, 100mb")
flag.BoolVar(&cleanup, "cleanup", false, "delete all objects in test bucket before running")
flag.BoolVar(&listObjects, "list", false, "list all objects in target bucket after sync")
flag.Parse()
minSize, err := parseSize(minSizeText)
if err != nil {
log.Fatal(err)
}
maxSize, err := parseSize(maxSizeText)
if err != nil {
log.Fatal(err)
}
if objectCount <= 0 {
log.Fatal("count must be greater than 0")
}
if maxSize < minSize {
log.Fatal("max-size must be greater than or equal to min-size")
}
opt := Options{
Bucket: bucket,
TargetName: "otterio-261004",
ObjectCount: objectCount,
MinSize: minSize,
MaxSize: maxSize,
Cleanup: cleanup,
ListObjects: listObjects,
}
accessKey := env("S3_ACCESS_KEY", "admin")
secretKey := env("S3_SECRET_KEY", "otterio123456")
region := env("S3_REGION", "us-east-1")
endpoints := []Endpoint{
{
Name: "otterio-260607",
Endpoint: env("OTTERIO_0607_ENDPOINT", "http://127.0.0.1:9102"),
AccessKey: accessKey,
SecretKey: secretKey,
Region: region,
},
{
Name: "otterio-261004",
Endpoint: env("OTTERIO_1004_ENDPOINT", "http://127.0.0.1:9101"),
AccessKey: accessKey,
SecretKey: secretKey,
Region: region,
},
}
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Minute)
defer cancel()
for i := range endpoints {
client, err := newS3Client(ctx, endpoints[i])
if err != nil {
log.Fatalf("[%s] create client: %v", endpoints[i].Name, err)
}
endpoints[i].Client = client
log.Printf("[%s] endpoint=%s", endpoints[i].Name, endpoints[i].Endpoint)
}
var target Endpoint
var sources []Endpoint
for _, ep := range endpoints {
if ep.Name == opt.TargetName {
target = ep
} else {
sources = append(sources, ep)
}
}
if target.Client == nil {
log.Fatalf("target endpoint %q not found", opt.TargetName)
}
log.Printf("[otterio] ensure target bucket: %s", opt.Bucket)
if err := ensureBucket(ctx, target.Client, opt.Bucket); err != nil {
log.Fatalf("[otterio] create bucket: %v", err)
}
if opt.Cleanup {
log.Printf("[cleanup] delete all objects in bucket %s on all endpoints", opt.Bucket)
for _, ep := range endpoints {
if err := deleteBucketObjects(ctx, ep.Client, opt.Bucket); err != nil {
log.Fatalf("[%s] cleanup failed: %v", ep.Name, err)
}
}
}
var allRecords []ObjectRecord
for _, source := range sources {
records, err := uploadRandomObjects(ctx, source, opt)
if err != nil {
log.Fatal(err)
}
allRecords = append(allRecords, records...)
}
for _, source := range sources {
for _, record := range allRecords {
if record.SourceName != source.Name {
continue
}
if err := syncOneObject(ctx, source, target, record, opt); err != nil {
log.Fatal(err)
}
}
}
verifiedCount := 0
for _, record := range allRecords {
if err := verifyTargetObject(ctx, target, opt, record); err != nil {
log.Fatal(err)
}
verifiedCount++
}
if opt.ListObjects {
if err := listTargetObjects(ctx, target, opt.Bucket); err != nil {
log.Fatal(err)
}
}
if err := printSyncSummary(ctx, target, sources, opt, allRecords, verifiedCount); err != nil {
log.Fatal(err)
}
}
最后
这轮先折腾到这里。有需要的朋友可以更新试试,也欢迎带上版本号、部署方式和复现步骤来提 issue。
朴素的对象存储,还是安安静静地干活就好了。
—EOF