Validate per-call Mock media before accepting approved execute
This commit is contained in:
@@ -59,7 +59,7 @@
|
||||
- 隔离 Mock AI 适配器仅消费显式冻结的模拟媒体与最终识别脚本,不手写或替代供应商协议;缺失/多余轮次、未批准的开场或助理音频明确拒绝。ASR-only 不执行开场与回复;full-AI 开场只允许一次,用户最终识别的已配置字面关键词先挂断、绝不播放该轮候选回复,未配置关键词不推断拒联。任务未配置最大轮数时严格使用共享通话流程的一轮边界,负数仍拒绝。Mock 输出不代表真实 ASR/LLM/TTS 已通过。
|
||||
- 新增 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**。
|
||||
- `ApprovedCallWorker` 与 `NewApprovedAgentServer` 强制把已授权执行和控制 RPC 绑定同一个任务栅栏、显式 Mock 模式及 Agent 进程期限;每通话先在无拨号/无业务写入前准备并核验签发 AI、合成媒体和交付目标,准备失败即明确拒绝且不发完成事实;签发拨号时限在异步工人发出接受回执前重验,超时或任务关闭同样在无拨号时明确拒绝;启动请求结束不取消在途通话,任务 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 结果代签。
|
||||
|
||||
@@ -76,7 +76,7 @@
|
||||
- Agent 录音交付隔离组件:`RecordingDelivery` 先确认结束,再依照录音是否实际生成分别上报唯一空录音结果或请求原授权并直传内存 WAV;录音生成失败保留通话真实结果、空录音对象及明确原因,不虚构上传事实。隔离测试通过本地 HTTP PUT 和假 Dispatcher RPC 覆盖成功无业务文件、OSS 明确失败后私有文件保存、恢复写入失败、未知 PUT 隔离、重启重领原目标、上传已确认后只重发原结果。再次调用不会隐式重新 PUT;正常已确认上传但尚未被 D 持久收讫的跨进程间隙仍受 K16 边界约束。此处未连接主入口批准执行媒体或 MQ;下述隔离链路另测本地真正的 D gRPC/SQLite。
|
||||
- 最终结果隔离组件:`ApprovedExecution` 已保留获批任务的 `caller_profile_id`;`FinalResultPayload` 只从批准的任务身份、确认时间及实际收到的用户媒体窗口生成当前唯一结果;最终 ASR 文本与字面关键词拒联写入完整转写,未播放的助手候选内容不冒充转写,未知 SIP 响应码保留 `null`,无录音保留 `{}`。单元测试直接核验当前 MQ Schema、字段缺失/乱序时窗与无应答负例;开场白发送失败不得计入已播放音频。此组件尚未接入主入口,也未证明真实 RTP 转写时间精度。
|
||||
- 本地隔离链路:`RunApprovedCall` 显式使用合成 `ApprovedMockPipeline`,让共享批准通话流程按 ASR-only 限制消费实际读出的 PCM Mock 媒体帧;`RecordingSession` 生成内存 WAV,`FinalResultPayload` 仅用最终模拟识别及实测采集时窗构造结果。Agent 经双向 TLS gRPC 向 Dispatcher 确认结束、领取原始资产签名授权,再对本地 HTTPS OSS Mock 单次 PUT;Dispatcher 将原上传事实与唯一最终结果 outbox 同事务提交。测试确认一次 PUT、一次结果、无业务文件,使用官方 SDK 生成的路径;Mock ASR 与 OSS 服务不验证真实供应商协议或签名。这不是主入口执行、RabbitMQ 投递或外部验收。
|
||||
- `ApprovedRecordedMockCall` 将显式批准的 ASR-only AI 快照、合成 PCM 媒体、每通话独立脚本、实际采集的有界 WAV、用户最终识别时窗和 `RecordingDelivery` 组合到一个可供 Agent 工人调用的 Mock runner;隔离测试验证先结束事实、一次本地 OSS PUT、唯一最终结果及正常路径无业务文件。无实际读入媒体时明确报告失败并以空录音、生成失败原因收口,不伪造上传或转写;缺失每通话交付器在启动媒体前报错。此 runner 尚未接到主 CLI 的 `ApprovedCallWorker`,测试 OSS/RPC 均为隔离替身,不能代表真实线路、供应商或 MQ 投递。
|
||||
- `ApprovedRecordedMockCall` 将显式批准的 ASR-only AI 快照、合成 PCM 媒体、每通话独立脚本、实际采集的有界 WAV、用户最终识别时窗和 `RecordingDelivery` 组合到一个可供 Agent 工人调用的 Mock runner;隔离测试验证先结束事实、一次本地 OSS PUT、唯一最终结果及正常路径无业务文件。无实际读入媒体时明确报告失败并以空录音、生成失败原因收口,不伪造上传或转写;缺失每通话交付器、交付器与签发 D/数字租户/事件不一致、私有恢复目录并非 `0700` 或模拟脚本不满足已批准 AI 时,均在发外呼接受回执前拒绝。此 runner 尚未接到主 CLI 的 `ApprovedCallWorker`,测试 OSS/RPC 均为隔离替身,不能代表真实线路、供应商或 MQ 投递。
|
||||
- 已验证:`go test ./... -count=1`、`go test -race ./... -count=1`、`go vet ./...`、`go build ./...`、`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`、`git diff --check`。尚未完成批准执行媒体录音到 Agent 交付组件的主入口接线、主入口 Agent↔Dispatcher 实际会话与录音执行接线(隔离 mTLS gRPC 已测)、真实执行时上传事实/最终结果交付及 MQ/端到端验收,不能宣称 P06 通过。
|
||||
|
||||
## 验收台账
|
||||
|
||||
@@ -144,9 +144,9 @@ func (s *Server) ExecuteApproved(ctx context.Context, req *agentpb.ExecuteApprov
|
||||
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):
|
||||
case errors.Is(err, agent.ErrTaskAdmissionClosed), errors.Is(err, agent.ErrTaskStopped), errors.Is(err, ErrApprovedCallPreparation):
|
||||
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")
|
||||
return nil, status.Error(codes.FailedPrecondition, "approved call was not ready before dial")
|
||||
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")
|
||||
|
||||
@@ -3,6 +3,7 @@ package rpc
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"os"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
@@ -28,59 +29,90 @@ type ApprovedRecordedMockCall struct {
|
||||
Delivery *agent.RecordingDelivery
|
||||
}
|
||||
|
||||
// Run ends the media flow before releasing its caller's task registration.
|
||||
// Reporting uses a separately bounded context so a task hangup cannot silently
|
||||
// suppress its confirmed end or the single final result.
|
||||
func (r *ApprovedRecordedMockCall) Run(ctx context.Context, approved ApprovedExecution) error {
|
||||
if r == nil || ctx == nil || r.Delivery == nil || r.ReportTimeout <= 0 {
|
||||
return ErrApprovedMockDeliveryUnavailable
|
||||
// Prepare validates the per-call synthetic media, approved AI settings and
|
||||
// recording destination before the Agent can acknowledge an originate. It has
|
||||
// no outbound or durable side effects. The returned function stays synchronous
|
||||
// until the Mock media flow has ended and hangup has completed.
|
||||
func (r *ApprovedRecordedMockCall) Prepare(approved ApprovedExecution) (func(context.Context) error, error) {
|
||||
if r == nil || r.Delivery == nil || r.ReportTimeout <= 0 {
|
||||
return nil, ErrApprovedMockDeliveryUnavailable
|
||||
}
|
||||
if approved.DispatcherID == "" || approved.TenantID <= 0 || approved.SourceEventID == "" || approved.CallID != approved.SourceEventID ||
|
||||
approved.TaskID == "" || approved.CallerProfileID == "" || approved.Callee == "" || approved.SelectedTrunkID == "" || r.MaxWAVBytes <= 44 || strings.TrimSpace(r.ReasonMessage) == "" {
|
||||
return errors.New("approved Mock call identity, media limit or observed reason is incomplete")
|
||||
return nil, errors.New("approved Mock call identity, media limit or observed reason is incomplete")
|
||||
}
|
||||
if r.Outcome != "answered" && r.Outcome != "no_answer" && r.Outcome != "failed" {
|
||||
return errors.New("approved Mock call outcome is invalid")
|
||||
return nil, errors.New("approved Mock call outcome is invalid")
|
||||
}
|
||||
if r.ExpectedRecording && (r.Delivery.Recovery == nil || strings.TrimSpace(r.Delivery.Recovery.Root) == "") {
|
||||
return ErrApprovedMockDeliveryUnavailable
|
||||
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 {
|
||||
return nil, ErrApprovedMockDeliveryUnavailable
|
||||
}
|
||||
if r.ExpectedRecording {
|
||||
if r.Delivery.Recovery == nil || strings.TrimSpace(r.Delivery.Recovery.Root) == "" {
|
||||
return nil, ErrApprovedMockDeliveryUnavailable
|
||||
}
|
||||
info, err := os.Stat(r.Delivery.Recovery.Root)
|
||||
if err != nil || !info.IsDir() || info.Mode().Perm() != 0700 {
|
||||
return nil, ErrApprovedMockDeliveryUnavailable
|
||||
}
|
||||
}
|
||||
media, err := callflow.NewRecordingSession(callflow.NewMemorySession(r.InboundPCM16), r.MaxWAVBytes)
|
||||
if err != nil {
|
||||
return err
|
||||
return nil, err
|
||||
}
|
||||
// The isolated memory session has no external SIP channel to hang up.
|
||||
hangup := func(context.Context) error { return nil }
|
||||
pipeline, err := ai.NewApprovedMockPipeline(approved.AI, r.Script, hangup)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
outcome, reason, expected := r.Outcome, r.ReasonMessage, r.ExpectedRecording
|
||||
delivery, timeout := r.Delivery, r.ReportTimeout
|
||||
return func(ctx context.Context) error {
|
||||
if ctx == nil {
|
||||
return errors.New("approved Mock call requires a context")
|
||||
}
|
||||
startedAt := time.Now().UTC()
|
||||
observed, runErr := RunApprovedCall(ctx, approved, media, hangup, pipeline)
|
||||
endedAt := time.Now().UTC()
|
||||
resultOutcome, resultReason := outcome, reason
|
||||
if runErr != nil {
|
||||
resultOutcome, resultReason = "failed", "isolated Mock media did not complete"
|
||||
}
|
||||
payload, err := callflow.FinalResultPayload(callflow.FinalCallFacts{
|
||||
TaskID: approved.TaskID, CallerProfileID: approved.CallerProfileID,
|
||||
Callee: approved.Callee, TrunkID: approved.SelectedTrunkID,
|
||||
StartedAt: startedAt, EndedAt: endedAt, Outcome: resultOutcome, ReasonMessage: resultReason,
|
||||
}, observed)
|
||||
if err != nil {
|
||||
return errors.Join(runErr, err)
|
||||
}
|
||||
completed := agent.CompletedRecording{ResultPayload: payload, Expected: expected}
|
||||
if expected {
|
||||
wav, durationMS, captureErr := media.WAV()
|
||||
completed.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), timeout)
|
||||
defer cancel()
|
||||
return errors.Join(runErr, delivery.Complete(reportCtx, completed))
|
||||
}, nil
|
||||
}
|
||||
|
||||
// Run is a direct local Mock entry. Agent dispatch uses Prepare so invalid
|
||||
// per-call configuration is rejected before the execute acknowledgment.
|
||||
func (r *ApprovedRecordedMockCall) Run(ctx context.Context, approved ApprovedExecution) error {
|
||||
if ctx == nil {
|
||||
return errors.New("approved Mock call requires a context")
|
||||
}
|
||||
prepared, err := r.Prepare(approved)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
startedAt := time.Now().UTC()
|
||||
observed, runErr := RunApprovedCall(ctx, approved, media, hangup, pipeline)
|
||||
endedAt := time.Now().UTC()
|
||||
outcome, reason := r.Outcome, r.ReasonMessage
|
||||
if runErr != nil {
|
||||
outcome, reason = "failed", "isolated Mock media did not complete"
|
||||
}
|
||||
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)
|
||||
}
|
||||
completed := agent.CompletedRecording{ResultPayload: payload, Expected: r.ExpectedRecording}
|
||||
if r.ExpectedRecording {
|
||||
wav, durationMS, captureErr := media.WAV()
|
||||
completed.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))
|
||||
return prepared(ctx)
|
||||
}
|
||||
|
||||
@@ -165,3 +165,30 @@ func TestApprovedRecordedMockCallRefusesMissingPerCallDelivery(t *testing.T) {
|
||||
t.Fatalf("missing per-call delivery was silently accepted: calls=%v err=%v", stub.calls, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestApprovedRecordedMockCallPreparationRejectsScriptBeforeDispatch(t *testing.T) {
|
||||
runner, stub, puts, _, approved := recordedMockFixture(t, bytes.Repeat([]byte{1, 0}, 1600), time.Second)
|
||||
runner.Script.Turns = nil // cannot cover the approved ASR-only turn
|
||||
prepared, err := runner.Prepare(approved)
|
||||
if err == nil || prepared != nil || len(stub.calls) != 0 || puts.Load() != 0 {
|
||||
t.Fatalf("invalid immutable Mock script reached the call or delivery path: calls=%v err=%v", stub.calls, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestApprovedRecordedMockCallPreparationRequiresBoundIdentityAndPrivateRecovery(t *testing.T) {
|
||||
runner, stub, puts, root, approved := recordedMockFixture(t, bytes.Repeat([]byte{1, 0}, 1600), time.Second)
|
||||
runner.Delivery.Call.TenantID = approved.TenantID + 1
|
||||
if prepared, err := runner.Prepare(approved); prepared != nil || !errors.Is(err, ErrApprovedMockDeliveryUnavailable) {
|
||||
t.Fatalf("recording would have been sent under another tenant: %v", err)
|
||||
}
|
||||
runner.Delivery.Call.TenantID = approved.TenantID
|
||||
if err := os.Chmod(root, 0775); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if prepared, err := runner.Prepare(approved); prepared != nil || !errors.Is(err, ErrApprovedMockDeliveryUnavailable) {
|
||||
t.Fatalf("recording started without private recovery directory: %v", err)
|
||||
}
|
||||
if len(stub.calls) != 0 || puts.Load() != 0 {
|
||||
t.Fatal("failed preparation contacted Dispatcher or OSS")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -15,7 +15,7 @@ func NewApprovedAgentServer(options ServerOptions, worker *ApprovedCallWorker) (
|
||||
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 {
|
||||
if worker.Lifecycle == nil || worker.Calls == nil || worker.Prepare == 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 {
|
||||
|
||||
@@ -30,12 +30,12 @@ func TestApprovedAgentServerRegistersCallsBeforeExecuteAckAndDrainsOnControl(t *
|
||||
}()
|
||||
failures := make(chan error, 1)
|
||||
worker := &ApprovedCallWorker{Lifecycle: process, Calls: calls, Now: func() time.Time { return now },
|
||||
Run: func(ctx context.Context, _ ApprovedExecution) error {
|
||||
Prepare: prepareWorkerRun(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{
|
||||
@@ -107,7 +107,7 @@ 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 },
|
||||
Prepare: prepareWorkerRun(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"),
|
||||
@@ -153,6 +153,7 @@ func TestExecuteApprovedReturnsDefinitiveNoDialRefusalWithoutRetry(t *testing.T)
|
||||
{"paused task", agent.ErrTaskAdmissionClosed, codes.FailedPrecondition},
|
||||
{"stopped task", agent.ErrTaskStopped, codes.FailedPrecondition},
|
||||
{"expired before dispatch", ErrApprovedDialExpired, codes.DeadlineExceeded},
|
||||
{"invalid per-call preparation", ErrApprovedCallPreparation, codes.FailedPrecondition},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
req := approvedTestRequest(t, now)
|
||||
|
||||
@@ -9,16 +9,20 @@ import (
|
||||
"git.ipao.vip/rogee/go-sip/internal/agent"
|
||||
)
|
||||
|
||||
var ErrApprovedDialExpired = errors.New("Dispatcher-approved dial deadline expired")
|
||||
var (
|
||||
ErrApprovedDialExpired = errors.New("Dispatcher-approved dial deadline expired")
|
||||
ErrApprovedCallPreparation = errors.New("approved call preparation rejected before dispatch")
|
||||
)
|
||||
|
||||
// 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.
|
||||
// ApprovedCallWorker keeps a call registered until its prepared synchronous
|
||||
// runner has finished physical hangup. Prepare must validate the per-call
|
||||
// snapshot, media and delivery without dialing or persisting business data;
|
||||
// the returned runner must not return before hangup. 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
|
||||
Prepare func(ApprovedExecution) (func(context.Context) error, error)
|
||||
OnFailure func(ApprovedExecution, error) error
|
||||
Now func() time.Time
|
||||
}
|
||||
@@ -28,7 +32,7 @@ type ApprovedCallWorker struct {
|
||||
// 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 {
|
||||
if w == nil || request == nil || w.Lifecycle == nil || w.Calls == nil || w.Prepare == nil || w.OnFailure == nil {
|
||||
return errors.New("approved call lifecycle, runner and failure reporting are required")
|
||||
}
|
||||
if err := request.Err(); err != nil {
|
||||
@@ -47,6 +51,13 @@ func (w *ApprovedCallWorker) Originate(request context.Context, approved Approve
|
||||
if !now().Before(approved.DialBefore) {
|
||||
return ErrApprovedDialExpired
|
||||
}
|
||||
prepared, err := w.Prepare(approved)
|
||||
if err != nil {
|
||||
return errors.Join(ErrApprovedCallPreparation, err)
|
||||
}
|
||||
if prepared == nil {
|
||||
return ErrApprovedCallPreparation
|
||||
}
|
||||
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 {
|
||||
@@ -72,7 +83,7 @@ func (w *ApprovedCallWorker) Originate(request context.Context, approved Approve
|
||||
return
|
||||
}
|
||||
issued <- nil // the Unary response may now acknowledge scheduled work
|
||||
runErr := w.Run(callContext, approved)
|
||||
runErr := prepared(callContext)
|
||||
// Physical execution resources are released independently of any later
|
||||
// failure reporting, OSS work or MQ publisher confirmation.
|
||||
release()
|
||||
|
||||
@@ -17,6 +17,12 @@ func workerTestExecution(now time.Time) ApprovedExecution {
|
||||
DialBefore: now.Add(time.Minute), MaxCallDuration: time.Minute}
|
||||
}
|
||||
|
||||
func prepareWorkerRun(run func(context.Context, ApprovedExecution) error) func(ApprovedExecution) (func(context.Context) error, error) {
|
||||
return func(call ApprovedExecution) (func(context.Context) error, error) {
|
||||
return func(ctx context.Context) error { return run(ctx, call) }, nil
|
||||
}
|
||||
}
|
||||
|
||||
func TestApprovedCallWorkerOutlivesInitiatingUnaryAndWaitsForPhysicalEnd(t *testing.T) {
|
||||
now := time.Date(2026, 9, 20, 10, 0, 0, 0, time.UTC)
|
||||
process, stopProcess := context.WithCancel(context.Background())
|
||||
@@ -35,12 +41,12 @@ func TestApprovedCallWorkerOutlivesInitiatingUnaryAndWaitsForPhysicalEnd(t *test
|
||||
}()
|
||||
failures := make(chan error, 1)
|
||||
worker := ApprovedCallWorker{Lifecycle: process, Calls: calls, Now: func() time.Time { return now },
|
||||
Run: func(ctx context.Context, _ ApprovedExecution) error {
|
||||
Prepare: prepareWorkerRun(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)
|
||||
@@ -93,7 +99,7 @@ func TestApprovedCallWorkerRejectsExpiredOrClosedTaskBeforeIssuing(t *testing.T)
|
||||
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 },
|
||||
Prepare: prepareWorkerRun(func(context.Context, ApprovedExecution) error { starts.Add(1); return nil }),
|
||||
OnFailure: func(ApprovedExecution, error) error { return nil },
|
||||
}
|
||||
approved := workerTestExecution(now)
|
||||
@@ -127,12 +133,12 @@ func TestApprovedCallWorkerProcessStopCancelsAndReportsFailureOnce(t *testing.T)
|
||||
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 {
|
||||
Prepare: prepareWorkerRun(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)
|
||||
@@ -175,7 +181,7 @@ func TestApprovedCallWorkerRefusesExpiredAuthorizationBeforeDispatchAck(t *testi
|
||||
}
|
||||
return now.Add(2 * time.Minute) // expired when the scheduled worker begins
|
||||
},
|
||||
Run: func(context.Context, ApprovedExecution) error { starts.Add(1); return nil },
|
||||
Prepare: prepareWorkerRun(func(context.Context, ApprovedExecution) error { starts.Add(1); return nil }),
|
||||
OnFailure: func(ApprovedExecution, error) error { reports.Add(1); return nil },
|
||||
}
|
||||
approved := workerTestExecution(now)
|
||||
@@ -190,3 +196,25 @@ func TestApprovedCallWorkerRefusesExpiredAuthorizationBeforeDispatchAck(t *testi
|
||||
t.Fatalf("refused call remained registered: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestApprovedCallWorkerPreparationRejectsBeforeDispatchAck(t *testing.T) {
|
||||
now := time.Date(2026, 9, 20, 10, 0, 0, 0, time.UTC)
|
||||
calls := &agent.TaskCalls{}
|
||||
cause := errors.New("synthetic Mock script does not match approved AI")
|
||||
var reports atomic.Int32
|
||||
worker := ApprovedCallWorker{Lifecycle: context.Background(), Calls: calls, Now: func() time.Time { return now },
|
||||
Prepare: func(ApprovedExecution) (func(context.Context) error, error) { return nil, cause },
|
||||
OnFailure: func(ApprovedExecution, error) error { reports.Add(1); return nil },
|
||||
}
|
||||
approved := workerTestExecution(now)
|
||||
if err := worker.Originate(context.Background(), approved); !errors.Is(err, ErrApprovedCallPreparation) || !errors.Is(err, cause) {
|
||||
t.Fatalf("invalid per-call Mock parameters were acknowledged: %v", err)
|
||||
}
|
||||
if reports.Load() != 0 {
|
||||
t.Fatal("no-dial preparation error fabricated a completed call")
|
||||
}
|
||||
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 preparation registered a call: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user