From a7ababeced69199c5f6ae11253796a4bb41c65f3 Mon Sep 17 00:00:00 2001 From: Rogee Date: Wed, 30 Sep 2026 12:13:32 +0800 Subject: [PATCH] Remove unreachable Agent upload and legacy Mock origination CLI --- cmd/sip-go-agent/agent_mock_originator.go | 28 ---- .../agent_mock_originator_test.go | 27 ---- cmd/sip-go-agent/upload.go | 139 ------------------ cmd/sip-go-agent/upload_attempt.go | 121 --------------- cmd/sip-go-agent/upload_attempt_test.go | 84 ----------- .../upload_failure_recovery_test.go | 75 ---------- cmd/sip-go-agent/upload_recovery.go | 105 ------------- cmd/sip-go-agent/upload_retry.go | 83 ----------- cmd/sip-go-agent/upload_retry_test.go | 70 --------- .../saas-dispatcher-implementation.md | 1 + 10 files changed, 1 insertion(+), 732 deletions(-) delete mode 100644 cmd/sip-go-agent/agent_mock_originator.go delete mode 100644 cmd/sip-go-agent/agent_mock_originator_test.go delete mode 100644 cmd/sip-go-agent/upload.go delete mode 100644 cmd/sip-go-agent/upload_attempt.go delete mode 100644 cmd/sip-go-agent/upload_attempt_test.go delete mode 100644 cmd/sip-go-agent/upload_failure_recovery_test.go delete mode 100644 cmd/sip-go-agent/upload_recovery.go delete mode 100644 cmd/sip-go-agent/upload_retry.go delete mode 100644 cmd/sip-go-agent/upload_retry_test.go diff --git a/cmd/sip-go-agent/agent_mock_originator.go b/cmd/sip-go-agent/agent_mock_originator.go deleted file mode 100644 index 25a7489..0000000 --- a/cmd/sip-go-agent/agent_mock_originator.go +++ /dev/null @@ -1,28 +0,0 @@ -package main - -import ( - "context" - "errors" - "log/slog" - - agentpb "git.ipao.vip/rogee/go-sip/gen/agent" -) - -// mockAuthorizedOriginator is deliberately isolated from SIP and Asterisk. -// Mixed/real modes have no adapter until separately authorized and verified. -func mockAuthorizedOriginator(mode string) func(context.Context, *agentpb.ExecuteAuthorizedRequest) error { - if mode != "mock" { - return nil - } - return func(ctx context.Context, request *agentpb.ExecuteAuthorizedRequest) error { - if err := ctx.Err(); err != nil { - return err - } - if request == nil || request.Binding == nil || request.Binding.ExecutionId == "" { - return errors.New("mock authorized origination has no execution identity") - } - slog.Info("isolated Mock Agent simulated authorized originate; no SIP sent", - "execution_id", request.Binding.ExecutionId, "trunk_id", request.SelectedTrunkId) - return nil - } -} diff --git a/cmd/sip-go-agent/agent_mock_originator_test.go b/cmd/sip-go-agent/agent_mock_originator_test.go deleted file mode 100644 index bc572ce..0000000 --- a/cmd/sip-go-agent/agent_mock_originator_test.go +++ /dev/null @@ -1,27 +0,0 @@ -package main - -import ( - "context" - "testing" - - agentpb "git.ipao.vip/rogee/go-sip/gen/agent" -) - -func TestAuthorizedOriginatorIsExplicitlyMockOnly(t *testing.T) { - if mockAuthorizedOriginator("mixed") != nil || mockAuthorizedOriginator("real") != nil || mockAuthorizedOriginator("") != nil { - t.Fatal("non-mock mode exposed an authorized mock originator") - } - originator := mockAuthorizedOriginator("mock") - if originator == nil { - t.Fatal("mock mode has no isolated originator") - } - request := &agentpb.ExecuteAuthorizedRequest{Binding: &agentpb.ExecutionBinding{ExecutionId: "execution-1"}, SelectedTrunkId: "trunk-1"} - ctx, cancel := context.WithCancel(context.Background()) - cancel() - if err := originator(ctx, request); err == nil { - t.Fatal("cancelled mock invocation appeared to complete") - } - if err := originator(context.Background(), request); err != nil { - t.Fatalf("isolated mock invocation: %v", err) - } -} diff --git a/cmd/sip-go-agent/upload.go b/cmd/sip-go-agent/upload.go deleted file mode 100644 index 4b6503c..0000000 --- a/cmd/sip-go-agent/upload.go +++ /dev/null @@ -1,139 +0,0 @@ -package main - -import ( - "context" - "crypto/sha256" - "encoding/hex" - "errors" - "fmt" - "path/filepath" - "strings" - "time" - - agentpb "git.ipao.vip/rogee/go-sip/gen/agent" - "git.ipao.vip/rogee/go-sip/internal/agent" - "git.ipao.vip/rogee/go-sip/internal/callruntime" - "git.ipao.vip/rogee/go-sip/internal/config" - "git.ipao.vip/rogee/go-sip/internal/rpc" -) - -func uploadCallRecordings(ctx context.Context, cfg config.Config, result callruntime.Result) ([]map[string]any, error) { - if strings.TrimSpace(cfg.DispatcherGRPCEndpoint) == "" { - return nil, nil - } - for name, value := range map[string]string{ - "AGENT_CALL_TENANT_ID": cfg.CallTenantID, - "AGENT_CALL_TENANT_KEY": cfg.CallTenantKey, - "AGENT_CALL_TASK_ID": cfg.CallTaskID, - "AGENT_CALL_TASK_ITEM_ID": cfg.CallTaskItemID, - } { - if strings.TrimSpace(value) == "" { - return nil, fmt.Errorf("%s is required when Dispatcher OSS upload is enabled", name) - } - } - if result.ChannelID == "" { - return nil, errors.New("call result has no channel ID for upload binding") - } - executionID := cfg.CallExecutionID - if executionID == "" { - executionID = result.ChannelID - } - binding := &agentpb.ExecutionBinding{ - TenantId: cfg.CallTenantID, - TenantKey: cfg.CallTenantKey, - ExecutionId: executionID, - TaskId: cfg.CallTaskID, - TaskItemId: cfg.CallTaskItemID, - TaskRevision: 1, - CallId: result.ChannelID, - AttemptId: result.ChannelID, - AgentVersionId: cfg.Version, - } - client, err := rpc.DialFromFiles(cfg.DispatcherGRPCEndpoint, cfg.MTLSCAFile, cfg.MTLSCertFile, cfg.MTLSKeyFile, cfg.DispatcherGRPCServerName) - if err != nil { - return nil, fmt.Errorf("dial Dispatcher gRPC service: %w", err) - } - defer client.Close() - spool, err := agent.NewSpool(cfg.SpoolRoot, time.Now) - if err != nil { - return nil, err - } - uploader := agent.UploadClient{Now: time.Now} - facts := append(append([]callruntime.RecordingFact(nil), result.InboundRecordings...), result.OutboundRecordings...) - if len(facts) == 0 { - return nil, errors.New("call produced no recordings for OSS upload") - } - uploaded := make([]map[string]any, 0, len(facts)) - for index, recording := range facts { - if recording.Path == "" || recording.Bytes <= 0 || recording.SHA256 == "" { - return nil, fmt.Errorf("recording %d is missing path, size or checksum", index+1) - } - assetID := recordingAssetID(recording.Segment, recording.Path, index+1) - asset := &agentpb.AssetDescriptor{ - Kind: agentpb.AssetKind_ASSET_KIND_RECORDING, - AssetId: assetID, - CallId: result.ChannelID, - ExecutionId: executionID, - Format: "wav", - SizeBytes: int64(recording.Bytes), - ChecksumSha256: recording.SHA256, - Channels: 1, - SampleRateHz: 16000, - DurationMs: recording.DurationMS, - } - record, err := uploadRecording(ctx, cfg, client, uploader, spool, binding, asset, recording.Path) - if err != nil { - return nil, err - } - uploaded = append(uploaded, map[string]any{ - "upload_id": record.UploadID, "asset_id": asset.AssetId, "object_key": record.ObjectKey, - "bytes": record.Result.SizeBytes, "sha256": record.Result.SHA256, - "status_code": record.Result.StatusCode, "notification_state": record.State, - }) - } - return uploaded, nil -} - -func recordingAssetID(segment, path string, index int) string { - raw := segment - if raw == "" { - raw = fmt.Sprintf("segment-%d", index) - } - raw += "-" + filepath.Base(path) - var builder strings.Builder - for _, r := range raw { - if r == '/' || r == '\\' || r == ' ' || r == '\t' || r == '\n' || r == '\r' { - builder.WriteByte('-') - continue - } - builder.WriteRune(r) - } - assetID := strings.Trim(builder.String(), "-") - if assetID == "" { - assetID = fmt.Sprintf("segment-%d", index) - } - if len([]byte(assetID)) > 120 { - digest := sha256.Sum256([]byte(assetID)) - assetID = "recording-" + hex.EncodeToString(digest[:16]) - } - return assetID -} - -func stableUploadID(binding *agentpb.ExecutionBinding, asset *agentpb.AssetDescriptor) string { - value := binding.ExecutionId + "\x00" + asset.AssetId + "\x00" + asset.ChecksumSha256 - digest := sha256.Sum256([]byte(value)) - return "upload-" + hex.EncodeToString(digest[:16]) -} - -func uploadMeta(cfg config.Config, phase, uploadID string) *agentpb.RequestMeta { - operationID := "recording-" + phase + "-" + uploadID - return &agentpb.RequestMeta{ - ProtocolVersion: "agent.v1", - RequestId: operationID, - TraceId: operationID, - OperationId: operationID, - IdempotencyKey: operationID, - AgentId: cfg.AgentID, - CellId: cfg.CellID, - } -} diff --git a/cmd/sip-go-agent/upload_attempt.go b/cmd/sip-go-agent/upload_attempt.go deleted file mode 100644 index e2e63b0..0000000 --- a/cmd/sip-go-agent/upload_attempt.go +++ /dev/null @@ -1,121 +0,0 @@ -package main - -import ( - "context" - "crypto/sha256" - "encoding/hex" - "errors" - "fmt" - "net/url" - "os" - "strings" - - agentpb "git.ipao.vip/rogee/go-sip/gen/agent" - "git.ipao.vip/rogee/go-sip/internal/agent" - "git.ipao.vip/rogee/go-sip/internal/config" - "google.golang.org/protobuf/proto" -) - -type recordingUploadRPC interface { - RequestUpload(context.Context, *agentpb.RequestUploadRequest) (*agentpb.RequestUploadResponse, error) - CompleteUpload(context.Context, *agentpb.CompleteUploadRequest) (*agentpb.CompleteUploadResponse, error) -} - -func uploadRecording(ctx context.Context, cfg config.Config, client recordingUploadRPC, uploader agent.UploadClient, spool *agent.Spool, binding *agentpb.ExecutionBinding, asset *agentpb.AssetDescriptor, path string) (returned agent.UploadAttempt, returnErr error) { - if binding == nil || asset == nil { - return returned, errors.New("upload binding and asset are required") - } - id := stableUploadID(binding, asset) - lock, err := spool.LockUpload(id) - if err != nil { - return returned, err - } - defer func() { returnErr = errors.Join(returnErr, lock.Close()) }() - identityBytes, err := (proto.MarshalOptions{Deterministic: true}).Marshal(&agentpb.RequestUploadRequest{Binding: binding, Asset: asset, UploadId: id}) - if err != nil { - return agent.UploadAttempt{}, err - } - digest := sha256.Sum256(identityBytes) - identity := hex.EncodeToString(digest[:]) - record, err := spool.LoadUploadAttempt(id) - if errors.Is(err, os.ErrNotExist) { - response, err := client.RequestUpload(ctx, &agentpb.RequestUploadRequest{Meta: uploadMeta(cfg, "request", id), Binding: binding, Asset: asset, UploadId: id}) - if err != nil { - return record, fmt.Errorf("request upload %s: %w", id, err) - } - if response.GetReceipt().GetResult() != agentpb.ResultCode_RESULT_CODE_ACCEPTED || response.GetGrant() == nil { - return record, fmt.Errorf("request upload %s has no accepted grant", id) - } - grant := response.Grant - if grant.UploadId != id { - return record, errors.New("upload grant identity mismatch") - } - parsed, err := url.Parse(grant.TargetUrl) - if err != nil || parsed.Host == "" { - return record, errors.New("upload grant has invalid target URL") - } - record = agent.UploadAttempt{UploadID: id, Identity: identity, State: "attempted", ObjectKey: grant.ObjectKey, Binding: binding, Asset: asset, RequestID: uploadMeta(cfg, "request", id).OperationId} - // Durable exclusive claim must precede any network PUT, including an attempt - // that ends in an ambiguous transport failure. - if err := spool.ClaimUpload(record); err != nil { - return record, err - } - uploader.AllowedHosts = map[string]struct{}{strings.ToLower(parsed.Host): {}} - result, err := uploader.UploadFile(ctx, grant, path) - if err != nil { - return record, fmt.Errorf("upload %s data plane: %w", id, err) - } - if result.SizeBytes != asset.SizeBytes || !strings.EqualFold(result.SHA256, asset.ChecksumSha256) { - return record, errors.New("upload result does not match asset") - } - if err := spool.RecordUploadResult(id, result); err != nil { - return record, err - } - record.State, record.Result = "uploaded", result - } else if err != nil { - return record, err - } else if record.Identity != identity { - return record, errors.New("persisted upload binding or asset mismatch") - } else if record.State == "attempted" { - return record, errors.New("prior PUT outcome unknown or failed; automatic re-upload is forbidden") - } - if record.State == "completed" { - return record, nil - } - return notifyUploadedRecordingLocked(ctx, cfg, client, spool, record) -} - -func notifyUploadedRecording(ctx context.Context, cfg config.Config, client recordingUploadRPC, spool *agent.Spool, record agent.UploadAttempt) (returned agent.UploadAttempt, returnErr error) { - lock, err := spool.LockUpload(record.UploadID) - if err != nil { - return returned, err - } - defer func() { returnErr = errors.Join(returnErr, lock.Close()) }() - current, err := spool.LoadUploadAttempt(record.UploadID) - if err != nil { - return returned, err - } - if current.State == "completed" { - return current, nil - } - return notifyUploadedRecordingLocked(ctx, cfg, client, spool, current) -} - -func notifyUploadedRecordingLocked(ctx context.Context, cfg config.Config, client recordingUploadRPC, spool *agent.Spool, record agent.UploadAttempt) (agent.UploadAttempt, error) { - if record.State != "uploaded" || record.Binding == nil || record.Asset == nil { - return record, errors.New("successful upload and original notification metadata are required") - } - id := record.UploadID - response, err := client.CompleteUpload(ctx, &agentpb.CompleteUploadRequest{Meta: uploadMeta(cfg, "complete", id), Binding: record.Binding, Asset: record.Asset, UploadId: id, UploadedSizeBytes: record.Result.SizeBytes, UploadedChecksumSha256: record.Result.SHA256}) - if err != nil { - return record, fmt.Errorf("notify upload %s: %w", id, err) - } - if response.GetReceipt().GetResult() != agentpb.ResultCode_RESULT_CODE_ACCEPTED || response.GetState() != agentpb.UploadState_UPLOAD_STATE_COMPLETED { - return record, errors.New("upload notification has not completed MQ delivery") - } - if err := spool.CompleteUploadNotification(id); err != nil { - return record, err - } - record.State = "completed" - return record, nil -} diff --git a/cmd/sip-go-agent/upload_attempt_test.go b/cmd/sip-go-agent/upload_attempt_test.go deleted file mode 100644 index 99d8e83..0000000 --- a/cmd/sip-go-agent/upload_attempt_test.go +++ /dev/null @@ -1,84 +0,0 @@ -package main - -import ( - "context" - "crypto/sha256" - "encoding/hex" - "errors" - "net/http" - "net/http/httptest" - "os" - "path/filepath" - "testing" - "time" - - agentpb "git.ipao.vip/rogee/go-sip/gen/agent" - "git.ipao.vip/rogee/go-sip/internal/agent" - "git.ipao.vip/rogee/go-sip/internal/config" -) - -type uploadRPCStub struct { - grant *agentpb.UploadGrant - requests, notifications int - pending bool -} - -func (s *uploadRPCStub) RequestUpload(_ context.Context, r *agentpb.RequestUploadRequest) (*agentpb.RequestUploadResponse, error) { - s.requests++ - s.grant.UploadId = r.UploadId - return &agentpb.RequestUploadResponse{Grant: s.grant, Receipt: &agentpb.OperationReceipt{Result: agentpb.ResultCode_RESULT_CODE_ACCEPTED}}, nil -} -func (s *uploadRPCStub) CompleteUpload(_ context.Context, _ *agentpb.CompleteUploadRequest) (*agentpb.CompleteUploadResponse, error) { - s.notifications++ - if s.pending { - return nil, errors.New("notification pending") - } - return &agentpb.CompleteUploadResponse{State: agentpb.UploadState_UPLOAD_STATE_COMPLETED, Receipt: &agentpb.OperationReceipt{Result: agentpb.ResultCode_RESULT_CODE_ACCEPTED}}, nil -} - -func TestRecordingNotificationRecoveryDoesNotPUTAgain(t *testing.T) { - puts := 0 - server := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { puts++; w.WriteHeader(http.StatusOK) })) - defer server.Close() - root := t.TempDir() - path := filepath.Join(root, "audio.wav") - if err := os.WriteFile(path, []byte("audio"), 0600); err != nil { - t.Fatal(err) - } - 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", 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) - } - uploader := agent.UploadClient{HTTPClient: server.Client()} - if _, err := uploadRecording(context.Background(), config.Config{}, remote, uploader, spool, binding, asset, path); err == nil { - t.Fatal("pending MQ notification reported complete") - } - remote.pending = false - restarted, err := agent.NewSpool(root, nil) - if err != nil { - t.Fatal(err) - } - if err := recoverUploadNotifications(context.Background(), config.Config{}, remote, restarted); err != nil { - t.Fatal(err) - } - record, err := restarted.LoadUploadAttempt(stableUploadID(binding, asset)) - if err != nil { - t.Fatal(err) - } - if record.State != "completed" || puts != 1 || remote.requests != 1 || remote.notifications != 2 { - t.Fatalf("state=%s PUT=%d grant=%d notification=%d", record.State, puts, remote.requests, remote.notifications) - } - if _, err := uploadRecording(context.Background(), config.Config{}, remote, uploader, restarted, binding, asset, path); err != nil { - t.Fatal(err) - } - if puts != 1 || remote.notifications != 2 { - t.Fatal("completed upload repeated side effects") - } - if _, err := os.Stat(path); err != nil { - t.Fatal("source recording removed") - } -} diff --git a/cmd/sip-go-agent/upload_failure_recovery_test.go b/cmd/sip-go-agent/upload_failure_recovery_test.go deleted file mode 100644 index de27434..0000000 --- a/cmd/sip-go-agent/upload_failure_recovery_test.go +++ /dev/null @@ -1,75 +0,0 @@ -package main - -import ( - "context" - "errors" - "net/http" - "testing" - - agentpb "git.ipao.vip/rogee/go-sip/gen/agent" - "git.ipao.vip/rogee/go-sip/internal/agent" - "git.ipao.vip/rogee/go-sip/internal/config" - "github.com/google/uuid" -) - -type mockFailureRecoveryClient struct { - reports int - fail bool -} - -func (c *mockFailureRecoveryClient) ReportExecutionEvent(_ context.Context, req *agentpb.ReportExecutionEventRequest) (*agentpb.ReportExecutionEventResponse, error) { - c.reports++ - if c.fail { - return nil, errors.New("Dispatcher unavailable") - } - return &agentpb.ReportExecutionEventResponse{Receipt: &agentpb.OperationReceipt{ - Result: agentpb.ResultCode_RESULT_CODE_ACCEPTED, FactId: req.Fact.FactId, ContentSha256: req.Fact.ContentSha256, - }}, nil -} - -func TestMockFailureRecoveryWaitsForActivationAndRejectsRealMode(t *testing.T) { - spool, err := agent.NewSpool(t.TempDir(), nil) - if err != nil { - t.Fatal(err) - } - active := &agentpb.RequestMeta{ - ProtocolVersion: "agent.v1", AgentId: "agent-a", CellId: "cell-a", BootId: "boot-a", - DispatcherEpoch: "epoch-a", SessionGeneration: 1, - } - remote := &mockFailureRecoveryClient{fail: true} - claimAndFail := func(id string) { - t.Helper() - if err := spool.ClaimUpload(agent.UploadAttempt{ - UploadID: id, Identity: "identity-" + id, State: "attempted", RequestID: uuid.NewString(), - Binding: &agentpb.ExecutionBinding{TenantId: "tenant-a", TenantKey: "tenant-a", ExecutionId: "execution-" + id}, - Asset: &agentpb.AssetDescriptor{AssetId: "recording-" + id}, - }); err != nil { - t.Fatal(err) - } - known, err := spool.ReportMockUploadFailure(context.Background(), remote, active, id, &agent.UploadHTTPError{StatusCode: http.StatusForbidden}) - if !known || err == nil { - t.Fatalf("failed Mock PUT was not retained for recovery: known=%t err=%v", known, err) - } - } - claimAndFail("upload-a") - cfg := config.Config{Mode: "mock"} - missingSession := func() (*agentpb.RequestMeta, error) { return nil, errors.New("not activated") } - if err := recoverMockFailureNotifications(context.Background(), cfg, remote, spool, missingSession); err == nil || remote.reports != 1 { - t.Fatalf("unactivated Agent reported a durable fact: reports=%d err=%v", remote.reports, err) - } - remote.fail = false - currentSession := func() (*agentpb.RequestMeta, error) { return active, nil } - if err := recoverMockFailureNotifications(context.Background(), cfg, remote, spool, currentSession); err != nil || remote.reports != 2 { - t.Fatalf("Mock failure was not recovered through the active session: reports=%d err=%v", remote.reports, err) - } - pending, err := spool.PendingUploadFailures() - if err != nil || len(pending) != 0 { - t.Fatalf("acknowledged failure was still pending: count=%d err=%v", len(pending), err) - } - remote.fail = true - claimAndFail("upload-b") - cfg.Mode = "real" - if err := recoverMockFailureNotifications(context.Background(), cfg, remote, spool, currentSession); err == nil || remote.reports != 3 { - t.Fatalf("Mock failure escaped into real mode: reports=%d err=%v", remote.reports, err) - } -} diff --git a/cmd/sip-go-agent/upload_recovery.go b/cmd/sip-go-agent/upload_recovery.go deleted file mode 100644 index 2bdb710..0000000 --- a/cmd/sip-go-agent/upload_recovery.go +++ /dev/null @@ -1,105 +0,0 @@ -package main - -import ( - "context" - "errors" - "fmt" - "log/slog" - "time" - - agentpb "git.ipao.vip/rogee/go-sip/gen/agent" - "git.ipao.vip/rogee/go-sip/internal/agent" - "git.ipao.vip/rogee/go-sip/internal/config" - "git.ipao.vip/rogee/go-sip/internal/rpc" -) - -func recoverUploadNotifications(ctx context.Context, cfg config.Config, client recordingUploadRPC, spool *agent.Spool) error { - pending, err := spool.PendingUploadNotifications() - if err != nil { - return err - } - var failures []error - for _, record := range pending { - if _, err := notifyUploadedRecording(ctx, cfg, client, spool, record); err != nil { - failures = append(failures, err) - } - } - return errors.Join(failures...) -} - -// recoverMockFailureNotifications never treats a prior boot's persisted fact -// as a current session. The report uses metadata from the newly activated -// session, while preserving the original immutable fact identity. -func recoverMockFailureNotifications(ctx context.Context, cfg config.Config, client agent.MockFailureFactClient, spool *agent.Spool, activeMeta func() (*agentpb.RequestMeta, error)) error { - pending, err := spool.PendingUploadFailures() - if err != nil || len(pending) == 0 { - return err - } - if cfg.Mode != "mock" { - return fmt.Errorf("pending Mock recording failure cannot be reported in %q mode", cfg.Mode) - } - if activeMeta == nil { - return errors.New("pending Mock recording failure requires an active Agent session") - } - meta, err := activeMeta() - if err != nil { - return fmt.Errorf("load active Agent session for Mock recording failure: %w", err) - } - return spool.RecoverMockUploadFailures(ctx, client, meta) -} - -// startUploadNotificationRecovery never requests a token or opens a source file. -// The worker owns only durable metadata-to-Dispatcher notifications. -func startUploadNotificationRecovery(ctx context.Context, cfg config.Config, spool *agent.Spool, activeMeta func() (*agentpb.RequestMeta, error)) (func(), error) { - pending, err := spool.PendingUploadNotifications() - if err != nil { - return nil, fmt.Errorf("read upload recovery journal: %w", err) - } - failureFacts, err := spool.PendingUploadFailures() - if err != nil { - return nil, fmt.Errorf("read Mock failure recovery journal: %w", err) - } - if len(failureFacts) > 0 && cfg.Mode != "mock" { - return nil, errors.New("pending Mock recording failure cannot start in mixed/real mode") - } - if cfg.DispatcherGRPCEndpoint == "" { - if len(pending) > 0 || len(failureFacts) > 0 { - return nil, errors.New("pending upload facts require Dispatcher gRPC endpoint") - } - return func() {}, nil - } - client, err := rpc.DialFromFiles(cfg.DispatcherGRPCEndpoint, cfg.MTLSCAFile, cfg.MTLSCertFile, cfg.MTLSKeyFile, cfg.DispatcherGRPCServerName) - if err != nil { - return nil, err - } - workerCtx, cancel := context.WithCancel(ctx) - done := make(chan struct{}) - go func() { - defer close(done) - ticker := time.NewTicker(time.Second) - defer ticker.Stop() - for { - attemptCtx, finish := context.WithTimeout(workerCtx, 10*time.Second) - err := errors.Join( - recoverUploadNotifications(attemptCtx, cfg, client, spool), - recoverMockFailureNotifications(attemptCtx, cfg, client, spool, activeMeta), - ) - finish() - if err != nil && workerCtx.Err() == nil { - slog.Error("upload notification recovery failed; original facts retained", "error", err) - } - select { - case <-workerCtx.Done(): - return - case <-ticker.C: - } - } - }() - return func() { - cancel() - <-done - if err := client.Close(); err != nil { - slog.Error("close upload notification client", "error", err) - } - }, nil -} diff --git a/cmd/sip-go-agent/upload_retry.go b/cmd/sip-go-agent/upload_retry.go deleted file mode 100644 index f110fd6..0000000 --- a/cmd/sip-go-agent/upload_retry.go +++ /dev/null @@ -1,83 +0,0 @@ -package main - -import ( - "context" - "crypto/sha256" - "encoding/hex" - "errors" - "net/url" - "strings" - - agentpb "git.ipao.vip/rogee/go-sip/gen/agent" - "git.ipao.vip/rogee/go-sip/internal/agent" - "git.ipao.vip/rogee/go-sip/internal/config" - "github.com/google/uuid" - "google.golang.org/protobuf/proto" -) - -// retryUploadRecording is reachable only through an explicit operator command. -// A supplied request ID is stable across redelivery and can authorize one PUT. -func retryUploadRecording(ctx context.Context, cfg config.Config, client recordingUploadRPC, uploader agent.UploadClient, spool *agent.Spool, uploadID, requestID, path string) (returned agent.UploadAttempt, returnErr error) { - parsedID, err := uuid.Parse(requestID) - if err != nil || parsedID.Version() != 4 || parsedID.Variant() != uuid.RFC4122 || parsedID.String() != requestID { - return returned, errors.New("explicit request ID must be a canonical UUID v4") - } - lock, err := spool.LockUpload(uploadID) - if err != nil { - return returned, err - } - defer func() { returnErr = errors.Join(returnErr, lock.Close()) }() - record, err := spool.LoadUploadAttempt(uploadID) - if err != nil { - return record, err - } - if record.State != "attempted" || record.Binding == nil || record.Asset == nil || record.RequestID == "" || record.RequestID == requestID { - return record, errors.New("only an unsuccessful upload may explicitly request a new grant") - } - raw, err := (proto.MarshalOptions{Deterministic: true}).Marshal(&agentpb.RequestUploadRequest{Binding: record.Binding, Asset: record.Asset, UploadId: uploadID}) - if err != nil { - return record, err - } - sum := sha256.Sum256(raw) - if record.Identity != hex.EncodeToString(sum[:]) || stableUploadID(record.Binding, record.Asset) != uploadID { - return record, errors.New("persisted upload identity mismatch") - } - meta := uploadMeta(cfg, "request", uploadID) - meta.RequestId = requestID - meta.OperationId = requestID - meta.IdempotencyKey = requestID - meta.TraceId = requestID - response, err := client.RequestUpload(ctx, &agentpb.RequestUploadRequest{Meta: meta, Binding: record.Binding, Asset: record.Asset, UploadId: uploadID}) - if err != nil { - return record, err - } - if response.GetReceipt().GetResult() != agentpb.ResultCode_RESULT_CODE_ACCEPTED || response.GetGrant() == nil { - return record, errors.New("explicit upload request was not granted") - } - grant := response.Grant - if grant.UploadId != uploadID || grant.ObjectKey != record.ObjectKey || grant.RequiredChecksumSha256 != record.Asset.ChecksumSha256 || grant.MaxBytes < record.Asset.SizeBytes { - return record, errors.New("replacement grant does not match the original asset") - } - target, err := url.Parse(grant.TargetUrl) - if err != nil || target.Host == "" { - return record, errors.New("replacement grant has invalid target") - } - if err := spool.ReserveUploadRetry(uploadID, requestID); err != nil { - return record, err - } - uploader.AllowedHosts = map[string]struct{}{strings.ToLower(target.Host): {}} - result, err := uploader.UploadFile(ctx, grant, path) - if err != nil { - return record, err - } - if result.SizeBytes != record.Asset.SizeBytes || result.SHA256 != record.Asset.ChecksumSha256 { - return record, errors.New("replacement upload result does not match original asset") - } - if err := spool.RecordUploadResult(uploadID, result); err != nil { - return record, err - } - record.RequestID = requestID - record.State = "uploaded" - record.Result = result - return notifyUploadedRecordingLocked(ctx, cfg, client, spool, record) -} diff --git a/cmd/sip-go-agent/upload_retry_test.go b/cmd/sip-go-agent/upload_retry_test.go deleted file mode 100644 index 902af06..0000000 --- a/cmd/sip-go-agent/upload_retry_test.go +++ /dev/null @@ -1,70 +0,0 @@ -package main - -import ( - "context" - "crypto/sha256" - "encoding/hex" - "io" - "net/http" - "net/http/httptest" - "os" - "path/filepath" - "testing" - "time" - - agentpb "git.ipao.vip/rogee/go-sip/gen/agent" - "git.ipao.vip/rogee/go-sip/internal/agent" - "git.ipao.vip/rogee/go-sip/internal/config" -) - -func TestExplicitUploadRetryIsOneNewRequestAndOnePUT(t *testing.T) { - puts := 0 - server := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { - puts++ - if _, err := io.Copy(io.Discard, r.Body); err != nil { - t.Error(err) - } - w.WriteHeader(http.StatusOK) - })) - defer server.Close() - root := t.TempDir() - path := filepath.Join(root, "audio.wav") - if err := os.WriteFile(path, []byte("audio"), 0600); err != nil { - t.Fatal(err) - } - 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", 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) - } - uploader := agent.UploadClient{HTTPClient: server.Client()} - if _, err := uploadRecording(context.Background(), config.Config{}, remote, uploader, spool, binding, asset, path); err == nil { - t.Fatal("expired token accepted") - } - if puts != 0 || remote.requests != 1 { - t.Fatal("expired grant caused PUT or auto renewal") - } - remote.grant.ExpiresAtUnixMs = time.Now().Add(15 * time.Minute).UnixMilli() - id := stableUploadID(binding, asset) - if _, err := uploadRecording(context.Background(), config.Config{}, remote, uploader, spool, binding, asset, path); err == nil { - t.Fatal("ordinary recovery retried a failed attempt") - } - if _, err := retryUploadRecording(context.Background(), config.Config{}, remote, uploader, spool, id, "11111111-1111-4111-8111-111111111111", path); err != nil { - t.Fatal(err) - } - if puts != 1 || remote.requests != 2 || remote.notifications != 1 { - t.Fatalf("PUT=%d grants=%d notifications=%d", puts, remote.requests, remote.notifications) - } - if _, err := retryUploadRecording(context.Background(), config.Config{}, remote, uploader, spool, id, "22222222-2222-4222-8222-222222222222", path); err == nil { - t.Fatal("completed upload retried") - } - if puts != 1 || remote.requests != 2 { - t.Fatal("duplicate retry repeated effects") - } - if _, err := os.Stat(path); err != nil { - t.Fatal("retry removed source file") - } -} diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index c309486..81f9e87 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -110,6 +110,7 @@ - HTTP 配置读取旧入口:移除旧字符串租户的三资源拼接读取、单独任务状态读取及旧 snapshot/changes 发现分支和专属测试;保留统一的凭据头、HTTP/JSON 错误处理和当前五类只读读取。当前任务/发现的 503、身份错配及分页测试继续执行,不从旧配置回退。 - HTTP 配置读取名称收敛:四份现行读取/发现 Go 文件改为 `snapshots`/`discovery` 的无代次路径,公开快照与 `ReadSIP`/`ReadTask`/`ReadTasks`/`ReadAllTasks` 只保留一套入口。原始 AI JSON、租户数字身份、五类 HTTP 路径及断连拒绝行为未改变;相关调用方和隔离测试同步更新。 - Dispatcher 名称收敛:现行 `Current*` 调度、控制、发现、线路选择及执行类型/错误改为 `Runtime`、`Bootstrap`、`DiscoveryFollower` 等唯一代码入口;18 份 Go 源码/测试文件移除代次路径,相关 Agent 控制及调用方同步更新。任务归属、固定 SIP 快照、白名单/时段、额度和结果防重规则不变,现行隔离执行测试继续通过。 +- Agent 根命令不可达分支:删除旧执行模型的本地录音上传、手工重试、后台通知恢复和旧 Mock originate 代码及专属测试;根命令仍只注册获批的 Mock Agent/Dispatcher。当前录音 RPC、单次 PUT、失败保留和结果恢复另由现行隔离测试覆盖,删除旧命令不改变现有业务数据。 ## 验收台账