From 369ee9cd867a5406aaadd2aeb5d665b7b6b4ae4f Mon Sep 17 00:00:00 2001 From: Rogee Date: Tue, 29 Sep 2026 22:42:30 +0800 Subject: [PATCH] Upload approved recordings directly from memory --- .../saas-dispatcher-implementation.md | 5 + internal/agent/upload.go | 86 +++++++++++---- internal/agent/upload_bytes_test.go | 101 ++++++++++++++++++ 3 files changed, 170 insertions(+), 22 deletions(-) create mode 100644 internal/agent/upload_bytes_test.go diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index 4676f09..e8bdbef 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -56,6 +56,11 @@ - 验证:`PATH=/tmp/sip-go-agent-tools/bin:$PATH bash scripts/check-current-contracts.sh`、`PATH=/tmp/sip-go-agent-tools/bin:$PATH bash scripts/check-proto.sh`、`go test ./... -count=1`、`go test -race ./internal/ai ./internal/callflow ./internal/configread ./internal/dispatcher ./internal/rpc -count=1`、`go vet ./...`、`go build ./...`、`git diff --check` 均通过。 - **后续边界:** 录音直传/OSS 失败恢复及最终结果属 P06;新主 CLI 接线、真实媒体完整联动及旧路径清理属 P07;隔离端到端和 A01–A12 属 P08。真实 Agent/Asterisk、SaaS、MQ、AI 供应商联调未开展,不能由 Mock 结果代签。 +## P06:录音/OSS/最终结果(进行中,未签收) + +- 已新增 Agent 内存录音单次 PUT 组件 `UploadClient.UploadBytes`:复用受限授权、大小和 SHA-256 核对,成功路径不写临时录音或通话结果文件。OSS 明确拒绝保留状态码且不自动重试;传输结果不明时返回专用错误、停止自动重试并隐藏带签名的 URL。空授权头拒绝而非使服务崩溃。这里只验证组件,尚未连接实际录音或最终结果。 +- 已验证:`go test ./internal/agent -count=1`、`go test -race ./internal/agent -count=1`、`go vet ./...`、`go build ./...`、`git diff --check`。失败双文件持久化、48 小时恢复、录音生成失败/无录音、Dispatcher 唯一结果 outbox、MQ/重启恢复和端到端回归均未完成,不能宣称 P06 通过。 + ## 验收台账 A01–A12 的行为验证及 K01–K16 的运行时验证待 P03–P08 逐项填充;不得用本地 Mock 冒充外部签收。 diff --git a/internal/agent/upload.go b/internal/agent/upload.go index 591ec33..81c8c4e 100644 --- a/internal/agent/upload.go +++ b/internal/agent/upload.go @@ -1,6 +1,7 @@ package agent import ( + "bytes" "context" "crypto/sha256" "encoding/hex" @@ -33,6 +34,7 @@ type UploadClient struct { var ErrUploadGrantExpired = errors.New("upload grant is expired") var ErrUploadGrantInvalid = errors.New("upload grant is invalid") var ErrUploadChecksumMismatch = errors.New("upload checksum mismatch") +var ErrUploadOutcomeUnknown = errors.New("upload outcome unknown") // UploadHTTPError records only the status, never the signed URL or OSS body. type UploadHTTPError struct{ StatusCode int } @@ -46,36 +48,58 @@ type UploadResult struct { ETag string } -func (c UploadClient) UploadFile(ctx context.Context, grant *agentpb.UploadGrant, path string) (result UploadResult, err error) { +// UploadBytes sends an in-memory recording directly; the success path never +// creates a recording file or a local call-result journal. +func (c UploadClient) UploadBytes(ctx context.Context, grant *agentpb.UploadGrant, recording []byte) (UploadResult, error) { + parsed, err := c.validateGrant(ctx, grant) + if err != nil { + return UploadResult{}, err + } + if len(recording) == 0 { + return UploadResult{}, errors.New("recording audio is empty") + } + sum := sha256.Sum256(recording) + return c.putValidated(ctx, grant, parsed, bytes.NewReader(recording), int64(len(recording)), hex.EncodeToString(sum[:])) +} + +func (c UploadClient) validateGrant(ctx context.Context, grant *agentpb.UploadGrant) (*url.URL, error) { if grant == nil { - return UploadResult{}, fmt.Errorf("%w: grant is required", ErrUploadGrantInvalid) + return nil, fmt.Errorf("%w: grant is required", ErrUploadGrantInvalid) } if grant.TargetUrl == "" || grant.UploadId == "" || grant.ObjectKey == "" { - return UploadResult{}, fmt.Errorf("%w: URL, ID and object key are required", ErrUploadGrantInvalid) + return nil, fmt.Errorf("%w: URL, ID and object key are required", ErrUploadGrantInvalid) } if grant.ExpiresAtUnixMs <= 0 { - return UploadResult{}, fmt.Errorf("%w: expiry is required", ErrUploadGrantInvalid) + return nil, fmt.Errorf("%w: expiry is required", ErrUploadGrantInvalid) } now := time.Now if c.Now != nil { now = c.Now } if !now().Before(time.UnixMilli(grant.ExpiresAtUnixMs)) { - return UploadResult{}, ErrUploadGrantExpired + return nil, ErrUploadGrantExpired } parsed, err := url.Parse(grant.TargetUrl) if err != nil || parsed.Host == "" { - return UploadResult{}, fmt.Errorf("%w: URL is invalid", ErrUploadGrantInvalid) + return nil, fmt.Errorf("%w: URL is invalid", ErrUploadGrantInvalid) } if parsed.Scheme != "https" && !(c.AllowInsecureHTTP && parsed.Scheme == "http") { - return UploadResult{}, fmt.Errorf("%w: URL must use HTTPS", ErrUploadGrantInvalid) + return nil, fmt.Errorf("%w: URL must use HTTPS", ErrUploadGrantInvalid) } if len(c.AllowedHosts) > 0 { if _, ok := c.AllowedHosts[strings.ToLower(parsed.Host)]; !ok { - return UploadResult{}, fmt.Errorf("%w: host %q is not allowed", ErrUploadGrantInvalid, parsed.Host) + return nil, fmt.Errorf("%w: host %q is not allowed", ErrUploadGrantInvalid, parsed.Host) } } if err := ctx.Err(); err != nil { + return nil, err + } + return parsed, nil +} + +func (c UploadClient) UploadFile(ctx context.Context, grant *agentpb.UploadGrant, path string) (result UploadResult, err error) { + parsed, err := c.validateGrant(ctx, grant) + if err != nil { return UploadResult{}, err } file, err := os.Open(path) @@ -105,21 +129,30 @@ func (c UploadClient) UploadFile(ctx context.Context, grant *agentpb.UploadGrant if err != nil { return UploadResult{}, err } - if grant.RequiredChecksumSha256 != "" && !strings.EqualFold(grant.RequiredChecksumSha256, digest) { - return UploadResult{}, fmt.Errorf("%w: asset does not match grant", ErrUploadChecksumMismatch) - } - if _, err := file.Seek(0, io.SeekStart); err != nil { return UploadResult{}, err } + return c.putValidated(ctx, grant, parsed, file, stat.Size(), digest) +} + +func (c UploadClient) putValidated(ctx context.Context, grant *agentpb.UploadGrant, parsed *url.URL, source io.Reader, size int64, digest string) (UploadResult, error) { + if grant.MaxBytes > 0 && size > grant.MaxBytes { + return UploadResult{}, fmt.Errorf("asset exceeds grant limit: %d > %d", size, grant.MaxBytes) + } + if grant.RequiredChecksumSha256 != "" && !strings.EqualFold(grant.RequiredChecksumSha256, digest) { + return UploadResult{}, fmt.Errorf("%w: asset does not match grant", ErrUploadChecksumMismatch) + } transmitted := &uploadChecksum{hash: sha256.New()} - body := io.TeeReader(io.LimitReader(file, stat.Size()), transmitted) + body := io.TeeReader(io.LimitReader(source, size), transmitted) req, err := http.NewRequestWithContext(ctx, http.MethodPut, parsed.String(), body) if err != nil { - return UploadResult{}, err + return UploadResult{}, fmt.Errorf("%w: cannot construct PUT request", ErrUploadGrantInvalid) } - req.ContentLength = stat.Size() + req.ContentLength = size for _, header := range grant.Headers { + if header == nil { + return UploadResult{}, fmt.Errorf("%w: header is missing", ErrUploadGrantInvalid) + } if strings.EqualFold(header.Name, "host") || strings.EqualFold(header.Name, "content-length") { return UploadResult{}, fmt.Errorf("%w: forbidden header", ErrUploadGrantInvalid) } @@ -133,13 +166,22 @@ func (c UploadClient) UploadFile(ctx context.Context, grant *agentpb.UploadGrant copyClient.CheckRedirect = func(_ *http.Request, _ []*http.Request) error { return http.ErrUseLastResponse } resp, err := copyClient.Do(req) if err != nil { - // net/http includes the entire signed URL in *url.Error. Retain the - // underlying transport cause without exposing the temporary token. + // A transport failure after sending bytes has an unknown OSS outcome. + // Never expose the signed URL, even if a custom transport includes it. var requestError *url.Error if errors.As(err, &requestError) { - return UploadResult{}, fmt.Errorf("upload PUT transport failure: %w", requestError.Err) + err = requestError.Err } - return UploadResult{}, err + kind := fmt.Sprintf("%T", err) + switch { + case errors.Is(err, context.Canceled): + kind = "canceled" + case errors.Is(err, context.DeadlineExceeded): + kind = "deadline" + case errors.Is(err, io.ErrUnexpectedEOF): + kind = "unexpected_eof" + } + return UploadResult{}, fmt.Errorf("%w: PUT transport failure (%s)", ErrUploadOutcomeUnknown, kind) } defer resp.Body.Close() maxResponse := c.MaxResponseBodySize @@ -151,11 +193,11 @@ func (c UploadClient) UploadFile(ctx context.Context, grant *agentpb.UploadGrant return UploadResult{}, &UploadHTTPError{StatusCode: resp.StatusCode} } if _, err := io.Copy(io.Discard, io.LimitReader(resp.Body, maxResponse)); err != nil { - return UploadResult{}, fmt.Errorf("read upload response: %w", err) + return UploadResult{}, fmt.Errorf("%w: read PUT response (%T)", ErrUploadOutcomeUnknown, err) } sentDigest, sentBytes := transmitted.result() - if sentBytes != stat.Size() || sentDigest != digest { - return UploadResult{}, fmt.Errorf("%w: transmitted bytes differ from validated asset", ErrUploadChecksumMismatch) + if sentBytes != size || sentDigest != digest { + return UploadResult{}, errors.Join(ErrUploadOutcomeUnknown, fmt.Errorf("%w: transmitted bytes differ from validated asset", ErrUploadChecksumMismatch)) } return UploadResult{StatusCode: resp.StatusCode, SizeBytes: sentBytes, SHA256: sentDigest, ETag: resp.Header.Get("ETag")}, nil } diff --git a/internal/agent/upload_bytes_test.go b/internal/agent/upload_bytes_test.go new file mode 100644 index 0000000..ca8f1b3 --- /dev/null +++ b/internal/agent/upload_bytes_test.go @@ -0,0 +1,101 @@ +package agent + +import ( + "bytes" + "context" + "crypto/sha256" + "encoding/hex" + "errors" + "fmt" + "io" + "net/http" + "net/http/httptest" + "os" + "strings" + "sync/atomic" + "testing" + "time" + + agentpb "git.ipao.vip/rogee/go-sip/gen/agent" +) + +func TestUploadBytesPutsMemoryRecordingOnceWithoutBusinessFile(t *testing.T) { + body := []byte("recording held only in memory") + sum := sha256.Sum256(body) + var requests atomic.Int32 + var received []byte + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + requests.Add(1) + if r.Method != http.MethodPut || r.ContentLength != int64(len(body)) || r.Header.Get("x-upload-token") != "mock-token" { + http.Error(w, "invalid presigned PUT", http.StatusBadRequest) + return + } + received, _ = io.ReadAll(r.Body) + w.Header().Set("ETag", "mock-etag") + })) + 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"}}} + 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) + } + entries, err := os.ReadDir(noFiles) + if err != nil || len(entries) != 0 { + t.Fatalf("normal path created a business file: entries=%v err=%v", entries, 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)} + _, 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) + } +} + +func TestUploadBytesDefinitiveHTTPFailureIsOneAttempt(t *testing.T) { + var requests atomic.Int32 + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + requests.Add(1) + w.WriteHeader(http.StatusServiceUnavailable) + _, _ = 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} + _, 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") { + t.Fatalf("failed OSS PUT was retried or leaked its response: requests=%d err=%v", requests.Load(), err) + } +} + +type ambiguousBytesTransport struct{ requests *atomic.Int32 } + +func (t ambiguousBytesTransport) RoundTrip(req *http.Request) (*http.Response, error) { + t.requests.Add(1) + return nil, fmt.Errorf("transport says URL=%s", req.URL.String()) +} + +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} + _, 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) + } +} + +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}} + _, 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) + } +}