Assemble isolated Agent session and recording worker
This commit is contained in:
@@ -0,0 +1,123 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"errors"
|
||||
"log"
|
||||
"maps"
|
||||
"net"
|
||||
"net/http"
|
||||
"os"
|
||||
"reflect"
|
||||
"slices"
|
||||
"strings"
|
||||
"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/config"
|
||||
"git.ipao.vip/rogee/go-sip/internal/rpc"
|
||||
"github.com/google/uuid"
|
||||
)
|
||||
|
||||
// newCurrentAgentServer binds the authenticated Agent session, task controls
|
||||
// and per-call Mock recording delivery. Serving the returned server requires
|
||||
// a separately verified mutual-TLS listener and a pinned local D connection.
|
||||
func newCurrentAgentServer(ctx context.Context, settings config.AgentEnvironment, scenario approvedMockScenario, appliedSIP map[string]int64, dispatcher agentpb.AgentControlServiceClient) (*rpc.Server, error) {
|
||||
if ctx == nil || ctx.Err() != nil || settings.AgentID == "" || settings.CellID == "" || settings.SessionPath == "" ||
|
||||
settings.RecoveryRoot == "" || len(settings.PeerFingerprints) == 0 || len(appliedSIP) == 0 ||
|
||||
scenario.MaxWAVBytes <= 44 || len(scenario.Script.Turns) == 0 || strings.TrimSpace(scenario.ReasonMessage) == "" {
|
||||
return nil, errors.New("current Agent requires an active process, explicit Mock media and deployment identity")
|
||||
}
|
||||
if dispatcher == nil {
|
||||
return nil, errors.New("current Agent requires a pinned Dispatcher transport")
|
||||
}
|
||||
value := reflect.ValueOf(dispatcher)
|
||||
switch value.Kind() {
|
||||
case reflect.Chan, reflect.Func, reflect.Interface, reflect.Map, reflect.Pointer, reflect.Slice:
|
||||
if value.IsNil() {
|
||||
return nil, errors.New("current Agent requires a pinned Dispatcher transport")
|
||||
}
|
||||
}
|
||||
root, err := os.Stat(settings.RecoveryRoot)
|
||||
if err != nil || !root.IsDir() || root.Mode().Perm() != 0700 {
|
||||
return nil, errors.New("current Agent requires an existing private 0700 recovery directory")
|
||||
}
|
||||
for trunk, revision := range appliedSIP {
|
||||
if strings.TrimSpace(trunk) == "" || revision <= 0 {
|
||||
return nil, errors.New("current Agent requires explicit applied Mock SIP revisions")
|
||||
}
|
||||
}
|
||||
loaded := maps.Clone(appliedSIP)
|
||||
pins := maps.Clone(settings.PeerFingerprints)
|
||||
scenario.InboundPCM16 = bytes.Clone(scenario.InboundPCM16)
|
||||
scenario.Script.OpeningPCM16 = bytes.Clone(scenario.Script.OpeningPCM16)
|
||||
scenario.Script.Turns = slices.Clone(scenario.Script.Turns)
|
||||
for index := range scenario.Script.Turns {
|
||||
scenario.Script.Turns[index].ReplyPCM16 = bytes.Clone(scenario.Script.Turns[index].ReplyPCM16)
|
||||
}
|
||||
uploadHTTP := localMockHTTPClient()
|
||||
var handler *rpc.Server
|
||||
worker := &rpc.ApprovedCallWorker{
|
||||
Lifecycle: ctx,
|
||||
Calls: &agent.TaskCalls{},
|
||||
Prepare: func(execution rpc.ApprovedExecution) (func(context.Context) error, error) {
|
||||
if handler == nil {
|
||||
return nil, errors.New("Agent session is unavailable")
|
||||
}
|
||||
delivery := &agent.RecordingDelivery{
|
||||
Call: agent.RecordingClient{
|
||||
Client: dispatcher, DispatcherID: execution.DispatcherID, TenantID: execution.TenantID,
|
||||
SourceEventID: execution.SourceEventID, Session: func(context.Context) (*agentpb.RequestMeta, error) { return handler.ActiveSessionMeta() },
|
||||
},
|
||||
Recovery: &agent.RecordingRecovery{
|
||||
Root: settings.RecoveryRoot,
|
||||
Upload: agent.UploadClient{HTTPClient: uploadHTTP, AllowInsecureHTTP: true},
|
||||
},
|
||||
}
|
||||
mock := &rpc.ApprovedRecordedMockCall{
|
||||
InboundPCM16: scenario.InboundPCM16, Script: scenario.Script,
|
||||
MaxWAVBytes: scenario.MaxWAVBytes, ExpectedRecording: scenario.ExpectedRecording,
|
||||
Outcome: scenario.Outcome, ReasonMessage: scenario.ReasonMessage,
|
||||
ReportTimeout: 15 * time.Minute, Delivery: delivery,
|
||||
}
|
||||
return mock.Prepare(execution)
|
||||
},
|
||||
OnFailure: func(execution rpc.ApprovedExecution, cause error) error {
|
||||
// The runner already tried to report termination and persisted any
|
||||
// failed upload. Never invent a second result or retry an unknown PUT.
|
||||
log.Printf("Agent Mock call requires inspection: event_id=%q task_id=%q cause_type=%T", execution.SourceEventID, execution.TaskID, cause)
|
||||
return nil
|
||||
},
|
||||
}
|
||||
handler, err = rpc.NewApprovedAgentServer(rpc.ServerOptions{
|
||||
Mode: "mock", StatePath: settings.SessionPath,
|
||||
Status: &agentpb.AgentStatus{AgentId: settings.AgentID, CellId: settings.CellID, BootId: uuid.NewString(), ProtocolVersion: "agent.v1"},
|
||||
LoadedSIP: func(context.Context) (map[string]int64, error) { return maps.Clone(loaded), nil },
|
||||
PeerCertificateFingerprints: pins, RequirePeerCertificate: true,
|
||||
}, worker)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return handler, nil
|
||||
}
|
||||
|
||||
// The isolated Mock uploader may reach localhost only, even if a grant or
|
||||
// redirect unexpectedly names a real OSS endpoint. It never logs signed URLs.
|
||||
func localMockHTTPClient() *http.Client {
|
||||
transport := http.DefaultTransport.(*http.Transport).Clone()
|
||||
transport.Proxy = nil
|
||||
transport.DialContext = func(ctx context.Context, network, address string) (net.Conn, error) {
|
||||
host, _, err := net.SplitHostPort(address)
|
||||
if err != nil {
|
||||
return nil, errors.New("Mock upload target is not local")
|
||||
}
|
||||
ip := net.ParseIP(host)
|
||||
if !strings.EqualFold(host, "localhost") && (ip == nil || !ip.IsLoopback()) {
|
||||
return nil, errors.New("Mock upload target is not local")
|
||||
}
|
||||
return (&net.Dialer{}).DialContext(ctx, network, address)
|
||||
}
|
||||
return &http.Client{Transport: transport, CheckRedirect: func(*http.Request, []*http.Request) error { return http.ErrUseLastResponse }}
|
||||
}
|
||||
@@ -0,0 +1,121 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
agentpb "git.ipao.vip/rogee/go-sip/gen/agent"
|
||||
"git.ipao.vip/rogee/go-sip/internal/ai"
|
||||
"git.ipao.vip/rogee/go-sip/internal/config"
|
||||
)
|
||||
|
||||
type isolatedAgentRecordingClient struct {
|
||||
agentpb.AgentControlServiceClient
|
||||
}
|
||||
|
||||
func currentAgentSetupFixture(t *testing.T) (config.AgentEnvironment, approvedMockScenario) {
|
||||
t.Helper()
|
||||
root := t.TempDir()
|
||||
if err := os.Chmod(root, 0700); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
settings := config.AgentEnvironment{
|
||||
AgentID: "agent-mock", CellID: "cell-mock", Listen: "127.0.0.1:0",
|
||||
SessionPath: filepath.Join(root, "session.json"), RecoveryRoot: root,
|
||||
DispatcherEndpoint: "127.0.0.1:39443", DispatcherServerName: "dispatcher.local",
|
||||
CAFile: filepath.Join(root, "ca.pem"), CertFile: filepath.Join(root, "agent.pem"),
|
||||
KeyFile: filepath.Join(root, "agent.key"), PeerFingerprints: map[string]struct{}{strings.Repeat("a", 64): {}},
|
||||
}
|
||||
scenario := approvedMockScenario{
|
||||
InboundPCM16: bytes.Repeat([]byte{1, 0}, 1600),
|
||||
Script: ai.ApprovedMockScript{Turns: []ai.ApprovedMockTurn{{Transcript: "synthetic ASR fixture"}}},
|
||||
MaxWAVBytes: 4096, ExpectedRecording: true, Outcome: "answered", ReasonMessage: "isolated Mock answered",
|
||||
}
|
||||
return settings, scenario
|
||||
}
|
||||
|
||||
func TestNewCurrentAgentServerBindsMockCallsToOneSessionAndRecoveryRoot(t *testing.T) {
|
||||
settings, scenario := currentAgentSetupFixture(t)
|
||||
server, err := newCurrentAgentServer(context.Background(), settings, scenario, map[string]int64{"trunk-mock": 8}, &isolatedAgentRecordingClient{})
|
||||
if err != nil || server == nil {
|
||||
t.Fatalf("valid isolated Agent could not be assembled: %v", err)
|
||||
}
|
||||
if _, err := os.Stat(settings.SessionPath); !os.IsNotExist(err) {
|
||||
t.Fatalf("server assembly opened durable session before activation: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestNewCurrentAgentServerRefusesUnsafeAdaptersBeforeResources(t *testing.T) {
|
||||
settings, scenario := currentAgentSetupFixture(t)
|
||||
for _, tc := range []struct {
|
||||
name string
|
||||
change func(*config.AgentEnvironment, *approvedMockScenario, *map[string]int64, *agentpb.AgentControlServiceClient)
|
||||
}{
|
||||
{"missing Agent identity", func(c *config.AgentEnvironment, _ *approvedMockScenario, _ *map[string]int64, _ *agentpb.AgentControlServiceClient) {
|
||||
c.AgentID = ""
|
||||
}},
|
||||
{"missing SIP evidence", func(_ *config.AgentEnvironment, _ *approvedMockScenario, sip *map[string]int64, _ *agentpb.AgentControlServiceClient) {
|
||||
*sip = nil
|
||||
}},
|
||||
{"unknown SIP revision", func(_ *config.AgentEnvironment, _ *approvedMockScenario, sip *map[string]int64, _ *agentpb.AgentControlServiceClient) {
|
||||
(*sip)["trunk-mock"] = 0
|
||||
}},
|
||||
{"missing recovery", func(c *config.AgentEnvironment, _ *approvedMockScenario, _ *map[string]int64, _ *agentpb.AgentControlServiceClient) {
|
||||
c.RecoveryRoot = filepath.Join(t.TempDir(), "missing")
|
||||
}},
|
||||
{"missing scenario", func(_ *config.AgentEnvironment, s *approvedMockScenario, _ *map[string]int64, _ *agentpb.AgentControlServiceClient) {
|
||||
s.MaxWAVBytes = 0
|
||||
}},
|
||||
{"missing Dispatcher transport", func(_ *config.AgentEnvironment, _ *approvedMockScenario, _ *map[string]int64, client *agentpb.AgentControlServiceClient) {
|
||||
*client = nil
|
||||
}},
|
||||
{"typed-nil Dispatcher transport", func(_ *config.AgentEnvironment, _ *approvedMockScenario, _ *map[string]int64, client *agentpb.AgentControlServiceClient) {
|
||||
var typed *isolatedAgentRecordingClient
|
||||
*client = typed
|
||||
}},
|
||||
} {
|
||||
t.Run(tc.name, func(t *testing.T) {
|
||||
cfg, media := settings, scenario
|
||||
loaded := map[string]int64{"trunk-mock": 8}
|
||||
var client agentpb.AgentControlServiceClient = &isolatedAgentRecordingClient{}
|
||||
tc.change(&cfg, &media, &loaded, &client)
|
||||
if server, err := newCurrentAgentServer(context.Background(), cfg, media, loaded, client); err == nil || server != nil {
|
||||
t.Fatalf("unsafe current Agent assembly was accepted: %v", err)
|
||||
}
|
||||
if _, err := os.Stat(settings.SessionPath); !os.IsNotExist(err) {
|
||||
t.Fatalf("rejected Agent assembly wrote durable session: %v", err)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestLocalMockHTTPClientRefusesExternalAndRedirectedTargets(t *testing.T) {
|
||||
client := localMockHTTPClient()
|
||||
if response, err := client.Get("https://oss.example.invalid/approved"); err == nil || response != nil || !strings.Contains(err.Error(), "not local") {
|
||||
t.Fatalf("Mock uploader contacted an external target: %v", err)
|
||||
}
|
||||
local := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
if r.URL.Path == "/redirect" {
|
||||
http.Redirect(w, r, "https://oss.example.invalid/approved", http.StatusFound)
|
||||
return
|
||||
}
|
||||
w.WriteHeader(http.StatusOK)
|
||||
}))
|
||||
defer local.Close()
|
||||
response, err := client.Get(local.URL + "/ok")
|
||||
if err != nil || response.StatusCode != http.StatusOK {
|
||||
t.Fatalf("local Mock endpoint was refused: response=%v err=%v", response, err)
|
||||
}
|
||||
_ = response.Body.Close()
|
||||
response, err = client.Get(local.URL + "/redirect")
|
||||
if err != nil || response.StatusCode != http.StatusFound {
|
||||
t.Fatalf("Mock uploader followed an external redirect: response=%v err=%v", response, err)
|
||||
}
|
||||
_ = response.Body.Close()
|
||||
}
|
||||
@@ -63,6 +63,7 @@
|
||||
- 新增内部任务级 `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 进程期限;每通话先在无拨号/无业务写入前准备并核验签发 AI、合成媒体和交付目标,准备失败即明确拒绝且不发完成事实;签发拨号时限在异步工人发出接受回执前重验,超时或任务关闭同样在无拨号时明确拒绝;启动请求结束不取消在途通话,任务 hangup 等待 runner 明确结束后才完成控制。隔离测试覆盖先执行回执后合成挂断、进程关闭、失败观察、重复身份防重及缺失/竞争适配器拒绝;未知执行不自动重拨。这里只证明调度与取消顺序,**未在主 CLI 连接真实媒体、录音、失败报告或出站 MQ**。
|
||||
- CLI 增加隔离 Mock 场景读取器:只接受显式给出的有界 JSON 文件,合成 PCM16、脚本、录音预期、结果事实和 WAV 上限均须逐项提供;缺失/未知字段、额外 JSON、非法 base64 和超大文件明确拒绝,错误不回显脚本内容。此读取器尚未连接 `agent` 主命令,不能把场景配置或 Mock 结果称为 SaaS 智能体值或真实音频验收。
|
||||
- Agent 组装辅助 `newCurrentAgentServer` 已把显式 Mock 部署身份、已加载 SIP 隔离事实、每通话合成脚本、同一个控制/执行栅栏以及 Agent→D 的会话元数据/录音交付绑定;拒绝缺失或 typed-nil 的 D 客户端、非法 SIP revision、未私有化的恢复目录和竞争配置,构造时不打开会话状态文件。共用的 Mock 上传 HTTP 客户端只连接本机且不跟随跳往外网的重定向,脚本、已加载事实与指纹均复制为启动快照。相关隔离/反例与 race 测试通过;**还未由实际 `agent` 命令调用,也未在此入口建立 mTLS 监听与已鉴权的 D 连接**,不算主运行链路通过。
|
||||
- 验证:`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 结果代签。
|
||||
|
||||
|
||||
Reference in New Issue
Block a user