Switch Dispatcher root to isolated Mock runtime
This commit is contained in:
@@ -0,0 +1,209 @@
|
||||
//go:build integration
|
||||
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/x509"
|
||||
"encoding/pem"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
agentpb "git.ipao.vip/rogee/go-sip/gen/agent"
|
||||
"git.ipao.vip/rogee/go-sip/internal/mq"
|
||||
"git.ipao.vip/rogee/go-sip/internal/rpc"
|
||||
"git.ipao.vip/rogee/go-sip/internal/tenant"
|
||||
amqp "github.com/rabbitmq/amqp091-go"
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/credentials"
|
||||
)
|
||||
|
||||
func TestCurrentDispatcherCommandStartsWithIsolatedMQHTTPAndAgent(t *testing.T) {
|
||||
brokerURL, adminURL := os.Getenv("RABBITMQ_URL"), os.Getenv("RABBITMQ_PROVISIONER_URL")
|
||||
if brokerURL == "" || adminURL == "" {
|
||||
t.Skip("requires the isolated RabbitMQ mock provisioner")
|
||||
}
|
||||
const dispatcherID = "c046b893-8628-4589-ae50-619d049248a6"
|
||||
adminConnection, err := amqp.Dial(adminURL)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer adminConnection.Close()
|
||||
admin, err := adminConnection.Channel()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer admin.Close()
|
||||
for _, exchange := range []string{mq.CommandsExchangeCurrent, mq.ResultsExchangeCurrent, mq.DeadLetterExchangeCurrent} {
|
||||
if err := admin.ExchangeDeclare(exchange, "topic", true, false, false, false, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
control, err := tenant.CurrentControlRoute(dispatcherID)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
result, err := tenant.CurrentResultRoute(dispatcherID)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := admin.QueueDeclare(control.Queue, true, false, false, false, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer func() { _, _ = admin.QueueDelete(control.Queue, false, false, false) }()
|
||||
if err := admin.QueueBind(control.Queue, control.BindingKey, control.Exchange, false, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := admin.QueueDeclare(result.Queue, true, false, false, false, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer func() { _, _ = admin.QueueDelete(result.Queue, false, false, false) }()
|
||||
if err := admin.QueueBind(result.Queue, result.BindingKey, result.Exchange, false, nil); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer func() { _ = admin.QueueUnbind(result.Queue, result.BindingKey, result.Exchange, nil) }()
|
||||
for _, queue := range []string{control.Queue, result.Queue} {
|
||||
if _, err := admin.QueuePurge(queue, false); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
|
||||
root := t.TempDir()
|
||||
ca, agentCert, agentKey, dispatcherCert, dispatcherKey, dispatcherLeaf := localCommandCertificates(t)
|
||||
agentBlock, _ := pem.Decode(agentCert)
|
||||
if agentBlock == nil {
|
||||
t.Fatal("isolated Agent certificate is invalid")
|
||||
}
|
||||
agentLeaf, err := x509.ParseCertificate(agentBlock.Bytes)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
agentTLS, err := rpc.NewServerTLSConfig(ca, agentCert, agentKey)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
agentServer := rpc.NewServer(rpc.ServerOptions{
|
||||
Mode: "mock", StatePath: filepath.Join(root, "agent-session.json"), ApprovedDispatcherID: dispatcherID,
|
||||
RequirePeerCertificate: true, PeerCertificateFingerprints: map[string]struct{}{rpc.CertificateFingerprint(dispatcherLeaf): {}},
|
||||
Status: &agentpb.AgentStatus{AgentId: "agent-mock", CellId: "cell-mock", BootId: "boot-mock", ProtocolVersion: "agent.v1"},
|
||||
LoadedSIP: func(context.Context) (map[string]int64, error) { return map[string]int64{"trunk-mock": 8}, nil },
|
||||
})
|
||||
agentListener, err := net.Listen("tcp", "127.0.0.1:0")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
agentGRPC := grpc.NewServer(grpc.Creds(credentials.NewTLS(agentTLS)))
|
||||
agentpb.RegisterAgentControlServiceServer(agentGRPC, agentServer)
|
||||
go func() { _ = agentGRPC.Serve(agentListener) }()
|
||||
defer agentGRPC.Stop()
|
||||
|
||||
sipBody, err := os.ReadFile(filepath.Join("..", "..", "contracts", "local", "examples", "config-read-sip.json"))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
taskBody, err := os.ReadFile(filepath.Join("..", "..", "contracts", "local", "examples", "task-discovery-end.json"))
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var sipReads, taskReads atomic.Int64
|
||||
saas := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
if r.Header.Get("X-DISPATCHER-id") != dispatcherID || r.Header.Get("X-DISPATCHER-SECRET-KEY") != "isolated-secret" {
|
||||
http.Error(w, "unapproved Dispatcher", http.StatusForbidden)
|
||||
return
|
||||
}
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
switch r.URL.Path {
|
||||
case "/internal/v1/dispatcher/sip":
|
||||
sipReads.Add(1)
|
||||
_, _ = w.Write(sipBody)
|
||||
case "/internal/v1/dispatcher/tasks":
|
||||
taskReads.Add(1)
|
||||
_, _ = w.Write(taskBody)
|
||||
default:
|
||||
http.Error(w, "unexpected Mock read", http.StatusNotFound)
|
||||
}
|
||||
}))
|
||||
defer saas.Close()
|
||||
|
||||
agentInventory := filepath.Join(root, "agent-endpoints.json")
|
||||
ossFile := filepath.Join(root, "oss.json")
|
||||
for filename, body := range map[string][]byte{
|
||||
agentInventory: []byte(fmt.Sprintf(`[{"agent_id":"agent-mock","cell_id":"cell-mock","address":%q,"server_name":"agent.local"}]`, agentListener.Addr().String())),
|
||||
ossFile: []byte(fmt.Sprintf(`{"dispatcher_id":%q,"oss":{"endpoint":"https://127.0.0.1:19445","region":"cn-mock","bucket":"mock-bucket","object_prefix":"approved","access_key_id_env":"MOCK_OSS_KEY_ID","access_key_secret_env":"MOCK_OSS_KEY_SECRET","max_asset_bytes":1048576}}`, dispatcherID)),
|
||||
filepath.Join(root, "ca.pem"): ca,
|
||||
filepath.Join(root, "dispatcher.pem"): dispatcherCert,
|
||||
filepath.Join(root, "dispatcher.key"): dispatcherKey,
|
||||
} {
|
||||
if err := os.WriteFile(filename, body, 0600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
reserved, err := net.Listen("tcp", "127.0.0.1:0")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
listen := reserved.Addr().String()
|
||||
_ = reserved.Close()
|
||||
database := filepath.Join(root, "dispatcher.sqlite")
|
||||
for name, value := range map[string]string{
|
||||
"DISPATCHER_ID": dispatcherID, "DISPATCHER_SECRET_KEY": "isolated-secret",
|
||||
"SAAS_BASE_URL": saas.URL, "RABBITMQ_URL": brokerURL,
|
||||
"DISPATCHER_SQLITE_PATH": database, "DISPATCHER_GRPC_LISTEN": listen,
|
||||
"DISPATCHER_AGENT_ENDPOINTS_FILE": agentInventory, "DISPATCHER_OSS_CONFIG_FILE": ossFile,
|
||||
"MTLS_CA_FILE": filepath.Join(root, "ca.pem"),
|
||||
"MTLS_CERT_FILE": filepath.Join(root, "dispatcher.pem"),
|
||||
"MTLS_KEY_FILE": filepath.Join(root, "dispatcher.key"),
|
||||
"MTLS_PEER_CERT_FINGERPRINTS": rpc.CertificateFingerprint(agentLeaf),
|
||||
"MOCK_OSS_KEY_ID": "isolated-id", "MOCK_OSS_KEY_SECRET": "isolated-secret",
|
||||
} {
|
||||
t.Setenv(name, value)
|
||||
}
|
||||
process, cancel := context.WithCancel(context.Background())
|
||||
command := newRootCommand()
|
||||
command.SetContext(process)
|
||||
command.SetOut(io.Discard)
|
||||
command.SetErr(io.Discard)
|
||||
command.SetArgs([]string{"dispatcher", "--mode", "mock"})
|
||||
finished := make(chan error, 1)
|
||||
go func() { finished <- command.Execute() }()
|
||||
deadline := time.After(8 * time.Second)
|
||||
for sipReads.Load() == 0 || taskReads.Load() == 0 {
|
||||
select {
|
||||
case err := <-finished:
|
||||
cancel()
|
||||
t.Fatalf("isolated Dispatcher exited before HTTP bootstrap: %v", err)
|
||||
case <-deadline:
|
||||
cancel()
|
||||
t.Fatal("isolated Dispatcher never read SIP and the empty task snapshot")
|
||||
case <-time.After(20 * time.Millisecond):
|
||||
}
|
||||
}
|
||||
if _, err := os.Stat(database); err != nil {
|
||||
cancel()
|
||||
t.Fatalf("successful isolated bootstrap did not retain SQLite: %v", err)
|
||||
}
|
||||
select {
|
||||
case err := <-finished:
|
||||
cancel()
|
||||
t.Fatalf("isolated Dispatcher stopped despite a live process: %v", err)
|
||||
default:
|
||||
}
|
||||
cancel()
|
||||
select {
|
||||
case err := <-finished:
|
||||
if err != nil && !errors.Is(err, context.Canceled) {
|
||||
t.Fatalf("isolated Dispatcher shutdown failed: %v", err)
|
||||
}
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("isolated Dispatcher did not stop after cancellation")
|
||||
}
|
||||
}
|
||||
@@ -12,23 +12,25 @@ import (
|
||||
)
|
||||
|
||||
func TestDispatcherHasNoBusinessHTTPFlags(t *testing.T) {
|
||||
command := newDispatcherCommand()
|
||||
for _, name := range []string{"control-listen", "control-token"} {
|
||||
command := newCurrentDispatcherCommand()
|
||||
for _, name := range []string{"control-listen", "control-token", "config", "db"} {
|
||||
if command.Flags().Lookup(name) != nil {
|
||||
t.Fatalf("obsolete HTTP flag remains: --%s", name)
|
||||
t.Fatalf("obsolete Dispatcher flag remains: --%s", name)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestDispatcherRequiresFileBeforeOpeningDatabase(t *testing.T) {
|
||||
func TestDispatcherRequiresDeploymentIdentityBeforeOpeningDatabase(t *testing.T) {
|
||||
path := filepath.Join(t.TempDir(), "untouched.db")
|
||||
t.Setenv("DISPATCHER_ID", "")
|
||||
t.Setenv("DISPATCHER_SQLITE_PATH", path)
|
||||
root := newRootCommand()
|
||||
root.SetArgs([]string{"dispatcher", "--mode", "mock", "--db", path})
|
||||
if err := root.Execute(); err == nil || !strings.Contains(err.Error(), "--config") {
|
||||
t.Fatalf("missing file did not fail early: %v", err)
|
||||
root.SetArgs([]string{"dispatcher", "--mode", "mock"})
|
||||
if err := root.Execute(); err == nil || !strings.Contains(err.Error(), "DISPATCHER_ID") {
|
||||
t.Fatalf("missing Dispatcher identity did not fail early: %v", err)
|
||||
}
|
||||
if _, err := os.Stat(path); !errors.Is(err, os.ErrNotExist) {
|
||||
t.Fatalf("database was touched before configuration validation: %v", err)
|
||||
t.Fatalf("database was touched before deployment validation: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -52,7 +52,7 @@ func newRootCommand() *cobra.Command {
|
||||
SilenceUsage: true,
|
||||
SilenceErrors: true,
|
||||
}
|
||||
root.AddCommand(newCurrentAgentCommand(), newDispatcherCommand())
|
||||
root.AddCommand(newCurrentAgentCommand(), newCurrentDispatcherCommand())
|
||||
return root
|
||||
}
|
||||
|
||||
|
||||
@@ -27,6 +27,14 @@ func TestRootAgentRejectsOldOneShotFlag(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestRootDispatcherRejectsOldConfigFlag(t *testing.T) {
|
||||
root := newRootCommand()
|
||||
root.SetArgs([]string{"dispatcher", "--config", "/nonexistent"})
|
||||
if err := root.Execute(); err == nil || !strings.Contains(err.Error(), "unknown flag") {
|
||||
t.Fatalf("obsolete Dispatcher configuration entry was admitted: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestWriteResultIsJSON(t *testing.T) {
|
||||
var b bytes.Buffer
|
||||
data, err := marshalResultForTest(map[string]string{"status": "ok"})
|
||||
|
||||
Vendored
+17
-17
@@ -1,21 +1,21 @@
|
||||
# Physical-host example; real deployment requires separate authorization.
|
||||
# Install dispatcher.json from deploys/config/dispatcher.json.example and replace
|
||||
# its example dispatcher_id with this installation's unique stable UUID v4.
|
||||
SIP_GO_AGENT_MODE=real
|
||||
DISPATCHER_DB=/var/lib/sip-go-agent/dispatcher/dispatcher.db
|
||||
# Required by the current V3 runtime; inject approved values out of band.
|
||||
# DISPATCHER_CONFIG_READ_BASE_URL=<injected-config-read-base-url>
|
||||
# DISPATCHER_SECRET_KEY=<injected-secret>
|
||||
# RABBITMQ_URL=<injected-broker-url>
|
||||
# Isolated local Mock only; not a production or real-call configuration.
|
||||
# Native Asterisk and capture-first diagnostics are mandatory for an actual
|
||||
# non-production validation host: deploys/test/nonprod-call-evidence.sh.
|
||||
# Start explicitly: sip-go-agent dispatcher --mode mock
|
||||
# Inject a globally unique, preapproved canonical UUID v4 for this Dispatcher.
|
||||
# DISPATCHER_ID=<injected-uuid-v4>
|
||||
# DISPATCHER_SECRET_KEY=<injected-from-controlled-environment>
|
||||
SAAS_BASE_URL=http://127.0.0.1:18080
|
||||
RABBITMQ_URL=amqp://127.0.0.1:5672
|
||||
DISPATCHER_SQLITE_PATH=/var/lib/sip-go-agent/dispatcher/current.sqlite
|
||||
DISPATCHER_GRPC_LISTEN=127.0.0.1:19443
|
||||
# The inventory must name exactly one local Agent with its certificate SAN.
|
||||
DISPATCHER_AGENT_ENDPOINTS_FILE=/etc/sip-go-agent/agent-endpoints.json
|
||||
# This file must bind the same Dispatcher ID and reference injected OSS keys;
|
||||
# credentials and temporary upload URLs do not belong in this example.
|
||||
DISPATCHER_OSS_CONFIG_FILE=/etc/sip-go-agent/oss.json
|
||||
MTLS_CA_FILE=/etc/sip-go-agent/pki/ca.pem
|
||||
MTLS_CERT_FILE=/etc/sip-go-agent/pki/dispatcher.pem
|
||||
MTLS_KEY_FILE=/etc/sip-go-agent/pki/dispatcher.key
|
||||
MTLS_SERVER_NAME=dispatcher.internal
|
||||
DISPATCHER_GRPC_LISTEN=127.0.0.1:19443
|
||||
DISPATCHER_ALLOWED_AGENT_IDS=agent-cell-a
|
||||
# MTLS_PEER_CERT_FINGERPRINTS=<approved-agent-certificate-fingerprint>
|
||||
# Inject only the credential variables explicitly referenced by dispatcher.json.
|
||||
# Never put actual values in this example or commit them:
|
||||
# GO_SIP_OSS_ACCESS_KEY_ID=<injected-at-runtime>
|
||||
# GO_SIP_OSS_ACCESS_KEY_SECRET=<injected-at-runtime>
|
||||
# Required SHA-256 fingerprint of the assigned Agent's certificate.
|
||||
# MTLS_PEER_CERT_FINGERPRINTS=<injected-64-character-hex-fingerprint>
|
||||
|
||||
@@ -32,12 +32,12 @@
|
||||
- 原 Buf STANDARD 的 `PACKAGE_VERSION_SUFFIX` 与已批准的无代次内部包名冲突;`buf.yaml` 仅对这一条规则作例外,其余 STANDARD 规则保持。隔离安装 Buf v1.50.0 / protoc-gen-go v1.36.12 / protoc-gen-go-grpc v1.5.1 于 `/tmp/sip-go-agent-tools/bin`,未改应用依赖。
|
||||
- `PATH=/tmp/sip-go-agent-tools/bin:$PATH sh scripts/check-proto.sh` 通过(lint/build/generate/新包测试/7 文件清单 hash);`git diff --check`、`go test ./internal/rpc ./internal/agent ./internal/dispatcher ./cmd/sip-go-agent -count=1` 通过。其他旧实现/契约入口、文件名碰撞和迁移数据安全仍待 P02/P07,**不能据此称全仓已无代次或整体完成**。
|
||||
|
||||
- 新增无实现代次字段的 Dispatcher 环境预检:必须显式给出规范 UUID v4 身份、只读 HTTP 地址及密钥、RabbitMQ 地址和 SQLite 路径;当前 Mock 仅接受本机 HTTP/MQ 目标,mixed/real 在任何资源操作前拒绝。预检不打开数据库或网络,缺失值不继承旧默认配置,错误不打印凭据。`go test ./internal/config -run '^TestLoadDispatcherEnvironment' -count=1` 通过;此预检尚未接入主 CLI,不能当作 P02/P03 完成。
|
||||
- Dispatcher 另有纯运行环境预检:唯一 Agent 端点清单和 OSS 配置文件须提供路径,本机 gRPC 监听及双向 TLS 文件/Agent 证书指纹须显式给出;缺失、非法指纹或非本机监听明确拒绝且不回显值。该检查不读取配置文件、不连接 Agent、不打开 SQLite/MQ;`go test ./internal/config -run '^TestLoadDispatcherRuntimeEnvironment' -count=1` 通过。当前仍未由 `dispatcher` 主命令调用,不能当作 Agent 加载或 OSS 签发已验收。
|
||||
- D 的部署文件读取已独立于旧代次 Schema:Agent 清单仅接受**一个本机端点**;OSS 配置必须匹配本 D 身份、指定本机 HTTPS 目标、15 分钟授权与显式资产上限,凭据只经受控环境变量引用。超大文件、重复/未知字段、多个 Agent、非本机目标、缺失凭据和错误 D 归属均拒绝;JSON 唯一键及凭据引用基础函数已移出旧配置文件。`go test ./internal/config -count=1` 通过。此处仅验证部署数据,不代表 Agent 实际加载、官方 SDK 已签发目标,亦未接入主 `dispatcher` 命令。
|
||||
- 独立的隔离 D 命令函数现已组装严格部署预检、当前 HTTP 读取、`store.OpenCurrent`、被动核验 SaaS 预建拓扑的 `mq.OpenCurrent`、D 身份绑定的 Agent 会话、批准执行/控制、实际加载版本检查、`CurrentRuntime` 和录音双向 TLS RPC;启动开库后先持久关闭新准入,不清理旧状态或擅自建队。负例测试确认 mixed/real、旧 `--config` 与非本机 Agent/错误 OSS 归属在打开 SQLite 前拒绝。**根 `dispatcher` 命令尚未切换;该函数尚未用实际本机 Agent、RabbitMQ 和 HTTP 完整启动,不构成主程序通过或外部验收。**
|
||||
- 新增无实现代次字段的 Dispatcher 环境预检:必须显式给出规范 UUID v4 身份、只读 HTTP 地址及密钥、RabbitMQ 地址和 SQLite 路径;当前 Mock 仅接受本机 HTTP/MQ 目标,mixed/real 在任何资源操作前拒绝。预检不打开数据库或网络,缺失值不继承旧默认配置,错误不打印凭据。`go test ./internal/config -run '^TestLoadDispatcherEnvironment' -count=1` 通过;该预检已接入根 `dispatcher` 命令;隔离启动不代表 P02/P03 整体验收。
|
||||
- Dispatcher 另有纯运行环境预检:唯一 Agent 端点清单和 OSS 配置文件须提供路径,本机 gRPC 监听及双向 TLS 文件/Agent 证书指纹须显式给出;缺失、非法指纹或非本机监听明确拒绝且不回显值。该检查不读取配置文件、不连接 Agent、不打开 SQLite/MQ;`go test ./internal/config -run '^TestLoadDispatcherRuntimeEnvironment' -count=1` 通过。该检查现由根 `dispatcher` 命令调用;本机 Mock 不能作为真实 Agent 加载或 OSS 签发的验收。
|
||||
- D 的部署文件读取已独立于旧代次 Schema:Agent 清单仅接受**一个本机端点**;OSS 配置必须匹配本 D 身份、指定本机 HTTPS 目标、15 分钟授权与显式资产上限,凭据只经受控环境变量引用。超大文件、重复/未知字段、多个 Agent、非本机目标、缺失凭据和错误 D 归属均拒绝;JSON 唯一键及凭据引用基础函数已移出旧配置文件。`go test ./internal/config -count=1` 通过。此处仅验证部署数据,不代表 Agent 实际加载、官方 SDK 已签发目标,主 `dispatcher` 命令现已读取这些部署文件;仅隔离 Mock 完成启动,未验证真实目标。
|
||||
- 独立的隔离 D 命令函数现已组装严格部署预检、当前 HTTP 读取、`store.OpenCurrent`、被动核验 SaaS 预建拓扑的 `mq.OpenCurrent`、D 身份绑定的 Agent 会话、批准执行/控制、实际加载版本检查、`CurrentRuntime` 和录音双向 TLS RPC;启动开库后先持久关闭新准入,不清理旧状态或擅自建队。负例测试确认 mixed/real、旧 `--config` 与非本机 Agent/错误 OSS 归属在打开 SQLite 前拒绝。**根 `dispatcher` 命令已切换;隔离脚本用临时本机 RabbitMQ、Mock HTTP 和本机双向 TLS 模拟 Agent 服务验证启动、读取与取消退出,未覆盖有任务的主进程执行、录音、MQ 最终结果或外部验收。**
|
||||
|
||||
- D 的 Agent 会话激活现在必须显式携带规范 UUID v4 的本 D 身份,并写入 Agent 授权绑定;预配置了对应 D 身份的隔离 Agent 会拒绝错误归属,不能再靠 D 自报空身份放行。`go test ./internal/dispatcher -count=1` 通过;主 `dispatcher` 命令尚未使用这条会话链路,不代表本机 D↔A 执行已验收。
|
||||
- D 的 Agent 会话激活现在必须显式携带规范 UUID v4 的本 D 身份,并写入 Agent 授权绑定;预配置了对应 D 身份的隔离 Agent 会拒绝错误归属,不能再靠 D 自报空身份放行。`go test ./internal/dispatcher -count=1` 通过;主 `dispatcher` 命令已通过本机模拟 Agent 的探测与激活;D↔A 执行仍未在主进程联测。
|
||||
- D 的已授权调用元数据现在只从该 Agent **最新且未过期**的会话产生,每次生成不同操作身份;未注册、会话缺失或过期均明确拒绝。`TestAgentCoordinatorProbesBeforeActivation` 覆盖这些边界,不能代替实际 D 主进程的执行与录音联测。
|
||||
- Agent 增加独立的无代次环境预检:Agent/Cell 与预授权 D 的身份、会话与私有恢复路径、隔离 Mock 场景和已加载 SIP 测试事实、D gRPC 目标及 mTLS 文件/指纹均须显式提供;仅允许本机监听和本机 D 端点,mixed/real 在访问文件或网络前拒绝。缺失项不继承旧 `FromEnv` 默认值,凭据/地址不回显;`go test ./internal/config -run '^TestLoadAgentEnvironment' -count=1` 通过。此检查已接入根命令可达的唯一 Agent Mock 入口;本机双向 TLS 监听可由受信 D 证书探测,另一张同 CA 证书被指纹门禁拒绝。本地测试用批准的 D 身份激活会话,错误 D UUID 即使携带受信证书也在写入会话日志前拒绝;尚未执行呼叫、验证 Agent→D 录音或真实 Asterisk 加载。
|
||||
|
||||
@@ -46,9 +46,9 @@
|
||||
## P03:HTTP 读取分批改造(未整体签收)
|
||||
|
||||
- `contract.ValidateCurrent` 与 `configread` 按当前 Schema 读取 SIP、provider、task、quota 和 cursor 任务发现;严格检查数字 tenant_id、本 D 归属及不可变配置。provider 凭据原值只保留在内存快照,不写日志;Agent 参数中的显式 0/false 保真;无旧 Schema/旧配置回退。
|
||||
- CLI 组装辅助 `dispatcherConfigurationClient` 将严格 Mock 环境预检与当前 HTTP Client 绑定;本机 HTTP 隔离测试实际读取 SIP,核对固定 `/internal/v1/dispatcher/sip`、D 身份/密钥 Header、revision 与不打开 SQLite;旧 `DISPATCHER_CONFIG_READ_BASE_URL` 不可充当缺失的当前地址。`go test ./cmd/sip-go-agent -run '^TestDispatcherConfigurationClient' -count=1` 通过;主 `dispatcher` 命令尚未调用该辅助,不能声称启动已切换。
|
||||
- CLI 组装辅助 `dispatcherConfigurationClient` 将严格 Mock 环境预检与当前 HTTP Client 绑定;本机 HTTP 隔离测试实际读取 SIP,核对固定 `/internal/v1/dispatcher/sip`、D 身份/密钥 Header、revision 与不打开 SQLite;旧 `DISPATCHER_CONFIG_READ_BASE_URL` 不可充当缺失的当前地址。`go test ./cmd/sip-go-agent -run '^TestDispatcherConfigurationClient' -count=1` 通过;主 `dispatcher` 命令已通过该辅助读取本机 Mock;真实 SaaS 的授权和连通尚未验证。
|
||||
- `store.OpenCurrent` 新建数字租户 SQLite 状态;旧表、旧版当前布局、残缺布局均在写入前拒绝并保留原记录;不执行旧数据迁移或自动清理。启动时完整发现同一快照一次提交,分页增量逐页持久提交后才推进**内存** cursor;失败关闭准入,重启重新取完整快照。HTTP 的旧 running 不能解除 MQ 暂停/终止,同 revision 异内容及跨任务 SIP/租户额度冲突拒绝。
|
||||
- `CurrentBootstrap` 先关闭准入,核验 SIP 全量与 Agent/Asterisk 已加载 revision、读取任务和 provider/额度,再排空 MQ 控制积压,最后依据已验证 SIP revision 开准入;有更新的持久 SIP 通知时保持关闭但控制与结果处理仍可继续。`CurrentDiscoveryFollower` 逐页绑定任务快照;HTTP 错误、失效或授权不一致只失败,不回退旧读取。**目前只在隔离运行组件中调用,尚未接入 `cmd/sip-go-agent/main.go`;provider 向真实 Agent 交付及真实加载尚待 P05/P07。**
|
||||
- `CurrentBootstrap` 先关闭准入,核验 SIP 全量与 Agent/Asterisk 已加载 revision、读取任务和 provider/额度,再排空 MQ 控制积压,最后依据已验证 SIP revision 开准入;有更新的持久 SIP 通知时保持关闭但控制与结果处理仍可继续。`CurrentDiscoveryFollower` 逐页绑定任务快照;HTTP 错误、失效或授权不一致只失败,不回退旧读取。**根 `dispatcher` 命令已接入该启动链路,空任务隔离快照通过;含任务的主进程 provider 交付与真实 Agent/Asterisk 加载仍未验证。**
|
||||
- TDD 与回归:`go test ./internal/configread ./internal/tenant ./internal/store ./internal/dispatcher -count=1`、`bash scripts/check-current-contracts.sh`、已提交 `a0118e3` 的干净归档测试通过;旧布局行/表原样保留由 `TestCurrentStoreRefusesPreviousCurrentLayoutBeforeModifyingDatabase` 覆盖。
|
||||
|
||||
## P04:隔离 MQ、控制与外呼接纳(仅项目内 Mock)
|
||||
@@ -56,7 +56,7 @@
|
||||
- 固定 `v1` 精确路由与 SaaS 共享结果队列;各 D 控制/任务队列均由 SaaS 预建,D 仅被动检查和消费。隔离 RabbitMQ 实测无 `configure` 权限、D1/D2 不串收、shared queue 实际收讫、断绑后 mandatory 失败不算交付;畸形控制积压拒绝并阻止启动准入。
|
||||
- 执行消息先写持久 inbox 才 ACK;白名单号码格式错误单条拒绝,不暂停其它任务;任务/线路规则不满足时原执行身份留在 SQLite、该任务后续积压留在 SaaS 队列,窗口开放后重新核验才发出 Mock 指令。发指令前两次时窗/SIP/准入检查与持久额度占用;未知 Agent RPC 保持未知占用、不自动重拨。结果 outbox 只在 mandatory/return/confirm 成功后标记已入队,不宣称 SaaS 已处理。
|
||||
- pause/stop/resume 控制先持久挡住该任务新呼叫,Agent 确认收到指令后在同一事务提交应用状态和回执;省略策略默认 hangup,重复控制仍执行、重复回执复用事件身份,stopped 同 ID 不可恢复,未接纳旧外呼静默 ACK。SIP 通知先持久关闭准入,旧/未知通话未确认终结、SaaS 新版尚未分发或 Agent/Asterisk 未加载时不重开;等待期间控制和 outbox 仍可处理。
|
||||
- TDD 与隔离验证:`go test ./internal/store ./internal/dispatcher ./internal/mq -count=1`、`bash scripts/check-current-mq-mock.sh`、`git diff --check` 通过;Mock 覆盖任务积压、恢复、控制、重复投递、发布失败、SIP revision 栅栏与格式错误。当前 Agent side effect 为注入的**假外呼**,主 CLI 仍旧;真实 SaaS/RabbitMQ、供应商、线路、录音与 `call.result` 均**未验收**,分别留 P05–P08。
|
||||
- TDD 与隔离验证:`go test ./internal/store ./internal/dispatcher ./internal/mq -count=1`、`bash scripts/check-current-mq-mock.sh`、`git diff --check` 通过;Mock 覆盖任务积压、恢复、控制、重复投递、发布失败、SIP revision 栅栏与格式错误。当前 Agent side effect 为注入的**假外呼**;根 D/Agent 命令均已切换为 Mock-only,但含任务的主进程执行、录音及 `call.result` 未贯通。真实 SaaS/RabbitMQ、供应商与线路均**未验收**。
|
||||
|
||||
## P05:Agent 快照与 SDK 隔离链路(项目内组件通过,主入口待 P07)
|
||||
|
||||
@@ -78,14 +78,14 @@
|
||||
|
||||
- 已新增 Agent 内存录音单次 PUT 组件 `UploadClient.UploadBytes`:复用受限授权、大小和 SHA-256 核对,成功路径不写临时录音或通话结果文件。OSS 明确拒绝保留状态码且不自动重试;传输结果不明时返回专用错误、停止自动重试并隐藏带签名的 URL。空授权头拒绝而非使服务崩溃。这里只验证组件,尚未连接实际录音或最终结果。
|
||||
- 内部 OSS 授权增加 Dispatcher 原始 bucket,官方 SDK 签发时返回获批 bucket;Agent 收到缺少 bucket 的授权会在 PUT 前拒绝,失败恢复的显式重申请若返回了不同 bucket,也在再次 PUT 前拒绝。两条红灯测试证明先前会错误上传;修复后全包测试及 Agent/RPC/OSS race 测试通过。这里只核验本地授权载体,不代表新主入口已完成签发或真实 OSS 已验证。
|
||||
- Agent 失败恢复隔离组件:仅在 OSS PUT 明确失败后,把原 bucket/object_key 对应的录音和通话信息两文件写入私有目录并同步落盘;双文件缺失或损坏明确报错、不伪造结果。完成保存后固定 48 小时窗口,按 1 分钟递增至最长 1 小时重试;到期保留原文件。PUT 前持久写入 in-flight,结果不明或进程重启不会二次 PUT;确认上传后先持久记录成功,再经注入的 Mock 回报通话结果,回报失败/重启仅重发原结果。启动扫描识别遗漏文件并提供不暴露原路径的稳定摘要;真实 Dispatcher 授权及回报尚未接线。
|
||||
- Agent 失败恢复隔离组件:仅在 OSS PUT 明确失败后,把原 bucket/object_key 对应的录音和通话信息两文件写入私有目录并同步落盘;双文件缺失或损坏明确报错、不伪造结果。完成保存后固定 48 小时窗口,按 1 分钟递增至最长 1 小时重试;到期保留原文件。PUT 前持久写入 in-flight,结果不明或进程重启不会二次 PUT;确认上传后先持久记录成功,再经注入的 Mock 回报通话结果,回报失败/重启仅重发原结果。启动扫描识别遗漏文件并提供不暴露原路径的稳定摘要;主程序录音授权与回报已接线,但尚未通过含任务的端到端联测。
|
||||
- Dispatcher 的无录音最终结果隔离组件:`CurrentStore.RecordCallResult` 仅在确认通话结束后,按持久任务快照校验任务、被叫、主叫和已选线路,并以源执行事件固定生成唯一最终结果身份;消息通过严格 MQ Schema 校验后与 outbox 在同一事务写入。同内容重投/重启只恢复原消息,冲突结果和 SQLite 写入失败均不会产生第二份结果。这里只验证隔离组件,Agent 实际回报尚未连通。
|
||||
- Dispatcher 的原始 OSS 目标及已上传结果隔离组件:新 SQLite 布局把一次通话的 upload_id、bucket、object_key、录音格式/时长/大小和 SHA-256 唯一绑定到已保留的执行;不保存临时 URL 或 TOKEN。旧布局拒绝启动并原样保留待交付 outbox,不自动迁移或清理。已签发录音目标不能通过空录音结果绕过上传;已有空录音最终结果不能再签发录音目标。`RecordUploadedCallResult` 仅接受与持久绑定完全一致的录音事实及 Agent 所报告的成功 PUT 状态,录音确认与唯一最终结果 outbox 同事务提交;丢失回报或 MQ 投递时重用原消息,已确认后拒绝再次签发 PUT 授权。Mock 证明的是本地状态约束,不是独立 OSS 校验或真实 Agent 身份验证。
|
||||
- Dispatcher 的终结与外呼回执顺序竞争隔离修复:原流程在 Agent 接受执行的 RPC 返回后才写入外呼回执,快速结束或 RPC 超时可能先到;现以一次 SQLite 事务在确认通话已结束后补齐原回执并释放占用,未知执行仅在确认结束后释放。迟到的执行响应、超时和重复结束不会产生第二份回执;注入 outbox 写入失败保留原占用。并发竞争及结束后立即生成唯一最终结果有单元测试;主入口真实 Agent 会话注入与通话执行仍未接线。
|
||||
- Dispatcher 录音事实 Unary RPC 隔离服务:`RequestRecordingUpload`、`ReportCallEnded`、`ReportCallResult` 均要求已配置本 D、核验 mTLS 指纹及当前 Agent 会话、数字租户和已保留的执行;复用官方 SDK 仅对原始录音签发固定 15 分钟授权,显式重申请仍用相同 bucket/object_key。结束事实可先于外呼响应而持久化原回执;录音结果核对 D 已存目标和 Agent 报告的成功 PUT,再与唯一结果 outbox 同事务提交。Mock 覆盖会话/租户拒绝、同资产重申请、上传前结果拒绝、坏 JSON、已上传与无录音结果及重复回报。隔离测试还通过本地双向 TLS 的 gRPC 实际传输:Agent `RecordingClient` 每次读取并克隆当前会话元数据,经受控 D 客户端领取授权、上报结束和唯一结果;未配置客户端明确拒绝。其分层测试不包含 OSS PUT;下述本地隔离链路另验证直传,主入口批准执行的媒体录音仍未接线。
|
||||
- 内存录音隔离组件:`RecordingSession` 仅复制共享通话流程实际读到和成功发送的 16-kHz PCM16,`EncodeMonoWAV` 直接在内存生成有界单声道 WAV;空音频、奇数字节、超过上限及未成功发送的音频都不能伪造成可上传录音。单元与 race 测试未产生业务文件。批准执行入口尚未接入该组件,且 Mock 中观测到的帧不等于真实 Asterisk 通话的全量媒体验收。
|
||||
- Dispatcher 的终结与外呼回执顺序竞争隔离修复:原流程在 Agent 接受执行的 RPC 返回后才写入外呼回执,快速结束或 RPC 超时可能先到;现以一次 SQLite 事务在确认通话已结束后补齐原回执并释放占用,未知执行仅在确认结束后释放。迟到的执行响应、超时和重复结束不会产生第二份回执;注入 outbox 写入失败保留原占用。并发竞争及结束后立即生成唯一最终结果有单元测试;主入口 Agent 会话注入已接线,实际通话执行与结果交付尚未联测。
|
||||
- Dispatcher 录音事实 Unary RPC 隔离服务:`RequestRecordingUpload`、`ReportCallEnded`、`ReportCallResult` 均要求已配置本 D、核验 mTLS 指纹及当前 Agent 会话、数字租户和已保留的执行;复用官方 SDK 仅对原始录音签发固定 15 分钟授权,显式重申请仍用相同 bucket/object_key。结束事实可先于外呼响应而持久化原回执;录音结果核对 D 已存目标和 Agent 报告的成功 PUT,再与唯一结果 outbox 同事务提交。Mock 覆盖会话/租户拒绝、同资产重申请、上传前结果拒绝、坏 JSON、已上传与无录音结果及重复回报。隔离测试还通过本地双向 TLS 的 gRPC 实际传输:Agent `RecordingClient` 每次读取并克隆当前会话元数据,经受控 D 客户端领取授权、上报结束和唯一结果;未配置客户端明确拒绝。其分层测试不包含 OSS PUT;下述本地隔离链路另验证直传,主入口批准执行的合成媒体录音已接线,但完整录音交付尚未在主进程联测。
|
||||
- 内存录音隔离组件:`RecordingSession` 仅复制共享通话流程实际读到和成功发送的 16-kHz PCM16,`EncodeMonoWAV` 直接在内存生成有界单声道 WAV;空音频、奇数字节、超过上限及未成功发送的音频都不能伪造成可上传录音。单元与 race 测试未产生业务文件。批准执行的合成 Mock runner 已接入该组件;Mock 中观测到的帧不等于真实 Asterisk 通话的全量媒体验收。
|
||||
- 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 转写时间精度。
|
||||
- 最终结果隔离组件:`ApprovedExecution` 已保留获批任务的 `caller_profile_id`;`FinalResultPayload` 只从批准的任务身份、确认时间及实际收到的用户媒体窗口生成当前唯一结果;最终 ASR 文本与字面关键词拒联写入完整转写,未播放的助手候选内容不冒充转写,未知 SIP 响应码保留 `null`,无录音保留 `{}`。单元测试直接核验当前 MQ Schema、字段缺失/乱序时窗与无应答负例;开场白发送失败不得计入已播放音频。此组件已接入 Agent 合成 runner,但未证明真实 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、唯一最终结果及正常路径无业务文件。无实际读入媒体时明确报告失败并以空录音、生成失败原因收口,不伪造上传或转写;缺失每通话交付器、交付器与签发 D/数字租户/事件不一致、私有恢复目录并非 `0700` 或模拟脚本不满足已批准 AI 时,均在发外呼接受回执前拒绝。此 runner 已接到隔离 Mock 主 CLI 的 `ApprovedCallWorker`,但还没有从主进程完成 D→A 执行、Agent→D 录音与 MQ 出站的端到端联测;现有 OSS/RPC 测试均为隔离替身,不能代表真实线路或供应商通过。
|
||||
- 已验证:`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 入口已组装隔离合成媒体和录音交付,但尚未通过主进程 D↔A 会话/执行/录音、实际上传事实与最终结果交付、MQ 发布及重启恢复的整体联测,不能宣称 P06 通过。
|
||||
|
||||
@@ -53,3 +53,4 @@ export RABBITMQ_URL="amqp://dispatcher_mock:${d_pw}@127.0.0.1:${port}/"
|
||||
export RABBITMQ_PROVISIONER_URL="amqp://saas_mock:${saas_pw}@127.0.0.1:${port}/"
|
||||
go test -tags=integration ./internal/mq -run '^TestCurrentBrokerSharedResultQueueAndNoConfigure$' -count=1 -v
|
||||
go test -tags=integration ./internal/dispatcher -run '^TestCurrentRuntimeIsolated' -count=1 -v
|
||||
go test -tags=integration ./cmd/sip-go-agent -run '^TestCurrentDispatcherCommandStartsWithIsolatedMQHTTPAndAgent$' -count=1 -v
|
||||
|
||||
Reference in New Issue
Block a user