Fence Agent task calls during pause and stop drain
This commit is contained in:
@@ -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 结果代签。
|
||||
|
||||
|
||||
@@ -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 != ""
|
||||
}
|
||||
@@ -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()
|
||||
}
|
||||
Reference in New Issue
Block a user