Bind approved call worker lifetime to Agent task controls

This commit is contained in:
2026-09-30 04:15:36 +08:00
parent 9b1f539538
commit 0e9362447b
6 changed files with 495 additions and 4 deletions
@@ -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 结果代签。
+12 -2
View File
@@ -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
}
+30
View File
@@ -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
}
+170
View File
@@ -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)
}
})
}
}
+88
View File
@@ -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
}
+192
View File
@@ -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)
}
}