diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index 281f81a..4f30559 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -57,6 +57,7 @@ - 对话控制隔离验证:共享 `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 的持久控制仍为权威,不能据此称真实通话已受控。 - 验证:`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/agent/task_calls.go b/internal/agent/task_calls.go new file mode 100644 index 0000000..dd8211f --- /dev/null +++ b/internal/agent/task_calls.go @@ -0,0 +1,133 @@ +package agent + +import ( + "context" + "errors" + "fmt" + "sync" +) + +var ( + ErrTaskAdmissionClosed = errors.New("Agent task admission is paused") + ErrTaskStopped = errors.New("Agent task is stopped") +) + +type TaskIdentity struct { + DispatcherID string + TenantID int64 + TaskID string +} + +type activeTaskCall struct { + cancel context.CancelFunc + done chan struct{} +} + +type taskCalls struct { + state string + active map[string]*activeTaskCall +} + +// TaskCalls keeps an in-process execution barrier for task-level controls. +// Dispatcher owns durable task state; an Agent restart never authorizes a new +// call without a fresh Dispatcher instruction. The caller must authenticate +// the active Dispatcher session before invoking Register or Apply. +type TaskCalls struct { + mu sync.Mutex + tasks map[TaskIdentity]*taskCalls +} + +func (r *TaskCalls) Register(task TaskIdentity, callID string, cancel context.CancelFunc) (func(), error) { + if r == nil || !validTaskIdentity(task) || callID == "" || cancel == nil { + return nil, errors.New("Agent call requires a task identity, call ID and cancellation") + } + r.mu.Lock() + defer r.mu.Unlock() + entry := r.entry(task) + switch entry.state { + case "stopped": + return nil, ErrTaskStopped + case "paused": + return nil, ErrTaskAdmissionClosed + } + if _, exists := entry.active[callID]; exists { + return nil, errors.New("Agent call is already active") + } + call := &activeTaskCall{cancel: cancel, done: make(chan struct{})} + entry.active[callID] = call + var once sync.Once + return func() { + once.Do(func() { + r.mu.Lock() + delete(entry.active, callID) + close(call.done) + r.mu.Unlock() + }) + }, nil +} + +// Apply closes admission before requesting hangup or waiting for drain. Its +// success means all calls active when the control arrived have finished; a +// timed-out wait leaves the pause/stop barrier in place for re-delivery. +func (r *TaskCalls) Apply(ctx context.Context, task TaskIdentity, action, policy string) error { + if r == nil || ctx == nil || !validTaskIdentity(task) { + return errors.New("Agent task control requires a context and task identity") + } + if err := ctx.Err(); err != nil { + return err + } + if (action != "pause" && action != "stop" && action != "resume") || + (action == "resume" && policy != "") || + (action != "resume" && policy != "hangup" && policy != "drain") { + return errors.New("Agent task control action or active-call policy is invalid") + } + r.mu.Lock() + entry := r.entry(task) + if entry.state == "stopped" && action != "stop" { + r.mu.Unlock() + return ErrTaskStopped + } + if action == "resume" { + entry.state = "" + r.mu.Unlock() + return nil + } + entry.state = "paused" + if action == "stop" { + entry.state = "stopped" + } + active := make([]*activeTaskCall, 0, len(entry.active)) + for _, call := range entry.active { + active = append(active, call) + } + r.mu.Unlock() + if policy == "hangup" { + for _, call := range active { + call.cancel() + } + } + for _, call := range active { + select { + case <-call.done: + case <-ctx.Done(): + return fmt.Errorf("Agent task control waiting for active calls: %w", ctx.Err()) + } + } + return nil +} + +func (r *TaskCalls) entry(task TaskIdentity) *taskCalls { + if r.tasks == nil { + r.tasks = make(map[TaskIdentity]*taskCalls) + } + entry := r.tasks[task] + if entry == nil { + entry = &taskCalls{active: make(map[string]*activeTaskCall)} + r.tasks[task] = entry + } + return entry +} + +func validTaskIdentity(task TaskIdentity) bool { + return task.DispatcherID != "" && task.TenantID > 0 && task.TaskID != "" +} diff --git a/internal/agent/task_calls_test.go b/internal/agent/task_calls_test.go new file mode 100644 index 0000000..5603e23 --- /dev/null +++ b/internal/agent/task_calls_test.go @@ -0,0 +1,120 @@ +package agent + +import ( + "context" + "errors" + "sync/atomic" + "testing" + "time" +) + +var taskFixture = TaskIdentity{DispatcherID: "c046b893-8628-4589-ae50-619d049248a6", TenantID: 42, TaskID: "task-a"} + +func TestTaskCallsPauseHangupWaitsForCallAndRefusesNewWork(t *testing.T) { + var calls TaskCalls + canceled := make(chan struct{}) + release, err := calls.Register(taskFixture, "call-a", func() { close(canceled) }) + if err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) + defer cancel() + completed := make(chan error, 1) + go func() { completed <- calls.Apply(ctx, taskFixture, "pause", "hangup") }() + select { + case <-canceled: + case <-ctx.Done(): + t.Fatal("pause did not request hangup") + } + if _, err := calls.Register(taskFixture, "call-b", func() {}); !errors.Is(err, ErrTaskAdmissionClosed) { + t.Fatalf("pause admitted another call: %v", err) + } + select { + case err := <-completed: + t.Fatalf("pause acknowledged before active call ended: %v", err) + default: + } + release() + if err := <-completed; err != nil { + t.Fatalf("pause did not wait for confirmed completion: %v", err) + } + if err := calls.Apply(ctx, taskFixture, "resume", ""); err != nil { + t.Fatalf("Dispatcher-approved resume was rejected: %v", err) + } + releaseAgain, err := calls.Register(taskFixture, "call-b", func() {}) + if err != nil { + t.Fatalf("resumed task did not admit work: %v", err) + } + releaseAgain() +} + +func TestTaskCallsStopDrainKeepsBarrierAfterTimeoutAndCannotResume(t *testing.T) { + var calls TaskCalls + var hangups atomic.Int32 + release, err := calls.Register(taskFixture, "call-a", func() { hangups.Add(1) }) + if err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithTimeout(context.Background(), 20*time.Millisecond) + defer cancel() + if err := calls.Apply(ctx, taskFixture, "stop", "drain"); !errors.Is(err, context.DeadlineExceeded) { + t.Fatalf("stop/drain acknowledged an active call: %v", err) + } + if hangups.Load() != 0 { + t.Fatal("drain hung up an active call") + } + if _, err := calls.Register(taskFixture, "call-b", func() {}); !errors.Is(err, ErrTaskStopped) { + t.Fatalf("stop barrier reopened after timeout: %v", err) + } + release() + if err := calls.Apply(context.Background(), taskFixture, "stop", "drain"); err != nil { + t.Fatalf("re-delivered stop did not complete after draining: %v", err) + } + if err := calls.Apply(context.Background(), taskFixture, "resume", ""); !errors.Is(err, ErrTaskStopped) { + t.Fatalf("stopped task resumed: %v", err) + } +} + +func TestTaskCallsPreserveTaskIsolationAndValidateControls(t *testing.T) { + var calls TaskCalls + for _, tc := range []struct{ action, policy string }{ + {"pause", ""}, {"resume", "hangup"}, {"restart", "hangup"}, {"stop", "unknown"}, + } { + if err := calls.Apply(context.Background(), taskFixture, tc.action, tc.policy); err == nil { + t.Fatalf("invalid control %q/%q was accepted", tc.action, tc.policy) + } + } + if _, err := calls.Register(TaskIdentity{}, "call", func() {}); err == nil { + t.Fatal("missing task identity was accepted") + } + if err := calls.Apply(context.Background(), TaskIdentity{}, "stop", "hangup"); err == nil { + t.Fatal("control without an approved task identity was accepted") + } + if _, err := calls.Register(taskFixture, "call", nil); err == nil { + t.Fatal("missing call cancellation was accepted") + } + other := taskFixture + other.TaskID = "task-b" + releaseOther, err := calls.Register(other, "call-a", func() {}) + if err != nil { + t.Fatal(err) + } + releaseA, err := calls.Register(taskFixture, "call-a", func() {}) + if err != nil { + t.Fatal(err) + } + if _, err := calls.Register(taskFixture, "call-a", func() {}); err == nil { + t.Fatal("duplicate live execution was admitted") + } + releaseA() + releaseA() // repeated terminal observation must not close the channel twice + if err := calls.Apply(context.Background(), taskFixture, "pause", "drain"); err != nil { + t.Fatalf("released call remained active: %v", err) + } + releaseOtherB, err := calls.Register(other, "call-b", func() {}) + if err != nil { + t.Fatalf("other task was fenced by a different task: %v", err) + } + releaseOtherB() + releaseOther() +}