feat(agent): prepare single native call with observed media delivery
This commit is contained in:
@@ -0,0 +1,143 @@
|
||||
package rpc
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"os"
|
||||
"strings"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"git.ipao.vip/rogee/go-sip/internal/agent"
|
||||
"git.ipao.vip/rogee/go-sip/internal/ai"
|
||||
"git.ipao.vip/rogee/go-sip/internal/asterisk"
|
||||
"git.ipao.vip/rogee/go-sip/internal/callflow"
|
||||
)
|
||||
|
||||
// ApprovedRecordedRealCall uses the signed Dispatcher instruction to create
|
||||
// exactly one native ARI call. Unlike the isolated Mock it never supplies a
|
||||
// synthetic transcript, media frame, outcome, or recording.
|
||||
type ApprovedRecordedRealCall struct {
|
||||
Loader asterisk.Loader
|
||||
MediaPayloadType uint8
|
||||
MaxWAVBytes int64
|
||||
ReportTimeout time.Duration
|
||||
Delivery *agent.RecordingDelivery
|
||||
originator func(context.Context, asterisk.NativeDial) (realMedia, error) // isolated tests only
|
||||
}
|
||||
|
||||
type realMedia struct {
|
||||
ctx context.Context
|
||||
session callflow.MediaSession
|
||||
close func() error
|
||||
}
|
||||
|
||||
// Prepare rejects missing signed identity and recovery capacity before the
|
||||
// Agent acknowledges the call. No SIP channel is created during preparation.
|
||||
func (r *ApprovedRecordedRealCall) Prepare(approved ApprovedExecution) (func(context.Context) error, error) {
|
||||
if r == nil || r.Delivery == nil || r.ReportTimeout <= 0 || r.MaxWAVBytes <= 44 || r.MediaPayloadType < 96 || r.MediaPayloadType > 127 {
|
||||
return nil, errors.New("real call requires bounded media and per-call delivery")
|
||||
}
|
||||
if approved.DispatcherID == "" || approved.TenantID <= 0 || approved.SourceEventID == "" || approved.CallID != approved.SourceEventID || approved.TaskID == "" ||
|
||||
approved.CallerProfileID == "" || approved.CallerID == "" || approved.Callee == "" || approved.DialedCallee == "" || approved.SelectedTrunkID == "" ||
|
||||
approved.RingTimeout <= 0 || approved.MaxCallDuration <= 0 || approved.DialBefore.IsZero() || (approved.AI.Mode != string(ai.ModeASROnly) && approved.AI.Mode != string(ai.ModeFullAI)) {
|
||||
return nil, errors.New("approved real call identity, timing or AI mode is incomplete")
|
||||
}
|
||||
if r.Delivery.Call.Client == nil || r.Delivery.Call.Session == nil || r.Delivery.Call.DispatcherID != approved.DispatcherID ||
|
||||
r.Delivery.Call.TenantID != approved.TenantID || r.Delivery.Call.SourceEventID != approved.SourceEventID || r.Delivery.Recovery == nil ||
|
||||
strings.TrimSpace(r.Delivery.Recovery.Root) == "" {
|
||||
return nil, errors.New("real recording recovery or Dispatcher delivery unavailable")
|
||||
}
|
||||
info, err := os.Stat(r.Delivery.Recovery.Root)
|
||||
if err != nil || !info.IsDir() || info.Mode().Perm() != 0700 {
|
||||
return nil, errors.New("real recording recovery directory must exist with mode 0700")
|
||||
}
|
||||
if r.originator == nil && (r.Loader.ConfigDir == "" || r.Loader.Asterisk == "" || r.Loader.LibraryDir == "") {
|
||||
return nil, errors.New("native Asterisk configuration required for real call")
|
||||
}
|
||||
if _, err := ai.NewCall(approved.AI, func(context.Context) error { return nil }); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
originate := r.originator
|
||||
if originate == nil {
|
||||
originate = func(ctx context.Context, dial asterisk.NativeDial) (realMedia, error) {
|
||||
call, err := r.Loader.Originate(ctx, dial)
|
||||
if err != nil {
|
||||
return realMedia{}, err
|
||||
}
|
||||
return realMedia{ctx: call.Context, session: call.Media, close: call.Close}, nil
|
||||
}
|
||||
}
|
||||
dial := asterisk.NativeDial{
|
||||
ExecutionID: approved.SourceEventID, TrunkID: approved.SelectedTrunkID,
|
||||
DialedCallee: approved.DialedCallee, CallerID: approved.CallerID,
|
||||
AnswerTimeout: approved.RingTimeout, MediaPayloadType: r.MediaPayloadType,
|
||||
}
|
||||
var started atomic.Bool
|
||||
return func(ctx context.Context) error {
|
||||
if !started.CompareAndSwap(false, true) {
|
||||
return errors.New("approved real call already started; refusing a second origination")
|
||||
}
|
||||
if ctx == nil || ctx.Err() != nil || !time.Now().Before(approved.DialBefore) {
|
||||
return errors.New("real call instruction expired or cancelled before origination")
|
||||
}
|
||||
call, err := originate(ctx, dial)
|
||||
if err != nil {
|
||||
return err // unknown origination is never retried or reported as a completed call
|
||||
}
|
||||
if call.ctx == nil || call.session == nil || call.close == nil {
|
||||
if call.close != nil {
|
||||
return errors.Join(errors.New("real call media unavailable"), call.close())
|
||||
}
|
||||
return errors.New("real call has no hangup function; outcome unknown")
|
||||
}
|
||||
capture, err := callflow.NewRecordingSession(call.session, r.MaxWAVBytes)
|
||||
if err != nil {
|
||||
return errors.Join(err, call.close())
|
||||
}
|
||||
startedAt := time.Now().UTC() // Originate returns only after the actual channel entered Stasis.
|
||||
hangup := func(context.Context) error { return call.close() }
|
||||
pipeline, err := ai.NewCall(approved.AI, hangup)
|
||||
if err != nil {
|
||||
return errors.Join(err, call.close())
|
||||
}
|
||||
observed, runErr := RunApprovedCall(call.ctx, approved, capture, hangup, pipeline)
|
||||
if err := call.close(); err != nil {
|
||||
return errors.Join(runErr, err) // do not report an unconfirmed end or release its occupancy
|
||||
}
|
||||
endedAt := time.Now().UTC()
|
||||
outcome, reason := "answered", "real call ended"
|
||||
if runErr != nil {
|
||||
outcome, reason = "failed", "real AI or media flow failed"
|
||||
}
|
||||
payload, err := callflow.FinalResultPayload(callflow.FinalCallFacts{
|
||||
TaskID: approved.TaskID, CallerProfileID: approved.CallerProfileID,
|
||||
Callee: approved.Callee, TrunkID: approved.SelectedTrunkID,
|
||||
StartedAt: startedAt, EndedAt: endedAt, Outcome: outcome, ReasonMessage: reason,
|
||||
}, observed)
|
||||
if err != nil {
|
||||
return errors.Join(runErr, err)
|
||||
}
|
||||
wav, durationMS, captureErr := capture.WAV()
|
||||
completed := agent.CompletedRecording{ResultPayload: payload, Expected: true, CaptureError: captureErr}
|
||||
if captureErr == nil {
|
||||
completed.RecordingID = "recording-" + approved.SourceEventID
|
||||
completed.UploadID = "upload-" + approved.SourceEventID
|
||||
completed.WAV, completed.DurationMS = wav, durationMS
|
||||
}
|
||||
reportCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), r.ReportTimeout)
|
||||
defer cancel()
|
||||
return errors.Join(runErr, r.Delivery.Complete(reportCtx, completed))
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (r *ApprovedRecordedRealCall) Run(ctx context.Context, approved ApprovedExecution) error {
|
||||
if ctx == nil {
|
||||
return errors.New("real call requires a context")
|
||||
}
|
||||
prepared, err := r.Prepare(approved)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return prepared(ctx)
|
||||
}
|
||||
@@ -0,0 +1,105 @@
|
||||
package rpc
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"io"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.ipao.vip/rogee/go-sip/internal/asterisk"
|
||||
"git.ipao.vip/rogee/go-sip/internal/callflow"
|
||||
"git.ipao.vip/rogee/go-sip/internal/media"
|
||||
)
|
||||
|
||||
type endedRealMedia struct{}
|
||||
|
||||
func (endedRealMedia) ReadPayload(context.Context) ([]byte, error) { return nil, io.EOF }
|
||||
func (endedRealMedia) SendPCM16(context.Context, []byte, int) error { return nil }
|
||||
func (endedRealMedia) Stats() media.RTPStats { return media.RTPStats{} }
|
||||
|
||||
func TestApprovedRecordedRealCallReportsOnlyEndedObservedCall(t *testing.T) {
|
||||
fixture, stub, puts, _, approved := recordedMockFixture(t, nil, time.Second)
|
||||
approved.CallerID, approved.DialedCallee, approved.RingTimeout = "BD93205882", "7089"+approved.Callee, time.Second
|
||||
var originate, hangup int
|
||||
runner := &ApprovedRecordedRealCall{
|
||||
MaxWAVBytes: 4096, MediaPayloadType: 118, ReportTimeout: time.Second,
|
||||
Delivery: fixture.Delivery,
|
||||
originator: func(ctx context.Context, dial asterisk.NativeDial) (realMedia, error) {
|
||||
originate++
|
||||
if dial.ExecutionID != approved.SourceEventID || dial.DialedCallee != approved.DialedCallee || dial.CallerID != approved.CallerID || dial.TrunkID != approved.SelectedTrunkID || dial.MediaPayloadType != 118 {
|
||||
t.Errorf("signed SIP identity or media profile changed: %+v", dial)
|
||||
}
|
||||
return realMedia{ctx: ctx, session: endedRealMedia{}, close: func() error { hangup++; return nil }}, nil
|
||||
},
|
||||
}
|
||||
err := runner.Run(context.Background(), approved)
|
||||
if err == nil || originate != 1 || hangup != 1 || strings.Join(stub.calls, ",") != "end,result" || puts.Load() != 0 {
|
||||
t.Fatalf("real call must end and report actual media failure once: err=%v originate=%d hangup=%d calls=%v puts=%d", err, originate, hangup, stub.calls, puts.Load())
|
||||
}
|
||||
var result struct {
|
||||
Outcome string `json:"outcome"`
|
||||
Recording map[string]any `json:"recording"`
|
||||
}
|
||||
if err := json.Unmarshal(stub.result, &result); err != nil || result.Outcome != "failed" || len(result.Recording) != 0 {
|
||||
t.Fatalf("empty media must not become an answered recording: result=%+v err=%v", result, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestApprovedRecordedRealCallDoesNotReportUnknownHangup(t *testing.T) {
|
||||
fixture, stub, puts, _, approved := recordedMockFixture(t, nil, time.Second)
|
||||
approved.CallerID, approved.DialedCallee, approved.RingTimeout = "BD93205882", "7089"+approved.Callee, time.Second
|
||||
var originate, hangup int
|
||||
runner := &ApprovedRecordedRealCall{
|
||||
MaxWAVBytes: 4096, MediaPayloadType: 118, ReportTimeout: time.Second,
|
||||
Delivery: fixture.Delivery,
|
||||
originator: func(ctx context.Context, _ asterisk.NativeDial) (realMedia, error) {
|
||||
originate++
|
||||
return realMedia{ctx: ctx, session: endedRealMedia{}, close: func() error { hangup++; return errors.New("ARI hangup outcome unknown") }}, nil
|
||||
},
|
||||
}
|
||||
if err := runner.Run(context.Background(), approved); err == nil || originate != 1 || hangup != 1 || len(stub.calls) != 0 || puts.Load() != 0 {
|
||||
t.Fatalf("unknown hangup must not forge call end or retry: err=%v originate=%d hangup=%d calls=%v puts=%d", err, originate, hangup, stub.calls, puts.Load())
|
||||
}
|
||||
}
|
||||
|
||||
func TestApprovedRecordedRealCallCannotOriginateTwiceOrReportUnknownOrigination(t *testing.T) {
|
||||
fixture, stub, _, _, approved := recordedMockFixture(t, nil, time.Second)
|
||||
approved.CallerID, approved.DialedCallee, approved.RingTimeout = "BD93205882", "7089"+approved.Callee, time.Second
|
||||
originate := 0
|
||||
runner := &ApprovedRecordedRealCall{
|
||||
MaxWAVBytes: 4096, MediaPayloadType: 118, ReportTimeout: time.Second,
|
||||
Delivery: fixture.Delivery,
|
||||
originator: func(context.Context, asterisk.NativeDial) (realMedia, error) {
|
||||
originate++
|
||||
return realMedia{}, errors.New("ARI request outcome unknown")
|
||||
},
|
||||
}
|
||||
prepared, err := runner.Prepare(approved)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := prepared(context.Background()); err == nil {
|
||||
t.Fatal("uncertain origination must not become a completed call")
|
||||
}
|
||||
if err := prepared(context.Background()); err == nil || originate != 1 || len(stub.calls) != 0 {
|
||||
t.Fatalf("same accepted instruction must not reoriginate or report: err=%v originate=%d calls=%v", err, originate, stub.calls)
|
||||
}
|
||||
}
|
||||
|
||||
func TestApprovedRecordedRealCallRejectsMissingDeliveryBeforeOrigination(t *testing.T) {
|
||||
_, _, _, _, approved := recordedMockFixture(t, nil, time.Second)
|
||||
called := false
|
||||
runner := &ApprovedRecordedRealCall{
|
||||
MaxWAVBytes: 4096, MediaPayloadType: 118, ReportTimeout: time.Second,
|
||||
originator: func(context.Context, asterisk.NativeDial) (realMedia, error) {
|
||||
called = true
|
||||
return realMedia{session: callflow.NewMemorySession(nil)}, nil
|
||||
},
|
||||
}
|
||||
if _, err := runner.Prepare(approved); err == nil || called {
|
||||
t.Fatalf("no validated delivery must block before originating: err=%v called=%v", err, called)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user