diff --git a/cmd/sip-go-agent/agent_real.go b/cmd/sip-go-agent/agent_real.go index 041e159..3ad85e9 100644 --- a/cmd/sip-go-agent/agent_real.go +++ b/cmd/sip-go-agent/agent_real.go @@ -19,6 +19,15 @@ import ( "github.com/google/uuid" ) +func reportDefiniteNonDialFailure(ctx context.Context, cause error, client agent.RecordingClient) error { + if !asterisk.NativeCallNeverSubmitted(cause) { + return nil // ARI submission may have happened: keep the unknown reservation + } + endCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 10*time.Second) + defer cancel() + return client.ReportEnded(endCtx) +} + // newRealAgentServer only exposes the signed Dispatcher execution path. There // is no scenario, synthetic recording, local dial entry, or Mock upload host. func newRealAgentServer(ctx context.Context, settings config.AgentEnvironment, dispatcher agentpb.AgentControlServiceClient) (*rpc.Server, error) { @@ -53,6 +62,14 @@ func newRealAgentServer(ctx context.Context, settings config.AgentEnvironment, d } loader := asterisk.Loader{ConfigDir: settings.AsteriskConfigDir, Asterisk: settings.AsteriskBin, LibraryDir: settings.AsteriskLibraryDir} var handler *rpc.Server + recordingClient := func(execution rpc.ApprovedExecution) agent.RecordingClient { + return agent.RecordingClient{ + Client: dispatcher, DispatcherID: execution.DispatcherID, TenantID: execution.TenantID, + SourceEventID: execution.SourceEventID, Session: func(context.Context) (*agentpb.RequestMeta, error) { + return handler.ActiveSessionMeta() + }, + } + } worker := &rpc.ApprovedCallWorker{ Lifecycle: ctx, Calls: &agent.TaskCalls{}, Prepare: func(execution rpc.ApprovedExecution) (func(context.Context) error, error) { @@ -60,10 +77,7 @@ func newRealAgentServer(ctx context.Context, settings config.AgentEnvironment, d return nil, errors.New("Agent session is unavailable") } delivery := &agent.RecordingDelivery{ - Call: agent.RecordingClient{ - Client: dispatcher, DispatcherID: execution.DispatcherID, TenantID: execution.TenantID, - SourceEventID: execution.SourceEventID, Session: func(context.Context) (*agentpb.RequestMeta, error) { return handler.ActiveSessionMeta() }, - }, + Call: recordingClient(execution), Recovery: &agent.RecordingRecovery{ Root: settings.RecoveryRoot, Upload: agent.UploadClient{AllowedHosts: map[string]struct{}{strings.ToLower(settings.OSSAllowedHost): {}}}, @@ -76,6 +90,14 @@ func newRealAgentServer(ctx context.Context, settings config.AgentEnvironment, d }, OnFailure: func(execution rpc.ApprovedExecution, cause error) error { log.Printf("Agent real call requires inspection: event_id=%q task_id=%q native_phase=%q ari_http_status=%d cause_type=%T", execution.SourceEventID, execution.TaskID, asterisk.NativeCallPhase(cause), asterisk.NativeCallHTTPStatus(cause), cause) + if !asterisk.NativeCallNeverSubmitted(cause) { + return nil + } + if err := reportDefiniteNonDialFailure(ctx, cause, recordingClient(execution)); err != nil { + log.Printf("Agent definite pre-dial end unconfirmed: event_id=%q error_type=%T", execution.SourceEventID, err) + return err + } + log.Printf("Agent definite pre-dial end confirmed: event_id=%q", execution.SourceEventID) return nil }, } diff --git a/cmd/sip-go-agent/agent_real_test.go b/cmd/sip-go-agent/agent_real_test.go index fed75a8..a7b6bc2 100644 --- a/cmd/sip-go-agent/agent_real_test.go +++ b/cmd/sip-go-agent/agent_real_test.go @@ -2,6 +2,7 @@ package main import ( "context" + "errors" "io" "net" "os" @@ -10,11 +11,46 @@ import ( "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/asterisk" "git.ipao.vip/rogee/go-sip/internal/rpc" "google.golang.org/grpc" "google.golang.org/grpc/credentials" ) +type definiteEndClient struct { + agentpb.AgentControlServiceClient + calls int + fail error +} + +func (c *definiteEndClient) ReportCallEnded(_ context.Context, req *agentpb.ReportCallEndedRequest, _ ...grpc.CallOption) (*agentpb.ReportCallEndedResponse, error) { + c.calls++ + if c.fail != nil { + return nil, c.fail + } + return &agentpb.ReportCallEndedResponse{Receipt: &agentpb.OperationReceipt{Result: agentpb.ResultCode_RESULT_CODE_APPLIED, FactId: req.SourceEventId}}, nil +} + +func TestOnlyDefinitelyUnissuedRealCallReportsConfirmedEnd(t *testing.T) { + backend := &definiteEndClient{} + client := agent.RecordingClient{Client: backend, DispatcherID: "dispatcher-test", TenantID: 1, SourceEventID: "call-pre-dial", Session: func(context.Context) (*agentpb.RequestMeta, error) { + return &agentpb.RequestMeta{AgentId: "agent-test", CellId: "cell-test", BootId: "boot-test", DispatcherEpoch: "epoch-test", SessionGeneration: 1}, nil + }} + before := &asterisk.NativeCallFailure{Phase: "rtp_listen", Cause: errors.New("no RTP socket")} + if err := reportDefiniteNonDialFailure(context.Background(), before, client); err != nil || backend.calls != 1 { + t.Fatalf("definite pre-dial failure was not reported: calls=%d err=%v", backend.calls, err) + } + unknown := &asterisk.NativeCallFailure{Phase: "originate", Cause: errors.New("response lost")} + if err := reportDefiniteNonDialFailure(context.Background(), unknown, client); err != nil || backend.calls != 1 { + t.Fatalf("unknown ARI outcome must stay occupied: calls=%d err=%v", backend.calls, err) + } + backend.fail = errors.New("no authenticated end receipt") + if err := reportDefiniteNonDialFailure(context.Background(), before, client); err == nil { + t.Fatal("failed end reporting must not pretend capacity was released") + } +} + func TestRealAgentCommandStartsPinnedServerWithoutMockFixturesOrDial(t *testing.T) { ca, agentCert, agentKey, dispatcherCert, dispatcherKey, dispatcherLeaf := localCommandCertificates(t) settings, _ := currentAgentSetupFixture(t) diff --git a/internal/asterisk/call.go b/internal/asterisk/call.go index 4a40d20..88eee26 100644 --- a/internal/asterisk/call.go +++ b/internal/asterisk/call.go @@ -46,6 +46,18 @@ func NativeCallPhase(err error) string { return "unclassified" } +// NativeCallNeverSubmitted is true only when the error occurred before any +// ARI originate request could have been sent. After submission, even a local +// error or successful cleanup does not prove the carrier never received SIP. +func NativeCallNeverSubmitted(err error) bool { + switch NativeCallPhase(err) { + case "validate", "ari_open", "rtp_listen", "subscribe": + return true + default: + return false + } +} + // NativeCallHTTPStatus extracts a numeric ARI HTTP status without retaining // a provider response body, credential or request URL in diagnostic logs. func NativeCallHTTPStatus(err error) int { diff --git a/internal/asterisk/call_test.go b/internal/asterisk/call_test.go index d1c2e87..a89ede9 100644 --- a/internal/asterisk/call_test.go +++ b/internal/asterisk/call_test.go @@ -101,6 +101,19 @@ type nativeHTTPStatusError struct{ code int } func (e nativeHTTPStatusError) Error() string { return "private ARI response" } func (e nativeHTTPStatusError) Code() int { return e.code } +func TestOnlyPreSubmitNativeFailureCanBeConfirmedNeverDialed(t *testing.T) { + for _, phase := range []string{"validate", "ari_open", "rtp_listen", "subscribe"} { + if !NativeCallNeverSubmitted(errors.Join(&NativeCallFailure{Phase: phase, Cause: errors.New("failed")}, errors.New("cleanup"))) { + t.Fatalf("%s cannot have created an outbound call", phase) + } + } + for _, phase := range []string{"originate", "answer_wait", "bridge_create", "external_media_create", "unclassified"} { + if NativeCallNeverSubmitted(&NativeCallFailure{Phase: phase, Cause: errors.New("unknown")}) { + t.Fatalf("%s might have originated", phase) + } + } +} + func TestNativeCallFailureReportsStageAndHTTPStatusWithoutResponseBody(t *testing.T) { cause := fmt.Errorf("request failed: %w", nativeHTTPStatusError{code: 404}) err := errors.Join(&NativeCallFailure{Phase: "originate", Cause: cause}, errors.New("cleanup failed"))