Upload approved recordings directly from memory
This commit is contained in:
@@ -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 冒充外部签收。
|
||||
|
||||
+64
-22
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user