diff --git a/cmd/sip-go-agent/upload_attempt_test.go b/cmd/sip-go-agent/upload_attempt_test.go index ae06de1..99d8e83 100644 --- a/cmd/sip-go-agent/upload_attempt_test.go +++ b/cmd/sip-go-agent/upload_attempt_test.go @@ -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) diff --git a/cmd/sip-go-agent/upload_retry_test.go b/cmd/sip-go-agent/upload_retry_test.go index 752cd18..902af06 100644 --- a/cmd/sip-go-agent/upload_retry_test.go +++ b/cmd/sip-go-agent/upload_retry_test.go @@ -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) diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index 70945e5..174e62f 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -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 身份验证。 diff --git a/gen/agent/agent.pb.go b/gen/agent/agent.pb.go index 1459fbd..69a5db9 100644 --- a/gen/agent/agent.pb.go +++ b/gen/agent/agent.pb.go @@ -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" + diff --git a/internal/agent/recording_retry.go b/internal/agent/recording_retry.go index dbd1d48..6c0d830 100644 --- a/internal/agent/recording_retry.go +++ b/internal/agent/recording_retry.go @@ -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 { diff --git a/internal/agent/recording_retry_test.go b/internal/agent/recording_retry_test.go index 79ec2e3..7e1b062 100644 --- a/internal/agent/recording_retry_test.go +++ b/internal/agent/recording_retry_test.go @@ -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) diff --git a/internal/agent/upload.go b/internal/agent/upload.go index 81c8c4e..80bc67a 100644 --- a/internal/agent/upload.go +++ b/internal/agent/upload.go @@ -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) diff --git a/internal/agent/upload_bytes_test.go b/internal/agent/upload_bytes_test.go index ca8f1b3..f63d56d 100644 --- a/internal/agent/upload_bytes_test.go +++ b/internal/agent/upload_bytes_test.go @@ -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) diff --git a/internal/agent/upload_integrity_test.go b/internal/agent/upload_integrity_test.go index c80d062..4971c46 100644 --- a/internal/agent/upload_integrity_test.go +++ b/internal/agent/upload_integrity_test.go @@ -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") } diff --git a/internal/agent/upload_test.go b/internal/agent/upload_test.go index bab7d06..999c84c 100644 --- a/internal/agent/upload_test.go +++ b/internal/agent/upload_test.go @@ -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) } diff --git a/internal/oss/aliyun.go b/internal/oss/aliyun.go index ca9320e..64f0e1f 100644 --- a/internal/oss/aliyun.go +++ b/internal/oss/aliyun.go @@ -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 diff --git a/internal/oss/aliyun_test.go b/internal/oss/aliyun_test.go index 18af23f..5e69cb4 100644 --- a/internal/oss/aliyun_test.go +++ b/internal/oss/aliyun_test.go @@ -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()) + } } } diff --git a/internal/rpc/server.go b/internal/rpc/server.go index 58abd71..504de21 100644 --- a/internal/rpc/server.go +++ b/internal/rpc/server.go @@ -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 } diff --git a/proto/agent/agent.proto b/proto/agent/agent.proto index c9a4a97..acf2560 100644 --- a/proto/agent/agent.proto +++ b/proto/agent/agent.proto @@ -495,6 +495,7 @@ message UploadGrant { string object_key = 5; string required_checksum_sha256 = 6; int64 max_bytes = 7; + string bucket = 8; } message RequestUploadResponse { diff --git a/proto/manifest.json b/proto/manifest.json index fd55702..cef59b5 100644 --- a/proto/manifest.json +++ b/proto/manifest.json @@ -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",