diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index 105c2cc..a9eecbe 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -57,8 +57,9 @@ - 对话控制隔离验证:共享 `callflow` 不再根据转写或回复内的硬编码词推断拒联,仅处理明确的关键词/拒联事实;ASR-only 不播报开场或 TTS 回复。配置的首语音/静默期限、整通话期限与最大轮数交给媒体控制器;回复按 Unicode 字符数分片,火山 TTS 的缓存音频块数超限明确失败而不重试;不支持的插话配置在准入前拒绝。该控制器尚未接入新主 CLI,不能宣称真实通话媒体已验证。 - Agent 隔离媒体入口 `rpc.RunApprovedCall` 仅接收签发的执行快照、媒体会话、挂断动作和显式每通话 AI pipeline;缺失或 typed nil 适配器在读媒体前拒绝,无隐式 SDK/Mock 回退。测试分别由已绑定快照构建 SDK pipeline 与注入受控合成脚本的 `ApprovedMockPipeline`。签发通话期限与 AI 总期限均约束整通电话;ASR-only 不合成开场,full-AI 开场完成才采集语音。父级期限不会被“未检测到语音”掩盖,空媒体 Mock 等到上下文真正结束。当前媒体只处理 16 kHz PCM16,SDK 虽接受、但媒体不能正确处理的 24 kHz ASR/TTS 配置在 Dispatcher 与 Agent 绑定时明确拒绝,不静默播放或识别错速音频。该入口尚未连接主 CLI,也未用于真实拨号。 - 隔离 Mock AI 适配器仅消费显式冻结的模拟媒体与最终识别脚本,不手写或替代供应商协议;缺失/多余轮次、未批准的开场或助理音频明确拒绝。ASR-only 不执行开场与回复;full-AI 开场只允许一次,用户最终识别的已配置字面关键词先挂断、绝不播放该轮候选回复,未配置关键词不推断拒联。任务未配置最大轮数时严格使用共享通话流程的一轮边界,负数仍拒绝。Mock 输出不代表真实 ASR/LLM/TTS 已通过。 -- 新增 Agent 任务级在途通话栅栏 `TaskCalls`:按 D/数字租户/任务隔离;暂停或停止先拒新通话,hangup 请求取消、drain 不挂断,两者均等待已登记通话结束才确认;等待超时仍保留关闭栅栏,停止后不能恢复,同一控制重复执行不另设去重。隔离测试验证两任务互不干扰、重复结束不破坏状态。该栅栏已由活跃会话校验的 Agent RPC 隔离单测调用,但尚未接入实际批准通话的登记/取消和主入口;Dispatcher 的持久控制仍为权威,不能据此称真实通话已受控。 -- 新增内部任务级 `ApplyApprovedTaskControl` Proto、生成物及 Dispatcher `ApprovedOriginator.SendControl`:D 每次发送前从已激活会话领取元数据,只生成独立 RPC 追踪 ID,不从 SaaS 控制虚构 command_id、幂等键或 revision;Agent 校验活跃 D 会话、数字租户、任务和 pause/stop 的 hangup/drain 政策,在 Mock 中等待已登记活动通话结束后才确认,超时保留本地停止栅栏。隔离测试覆盖拒绝伪 D/旧会话、非法字段与政策、无适配器、未知 RPC 结果不自动重试及重复控制的显式重送;本地双向 TLS gRPC 还验证了受信 D 证书经真实生成 Stub 传送控制后,先挂断再等已登记通话结束才确认,不受信的客户端证书与旧会话均拒绝且不能重新准入。**尚无实际批准通话登记或新主 CLI 接线**,不能据此称真实通话控制已验收。 +- 新增 Agent 任务级在途通话栅栏 `TaskCalls`:按 D/数字租户/任务隔离;暂停或停止先拒新通话,hangup 请求取消、drain 不挂断,两者均等待已登记通话结束才确认;等待超时仍保留关闭栅栏,停止后不能恢复,同一控制重复执行不另设去重。隔离测试验证两任务互不干扰、重复结束不破坏状态。该栅栏现由严格 Agent Mock 组装将已授权执行的**合成**在途通话登记到同一个控制器;主 CLI 仍未使用实际媒体入口,Dispatcher 的持久控制仍为权威,不能据此称真实通话已受控。 +- 新增内部任务级 `ApplyApprovedTaskControl` Proto、生成物及 Dispatcher `ApprovedOriginator.SendControl`:D 每次发送前从已激活会话领取元数据,只生成独立 RPC 追踪 ID,不从 SaaS 控制虚构 command_id、幂等键或 revision;Agent 校验活跃 D 会话、数字租户、任务和 pause/stop 的 hangup/drain 政策,在 Mock 中等待已登记活动通话结束后才确认,超时保留本地停止栅栏。隔离测试覆盖拒绝伪 D/旧会话、非法字段与政策、无适配器、未知 RPC 结果不自动重试及重复控制的显式重送;本地双向 TLS gRPC 还验证了受信 D 证书经真实生成 Stub 传送控制后,先挂断再等已登记通话结束才确认,不受信的客户端证书与旧会话均拒绝且不能重新准入。**仅隔离合成 runner 已经登记,实际媒体与新主 CLI 尚未接线**,不能据此称真实通话控制已验收。 +- `ApprovedCallWorker` 与 `NewApprovedAgentServer` 强制把已授权执行和控制 RPC 绑定同一个任务栅栏、显式 Mock 模式及 Agent 进程期限;签发拨号时限在异步工人发出接受回执前重验,超时或任务关闭在无拨号时明确拒绝;启动请求结束不取消在途通话,任务 hangup 等待 runner 明确结束后才完成控制。隔离测试覆盖先执行回执后合成挂断、进程关闭、失败观察、重复身份防重及缺失/竞争适配器拒绝;未知执行不自动重拨。这里只证明调度与取消顺序,**未在主 CLI 连接真实媒体、录音、失败报告或出站 MQ**。 - 验证:`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 结果代签。 diff --git a/internal/rpc/approved_execution.go b/internal/rpc/approved_execution.go index 56af6ab..ec5a071 100644 --- a/internal/rpc/approved_execution.go +++ b/internal/rpc/approved_execution.go @@ -13,6 +13,7 @@ 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/ai" "git.ipao.vip/rogee/go-sip/internal/configread" "google.golang.org/grpc/codes" @@ -139,8 +140,17 @@ func (s *Server) ExecuteApproved(ctx context.Context, req *agentpb.ExecuteApprov DialBefore: time.UnixMilli(req.DialBeforeUnixMs), AI: bound, } if err := s.mockApprovedOriginate(ctx, approved); err != nil { - log.Printf("Agent approved execution outcome unknown: event_id=%q task_id=%q cause_type=%T", req.SourceEventId, req.TaskId, err) - return nil, status.Error(codes.Unavailable, "approved mock execution outcome unknown") + switch { + case errors.Is(err, ErrApprovedDialExpired): + log.Printf("Agent approved execution refused before dial: event_id=%q task_id=%q cause_type=%T", req.SourceEventId, req.TaskId, err) + return nil, status.Error(codes.DeadlineExceeded, "Dispatcher dial authorization expired before issuing") + case errors.Is(err, agent.ErrTaskAdmissionClosed), errors.Is(err, agent.ErrTaskStopped): + log.Printf("Agent approved execution refused before dial: event_id=%q task_id=%q cause_type=%T", req.SourceEventId, req.TaskId, err) + return nil, status.Error(codes.FailedPrecondition, "approved task no longer admits calls") + default: + log.Printf("Agent approved execution outcome unknown: event_id=%q task_id=%q cause_type=%T", req.SourceEventId, req.TaskId, err) + return nil, status.Error(codes.Unavailable, "approved mock execution outcome unknown") + } } return &agentpb.ExecuteApprovedResponse{CallId: req.CallId, Accepted: true}, nil } diff --git a/internal/rpc/approved_server.go b/internal/rpc/approved_server.go new file mode 100644 index 0000000..016f7ed --- /dev/null +++ b/internal/rpc/approved_server.go @@ -0,0 +1,30 @@ +package rpc + +import "errors" + +var ErrApprovedWorkerRequired = errors.New("approved call worker is required") + +// NewApprovedAgentServer binds the control RPC and one-shot originate callback +// to the same task-call registry. The separate legacy hooks are not accepted +// on this current Mock-only entry: otherwise a control could confirm while a +// call started by a different adapter remains active. +func NewApprovedAgentServer(options ServerOptions, worker *ApprovedCallWorker) (*Server, error) { + if worker == nil { + return nil, ErrApprovedWorkerRequired + } + if options.Mode != "mock" || options.StatePath == "" || options.Status == nil || options.Status.AgentId == "" || options.Status.CellId == "" || options.Status.BootId == "" || options.LoadedSIP == nil { + return nil, errors.New("approved Agent requires explicit Mock mode, durable state, identity and observed SIP") + } + if worker.Lifecycle == nil || worker.Calls == nil || worker.Run == nil || worker.OnFailure == nil { + return nil, errors.New("approved Agent requires a process lifecycle, task calls, runner and failure reporting") + } + if options.MockApprovedOriginate != nil || options.ApprovedTaskCalls != nil { + return nil, errors.New("approved Agent cannot use competing call or control adapters") + } + // Copy the configuration so later mutation of worker fields cannot cause + // ExecuteApproved and task controls to observe different registries. + frozen := *worker + options.MockApprovedOriginate = frozen.Originate + options.ApprovedTaskCalls = frozen.Calls + return NewServer(options), nil +} diff --git a/internal/rpc/approved_server_test.go b/internal/rpc/approved_server_test.go new file mode 100644 index 0000000..1b1a089 --- /dev/null +++ b/internal/rpc/approved_server_test.go @@ -0,0 +1,170 @@ +package rpc + +import ( + "context" + "errors" + "path/filepath" + "testing" + "time" + + agentpb "git.ipao.vip/rogee/go-sip/gen/agent" + "git.ipao.vip/rogee/go-sip/internal/agent" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" +) + +func TestApprovedAgentServerRegistersCallsBeforeExecuteAckAndDrainsOnControl(t *testing.T) { + now := time.Date(2026, 9, 20, 10, 0, 0, 0, time.UTC) + req := approvedTestRequest(t, now) + process, stopProcess := context.WithCancel(context.Background()) + defer stopProcess() + calls := &agent.TaskCalls{} + started := make(chan context.Context, 1) + hangupFinished := make(chan struct{}) + defer func() { + select { + case <-hangupFinished: + default: + close(hangupFinished) + } + }() + failures := make(chan error, 1) + worker := &ApprovedCallWorker{Lifecycle: process, Calls: calls, Now: func() time.Time { return now }, + Run: func(ctx context.Context, _ ApprovedExecution) error { + started <- ctx + <-ctx.Done() + <-hangupFinished + return nil + }, + OnFailure: func(_ ApprovedExecution, err error) error { failures <- err; return nil }, + } + server, err := NewApprovedAgentServer(ServerOptions{ + Mode: "mock", StatePath: filepath.Join(t.TempDir(), "agent-session.json"), + Now: func() time.Time { return now }, + Status: &agentpb.AgentStatus{AgentId: "agent-1", CellId: "cell-1", BootId: "boot-1"}, + LoadedSIP: func(context.Context) (map[string]int64, error) { + return map[string]int64{req.SelectedTrunkId: req.SipRevision}, nil + }, + }, worker) + if err != nil { + t.Fatal(err) + } + if _, err := server.ActivateAgent(context.Background(), &agentpb.ActivateAgentRequest{ + Meta: testMeta("activate-approved", "", 0), + Binding: &agentpb.AgentBinding{AgentId: "agent-1", CellId: "cell-1", ExpectedBootId: "boot-1", + DispatcherEpoch: "epoch-1", SessionGeneration: 1, DispatcherId: req.DispatcherId}, + ActivationOperationId: "activate-approved", + }); err != nil { + t.Fatal(err) + } + request, closeRequest := context.WithCancel(context.Background()) + response, err := server.ExecuteApproved(request, req) + if err != nil || response == nil || !response.Accepted { + t.Fatalf("registered call did not receive dispatch ACK: response=%v err=%v", response, err) + } + closeRequest() // the outbound gRPC request lifetime has ended + var callContext context.Context + select { + case callContext = <-started: + case <-time.After(2 * time.Second): + t.Fatal("approved worker did not begin after ACK") + } + if callContext.Err() != nil { + t.Fatal("active call inherited the completed unary context") + } + control := &agentpb.ApplyApprovedTaskControlRequest{ + Meta: testMeta("control", "", 1), DispatcherId: req.DispatcherId, + TenantId: req.TenantId, TaskId: req.TaskId, + Action: agentpb.ControlAction_CONTROL_ACTION_STOP, ActiveCallPolicy: agentpb.ActiveCallPolicy_ACTIVE_CALL_POLICY_HANGUP, + } + controlResult := make(chan error, 1) + go func() { _, err := server.ApplyApprovedTaskControl(context.Background(), control); controlResult <- err }() + select { + case <-callContext.Done(): + case <-time.After(2 * time.Second): + t.Fatal("stop did not reach the registered call") + } + select { + case err := <-controlResult: + t.Fatalf("stop was acknowledged before media hangup: %v", err) + default: + } + close(hangupFinished) + if err := <-controlResult; err != nil { + t.Fatalf("stop did not wait for the actual runner: %v", err) + } + if _, err := server.ExecuteApproved(context.Background(), req); status.Code(err) != codes.FailedPrecondition { + t.Fatalf("completed or unknown durable identity was redialed: %v", err) + } + select { + case err := <-failures: + t.Fatalf("successful controlled hangup fabricated an execution failure: %v", err) + default: + } +} + +func TestApprovedAgentServerRefusesMissingOrCompetingAdapters(t *testing.T) { + process := context.Background() + calls := &agent.TaskCalls{} + worker := &ApprovedCallWorker{Lifecycle: process, Calls: calls, + Run: func(context.Context, ApprovedExecution) error { return nil }, + OnFailure: func(ApprovedExecution, error) error { return nil }, + } + valid := ServerOptions{Mode: "mock", StatePath: filepath.Join(t.TempDir(), "agent-session.json"), + Status: &agentpb.AgentStatus{AgentId: "agent-1", CellId: "cell-1", BootId: "boot-1"}, + LoadedSIP: func(context.Context) (map[string]int64, error) { return map[string]int64{"trunk-mock": 8}, nil }, + } + for _, tc := range []struct { + name string + edit func(*ServerOptions, *ApprovedCallWorker) + }{ + {"missing state", func(o *ServerOptions, _ *ApprovedCallWorker) { o.StatePath = "" }}, + {"non-Mock mode", func(o *ServerOptions, _ *ApprovedCallWorker) { o.Mode = "real" }}, + {"implicit Mock default", func(o *ServerOptions, _ *ApprovedCallWorker) { o.Mode = "" }}, + {"missing SIP observation", func(o *ServerOptions, _ *ApprovedCallWorker) { o.LoadedSIP = nil }}, + {"missing Agent identity", func(o *ServerOptions, _ *ApprovedCallWorker) { o.Status = nil }}, + {"competing originator", func(o *ServerOptions, _ *ApprovedCallWorker) { + o.MockApprovedOriginate = func(context.Context, ApprovedExecution) error { return nil } + }}, + {"competing control registry", func(o *ServerOptions, _ *ApprovedCallWorker) { o.ApprovedTaskCalls = &agent.TaskCalls{} }}, + {"missing worker process", func(_ *ServerOptions, w *ApprovedCallWorker) { w.Lifecycle = nil }}, + {"missing failure handler", func(_ *ServerOptions, w *ApprovedCallWorker) { w.OnFailure = nil }}, + } { + t.Run(tc.name, func(t *testing.T) { + options, isolated := valid, *worker + tc.edit(&options, &isolated) + if server, err := NewApprovedAgentServer(options, &isolated); err == nil || server != nil { + t.Fatalf("unsafe approved server setup was admitted: %v", err) + } + }) + } + if server, err := NewApprovedAgentServer(valid, nil); err == nil || server != nil || !errors.Is(err, ErrApprovedWorkerRequired) { + t.Fatalf("missing approved call worker was admitted: %v", err) + } +} + +func TestExecuteApprovedReturnsDefinitiveNoDialRefusalWithoutRetry(t *testing.T) { + now := time.Date(2026, 9, 20, 10, 0, 0, 0, time.UTC) + for _, tc := range []struct { + name string + cause error + want codes.Code + }{ + {"paused task", agent.ErrTaskAdmissionClosed, codes.FailedPrecondition}, + {"stopped task", agent.ErrTaskStopped, codes.FailedPrecondition}, + {"expired before dispatch", ErrApprovedDialExpired, codes.DeadlineExceeded}, + } { + t.Run(tc.name, func(t *testing.T) { + req := approvedTestRequest(t, now) + attempts := 0 + server := activatedApprovedServer(t, now, filepath.Join(t.TempDir(), "agent-session.json"), req.DispatcherId, 1, + func(context.Context, ApprovedExecution) error { attempts++; return tc.cause }) + if _, err := server.ExecuteApproved(context.Background(), req); status.Code(err) != tc.want { + t.Fatalf("known no-dial was treated as ambiguous: %v", err) + } + if _, err := server.ExecuteApproved(context.Background(), req); status.Code(err) != codes.FailedPrecondition || attempts != 1 { + t.Fatalf("durable no-redial guard was bypassed: attempts=%d err=%v", attempts, err) + } + }) + } +} diff --git a/internal/rpc/approved_worker.go b/internal/rpc/approved_worker.go new file mode 100644 index 0000000..b3e4eaf --- /dev/null +++ b/internal/rpc/approved_worker.go @@ -0,0 +1,88 @@ +package rpc + +import ( + "context" + "errors" + "log" + "time" + + "git.ipao.vip/rogee/go-sip/internal/agent" +) + +var ErrApprovedDialExpired = errors.New("Dispatcher-approved dial deadline expired") + +// ApprovedCallWorker keeps a call registered until its synchronous runner has +// finished the physical call. Run must not return before hangup; it may report +// call results independently, but must never start a second originate. The +// process lifecycle, not the initiating Unary request, owns the call context. +type ApprovedCallWorker struct { + Lifecycle context.Context + Calls *agent.TaskCalls + Run func(context.Context, ApprovedExecution) error + OnFailure func(ApprovedExecution, error) error + Now func() time.Time +} + +// Originate returns after registration and launch, allowing Dispatcher to +// acknowledge issuance while the call remains cancellable by task control. +// If scheduling delays the worker past the signed deadline it reports failure +// without dialing; an unknown durable execution is never automatically retried. +func (w *ApprovedCallWorker) Originate(request context.Context, approved ApprovedExecution) error { + if w == nil || request == nil || w.Lifecycle == nil || w.Calls == nil || w.Run == nil || w.OnFailure == nil { + return errors.New("approved call lifecycle, runner and failure reporting are required") + } + if err := request.Err(); err != nil { + return err + } + if err := w.Lifecycle.Err(); err != nil { + return err + } + if approved.DispatcherID == "" || approved.TenantID <= 0 || approved.TaskID == "" || approved.CallID == "" || approved.MaxCallDuration <= 0 || approved.DialBefore.IsZero() { + return errors.New("approved call identity, duration or dial deadline is incomplete") + } + now := w.Now + if now == nil { + now = time.Now + } + if !now().Before(approved.DialBefore) { + return ErrApprovedDialExpired + } + callContext, cancel := context.WithTimeout(w.Lifecycle, approved.MaxCallDuration) + release, err := w.Calls.Register(agent.TaskIdentity{DispatcherID: approved.DispatcherID, TenantID: approved.TenantID, TaskID: approved.TaskID}, approved.CallID, cancel) + if err != nil { + cancel() + return err + } + issued := make(chan error, 1) + go func() { + // Deferred cleanup also runs if a broken runner panics; panics are not + // swallowed, and the process does not claim a successful outcome. + defer cancel() + defer release() + if err := callContext.Err(); err != nil { + issued <- err + return + } + if err := request.Err(); err != nil { + issued <- err + return + } + if !now().Before(approved.DialBefore) { + issued <- ErrApprovedDialExpired + return + } + issued <- nil // the Unary response may now acknowledge scheduled work + runErr := w.Run(callContext, approved) + // Physical execution resources are released independently of any later + // failure reporting, OSS work or MQ publisher confirmation. + release() + if runErr == nil { + return + } + log.Printf("Agent approved call ended with error: event_id=%q task_id=%q cause_type=%T", approved.SourceEventID, approved.TaskID, runErr) + if err := w.OnFailure(approved, runErr); err != nil { + log.Printf("Agent approved call failure reporting failed: event_id=%q task_id=%q cause_type=%T", approved.SourceEventID, approved.TaskID, err) + } + }() + return <-issued +} diff --git a/internal/rpc/approved_worker_test.go b/internal/rpc/approved_worker_test.go new file mode 100644 index 0000000..5277b4a --- /dev/null +++ b/internal/rpc/approved_worker_test.go @@ -0,0 +1,192 @@ +package rpc + +import ( + "context" + "errors" + "strings" + "sync/atomic" + "testing" + "time" + + "git.ipao.vip/rogee/go-sip/internal/agent" +) + +func workerTestExecution(now time.Time) ApprovedExecution { + return ApprovedExecution{DispatcherID: "c046b893-8628-4589-ae50-619d049248a6", TenantID: 42, + TaskID: "task-a", SourceEventID: "event-a", CallID: "call-a", + DialBefore: now.Add(time.Minute), MaxCallDuration: time.Minute} +} + +func TestApprovedCallWorkerOutlivesInitiatingUnaryAndWaitsForPhysicalEnd(t *testing.T) { + now := time.Date(2026, 9, 20, 10, 0, 0, 0, time.UTC) + process, stopProcess := context.WithCancel(context.Background()) + defer stopProcess() + request, closeRequest := context.WithCancel(context.Background()) + defer closeRequest() + calls := &agent.TaskCalls{} + started := make(chan context.Context, 1) + finish := make(chan struct{}) + defer func() { + select { + case <-finish: + default: + close(finish) + } + }() + failures := make(chan error, 1) + worker := ApprovedCallWorker{Lifecycle: process, Calls: calls, Now: func() time.Time { return now }, + Run: func(ctx context.Context, _ ApprovedExecution) error { + started <- ctx + <-ctx.Done() + <-finish // simulated hangup must finish before TaskCalls releases the call + return nil + }, + OnFailure: func(_ ApprovedExecution, err error) error { failures <- err; return nil }, + } + approved := workerTestExecution(now) + if err := worker.Originate(request, approved); err != nil { + t.Fatal(err) + } + var callContext context.Context + select { + case callContext = <-started: + case <-time.After(2 * time.Second): + t.Fatal("accepted call did not start") + } + closeRequest() // the unary response must not cancel the active call + if callContext.Err() != nil { + t.Fatal("call inherited the completed unary request deadline") + } + task := agent.TaskIdentity{DispatcherID: approved.DispatcherID, TenantID: approved.TenantID, TaskID: approved.TaskID} + if _, err := calls.Register(task, approved.CallID, func() {}); err == nil { + t.Fatal("accepted call was not registered before the RPC returned") + } + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + controlDone := make(chan error, 1) + go func() { controlDone <- calls.Apply(ctx, task, "pause", "hangup") }() + select { + case <-callContext.Done(): + case <-ctx.Done(): + t.Fatal("Agent control did not cancel the active call") + } + select { + case err := <-controlDone: + t.Fatalf("control was acknowledged before physical hangup finished: %v", err) + default: + } + close(finish) + if err := <-controlDone; err != nil { + t.Fatalf("control did not wait for the Agent runner: %v", err) + } + select { + case err := <-failures: + t.Fatalf("successful cancellation generated a false failure: %v", err) + default: + } +} + +func TestApprovedCallWorkerRejectsExpiredOrClosedTaskBeforeIssuing(t *testing.T) { + now := time.Date(2026, 9, 20, 10, 0, 0, 0, time.UTC) + process, stopProcess := context.WithCancel(context.Background()) + defer stopProcess() + calls := &agent.TaskCalls{} + var starts atomic.Int32 + worker := ApprovedCallWorker{Lifecycle: process, Calls: calls, Now: func() time.Time { return now }, + Run: func(context.Context, ApprovedExecution) error { starts.Add(1); return nil }, + OnFailure: func(ApprovedExecution, error) error { return nil }, + } + approved := workerTestExecution(now) + approved.DialBefore = now.Add(-time.Millisecond) + if err := worker.Originate(context.Background(), approved); !errors.Is(err, ErrApprovedDialExpired) { + t.Fatalf("expired authorization was admitted: %v", err) + } + approved.DialBefore = now.Add(time.Minute) + task := agent.TaskIdentity{DispatcherID: approved.DispatcherID, TenantID: approved.TenantID, TaskID: approved.TaskID} + if err := calls.Apply(context.Background(), task, "pause", "hangup"); err != nil { + t.Fatal(err) + } + if err := worker.Originate(context.Background(), approved); !errors.Is(err, agent.ErrTaskAdmissionClosed) { + t.Fatalf("paused task was admitted: %v", err) + } + if starts.Load() != 0 { + t.Fatal("Agent issued a call despite expired authorization or task pause") + } + worker.OnFailure = nil + if err := worker.Originate(context.Background(), approved); err == nil || !strings.Contains(err.Error(), "failure") { + t.Fatalf("missing failure observability was accepted: %v", err) + } +} + +func TestApprovedCallWorkerProcessStopCancelsAndReportsFailureOnce(t *testing.T) { + now := time.Date(2026, 9, 20, 10, 0, 0, 0, time.UTC) + process, stopProcess := context.WithCancel(context.Background()) + defer stopProcess() + calls := &agent.TaskCalls{} + started := make(chan struct{}) + reported := make(chan error, 1) + var attempts atomic.Int32 + worker := ApprovedCallWorker{Lifecycle: process, Calls: calls, Now: func() time.Time { return now }, + Run: func(ctx context.Context, _ ApprovedExecution) error { + attempts.Add(1) + close(started) + <-ctx.Done() + return ctx.Err() + }, + OnFailure: func(_ ApprovedExecution, err error) error { reported <- err; return nil }, + } + approved := workerTestExecution(now) + if err := worker.Originate(context.Background(), approved); err != nil { + t.Fatal(err) + } + select { + case <-started: + case <-time.After(2 * time.Second): + t.Fatal("accepted worker did not start") + } + stopProcess() + select { + case err := <-reported: + if !errors.Is(err, context.Canceled) { + t.Fatalf("lost process-shutdown failure reason: %v", err) + } + case <-time.After(2 * time.Second): + t.Fatal("worker failed silently after process shutdown") + } + task := agent.TaskIdentity{DispatcherID: approved.DispatcherID, TenantID: approved.TenantID, TaskID: approved.TaskID} + if err := calls.Apply(context.Background(), task, "stop", "drain"); err != nil { + t.Fatalf("failed worker was not released: %v", err) + } + if attempts.Load() != 1 { + t.Fatalf("worker was replayed after an unknown outcome: %d", attempts.Load()) + } +} + +func TestApprovedCallWorkerRefusesExpiredAuthorizationBeforeDispatchAck(t *testing.T) { + now := time.Date(2026, 9, 20, 10, 0, 0, 0, time.UTC) + process, stopProcess := context.WithCancel(context.Background()) + defer stopProcess() + calls := &agent.TaskCalls{} + var checks, starts, reports atomic.Int32 + worker := ApprovedCallWorker{Lifecycle: process, Calls: calls, + Now: func() time.Time { + if checks.Add(1) == 1 { + return now // approved when entering the worker + } + return now.Add(2 * time.Minute) // expired when the scheduled worker begins + }, + Run: func(context.Context, ApprovedExecution) error { starts.Add(1); return nil }, + OnFailure: func(ApprovedExecution, error) error { reports.Add(1); return nil }, + } + approved := workerTestExecution(now) + if err := worker.Originate(context.Background(), approved); !errors.Is(err, ErrApprovedDialExpired) { + t.Fatalf("expired scheduled call received an execution ACK: %v", err) + } + if starts.Load() != 0 || reports.Load() != 0 { + t.Fatal("no-dial pre-dispatch expiry started work or fabricated a call result") + } + task := agent.TaskIdentity{DispatcherID: approved.DispatcherID, TenantID: approved.TenantID, TaskID: approved.TaskID} + if err := calls.Apply(context.Background(), task, "stop", "drain"); err != nil { + t.Fatalf("refused call remained registered: %v", err) + } +}