Carry approved OSS bucket through upload grants
This commit is contained in:
@@ -48,7 +48,7 @@ func TestRecordingNotificationRecoveryDoesNotPUTAgain(t *testing.T) {
|
||||
sum := sha256.Sum256([]byte("audio"))
|
||||
binding := &agentpb.ExecutionBinding{TenantId: "tenant-a", TenantKey: "tenant-a", ExecutionId: "execution-a", CallId: "call-a"}
|
||||
asset := &agentpb.AssetDescriptor{AssetId: "recording-a", SizeBytes: 5, ChecksumSha256: hex.EncodeToString(sum[:])}
|
||||
remote := &uploadRPCStub{pending: true, grant: &agentpb.UploadGrant{ObjectKey: "recording-a", TargetUrl: server.URL, MaxBytes: 5, RequiredChecksumSha256: asset.ChecksumSha256, ExpiresAtUnixMs: time.Now().Add(time.Minute).UnixMilli()}}
|
||||
remote := &uploadRPCStub{pending: true, grant: &agentpb.UploadGrant{ObjectKey: "recording-a", Bucket: "mock-bucket", TargetUrl: server.URL, MaxBytes: 5, RequiredChecksumSha256: asset.ChecksumSha256, ExpiresAtUnixMs: time.Now().Add(time.Minute).UnixMilli()}}
|
||||
spool, err := agent.NewSpool(root, nil)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
|
||||
@@ -35,7 +35,7 @@ func TestExplicitUploadRetryIsOneNewRequestAndOnePUT(t *testing.T) {
|
||||
sum := sha256.Sum256([]byte("audio"))
|
||||
binding := &agentpb.ExecutionBinding{TenantId: "tenant-a", TenantKey: "tenant-a", ExecutionId: "execution-a", CallId: "call-a"}
|
||||
asset := &agentpb.AssetDescriptor{AssetId: "recording-a", SizeBytes: 5, ChecksumSha256: hex.EncodeToString(sum[:])}
|
||||
remote := &uploadRPCStub{grant: &agentpb.UploadGrant{ObjectKey: "recording-a", TargetUrl: server.URL, MaxBytes: 5, RequiredChecksumSha256: asset.ChecksumSha256, ExpiresAtUnixMs: time.Now().Add(-time.Minute).UnixMilli()}}
|
||||
remote := &uploadRPCStub{grant: &agentpb.UploadGrant{ObjectKey: "recording-a", Bucket: "mock-bucket", TargetUrl: server.URL, MaxBytes: 5, RequiredChecksumSha256: asset.ChecksumSha256, ExpiresAtUnixMs: time.Now().Add(-time.Minute).UnixMilli()}}
|
||||
spool, err := agent.NewSpool(root, nil)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
|
||||
@@ -59,6 +59,7 @@
|
||||
## P06:录音/OSS/最终结果(进行中,未签收)
|
||||
|
||||
- 已新增 Agent 内存录音单次 PUT 组件 `UploadClient.UploadBytes`:复用受限授权、大小和 SHA-256 核对,成功路径不写临时录音或通话结果文件。OSS 明确拒绝保留状态码且不自动重试;传输结果不明时返回专用错误、停止自动重试并隐藏带签名的 URL。空授权头拒绝而非使服务崩溃。这里只验证组件,尚未连接实际录音或最终结果。
|
||||
- 内部 OSS 授权增加 Dispatcher 原始 bucket,官方 SDK 签发时返回获批 bucket;Agent 收到缺少 bucket 的授权会在 PUT 前拒绝,失败恢复的显式重申请若返回了不同 bucket,也在再次 PUT 前拒绝。两条红灯测试证明先前会错误上传;修复后全包测试及 Agent/RPC/OSS race 测试通过。这里只核验本地授权载体,不代表新主入口已完成签发或真实 OSS 已验证。
|
||||
- Agent 失败恢复隔离组件:仅在 OSS PUT 明确失败后,把原 bucket/object_key 对应的录音和通话信息两文件写入私有目录并同步落盘;双文件缺失或损坏明确报错、不伪造结果。完成保存后固定 48 小时窗口,按 1 分钟递增至最长 1 小时重试;到期保留原文件。PUT 前持久写入 in-flight,结果不明或进程重启不会二次 PUT;确认上传后先持久记录成功,再经注入的 Mock 回报通话结果,回报失败/重启仅重发原结果。启动扫描识别遗漏文件并提供不暴露原路径的稳定摘要;真实 Dispatcher 授权及回报尚未接线。
|
||||
- Dispatcher 的无录音最终结果隔离组件:`CurrentStore.RecordCallResult` 仅在确认通话结束后,按持久任务快照校验任务、被叫、主叫和已选线路,并以源执行事件固定生成唯一最终结果身份;消息通过严格 MQ Schema 校验后与 outbox 在同一事务写入。同内容重投/重启只恢复原消息,冲突结果和 SQLite 写入失败均不会产生第二份结果。这里只验证隔离组件,Agent 实际回报尚未连通。
|
||||
- Dispatcher 的原始 OSS 目标及已上传结果隔离组件:新 SQLite 布局把一次通话的 upload_id、bucket、object_key、录音格式/时长/大小和 SHA-256 唯一绑定到已保留的执行;不保存临时 URL 或 TOKEN。旧布局拒绝启动并原样保留待交付 outbox,不自动迁移或清理。已签发录音目标不能通过空录音结果绕过上传;已有空录音最终结果不能再签发录音目标。`RecordUploadedCallResult` 仅接受与持久绑定完全一致的录音事实及 Agent 所报告的成功 PUT 状态,录音确认与唯一最终结果 outbox 同事务提交;丢失回报或 MQ 投递时重用原消息,已确认后拒绝再次签发 PUT 授权。Mock 证明的是本地状态约束,不是独立 OSS 校验或真实 Agent 身份验证。
|
||||
|
||||
+11
-2
@@ -4006,6 +4006,7 @@ type UploadGrant struct {
|
||||
ObjectKey string `protobuf:"bytes,5,opt,name=object_key,json=objectKey,proto3" json:"object_key,omitempty"`
|
||||
RequiredChecksumSha256 string `protobuf:"bytes,6,opt,name=required_checksum_sha256,json=requiredChecksumSha256,proto3" json:"required_checksum_sha256,omitempty"`
|
||||
MaxBytes int64 `protobuf:"varint,7,opt,name=max_bytes,json=maxBytes,proto3" json:"max_bytes,omitempty"`
|
||||
Bucket string `protobuf:"bytes,8,opt,name=bucket,proto3" json:"bucket,omitempty"`
|
||||
unknownFields protoimpl.UnknownFields
|
||||
sizeCache protoimpl.SizeCache
|
||||
}
|
||||
@@ -4089,6 +4090,13 @@ func (x *UploadGrant) GetMaxBytes() int64 {
|
||||
return 0
|
||||
}
|
||||
|
||||
func (x *UploadGrant) GetBucket() string {
|
||||
if x != nil {
|
||||
return x.Bucket
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
type RequestUploadResponse struct {
|
||||
state protoimpl.MessageState `protogen:"open.v1"`
|
||||
Receipt *OperationReceipt `protobuf:"bytes,1,opt,name=receipt,proto3" json:"receipt,omitempty"`
|
||||
@@ -4587,7 +4595,7 @@ const file_agent_agent_proto_rawDesc = "" +
|
||||
"\x04meta\x18\x01 \x01(\v2\x12.agent.RequestMetaR\x04meta\x121\n" +
|
||||
"\abinding\x18\x02 \x01(\v2\x17.agent.ExecutionBindingR\abinding\x12,\n" +
|
||||
"\x05asset\x18\x03 \x01(\v2\x16.agent.AssetDescriptorR\x05asset\x12\x1b\n" +
|
||||
"\tupload_id\x18\x04 \x01(\tR\buploadId\"\x95\x02\n" +
|
||||
"\tupload_id\x18\x04 \x01(\tR\buploadId\"\xad\x02\n" +
|
||||
"\vUploadGrant\x12\x1b\n" +
|
||||
"\tupload_id\x18\x01 \x01(\tR\buploadId\x12\x1d\n" +
|
||||
"\n" +
|
||||
@@ -4597,7 +4605,8 @@ const file_agent_agent_proto_rawDesc = "" +
|
||||
"\n" +
|
||||
"object_key\x18\x05 \x01(\tR\tobjectKey\x128\n" +
|
||||
"\x18required_checksum_sha256\x18\x06 \x01(\tR\x16requiredChecksumSha256\x12\x1b\n" +
|
||||
"\tmax_bytes\x18\a \x01(\x03R\bmaxBytes\"\x9e\x01\n" +
|
||||
"\tmax_bytes\x18\a \x01(\x03R\bmaxBytes\x12\x16\n" +
|
||||
"\x06bucket\x18\b \x01(\tR\x06bucket\"\x9e\x01\n" +
|
||||
"\x15RequestUploadResponse\x121\n" +
|
||||
"\areceipt\x18\x01 \x01(\v2\x17.agent.OperationReceiptR\areceipt\x12(\n" +
|
||||
"\x05grant\x18\x02 \x01(\v2\x12.agent.UploadGrantR\x05grant\x12(\n" +
|
||||
|
||||
@@ -111,7 +111,7 @@ func (r *RecordingRecovery) Retry(ctx context.Context, bucket, objectKey string,
|
||||
}
|
||||
|
||||
func (r *RecordingRecovery) validateTarget(ctx context.Context, entry RecordingRecoveryEntry, target RecoveryTarget) error {
|
||||
if target.Grant == nil || target.Bucket != entry.Bucket || target.Grant.UploadId != entry.UploadID || target.Grant.ObjectKey != entry.ObjectKey || target.Grant.MaxBytes < entry.SizeBytes || target.Grant.MaxBytes <= 0 || !strings.EqualFold(target.Grant.RequiredChecksumSha256, entry.SHA256) {
|
||||
if target.Grant == nil || target.Bucket != entry.Bucket || target.Grant.Bucket != entry.Bucket || target.Grant.UploadId != entry.UploadID || target.Grant.ObjectKey != entry.ObjectKey || target.Grant.MaxBytes < entry.SizeBytes || target.Grant.MaxBytes <= 0 || !strings.EqualFold(target.Grant.RequiredChecksumSha256, entry.SHA256) {
|
||||
return errors.New("recovery grant target does not match the original OSS asset")
|
||||
}
|
||||
if _, err := r.Upload.validateGrant(ctx, target.Grant); err != nil {
|
||||
|
||||
@@ -18,7 +18,7 @@ import (
|
||||
|
||||
func recoveryTarget(now time.Time, entry RecordingRecoveryEntry, signedURL string) RecoveryTarget {
|
||||
return RecoveryTarget{Bucket: entry.Bucket, Grant: &agentpb.UploadGrant{
|
||||
UploadId: entry.UploadID, ObjectKey: entry.ObjectKey, TargetUrl: signedURL,
|
||||
UploadId: entry.UploadID, Bucket: entry.Bucket, ObjectKey: entry.ObjectKey, TargetUrl: signedURL,
|
||||
ExpiresAtUnixMs: now.Add(15 * time.Minute).UnixMilli(), MaxBytes: entry.SizeBytes,
|
||||
RequiredChecksumSha256: entry.SHA256,
|
||||
}}
|
||||
@@ -168,6 +168,33 @@ func TestRecordingRecoveryRejectsChangedBucketBeforePUT(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestRecordingRecoveryRejectsChangedSignedGrantBucketBeforePUT(t *testing.T) {
|
||||
now := time.Date(2026, 10, 1, 10, 0, 0, 0, time.UTC)
|
||||
var puts atomic.Int32
|
||||
server := httptest.NewServer(http.HandlerFunc(func(http.ResponseWriter, *http.Request) { puts.Add(1) }))
|
||||
defer server.Close()
|
||||
recovery := RecordingRecovery{Root: privateRecoveryRoot(t), Now: func() time.Time { return now }, Upload: UploadClient{AllowInsecureHTTP: true, Now: func() time.Time { return now }}}
|
||||
entry, err := recovery.SaveFailure([]byte("RIFF original audio"), failedRecordingFixture(), &UploadHTTPError{StatusCode: 503})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
now = now.Add(time.Minute)
|
||||
_, err = recovery.Retry(context.Background(), entry.Bucket, entry.ObjectKey,
|
||||
func(context.Context, RecordingRecoveryEntry) (RecoveryTarget, error) {
|
||||
target := recoveryTarget(now, entry, server.URL+"/original?signature=DO_NOT_LOG")
|
||||
target.Grant.Bucket = "changed-bucket"
|
||||
return target, nil
|
||||
},
|
||||
func(context.Context, RecordingRecoveryEntry) error { return errors.New("unexpected report") })
|
||||
if err == nil || !strings.Contains(err.Error(), "target") || puts.Load() != 0 {
|
||||
t.Fatalf("changed signed-grant bucket reached OSS: puts=%d err=%v", puts.Load(), err)
|
||||
}
|
||||
pending, err := recovery.Load(entry.Bucket, entry.ObjectKey)
|
||||
if err != nil || pending.State != "retry_pending" {
|
||||
t.Fatalf("changed grant moved original recording: state=%s err=%v", pending.State, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRecordingRecoveryPersistsInFlightBeforeUnknownPUT(t *testing.T) {
|
||||
now := time.Date(2026, 10, 1, 10, 0, 0, 0, time.UTC)
|
||||
root := privateRecoveryRoot(t)
|
||||
|
||||
@@ -66,8 +66,8 @@ func (c UploadClient) validateGrant(ctx context.Context, grant *agentpb.UploadGr
|
||||
if grant == nil {
|
||||
return nil, fmt.Errorf("%w: grant is required", ErrUploadGrantInvalid)
|
||||
}
|
||||
if grant.TargetUrl == "" || grant.UploadId == "" || grant.ObjectKey == "" {
|
||||
return nil, fmt.Errorf("%w: URL, ID and object key are required", ErrUploadGrantInvalid)
|
||||
if grant.TargetUrl == "" || grant.UploadId == "" || grant.ObjectKey == "" || grant.Bucket == "" {
|
||||
return nil, fmt.Errorf("%w: URL, ID, bucket and object key are required", ErrUploadGrantInvalid)
|
||||
}
|
||||
if grant.ExpiresAtUnixMs <= 0 {
|
||||
return nil, fmt.Errorf("%w: expiry is required", ErrUploadGrantInvalid)
|
||||
|
||||
@@ -36,7 +36,7 @@ func TestUploadBytesPutsMemoryRecordingOnceWithoutBusinessFile(t *testing.T) {
|
||||
defer server.Close()
|
||||
noFiles := t.TempDir()
|
||||
t.Setenv("TMPDIR", noFiles)
|
||||
grant := &agentpb.UploadGrant{UploadId: "upload-mem", ObjectKey: "recordings/call-1.wav", TargetUrl: server.URL + "/file?signature=DO_NOT_LOG", ExpiresAtUnixMs: time.Now().Add(time.Minute).UnixMilli(), MaxBytes: int64(len(body)), RequiredChecksumSha256: hex.EncodeToString(sum[:]), Headers: []*agentpb.Header{{Name: "x-upload-token", Value: "mock-token"}}}
|
||||
grant := &agentpb.UploadGrant{UploadId: "upload-mem", Bucket: "mock-bucket", ObjectKey: "recordings/call-1.wav", TargetUrl: server.URL + "/file?signature=DO_NOT_LOG", ExpiresAtUnixMs: time.Now().Add(time.Minute).UnixMilli(), MaxBytes: int64(len(body)), RequiredChecksumSha256: hex.EncodeToString(sum[:]), Headers: []*agentpb.Header{{Name: "x-upload-token", Value: "mock-token"}}}
|
||||
result, err := (UploadClient{AllowInsecureHTTP: true}).UploadBytes(context.Background(), grant, body)
|
||||
if err != nil || requests.Load() != 1 || !bytes.Equal(received, body) || result.SHA256 != grant.RequiredChecksumSha256 || result.SizeBytes != int64(len(body)) || result.ETag != "mock-etag" {
|
||||
t.Fatalf("normal path must PUT original bytes once: result=%+v requests=%d err=%v", result, requests.Load(), err)
|
||||
@@ -47,12 +47,24 @@ func TestUploadBytesPutsMemoryRecordingOnceWithoutBusinessFile(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestUploadBytesRequiresOriginalBucketBeforePUT(t *testing.T) {
|
||||
var requests atomic.Int32
|
||||
server := httptest.NewServer(http.HandlerFunc(func(http.ResponseWriter, *http.Request) { requests.Add(1) }))
|
||||
defer server.Close()
|
||||
body := []byte("recording")
|
||||
sum := sha256.Sum256(body)
|
||||
grant := &agentpb.UploadGrant{UploadId: "upload-no-bucket", ObjectKey: "recordings/call.wav", TargetUrl: server.URL, ExpiresAtUnixMs: time.Now().Add(time.Minute).UnixMilli(), MaxBytes: int64(len(body)), RequiredChecksumSha256: hex.EncodeToString(sum[:])}
|
||||
if _, err := (UploadClient{AllowInsecureHTTP: true}).UploadBytes(context.Background(), grant, body); !errors.Is(err, ErrUploadGrantInvalid) || requests.Load() != 0 {
|
||||
t.Fatalf("missing D-owned original bucket sent a PUT: attempts=%d err=%v", requests.Load(), err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUploadBytesRejectsUnchangedGrantMismatchBeforePUT(t *testing.T) {
|
||||
var requests atomic.Int32
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { requests.Add(1) }))
|
||||
defer server.Close()
|
||||
data := []byte("original recording")
|
||||
grant := &agentpb.UploadGrant{UploadId: "upload-mismatch", ObjectKey: "recordings/call-2.wav", TargetUrl: server.URL, ExpiresAtUnixMs: time.Now().Add(time.Minute).UnixMilli(), MaxBytes: int64(len(data)), RequiredChecksumSha256: strings.Repeat("f", 64)}
|
||||
grant := &agentpb.UploadGrant{UploadId: "upload-mismatch", Bucket: "mock-bucket", ObjectKey: "recordings/call-2.wav", TargetUrl: server.URL, ExpiresAtUnixMs: time.Now().Add(time.Minute).UnixMilli(), MaxBytes: int64(len(data)), RequiredChecksumSha256: strings.Repeat("f", 64)}
|
||||
_, err := (UploadClient{AllowInsecureHTTP: true}).UploadBytes(context.Background(), grant, data)
|
||||
if !errors.Is(err, ErrUploadChecksumMismatch) || requests.Load() != 0 || string(data) != "original recording" {
|
||||
t.Fatalf("mismatched grant caused a PUT or changed recording: requests=%d err=%v", requests.Load(), err)
|
||||
@@ -67,7 +79,7 @@ func TestUploadBytesDefinitiveHTTPFailureIsOneAttempt(t *testing.T) {
|
||||
_, _ = io.WriteString(w, "DO_NOT_LOG_RESPONSE")
|
||||
}))
|
||||
defer server.Close()
|
||||
grant := &agentpb.UploadGrant{UploadId: "upload-fail", ObjectKey: "recordings/call-3.wav", TargetUrl: server.URL, ExpiresAtUnixMs: time.Now().Add(time.Minute).UnixMilli(), MaxBytes: 100}
|
||||
grant := &agentpb.UploadGrant{UploadId: "upload-fail", Bucket: "mock-bucket", ObjectKey: "recordings/call-3.wav", TargetUrl: server.URL, ExpiresAtUnixMs: time.Now().Add(time.Minute).UnixMilli(), MaxBytes: 100}
|
||||
_, err := (UploadClient{AllowInsecureHTTP: true}).UploadBytes(context.Background(), grant, []byte("recording"))
|
||||
var rejected *UploadHTTPError
|
||||
if !errors.As(err, &rejected) || rejected.StatusCode != http.StatusServiceUnavailable || requests.Load() != 1 || strings.Contains(err.Error(), "DO_NOT_LOG_RESPONSE") {
|
||||
@@ -85,7 +97,7 @@ func (t ambiguousBytesTransport) RoundTrip(req *http.Request) (*http.Response, e
|
||||
func TestUploadBytesAmbiguousTransportFailureIsNotRetriedOrLeaked(t *testing.T) {
|
||||
var requests atomic.Int32
|
||||
client := &http.Client{Transport: ambiguousBytesTransport{requests: &requests}}
|
||||
grant := &agentpb.UploadGrant{UploadId: "upload-unknown", ObjectKey: "recordings/call-4.wav", TargetUrl: "https://oss.example.invalid/file?signature=DO_NOT_LOG", ExpiresAtUnixMs: time.Now().Add(time.Minute).UnixMilli(), MaxBytes: 100}
|
||||
grant := &agentpb.UploadGrant{UploadId: "upload-unknown", Bucket: "mock-bucket", ObjectKey: "recordings/call-4.wav", TargetUrl: "https://oss.example.invalid/file?signature=DO_NOT_LOG", ExpiresAtUnixMs: time.Now().Add(time.Minute).UnixMilli(), MaxBytes: 100}
|
||||
_, err := (UploadClient{HTTPClient: client}).UploadBytes(context.Background(), grant, []byte("recording"))
|
||||
if !errors.Is(err, ErrUploadOutcomeUnknown) || requests.Load() != 1 || strings.Contains(err.Error(), "DO_NOT_LOG") || strings.Contains(err.Error(), "signature=") {
|
||||
t.Fatalf("ambiguous upload must halt without token disclosure: requests=%d err=%v", requests.Load(), err)
|
||||
@@ -93,7 +105,7 @@ func TestUploadBytesAmbiguousTransportFailureIsNotRetriedOrLeaked(t *testing.T)
|
||||
}
|
||||
|
||||
func TestUploadBytesRejectsNilGrantHeaderBeforePUT(t *testing.T) {
|
||||
grant := &agentpb.UploadGrant{UploadId: "upload-bad-header", ObjectKey: "recordings/call-5.wav", TargetUrl: "https://oss.example.invalid/file", ExpiresAtUnixMs: time.Now().Add(time.Minute).UnixMilli(), MaxBytes: 100, Headers: []*agentpb.Header{nil}}
|
||||
grant := &agentpb.UploadGrant{UploadId: "upload-bad-header", Bucket: "mock-bucket", ObjectKey: "recordings/call-5.wav", TargetUrl: "https://oss.example.invalid/file", ExpiresAtUnixMs: time.Now().Add(time.Minute).UnixMilli(), MaxBytes: 100, Headers: []*agentpb.Header{nil}}
|
||||
_, err := (UploadClient{}).UploadBytes(context.Background(), grant, []byte("recording"))
|
||||
if !errors.Is(err, ErrUploadGrantInvalid) {
|
||||
t.Fatalf("malformed grant header must be rejected without a panic: %v", err)
|
||||
|
||||
@@ -36,7 +36,7 @@ func TestUploadDoesNotReportPreReadDigestForChangedBytes(t *testing.T) {
|
||||
}
|
||||
return &http.Response{StatusCode: 200, Header: make(http.Header), Body: io.NopCloser(strings.NewReader(""))}, nil
|
||||
})}
|
||||
grant := &agentpb.UploadGrant{UploadId: "upload-a", ObjectKey: "recording", TargetUrl: "https://oss.invalid/object", MaxBytes: 8, RequiredChecksumSha256: hex.EncodeToString(sum[:]), ExpiresAtUnixMs: time.Now().Add(time.Minute).UnixMilli()}
|
||||
grant := &agentpb.UploadGrant{UploadId: "upload-a", Bucket: "mock-bucket", ObjectKey: "recording", TargetUrl: "https://oss.invalid/object", MaxBytes: 8, RequiredChecksumSha256: hex.EncodeToString(sum[:]), ExpiresAtUnixMs: time.Now().Add(time.Minute).UnixMilli()}
|
||||
if _, err := (UploadClient{HTTPClient: client}).UploadFile(context.Background(), grant, path); err == nil {
|
||||
t.Fatal("reported original checksum after sending changed bytes")
|
||||
}
|
||||
|
||||
@@ -24,7 +24,7 @@ func TestUploadTransportFailureIsSingleAttemptAndRedactsSignedURL(t *testing.T)
|
||||
}
|
||||
var attempts int
|
||||
client := &http.Client{Transport: uploadFailureTransport{calls: &attempts}}
|
||||
grant := &agentpb.UploadGrant{UploadId: "upload-error", ObjectKey: "recording", TargetUrl: "https://oss.example.invalid/object?signature=DO_NOT_LOG", MaxBytes: 5, ExpiresAtUnixMs: time.Now().Add(time.Minute).UnixMilli()}
|
||||
grant := &agentpb.UploadGrant{UploadId: "upload-error", Bucket: "mock-bucket", ObjectKey: "recording", TargetUrl: "https://oss.example.invalid/object?signature=DO_NOT_LOG", MaxBytes: 5, ExpiresAtUnixMs: time.Now().Add(time.Minute).UnixMilli()}
|
||||
_, err := (UploadClient{HTTPClient: client}).UploadFile(context.Background(), grant, path)
|
||||
if err == nil {
|
||||
t.Fatal("transport failure hidden")
|
||||
@@ -63,7 +63,7 @@ func TestUploadClientUsesGrantAndVerifiesChecksum(t *testing.T) {
|
||||
if err := os.WriteFile(assetPath, body, 0o600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
grant := &agentpb.UploadGrant{UploadId: "upload-1", TargetUrl: server.URL, ObjectKey: "recording-1", ExpiresAtUnixMs: time.Unix(101, 0).UnixMilli(), Headers: []*agentpb.Header{{Name: "x-upload-token", Value: "mock-token"}}, RequiredChecksumSha256: hex.EncodeToString(digest[:]), MaxBytes: int64(len(body))}
|
||||
grant := &agentpb.UploadGrant{UploadId: "upload-1", Bucket: "mock-bucket", TargetUrl: server.URL, ObjectKey: "recording-1", ExpiresAtUnixMs: time.Unix(101, 0).UnixMilli(), Headers: []*agentpb.Header{{Name: "x-upload-token", Value: "mock-token"}}, RequiredChecksumSha256: hex.EncodeToString(digest[:]), MaxBytes: int64(len(body))}
|
||||
result, err := (UploadClient{AllowInsecureHTTP: true, Now: func() time.Time { return time.Unix(100, 0) }, AllowedHosts: map[string]struct{}{strings.TrimPrefix(server.URL, "http://"): {}}}).UploadFile(context.Background(), grant, assetPath)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
@@ -81,7 +81,7 @@ func TestUploadClientFailsClosedForGrantMismatchAndHTTP(t *testing.T) {
|
||||
if err := os.WriteFile(assetPath, []byte("bytes"), 0o600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
grant := &agentpb.UploadGrant{UploadId: "upload-2", TargetUrl: "http://127.0.0.1:1/upload", ObjectKey: "recording-2", ExpiresAtUnixMs: time.Now().Add(time.Hour).UnixMilli(), RequiredChecksumSha256: strings.Repeat("a", 64), MaxBytes: 1024}
|
||||
grant := &agentpb.UploadGrant{UploadId: "upload-2", Bucket: "mock-bucket", TargetUrl: "http://127.0.0.1:1/upload", ObjectKey: "recording-2", ExpiresAtUnixMs: time.Now().Add(time.Hour).UnixMilli(), RequiredChecksumSha256: strings.Repeat("a", 64), MaxBytes: 1024}
|
||||
if _, err := (UploadClient{AllowInsecureHTTP: true}).UploadFile(context.Background(), grant, assetPath); !errors.Is(err, ErrUploadChecksumMismatch) {
|
||||
t.Fatalf("expected typed checksum mismatch, got %v", err)
|
||||
}
|
||||
@@ -103,7 +103,7 @@ func TestUploadClientHTTPFailureHasStatusWithoutResponseBody(t *testing.T) {
|
||||
}))
|
||||
defer server.Close()
|
||||
sum := sha256.Sum256(body)
|
||||
grant := &agentpb.UploadGrant{UploadId: "upload-forbidden", TargetUrl: server.URL, ObjectKey: "recording-forbidden", ExpiresAtUnixMs: time.Now().Add(time.Minute).UnixMilli(), RequiredChecksumSha256: hex.EncodeToString(sum[:]), MaxBytes: 1024}
|
||||
grant := &agentpb.UploadGrant{UploadId: "upload-forbidden", Bucket: "mock-bucket", TargetUrl: server.URL, ObjectKey: "recording-forbidden", ExpiresAtUnixMs: time.Now().Add(time.Minute).UnixMilli(), RequiredChecksumSha256: hex.EncodeToString(sum[:]), MaxBytes: 1024}
|
||||
_, err := (UploadClient{AllowInsecureHTTP: true}).UploadFile(context.Background(), grant, assetPath)
|
||||
var response *UploadHTTPError
|
||||
if !errors.As(err, &response) || response.StatusCode != http.StatusForbidden || strings.Contains(err.Error(), "DO_NOT_LOG_SECRET") {
|
||||
@@ -116,7 +116,7 @@ func TestUploadClientEnforcesSizeAndHost(t *testing.T) {
|
||||
if err := os.WriteFile(assetPath, []byte("bytes"), 0o600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
grant := &agentpb.UploadGrant{UploadId: "upload-3", TargetUrl: "https://oss.example.invalid/upload", ObjectKey: "recording-3", ExpiresAtUnixMs: time.Now().Add(time.Hour).UnixMilli(), MaxBytes: 1}
|
||||
grant := &agentpb.UploadGrant{UploadId: "upload-3", Bucket: "mock-bucket", TargetUrl: "https://oss.example.invalid/upload", ObjectKey: "recording-3", ExpiresAtUnixMs: time.Now().Add(time.Hour).UnixMilli(), MaxBytes: 1}
|
||||
if _, err := (UploadClient{}).UploadFile(context.Background(), grant, assetPath); err == nil {
|
||||
t.Fatal("expected size rejection")
|
||||
}
|
||||
@@ -131,7 +131,7 @@ func TestUploadClientRejectsExpiredGrant(t *testing.T) {
|
||||
if err := os.WriteFile(assetPath, []byte("bytes"), 0o600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
grant := &agentpb.UploadGrant{UploadId: "upload-expired", TargetUrl: "https://oss.example.invalid/upload", ObjectKey: "recording-expired", ExpiresAtUnixMs: time.Unix(100, 0).UnixMilli(), MaxBytes: 1024}
|
||||
grant := &agentpb.UploadGrant{UploadId: "upload-expired", Bucket: "mock-bucket", TargetUrl: "https://oss.example.invalid/upload", ObjectKey: "recording-expired", ExpiresAtUnixMs: time.Unix(100, 0).UnixMilli(), MaxBytes: 1024}
|
||||
if _, err := (UploadClient{Now: func() time.Time { return time.Unix(100, 0) }}).UploadFile(context.Background(), grant, assetPath); !errors.Is(err, ErrUploadGrantExpired) {
|
||||
t.Fatalf("expected typed expired grant rejection, got %v", err)
|
||||
}
|
||||
|
||||
@@ -115,6 +115,7 @@ func (c *Client) Grant(ctx context.Context, uploadID, objectKey, checksum string
|
||||
Headers: headers,
|
||||
ExpiresAtUnixMs: presigned.Expiration.UnixMilli(),
|
||||
ObjectKey: objectKey,
|
||||
Bucket: c.config.Bucket,
|
||||
RequiredChecksumSha256: checksum,
|
||||
MaxBytes: maxBytes,
|
||||
}, nil
|
||||
|
||||
@@ -37,6 +37,9 @@ func TestGrantLifetimeIsFixed(t *testing.T) {
|
||||
if grant.MaxBytes != 100 {
|
||||
t.Fatal("requested object size was changed")
|
||||
}
|
||||
if grant.GetBucket() != cfg.Bucket {
|
||||
t.Fatalf("Agent was not given the Dispatcher-owned original OSS bucket: %q", grant.GetBucket())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -1009,7 +1009,7 @@ func (s *Server) RequestUpload(ctx context.Context, req *agentpb.RequestUploadRe
|
||||
}
|
||||
return &agentpb.RequestUploadResponse{Receipt: s.receipt(req.Meta, agentpb.ResultCode_RESULT_CODE_ACCEPTED, agentpb.FailureCode_FAILURE_CODE_UNSPECIFIED, "duplicate upload request", false), Grant: proto.Clone(existing.grant).(*agentpb.UploadGrant), State: existing.state}, nil
|
||||
}
|
||||
grant := &agentpb.UploadGrant{UploadId: req.UploadId, TargetUrl: "https://oss.mock.invalid/upload/" + req.UploadId, ExpiresAtUnixMs: s.now().Add(5 * time.Minute).UnixMilli(), ObjectKey: req.Asset.AssetId, RequiredChecksumSha256: req.Asset.ChecksumSha256, MaxBytes: s.uploadPolicy.MaxAssetBytes}
|
||||
grant := &agentpb.UploadGrant{UploadId: req.UploadId, TargetUrl: "https://oss.mock.invalid/upload/" + req.UploadId, ExpiresAtUnixMs: s.now().Add(5 * time.Minute).UnixMilli(), ObjectKey: req.Asset.AssetId, Bucket: "mock-bucket", RequiredChecksumSha256: req.Asset.ChecksumSha256, MaxBytes: s.uploadPolicy.MaxAssetBytes}
|
||||
if grant.MaxBytes == 0 {
|
||||
grant.MaxBytes = req.Asset.SizeBytes
|
||||
}
|
||||
|
||||
@@ -495,6 +495,7 @@ message UploadGrant {
|
||||
string object_key = 5;
|
||||
string required_checksum_sha256 = 6;
|
||||
int64 max_bytes = 7;
|
||||
string bucket = 8;
|
||||
}
|
||||
|
||||
message RequestUploadResponse {
|
||||
|
||||
+4
-4
@@ -34,13 +34,13 @@
|
||||
},
|
||||
{
|
||||
"path": "proto/agent/agent.proto",
|
||||
"bytes": 13260,
|
||||
"sha256": "28e8702ca800ee696ef924eb1d79a2d60ae6217ee7c1d91c1921125c03661b43"
|
||||
"bytes": 13281,
|
||||
"sha256": "dc086c1963fa898d39f1f9e10de85abf3c4449834a84f0a87ec7b5e7aa281f91"
|
||||
},
|
||||
{
|
||||
"path": "gen/agent/agent.pb.go",
|
||||
"bytes": 163063,
|
||||
"sha256": "b66939cfb70bc659813238bb4a74c7e633e5551efbe85a00eeafdcf39ff8402d"
|
||||
"bytes": 163322,
|
||||
"sha256": "a9da5dc255d2c5d3af736fcdd31909ebd59409f07bb4cef4de578d7baaf40563"
|
||||
},
|
||||
{
|
||||
"path": "gen/agent/agent_grpc.pb.go",
|
||||
|
||||
Reference in New Issue
Block a user