From a52fe5a741ace0ebf0664da3da8863b01b73a39a Mon Sep 17 00:00:00 2001 From: Rogee Date: Sat, 3 Oct 2026 12:45:11 +0800 Subject: [PATCH] Apply approved SIP snapshots through isolated Dispatcher and Agent channel --- AGENTS.md | 4 +- cmd/sip-go-agent/agent_command.go | 47 +-- cmd/sip-go-agent/agent_sip.go | 40 +++ cmd/sip-go-agent/dispatcher_command.go | 5 +- cmd/sip-go-agent/dispatcher_sip.go | 162 ++++++++++ deploys/README.md | 9 +- docs/evidence/sip-user-service-nonprod.md | 10 +- gen/agent/agent.pb.go | 270 ++++++++++++----- gen/agent/agent_grpc.pb.go | 40 +++ internal/asterisk/loader.go | 284 ++++++++++++++++++ internal/asterisk/loader_test.go | 100 ++++++ internal/config/agent_runtime.go | 31 +- internal/config/agent_runtime_test.go | 17 ++ internal/config/dispatcher_runtime.go | 11 +- internal/config/dispatcher_runtime_test.go | 11 + internal/config/runtime.go | 6 +- internal/dispatcher/approved_originator.go | 57 +++- .../dispatcher/approved_originator_test.go | 57 ++++ internal/dispatcher/config.go | 6 + internal/dispatcher/config_test.go | 17 +- internal/dispatcher/sip_only.go | 83 +++++ internal/dispatcher/sip_only_test.go | 67 +++++ internal/dispatcher/sip_reload.go | 6 + internal/dispatcher/sip_runtime_test.go | 19 +- internal/rpc/server.go | 7 +- internal/rpc/service_test.go | 2 +- internal/rpc/sip_apply.go | 70 +++++ internal/rpc/sip_apply_test.go | 55 ++++ internal/store/sip.go | 22 ++ internal/store/sip_change_test.go | 27 ++ proto/agent/agent.proto | 11 + proto/manifest.json | 12 +- 32 files changed, 1437 insertions(+), 128 deletions(-) create mode 100644 cmd/sip-go-agent/agent_sip.go create mode 100644 cmd/sip-go-agent/dispatcher_sip.go create mode 100644 internal/asterisk/loader.go create mode 100644 internal/asterisk/loader_test.go create mode 100644 internal/dispatcher/sip_only.go create mode 100644 internal/dispatcher/sip_only_test.go create mode 100644 internal/rpc/sip_apply.go create mode 100644 internal/rpc/sip_apply_test.go diff --git a/AGENTS.md b/AGENTS.md index deaa4b4..2329a1c 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -80,7 +80,7 @@ - P01–P08 及 K01–K16 已完成**项目内隔离 Mock** 核验;本轮把分散的人类可读契约、文档与第三方对接合为唯一当前规范,不重审已确认规则。若未来另获启动开发/审查子 Agent 授权,必须按使用者指定的 `gpt-5.6-luna`、`max` 思考和 `fast: true` 逐项核验并显式配置,不静默换模型、降档或关闭 fast。 - 唯一现行 SaaS↔Dispatcher 业务规范是 [`docs/thirds/saas-dispatcher.md`](docs/thirds/saas-dispatcher.md);当前项目内 Schema、拓扑、正反例及来源/hash 在 [`contracts/local/`](contracts/local/);内部 Agent RPC 在 [`proto/agent/agent.proto`](proto/agent/agent.proto)。本地验收与外部缺口见 [`docs/evidence/saas-dispatcher-p08-acceptance.md`](docs/evidence/saas-dispatcher-p08-acceptance.md)。Markdown 不代替机器合同或外部签收,也不另外维护一份平行字段定义。 - 已由当前合同来源清单固定哈希的历史提案与计划保留**原字节**于 [`docs/archive/sources/`](docs/archive/sources/README.md),使用者原有未提交的两份旧对接文档也按原字节归档;旧上游 v1 在 [`docs/archive/upstream/`](docs/archive/upstream/README.md) 可离线校验,但不嵌入运行合同。旧 F/W 工作包、旧 MQ-only 合同及归档不作为当前运行入口。固定 MQ `v1`、HTTP `/internal/v1/dispatcher/...` 和业务 revision 是现行通信规则,不是自有实现代次;不得为历史路径新建兼容或回退。 -- 当前本地范围仅**单节点、单 Dispatcher、单 Agent、单 Cell、单租户和隔离 Mock**。根命令只接受显式 `agent`/`dispatcher` 业务子命令,mixed/real 启动即拒绝;没有真实 SaaS、management、RabbitMQ、OSS、AI 供应商、Asterisk、SIP 线路、ECS 或生产签收。生产发布包仍为 `production_approval=false`。本机 Mock HTTPS、双向 TLS、RabbitMQ 和哈希通过均不授权真实呼叫。 +- 当前业务范围仍仅**单节点、单 Dispatcher、单 Agent、单 Cell、单租户和隔离 Mock**。根命令只接受显式 `agent`/`dispatcher`;mixed/real 启动即拒绝。另有严格隔离的 `--mode sip-only`:只允许 Dispatcher 读完整 SIP、持久接纳归属 `sip.config`、经已激活双向 TLS Agent 会话将完整快照应用到原生 Asterisk,并从运行态核对版本;不发现任务、不启动业务呼叫、不开放准入、不处理其他业务控制。测试机已仅凭使用者批准写入三条免鉴权、无需 REGISTER 的 UDP 历史登记线路并核对运行态,无拨号;这不证明供应商线路可用或生产签收。没有真实 SaaS、management、OSS、AI 供应商或生产签收;生产发布包仍为 `production_approval=false`。任何本机 Mock 或 SIP-only 核验均不授权真实呼叫。 - 开发按 TDD 分批,小步提交;不得覆盖使用者现存修改/未跟踪文件,不自动清理、迁移或覆盖任何现存 SQLite、spool、outbox 和 Agent 恢复文件。旧 `.executions` 及恢复根目录中旧 `.uploads`、`.upload-locks`、逐执行 `state.json` 的发现须只读失败关闭,现存未交付事实由使用者确认处置。真实云账号、EIP、线路、拨号、生产部署和共享数据操作分别需要明确授权。 ## SaaS、Dispatcher 与 Agent 的现行边界 @@ -89,7 +89,7 @@ - 呼叫、控制、必要回执和**每通话唯一最终结果**经固定 `v1` RabbitMQ Topic/队列,任务与控制队列由 SaaS 预建,Dispatcher 不可自行建/删/绑定;stop 时先停该任务消费者并关闭持久准入,再 purge 仅该任务队列的待投递消息,失败不回成功;SaaS 停止继续投递已 stop 任务,后续误投递不执行。结果进入指定共享 durable 队列。独立 D UUID 和接收队列不能广播后正文过滤;`tenant_key` 原值保留,任务只能由归属 D 执行。入站先校验和持久 inbox 后 ACK,状态/outbox 同事务;persistent、mandatory、无 return、publisher confirm 成功才记交付,confirm **不是** SaaS 应用收讫。失败/确认丢失与重启只重发同身份消息,不重复拨号或捏造结果。 - `call.execute` 只带获批 `task_id/callee`,调用线路、主叫、AI 和时限由该任务快照固定;`task.control` 的 start/pause/resume/stop 无旧 CAS/版本字段,控制回 task/action/status 处理结果,stop 不可恢复。任务启动/重启完整读取列表,运行中不定时轮询任务列表;新建任务由 start 读取任务配置和租户额度,暂停后修改的任务由 resume 重读;全局 AI 服务商列表每次启动只读取一次,本进程全部任务复用,变更须重启后才生效;没有 edit 事件。启动时全量读取并核验 SIP,运行中只由 `sip.config` 通知触发全量 SIP 读取,无常规定时 SIP 轮询;任务 start/resume 使用已生效 SIP 快照,不自行拉取或向 Agent 核验 SIP。SIP 修订变化待旧呼叫结束且 Agent 已加载后,仅在本地重新绑定任务快照,不重读任务列表;待处理期间新准入关闭。发现分页同快照先全部校验后提交,MQ 控制持久状态高于偶发 HTTP 状态;冲突关准入、resume 须重新核验。历史命令不因消息年龄过期,但实际呼出前仍检查任务/白名单/时段/授权/额度和有效通话上限;任务结束不删未确认结果、恢复、幂等或未知占用。 - 独立 Dispatcher 的 SQLite 是任务、额度、inbox/outbox 的权威数据;Agent 无业务数据库,录音、执行与上传恢复只写受控私有文件。额度包含未知占用,新 boot/租约到期不得自动清除未知执行;不实现双活数据库、自动跨机热备、多 D 共享额度或第二租户公平。本轮不借本机 D1/D2 隔离夹具宣称多 D 运行。不得建立旧表/旧消息/旧 HTTP 执行兼容通道。 -- Dispatcher↔Agent 复用 Unary gRPC 和受控 Endpoint;Agent 预绑定 D UUID 与服务端证书指纹,激活/会话代际、peer mTLS/SAN/SNI 和已签发期限须核对,新 boot 不清未知占用。Agent 不自行向 SaaS 取任务/AI/OSS 授权;Dispatcher 只用已经核验的 Agent `GetLoadedSIP` revision 开执行准入。本机 Mock 的加载回报不证明 Asterisk 已实际加载,SIP 配置的唯一编辑/审批面仍是 management。 +- Dispatcher↔Agent 复用 Unary gRPC 和受控 Endpoint;Agent 预绑定 D UUID 与服务端证书指纹,激活/会话代际、peer mTLS/SAN/SNI 和已签发期限须核对,新 boot 不清未知占用。Agent 不自行向 SaaS 取任务/AI/OSS 授权;Dispatcher 只用已经核验的 Agent `GetLoadedSIP` revision 开执行准入。本机 Mock 的加载回报不证明 Asterisk 已实际加载;仅隔离 SIP-only Agent 在配置原子写入、PJSIP reload 和运行态 endpoint/AOR/UDP transport 一致后持久标记 revision,每次加载查询重新核验。只支持明确的 UDP、IP/none 鉴权、无需 REGISTER 的 IPv4/PCMA 线路;未知字段和其他传输/鉴权/注册方式拒绝,不热更静态 transport。SIP 配置的唯一编辑/审批面仍是 management。 - AI 使用任务内不可变授权快照:仅经获批准百炼/火山 ASR、OpenAI 兼容 LLM、火山 TTS 能表达的参数进入每通话实例;ASR-only 不启动 LLM/TTS,完整 AI 不借旧语音测试的 LLM/TTS。只有最终用户 ASR 文本的明确字面关键词可触发拒联/挂断;不由 SDK 默认值、环境、CLI、metadata 或宽松 Schema 改写业务参数,不因 SDK 重试产生第二次发起/收费或重播。日志只存脱敏版本/摘要/计数,不存密钥、prompt、完整对话或音频。 - Agent 录音经受控双向 TLS 向 D 领取短期 OSS 上传授权,每次尝试只作**一次 HTTPS PUT**;正常上传成功不写录音/结果业务文件。首次明确失败须先完整保存录音与结果两份恢复文件,才从该时刻启动 48 小时重试;按 1、2、4、8、16、32、60 分钟及其后每 60 分钟的固定节奏显式重新申请授权,同一 OSS 目标、同一消息身份。PUT 结果未知不得盲目重传;48 小时届满仍失败时保留文件待人工,**不伪造最终结果或自动清理**。D 不转发文件,已确认结束的通话及时释放执行占用;未知执行仍占用。只有真实终结后才通过唯一 `call.execute.result` 回报录音路径、最终转写和拒联事实;无录音或录音生成失败以空 `recording={}` 和真实结果收口,生成失败须说明原因。不能恢复的录音不声称零丢失,也不伪造 OSS/SaaS 应用回执。凭据/TOKEN/签名 URL 不写入样例、日志、源码或证据。 diff --git a/cmd/sip-go-agent/agent_command.go b/cmd/sip-go-agent/agent_command.go index 8bf3598..0afabde 100644 --- a/cmd/sip-go-agent/agent_command.go +++ b/cmd/sip-go-agent/agent_command.go @@ -2,6 +2,7 @@ package main import ( "bytes" + "crypto/tls" "encoding/json" "errors" "fmt" @@ -28,14 +29,6 @@ func newAgentCommand() *cobra.Command { if err != nil { return err } - scenario, err := loadApprovedMockScenario(settings.MockScenarioFile) - if err != nil { - return err - } - applied, err := loadMockAppliedSIP(settings.MockAppliedSIPFile) - if err != nil { - return err - } ca, err := readAgentPEM("MTLS_CA_FILE", settings.CAFile) if err != nil { return err @@ -52,22 +45,38 @@ func newAgentCommand() *cobra.Command { if err != nil { return errors.New("Agent mTLS listener certificate is invalid") } - clientTLS, err := rpc.NewClientTLSConfig(ca, cert, key, settings.DispatcherServerName) - if err != nil { - return errors.New("Agent mTLS Dispatcher certificate configuration is invalid") + var handler *rpc.Server + if mode == "sip-only" { + handler, err = newSIPOnlyAgentServer(settings) + } else { + var scenario approvedMockScenario + scenario, err = loadApprovedMockScenario(settings.MockScenarioFile) + if err != nil { + return err + } + var applied map[string]int64 + applied, err = loadMockAppliedSIP(settings.MockAppliedSIPFile) + if err != nil { + return err + } + var clientTLS *tls.Config + clientTLS, err = rpc.NewClientTLSConfig(ca, cert, key, settings.DispatcherServerName) + if err != nil { + return errors.New("Agent mTLS Dispatcher certificate configuration is invalid") + } + connection, connectErr := grpc.NewClient(settings.DispatcherEndpoint, grpc.WithTransportCredentials(credentials.NewTLS(clientTLS))) + if connectErr != nil { + return errors.New("Agent cannot create pinned Dispatcher connection") + } + defer connection.Close() + handler, err = newAgentServer(cmd.Context(), settings, scenario, applied, agentpb.NewAgentControlServiceClient(connection)) } - connection, err := grpc.NewClient(settings.DispatcherEndpoint, grpc.WithTransportCredentials(credentials.NewTLS(clientTLS))) - if err != nil { - return errors.New("Agent cannot create pinned Dispatcher connection") - } - defer connection.Close() - handler, err := newAgentServer(cmd.Context(), settings, scenario, applied, agentpb.NewAgentControlServiceClient(connection)) if err != nil { return err } listener, err := net.Listen("tcp", settings.Listen) if err != nil { - return fmt.Errorf("Agent cannot listen on configured Mock address: %w", err) + return fmt.Errorf("Agent cannot listen on configured local address: %w", err) } defer listener.Close() server := grpc.NewServer(grpc.Creds(credentials.NewTLS(serverTLS))) @@ -91,7 +100,7 @@ func newAgentCommand() *cobra.Command { return nil }, } - command.Flags().StringVar(&mode, "mode", "mock", "isolated Mock mode only") + command.Flags().StringVar(&mode, "mode", "mock", "isolated Mock or SIP-only mode") return command } diff --git a/cmd/sip-go-agent/agent_sip.go b/cmd/sip-go-agent/agent_sip.go new file mode 100644 index 0000000..463711c --- /dev/null +++ b/cmd/sip-go-agent/agent_sip.go @@ -0,0 +1,40 @@ +package main + +import ( + "errors" + "os" + "path/filepath" + + agentpb "git.ipao.vip/rogee/go-sip/gen/agent" + "git.ipao.vip/rogee/go-sip/internal/asterisk" + "git.ipao.vip/rogee/go-sip/internal/config" + "git.ipao.vip/rogee/go-sip/internal/rpc" + "github.com/google/uuid" +) + +func newSIPOnlyAgentServer(settings config.AgentEnvironment) (*rpc.Server, error) { + if settings.AgentID == "" || settings.CellID == "" || settings.DispatcherID == "" || settings.SessionPath == "" || len(settings.PeerFingerprints) == 0 { + return nil, errors.New("SIP-only Agent requires bound identities, persistent sessions and pinned peer certificates") + } + for _, path := range []string{settings.AsteriskConfigDir, settings.AsteriskBin, settings.AsteriskLibraryDir} { + if !filepath.IsAbs(path) { + return nil, errors.New("native Asterisk paths must be explicit absolute paths") + } + } + if err := rejectLegacyAgentSpool(settings.RecoveryRoot); err != nil { + return nil, err + } + if _, err := os.Lstat(filepath.Join(settings.RecoveryRoot, ".executions")); err == nil { + return nil, errors.New("legacy Agent execution state requires operator disposition") + } else if !errors.Is(err, os.ErrNotExist) { + return nil, err + } + loader := asterisk.Loader{ConfigDir: settings.AsteriskConfigDir, Asterisk: settings.AsteriskBin, LibraryDir: settings.AsteriskLibraryDir} + return rpc.NewServer(rpc.ServerOptions{ + Mode: "sip-only", StatePath: settings.SessionPath, + ApprovedDispatcherID: settings.DispatcherID, + Status: &agentpb.AgentStatus{AgentId: settings.AgentID, CellId: settings.CellID, BootId: uuid.NewString(), ProtocolVersion: "agent.v1"}, + RequirePeerCertificate: true, PeerCertificateFingerprints: settings.PeerFingerprints, + LoadedSIP: loader.LoadedSIP, ApplySIP: loader.Apply, + }), nil +} diff --git a/cmd/sip-go-agent/dispatcher_command.go b/cmd/sip-go-agent/dispatcher_command.go index 38d333b..0fd7bb7 100644 --- a/cmd/sip-go-agent/dispatcher_command.go +++ b/cmd/sip-go-agent/dispatcher_command.go @@ -30,7 +30,7 @@ func newDispatcherCommand() *cobra.Command { return runDispatcher(cmd.Context(), mode) }, } - command.Flags().StringVar(&mode, "mode", "mock", "isolated Mock mode only") + command.Flags().StringVar(&mode, "mode", "mock", "isolated Mock or SIP-only mode") return command } @@ -41,6 +41,9 @@ func runDispatcher(ctx context.Context, mode string) (result error) { if err != nil { return err } + if mode == "sip-only" { + return runSIPOnlyDispatcher(ctx, settings) + } agentEndpoint, err := config.LoadMockAgentEndpoint(settings.AgentEndpointsFile) if err != nil { return err diff --git a/cmd/sip-go-agent/dispatcher_sip.go b/cmd/sip-go-agent/dispatcher_sip.go new file mode 100644 index 0000000..8d55fcf --- /dev/null +++ b/cmd/sip-go-agent/dispatcher_sip.go @@ -0,0 +1,162 @@ +package main + +import ( + "context" + "errors" + "fmt" + "log/slog" + "time" + + agentpb "git.ipao.vip/rogee/go-sip/gen/agent" + "git.ipao.vip/rogee/go-sip/internal/config" + "git.ipao.vip/rogee/go-sip/internal/dispatcher" + "git.ipao.vip/rogee/go-sip/internal/mq" + "git.ipao.vip/rogee/go-sip/internal/rpc" + "git.ipao.vip/rogee/go-sip/internal/store" + "github.com/google/uuid" + "google.golang.org/grpc" + "google.golang.org/grpc/credentials" +) + +// runSIPOnlyDispatcher owns no business queue, task discovery, recording +// listener, OSS signer, call executor or call admission. It consumes only the +// assigned control queue; unrelated controls are requeued and stop this lane. +func runSIPOnlyDispatcher(ctx context.Context, settings config.DispatcherRuntimeEnvironment) (result error) { + endpoints, err := config.LoadAgentEndpoints(settings.AgentEndpointsFile) + if err != nil { + return err + } + if len(endpoints) != 1 { + return errors.New("SIP-only Dispatcher requires exactly one assigned Agent") + } + endpoint := endpoints[0] + ca, err := readAgentPEM("MTLS_CA_FILE", settings.CAFile) + if err != nil { + return err + } + cert, err := readAgentPEM("MTLS_CERT_FILE", settings.CertFile) + if err != nil { + return err + } + key, err := readAgentPEM("MTLS_KEY_FILE", settings.KeyFile) + if err != nil { + return err + } + clientTLS, err := rpc.NewClientTLSConfig(ca, cert, key, endpoint.ServerName) + if err != nil { + return errors.New("SIP-only Dispatcher mTLS Agent certificate is invalid") + } + httpClient, err := localMockHTTPClient(ca) + if err != nil { + return err + } + httpClient.Timeout = 10 * time.Second + _, reader, err := dispatcherConfigurationClient("sip-only", httpClient) + if err != nil { + return err + } + db, err := store.Open(settings.SQLitePath) + if err != nil { + return err + } + defer func() { result = errors.Join(result, db.Close()) }() + if err := db.CloseAdmission(settings.DispatcherID); err != nil { + return err + } + broker, err := mq.Open(settings.RabbitMQURL, settings.DispatcherID, 1) + if err != nil { + return err + } + defer func() { result = errors.Join(result, broker.Close()) }() + connection, err := grpc.NewClient(endpoint.Address, grpc.WithTransportCredentials(credentials.NewTLS(clientTLS))) + if err != nil { + return errors.New("SIP-only Dispatcher cannot create pinned Agent connection") + } + defer func() { result = errors.Join(result, connection.Close()) }() + agent := agentpb.NewAgentControlServiceClient(connection) + coordinator := dispatcher.NewAgentCoordinator(time.Now) + if err := coordinator.Register(endpoint.AgentID, agent); err != nil { + return err + } + probeCtx, cancel := context.WithTimeout(ctx, 5*time.Second) + status, err := coordinator.Probe(probeCtx, endpoint.AgentID, endpoint.CellID) + if err != nil { + cancel() + return fmt.Errorf("probe SIP-only assigned Agent: %w", err) + } + epoch, err := uuid.NewRandom() + if err != nil { + cancel() + return err + } + session, err := coordinator.Activate(probeCtx, settings.DispatcherID, endpoint.AgentID, endpoint.CellID, status.BootId, epoch.String(), 0) + cancel() + if err != nil { + return fmt.Errorf("activate SIP-only Agent session: %w", err) + } + originator := &dispatcher.ApprovedOriginator{DispatcherID: settings.DispatcherID, Client: agent, Meta: func(callCtx context.Context) (*agentpb.RequestMeta, error) { + return coordinator.ApprovedMeta(callCtx, endpoint.AgentID) + }} + lane := dispatcher.SIPOnly{DispatcherID: settings.DispatcherID, Store: db, Client: reader, ApplySIP: originator.ApplySIP, VerifySIP: originator.VerifySIP} + if _, err := broker.DrainControlPredeclared(ctx, broker.ControlQueue(), lane.HandleNotification); err != nil { + return fmt.Errorf("drain assigned SIP controls before startup: %w", err) + } + if err := lane.Sync(ctx); err != nil { + return fmt.Errorf("initial native SIP-only sync: %w", err) + } + serveCtx, stop := context.WithCancel(ctx) + defer stop() + renewErrors := make(chan error, 1) + go func() { + renewErr := maintainAgentSession(serveCtx, session, 5*time.Minute, func(c context.Context, previous dispatcher.AgentSession) (dispatcher.AgentSession, error) { + return coordinator.Renew(c, settings.DispatcherID, previous) + }) + renewErrors <- renewErr + stop() + }() + consumer, err := broker.StartPredeclaredConsumer(serveCtx, broker.ControlQueue(), lane.HandleNotification) + if err != nil { + return err + } + defer func() { + waitCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + result = errors.Join(result, consumer.Stop(waitCtx)) + }() + consumerErrors := make(chan error, 1) + go func() { consumerErrors <- consumer.Wait(context.Background()); stop() }() + ticker := time.NewTicker(time.Second) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return nil + case err := <-renewErrors: + if err != nil { + return fmt.Errorf("SIP-only Agent session renewal failed: %w", err) + } + return errors.New("SIP-only Agent session stopped unexpectedly") + case err := <-consumerErrors: + if err != nil { + return fmt.Errorf("SIP-only control queue stopped: %w", err) + } + return errors.New("SIP-only control queue stopped unexpectedly") + case err := <-broker.Done(): + return fmt.Errorf("SIP-only broker connection lost: %v", err) + case <-ticker.C: + _, pending, err := db.SIPState(settings.DispatcherID) + if err != nil { + return err + } + if pending == 0 { + continue + } + readCtx, cancel := context.WithTimeout(serveCtx, 20*time.Second) + err = lane.Sync(readCtx) + cancel() + if err != nil { + slog.Warn("SIP-only revision remains durably pending", "dispatcher_id", settings.DispatcherID, "pending_revision", pending, "error", err) + } + } + } +} diff --git a/deploys/README.md b/deploys/README.md index 3503ef7..762ff84 100644 --- a/deploys/README.md +++ b/deploys/README.md @@ -41,8 +41,13 @@ artifact are injected separately. Asterisk is built or installed with the scripts under [`cell/`](cell/), and its management-owned configuration is never overwritten. -The current business binary is **isolated Mock-only**; mixed/real startup is -rejected, and the current release is not approved for production installation. +The business call runtime remains **isolated Mock-only**; mixed/real startup is +rejected. A separate `agent --mode sip-only` / `dispatcher --mode sip-only` lane +can update native Asterisk endpoint/AOR objects from a complete approved SIP +snapshot and durable `sip.config` notification. It requires preinstalled, +management-reviewed UDP transport plus a private endpoint include, and never +opens call admission or enables dialing. SIP-only does not make the current +release production-approved. `config/agent-endpoints.example.json` is the only current local configuration example. Earlier Dispatcher and static Cell JSON examples are preserved byte-for-byte under [`docs/archive/deployment-examples/`](../docs/archive/deployment-examples/) diff --git a/docs/evidence/sip-user-service-nonprod.md b/docs/evidence/sip-user-service-nonprod.md index 6df0b82..682d376 100644 --- a/docs/evidence/sip-user-service-nonprod.md +++ b/docs/evidence/sip-user-service-nonprod.md @@ -6,14 +6,20 @@ The test host's address and raw logs are omitted from committed evidence. - Fresh Debian 13 amd64 host, root directory mode `0755 root:root` (read-only inspection). - Native Asterisk 22.10.1 stage SHA-256: `68006a1a8efed288be4ca4a2ae3cb9554a31d733eac08eaacf4c646c95faf74d`. - Installed under `rogee`'s home without sudo; `systemctl --user` reported `go-sip-asterisk.service` **enabled and active**, and the live CLI returned Asterisk 22.10.1. -- `systemctl --user reload` succeeded and `res_pjsip.so` reported running. **No PJSIP endpoints were configured or loaded**; this does not verify adding/changing a real trunk or reloading it during an active call. +- Initial installation: `systemctl --user reload` succeeded and `res_pjsip.so` reported running. **At that time no PJSIP endpoints were configured**. Later authorized three-line loading is documented below; reload during an active call remains unverified. - User lingering was **disabled** by user choice: reboot persistence is **not verified**. The production installer now refuses installation without lingering; `--nonprod` is explicitly session-scoped. - An earlier user-install draft generated a template-only `[directories](!)` stanza and missed a startup crash. Corrected to `[directories]` and a live CLI readiness check; the corrected installer was checked for platform/lingering and no-overwrite behavior, but a fresh-install run of that final revision remains **unverified**. - Prior Debian 12 nonproduction host: a native, same-OS Asterisk 22.10.1 stage was built and its checksum checked. The default installer rejected Debian 12; `--nonprod` installed the stage. Its system service could not be accepted and the host was subsequently reinstalled; do not count that as a successful Cell deployment. +## Authorized native SIP load — 2026-10-03 + +- After confirming **zero active calls**, a fixed UDP transport and Agent-owned include were installed in the `rogee` user config. Restarting the user service took time; an early CLI query found it unready, so no loaded revision was claimed then. Once ready, `go-sip-udp` was observed on `0.0.0.0:5060`. +- The project's native renderer and loader applied a complete approved, registration-free, IP-auth three-trunk snapshot and inspected live PJSIP endpoint, AOR/contact, PCMA and transport state. **Three endpoints loaded, revision 9 recorded, zero active calls.** The initial interrupted observation recorded no revision; recovery required both the exact approved file and matching Asterisk runtime state. No dial, REGISTER, provider authorization, media, or active-call reload was attempted. +- This is a direct native-loader probe, **not yet** proof of Dispatcher→Agent RPC, RabbitMQ `sip.config`, or external SaaS. Only counts, status and revision are retained here; service addresses, caller identities, keys and raw host logs are omitted. + ## Still required before a real call or production acceptance -1. Implement and prove real Dispatcher→Agent SIP apply/reload and Agent-observed loaded state (current Go call runtime is still Mock-only). Validate endpoint add/edit against Asterisk while a call is active; explicitly reject unsupported authentication/REGISTER and transport changes. +1. Prove the isolated SIP-only Dispatcher→Agent RPC and durable `sip.config` notification path end to end against real Asterisk. Verify endpoint updates while a call is active separately; the business call runtime remains Mock-only and is not authorized to dial. Unsupported authentication/REGISTER and transport changes must continue to fail closed. 2. Provide approved real line configuration and arrange each whitelist trial (trunk, original number, time, attempt count). The fixed Asia/Shanghai 09:00–20:00 gate and per-number daily cap remain mandatory. 3. Establish RabbitMQ/OSS/AI real integrations and nonproduction call-evidence capture. `deploys/test/nonprod-call-evidence.sh` requires root or the required capture capabilities; this host's `rogee` currently has no sudo. If tcpdump, logger, or ARI/PJSIP state is unavailable, do not dial. 4. Enable `rogee` user lingering and verify reboot-persistent `enabled+active` before claiming production readiness. Complete the host/network/dependency diagnostics and external signoffs separately; local `make check` cannot replace them. diff --git a/gen/agent/agent.pb.go b/gen/agent/agent.pb.go index 9de7c68..9047077 100644 --- a/gen/agent/agent.pb.go +++ b/gen/agent/agent.pb.go @@ -2107,6 +2107,102 @@ func (x *GetLoadedSIPResponse) GetTrunkRevision() map[string]int64 { return nil } +type ApplySIPRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + Meta *RequestMeta `protobuf:"bytes,1,opt,name=meta,proto3" json:"meta,omitempty"` + ApprovedSnapshotJson []byte `protobuf:"bytes,2,opt,name=approved_snapshot_json,json=approvedSnapshotJson,proto3" json:"approved_snapshot_json,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ApplySIPRequest) Reset() { + *x = ApplySIPRequest{} + mi := &file_agent_agent_proto_msgTypes[20] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ApplySIPRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ApplySIPRequest) ProtoMessage() {} + +func (x *ApplySIPRequest) ProtoReflect() protoreflect.Message { + mi := &file_agent_agent_proto_msgTypes[20] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ApplySIPRequest.ProtoReflect.Descriptor instead. +func (*ApplySIPRequest) Descriptor() ([]byte, []int) { + return file_agent_agent_proto_rawDescGZIP(), []int{20} +} + +func (x *ApplySIPRequest) GetMeta() *RequestMeta { + if x != nil { + return x.Meta + } + return nil +} + +func (x *ApplySIPRequest) GetApprovedSnapshotJson() []byte { + if x != nil { + return x.ApprovedSnapshotJson + } + return nil +} + +type ApplySIPResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + TrunkRevision map[string]int64 `protobuf:"bytes,1,rep,name=trunk_revision,json=trunkRevision,proto3" json:"trunk_revision,omitempty" protobuf_key:"bytes,1,opt,name=key" protobuf_val:"varint,2,opt,name=value"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ApplySIPResponse) Reset() { + *x = ApplySIPResponse{} + mi := &file_agent_agent_proto_msgTypes[21] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ApplySIPResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ApplySIPResponse) ProtoMessage() {} + +func (x *ApplySIPResponse) ProtoReflect() protoreflect.Message { + mi := &file_agent_agent_proto_msgTypes[21] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ApplySIPResponse.ProtoReflect.Descriptor instead. +func (*ApplySIPResponse) Descriptor() ([]byte, []int) { + return file_agent_agent_proto_rawDescGZIP(), []int{21} +} + +func (x *ApplySIPResponse) GetTrunkRevision() map[string]int64 { + if x != nil { + return x.TrunkRevision + } + return nil +} + // Current task-level control has no external command ID, revision CAS or // execution binding. The authenticated session must match dispatcher_id. type ApplyApprovedTaskControlRequest struct { @@ -2123,7 +2219,7 @@ type ApplyApprovedTaskControlRequest struct { func (x *ApplyApprovedTaskControlRequest) Reset() { *x = ApplyApprovedTaskControlRequest{} - mi := &file_agent_agent_proto_msgTypes[20] + mi := &file_agent_agent_proto_msgTypes[22] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -2135,7 +2231,7 @@ func (x *ApplyApprovedTaskControlRequest) String() string { func (*ApplyApprovedTaskControlRequest) ProtoMessage() {} func (x *ApplyApprovedTaskControlRequest) ProtoReflect() protoreflect.Message { - mi := &file_agent_agent_proto_msgTypes[20] + mi := &file_agent_agent_proto_msgTypes[22] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -2148,7 +2244,7 @@ func (x *ApplyApprovedTaskControlRequest) ProtoReflect() protoreflect.Message { // Deprecated: Use ApplyApprovedTaskControlRequest.ProtoReflect.Descriptor instead. func (*ApplyApprovedTaskControlRequest) Descriptor() ([]byte, []int) { - return file_agent_agent_proto_rawDescGZIP(), []int{20} + return file_agent_agent_proto_rawDescGZIP(), []int{22} } func (x *ApplyApprovedTaskControlRequest) GetMeta() *RequestMeta { @@ -2202,7 +2298,7 @@ type ApplyApprovedTaskControlResponse struct { func (x *ApplyApprovedTaskControlResponse) Reset() { *x = ApplyApprovedTaskControlResponse{} - mi := &file_agent_agent_proto_msgTypes[21] + mi := &file_agent_agent_proto_msgTypes[23] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -2214,7 +2310,7 @@ func (x *ApplyApprovedTaskControlResponse) String() string { func (*ApplyApprovedTaskControlResponse) ProtoMessage() {} func (x *ApplyApprovedTaskControlResponse) ProtoReflect() protoreflect.Message { - mi := &file_agent_agent_proto_msgTypes[21] + mi := &file_agent_agent_proto_msgTypes[23] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -2227,7 +2323,7 @@ func (x *ApplyApprovedTaskControlResponse) ProtoReflect() protoreflect.Message { // Deprecated: Use ApplyApprovedTaskControlResponse.ProtoReflect.Descriptor instead. func (*ApplyApprovedTaskControlResponse) Descriptor() ([]byte, []int) { - return file_agent_agent_proto_rawDescGZIP(), []int{21} + return file_agent_agent_proto_rawDescGZIP(), []int{23} } func (x *ApplyApprovedTaskControlResponse) GetAccepted() bool { @@ -2253,7 +2349,7 @@ type UploadGrant struct { func (x *UploadGrant) Reset() { *x = UploadGrant{} - mi := &file_agent_agent_proto_msgTypes[22] + mi := &file_agent_agent_proto_msgTypes[24] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -2265,7 +2361,7 @@ func (x *UploadGrant) String() string { func (*UploadGrant) ProtoMessage() {} func (x *UploadGrant) ProtoReflect() protoreflect.Message { - mi := &file_agent_agent_proto_msgTypes[22] + mi := &file_agent_agent_proto_msgTypes[24] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -2278,7 +2374,7 @@ func (x *UploadGrant) ProtoReflect() protoreflect.Message { // Deprecated: Use UploadGrant.ProtoReflect.Descriptor instead. func (*UploadGrant) Descriptor() ([]byte, []int) { - return file_agent_agent_proto_rawDescGZIP(), []int{22} + return file_agent_agent_proto_rawDescGZIP(), []int{24} } func (x *UploadGrant) GetUploadId() string { @@ -2353,7 +2449,7 @@ type RequestRecordingUploadRequest struct { func (x *RequestRecordingUploadRequest) Reset() { *x = RequestRecordingUploadRequest{} - mi := &file_agent_agent_proto_msgTypes[23] + mi := &file_agent_agent_proto_msgTypes[25] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -2365,7 +2461,7 @@ func (x *RequestRecordingUploadRequest) String() string { func (*RequestRecordingUploadRequest) ProtoMessage() {} func (x *RequestRecordingUploadRequest) ProtoReflect() protoreflect.Message { - mi := &file_agent_agent_proto_msgTypes[23] + mi := &file_agent_agent_proto_msgTypes[25] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -2378,7 +2474,7 @@ func (x *RequestRecordingUploadRequest) ProtoReflect() protoreflect.Message { // Deprecated: Use RequestRecordingUploadRequest.ProtoReflect.Descriptor instead. func (*RequestRecordingUploadRequest) Descriptor() ([]byte, []int) { - return file_agent_agent_proto_rawDescGZIP(), []int{23} + return file_agent_agent_proto_rawDescGZIP(), []int{25} } func (x *RequestRecordingUploadRequest) GetMeta() *RequestMeta { @@ -2432,7 +2528,7 @@ type RequestRecordingUploadResponse struct { func (x *RequestRecordingUploadResponse) Reset() { *x = RequestRecordingUploadResponse{} - mi := &file_agent_agent_proto_msgTypes[24] + mi := &file_agent_agent_proto_msgTypes[26] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -2444,7 +2540,7 @@ func (x *RequestRecordingUploadResponse) String() string { func (*RequestRecordingUploadResponse) ProtoMessage() {} func (x *RequestRecordingUploadResponse) ProtoReflect() protoreflect.Message { - mi := &file_agent_agent_proto_msgTypes[24] + mi := &file_agent_agent_proto_msgTypes[26] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -2457,7 +2553,7 @@ func (x *RequestRecordingUploadResponse) ProtoReflect() protoreflect.Message { // Deprecated: Use RequestRecordingUploadResponse.ProtoReflect.Descriptor instead. func (*RequestRecordingUploadResponse) Descriptor() ([]byte, []int) { - return file_agent_agent_proto_rawDescGZIP(), []int{24} + return file_agent_agent_proto_rawDescGZIP(), []int{26} } func (x *RequestRecordingUploadResponse) GetGrant() *UploadGrant { @@ -2479,7 +2575,7 @@ type ReportCallEndedRequest struct { func (x *ReportCallEndedRequest) Reset() { *x = ReportCallEndedRequest{} - mi := &file_agent_agent_proto_msgTypes[25] + mi := &file_agent_agent_proto_msgTypes[27] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -2491,7 +2587,7 @@ func (x *ReportCallEndedRequest) String() string { func (*ReportCallEndedRequest) ProtoMessage() {} func (x *ReportCallEndedRequest) ProtoReflect() protoreflect.Message { - mi := &file_agent_agent_proto_msgTypes[25] + mi := &file_agent_agent_proto_msgTypes[27] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -2504,7 +2600,7 @@ func (x *ReportCallEndedRequest) ProtoReflect() protoreflect.Message { // Deprecated: Use ReportCallEndedRequest.ProtoReflect.Descriptor instead. func (*ReportCallEndedRequest) Descriptor() ([]byte, []int) { - return file_agent_agent_proto_rawDescGZIP(), []int{25} + return file_agent_agent_proto_rawDescGZIP(), []int{27} } func (x *ReportCallEndedRequest) GetMeta() *RequestMeta { @@ -2544,7 +2640,7 @@ type ReportCallEndedResponse struct { func (x *ReportCallEndedResponse) Reset() { *x = ReportCallEndedResponse{} - mi := &file_agent_agent_proto_msgTypes[26] + mi := &file_agent_agent_proto_msgTypes[28] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -2556,7 +2652,7 @@ func (x *ReportCallEndedResponse) String() string { func (*ReportCallEndedResponse) ProtoMessage() {} func (x *ReportCallEndedResponse) ProtoReflect() protoreflect.Message { - mi := &file_agent_agent_proto_msgTypes[26] + mi := &file_agent_agent_proto_msgTypes[28] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -2569,7 +2665,7 @@ func (x *ReportCallEndedResponse) ProtoReflect() protoreflect.Message { // Deprecated: Use ReportCallEndedResponse.ProtoReflect.Descriptor instead. func (*ReportCallEndedResponse) Descriptor() ([]byte, []int) { - return file_agent_agent_proto_rawDescGZIP(), []int{26} + return file_agent_agent_proto_rawDescGZIP(), []int{28} } func (x *ReportCallEndedResponse) GetReceipt() *OperationReceipt { @@ -2592,7 +2688,7 @@ type UploadObservation struct { func (x *UploadObservation) Reset() { *x = UploadObservation{} - mi := &file_agent_agent_proto_msgTypes[27] + mi := &file_agent_agent_proto_msgTypes[29] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -2604,7 +2700,7 @@ func (x *UploadObservation) String() string { func (*UploadObservation) ProtoMessage() {} func (x *UploadObservation) ProtoReflect() protoreflect.Message { - mi := &file_agent_agent_proto_msgTypes[27] + mi := &file_agent_agent_proto_msgTypes[29] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -2617,7 +2713,7 @@ func (x *UploadObservation) ProtoReflect() protoreflect.Message { // Deprecated: Use UploadObservation.ProtoReflect.Descriptor instead. func (*UploadObservation) Descriptor() ([]byte, []int) { - return file_agent_agent_proto_rawDescGZIP(), []int{27} + return file_agent_agent_proto_rawDescGZIP(), []int{29} } func (x *UploadObservation) GetUploadId() string { @@ -2669,7 +2765,7 @@ type ReportCallResultRequest struct { func (x *ReportCallResultRequest) Reset() { *x = ReportCallResultRequest{} - mi := &file_agent_agent_proto_msgTypes[28] + mi := &file_agent_agent_proto_msgTypes[30] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -2681,7 +2777,7 @@ func (x *ReportCallResultRequest) String() string { func (*ReportCallResultRequest) ProtoMessage() {} func (x *ReportCallResultRequest) ProtoReflect() protoreflect.Message { - mi := &file_agent_agent_proto_msgTypes[28] + mi := &file_agent_agent_proto_msgTypes[30] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -2694,7 +2790,7 @@ func (x *ReportCallResultRequest) ProtoReflect() protoreflect.Message { // Deprecated: Use ReportCallResultRequest.ProtoReflect.Descriptor instead. func (*ReportCallResultRequest) Descriptor() ([]byte, []int) { - return file_agent_agent_proto_rawDescGZIP(), []int{28} + return file_agent_agent_proto_rawDescGZIP(), []int{30} } func (x *ReportCallResultRequest) GetMeta() *RequestMeta { @@ -2748,7 +2844,7 @@ type ReportCallResultResponse struct { func (x *ReportCallResultResponse) Reset() { *x = ReportCallResultResponse{} - mi := &file_agent_agent_proto_msgTypes[29] + mi := &file_agent_agent_proto_msgTypes[31] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -2760,7 +2856,7 @@ func (x *ReportCallResultResponse) String() string { func (*ReportCallResultResponse) ProtoMessage() {} func (x *ReportCallResultResponse) ProtoReflect() protoreflect.Message { - mi := &file_agent_agent_proto_msgTypes[29] + mi := &file_agent_agent_proto_msgTypes[31] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -2773,7 +2869,7 @@ func (x *ReportCallResultResponse) ProtoReflect() protoreflect.Message { // Deprecated: Use ReportCallResultResponse.ProtoReflect.Descriptor instead. func (*ReportCallResultResponse) Descriptor() ([]byte, []int) { - return file_agent_agent_proto_rawDescGZIP(), []int{29} + return file_agent_agent_proto_rawDescGZIP(), []int{31} } func (x *ReportCallResultResponse) GetReceipt() *OperationReceipt { @@ -2943,6 +3039,14 @@ const file_agent_agent_proto_rawDesc = "" + "\x0etrunk_revision\x18\x01 \x03(\v2..agent.GetLoadedSIPResponse.TrunkRevisionEntryR\rtrunkRevision\x1a@\n" + "\x12TrunkRevisionEntry\x12\x10\n" + "\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n" + + "\x05value\x18\x02 \x01(\x03R\x05value:\x028\x01\"o\n" + + "\x0fApplySIPRequest\x12&\n" + + "\x04meta\x18\x01 \x01(\v2\x12.agent.RequestMetaR\x04meta\x124\n" + + "\x16approved_snapshot_json\x18\x02 \x01(\fR\x14approvedSnapshotJson\"\xa7\x01\n" + + "\x10ApplySIPResponse\x12Q\n" + + "\x0etrunk_revision\x18\x01 \x03(\v2*.agent.ApplySIPResponse.TrunkRevisionEntryR\rtrunkRevision\x1a@\n" + + "\x12TrunkRevisionEntry\x12\x10\n" + + "\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n" + "\x05value\x18\x02 \x01(\x03R\x05value:\x028\x01\"\x99\x02\n" + "\x1fApplyApprovedTaskControlRequest\x12&\n" + "\x04meta\x18\x01 \x01(\v2\x12.agent.RequestMetaR\x04meta\x12#\n" + @@ -3041,13 +3145,14 @@ const file_agent_agent_proto_rawDesc = "" + "\tAssetKind\x12\x1a\n" + "\x16ASSET_KIND_UNSPECIFIED\x10\x00\x12\x18\n" + "\x14ASSET_KIND_RECORDING\x10\x01\x12\x19\n" + - "\x15ASSET_KIND_TRANSCRIPT\x10\x022\xc6\x05\n" + + "\x15ASSET_KIND_TRANSCRIPT\x10\x022\x83\x06\n" + "\x13AgentControlService\x12M\n" + "\x0eGetAgentStatus\x12\x1c.agent.GetAgentStatusRequest\x1a\x1d.agent.GetAgentStatusResponse\x12J\n" + "\rActivateAgent\x12\x1b.agent.ActivateAgentRequest\x1a\x1c.agent.ActivateAgentResponse\x12P\n" + "\x0fExecuteApproved\x12\x1d.agent.ExecuteApprovedRequest\x1a\x1e.agent.ExecuteApprovedResponse\x12k\n" + "\x18ApplyApprovedTaskControl\x12&.agent.ApplyApprovedTaskControlRequest\x1a'.agent.ApplyApprovedTaskControlResponse\x12G\n" + - "\fGetLoadedSIP\x12\x1a.agent.GetLoadedSIPRequest\x1a\x1b.agent.GetLoadedSIPResponse\x12e\n" + + "\fGetLoadedSIP\x12\x1a.agent.GetLoadedSIPRequest\x1a\x1b.agent.GetLoadedSIPResponse\x12;\n" + + "\bApplySIP\x12\x16.agent.ApplySIPRequest\x1a\x17.agent.ApplySIPResponse\x12e\n" + "\x16RequestRecordingUpload\x12$.agent.RequestRecordingUploadRequest\x1a%.agent.RequestRecordingUploadResponse\x12P\n" + "\x0fReportCallEnded\x12\x1d.agent.ReportCallEndedRequest\x1a\x1e.agent.ReportCallEndedResponse\x12S\n" + "\x10ReportCallResult\x12\x1e.agent.ReportCallResultRequest\x1a\x1f.agent.ReportCallResultResponseB-Z+git.ipao.vip/rogee/go-sip/gen/agent;agentpbb\x06proto3" @@ -3065,7 +3170,7 @@ func file_agent_agent_proto_rawDescGZIP() []byte { } var file_agent_agent_proto_enumTypes = make([]protoimpl.EnumInfo, 7) -var file_agent_agent_proto_msgTypes = make([]protoimpl.MessageInfo, 31) +var file_agent_agent_proto_msgTypes = make([]protoimpl.MessageInfo, 34) var file_agent_agent_proto_goTypes = []any{ (ResultCode)(0), // 0: agent.ResultCode (FailureCode)(0), // 1: agent.FailureCode @@ -3094,17 +3199,20 @@ var file_agent_agent_proto_goTypes = []any{ (*ExecuteApprovedResponse)(nil), // 24: agent.ExecuteApprovedResponse (*GetLoadedSIPRequest)(nil), // 25: agent.GetLoadedSIPRequest (*GetLoadedSIPResponse)(nil), // 26: agent.GetLoadedSIPResponse - (*ApplyApprovedTaskControlRequest)(nil), // 27: agent.ApplyApprovedTaskControlRequest - (*ApplyApprovedTaskControlResponse)(nil), // 28: agent.ApplyApprovedTaskControlResponse - (*UploadGrant)(nil), // 29: agent.UploadGrant - (*RequestRecordingUploadRequest)(nil), // 30: agent.RequestRecordingUploadRequest - (*RequestRecordingUploadResponse)(nil), // 31: agent.RequestRecordingUploadResponse - (*ReportCallEndedRequest)(nil), // 32: agent.ReportCallEndedRequest - (*ReportCallEndedResponse)(nil), // 33: agent.ReportCallEndedResponse - (*UploadObservation)(nil), // 34: agent.UploadObservation - (*ReportCallResultRequest)(nil), // 35: agent.ReportCallResultRequest - (*ReportCallResultResponse)(nil), // 36: agent.ReportCallResultResponse - nil, // 37: agent.GetLoadedSIPResponse.TrunkRevisionEntry + (*ApplySIPRequest)(nil), // 27: agent.ApplySIPRequest + (*ApplySIPResponse)(nil), // 28: agent.ApplySIPResponse + (*ApplyApprovedTaskControlRequest)(nil), // 29: agent.ApplyApprovedTaskControlRequest + (*ApplyApprovedTaskControlResponse)(nil), // 30: agent.ApplyApprovedTaskControlResponse + (*UploadGrant)(nil), // 31: agent.UploadGrant + (*RequestRecordingUploadRequest)(nil), // 32: agent.RequestRecordingUploadRequest + (*RequestRecordingUploadResponse)(nil), // 33: agent.RequestRecordingUploadResponse + (*ReportCallEndedRequest)(nil), // 34: agent.ReportCallEndedRequest + (*ReportCallEndedResponse)(nil), // 35: agent.ReportCallEndedResponse + (*UploadObservation)(nil), // 36: agent.UploadObservation + (*ReportCallResultRequest)(nil), // 37: agent.ReportCallResultRequest + (*ReportCallResultResponse)(nil), // 38: agent.ReportCallResultResponse + nil, // 39: agent.GetLoadedSIPResponse.TrunkRevisionEntry + nil, // 40: agent.ApplySIPResponse.TrunkRevisionEntry } var file_agent_agent_proto_depIdxs = []int32{ 1, // 0: agent.Failure.code:type_name -> agent.FailureCode @@ -3129,40 +3237,44 @@ var file_agent_agent_proto_depIdxs = []int32{ 9, // 19: agent.ActivateAgentResponse.failure:type_name -> agent.Failure 7, // 20: agent.ExecuteApprovedRequest.meta:type_name -> agent.RequestMeta 7, // 21: agent.GetLoadedSIPRequest.meta:type_name -> agent.RequestMeta - 37, // 22: agent.GetLoadedSIPResponse.trunk_revision:type_name -> agent.GetLoadedSIPResponse.TrunkRevisionEntry - 7, // 23: agent.ApplyApprovedTaskControlRequest.meta:type_name -> agent.RequestMeta - 4, // 24: agent.ApplyApprovedTaskControlRequest.action:type_name -> agent.ControlAction - 5, // 25: agent.ApplyApprovedTaskControlRequest.active_call_policy:type_name -> agent.ActiveCallPolicy - 18, // 26: agent.UploadGrant.headers:type_name -> agent.Header - 7, // 27: agent.RequestRecordingUploadRequest.meta:type_name -> agent.RequestMeta - 17, // 28: agent.RequestRecordingUploadRequest.asset:type_name -> agent.AssetDescriptor - 29, // 29: agent.RequestRecordingUploadResponse.grant:type_name -> agent.UploadGrant - 7, // 30: agent.ReportCallEndedRequest.meta:type_name -> agent.RequestMeta - 10, // 31: agent.ReportCallEndedResponse.receipt:type_name -> agent.OperationReceipt - 7, // 32: agent.ReportCallResultRequest.meta:type_name -> agent.RequestMeta - 34, // 33: agent.ReportCallResultRequest.upload:type_name -> agent.UploadObservation - 10, // 34: agent.ReportCallResultResponse.receipt:type_name -> agent.OperationReceipt - 19, // 35: agent.AgentControlService.GetAgentStatus:input_type -> agent.GetAgentStatusRequest - 21, // 36: agent.AgentControlService.ActivateAgent:input_type -> agent.ActivateAgentRequest - 23, // 37: agent.AgentControlService.ExecuteApproved:input_type -> agent.ExecuteApprovedRequest - 27, // 38: agent.AgentControlService.ApplyApprovedTaskControl:input_type -> agent.ApplyApprovedTaskControlRequest - 25, // 39: agent.AgentControlService.GetLoadedSIP:input_type -> agent.GetLoadedSIPRequest - 30, // 40: agent.AgentControlService.RequestRecordingUpload:input_type -> agent.RequestRecordingUploadRequest - 32, // 41: agent.AgentControlService.ReportCallEnded:input_type -> agent.ReportCallEndedRequest - 35, // 42: agent.AgentControlService.ReportCallResult:input_type -> agent.ReportCallResultRequest - 20, // 43: agent.AgentControlService.GetAgentStatus:output_type -> agent.GetAgentStatusResponse - 22, // 44: agent.AgentControlService.ActivateAgent:output_type -> agent.ActivateAgentResponse - 24, // 45: agent.AgentControlService.ExecuteApproved:output_type -> agent.ExecuteApprovedResponse - 28, // 46: agent.AgentControlService.ApplyApprovedTaskControl:output_type -> agent.ApplyApprovedTaskControlResponse - 26, // 47: agent.AgentControlService.GetLoadedSIP:output_type -> agent.GetLoadedSIPResponse - 31, // 48: agent.AgentControlService.RequestRecordingUpload:output_type -> agent.RequestRecordingUploadResponse - 33, // 49: agent.AgentControlService.ReportCallEnded:output_type -> agent.ReportCallEndedResponse - 36, // 50: agent.AgentControlService.ReportCallResult:output_type -> agent.ReportCallResultResponse - 43, // [43:51] is the sub-list for method output_type - 35, // [35:43] is the sub-list for method input_type - 35, // [35:35] is the sub-list for extension type_name - 35, // [35:35] is the sub-list for extension extendee - 0, // [0:35] is the sub-list for field type_name + 39, // 22: agent.GetLoadedSIPResponse.trunk_revision:type_name -> agent.GetLoadedSIPResponse.TrunkRevisionEntry + 7, // 23: agent.ApplySIPRequest.meta:type_name -> agent.RequestMeta + 40, // 24: agent.ApplySIPResponse.trunk_revision:type_name -> agent.ApplySIPResponse.TrunkRevisionEntry + 7, // 25: agent.ApplyApprovedTaskControlRequest.meta:type_name -> agent.RequestMeta + 4, // 26: agent.ApplyApprovedTaskControlRequest.action:type_name -> agent.ControlAction + 5, // 27: agent.ApplyApprovedTaskControlRequest.active_call_policy:type_name -> agent.ActiveCallPolicy + 18, // 28: agent.UploadGrant.headers:type_name -> agent.Header + 7, // 29: agent.RequestRecordingUploadRequest.meta:type_name -> agent.RequestMeta + 17, // 30: agent.RequestRecordingUploadRequest.asset:type_name -> agent.AssetDescriptor + 31, // 31: agent.RequestRecordingUploadResponse.grant:type_name -> agent.UploadGrant + 7, // 32: agent.ReportCallEndedRequest.meta:type_name -> agent.RequestMeta + 10, // 33: agent.ReportCallEndedResponse.receipt:type_name -> agent.OperationReceipt + 7, // 34: agent.ReportCallResultRequest.meta:type_name -> agent.RequestMeta + 36, // 35: agent.ReportCallResultRequest.upload:type_name -> agent.UploadObservation + 10, // 36: agent.ReportCallResultResponse.receipt:type_name -> agent.OperationReceipt + 19, // 37: agent.AgentControlService.GetAgentStatus:input_type -> agent.GetAgentStatusRequest + 21, // 38: agent.AgentControlService.ActivateAgent:input_type -> agent.ActivateAgentRequest + 23, // 39: agent.AgentControlService.ExecuteApproved:input_type -> agent.ExecuteApprovedRequest + 29, // 40: agent.AgentControlService.ApplyApprovedTaskControl:input_type -> agent.ApplyApprovedTaskControlRequest + 25, // 41: agent.AgentControlService.GetLoadedSIP:input_type -> agent.GetLoadedSIPRequest + 27, // 42: agent.AgentControlService.ApplySIP:input_type -> agent.ApplySIPRequest + 32, // 43: agent.AgentControlService.RequestRecordingUpload:input_type -> agent.RequestRecordingUploadRequest + 34, // 44: agent.AgentControlService.ReportCallEnded:input_type -> agent.ReportCallEndedRequest + 37, // 45: agent.AgentControlService.ReportCallResult:input_type -> agent.ReportCallResultRequest + 20, // 46: agent.AgentControlService.GetAgentStatus:output_type -> agent.GetAgentStatusResponse + 22, // 47: agent.AgentControlService.ActivateAgent:output_type -> agent.ActivateAgentResponse + 24, // 48: agent.AgentControlService.ExecuteApproved:output_type -> agent.ExecuteApprovedResponse + 30, // 49: agent.AgentControlService.ApplyApprovedTaskControl:output_type -> agent.ApplyApprovedTaskControlResponse + 26, // 50: agent.AgentControlService.GetLoadedSIP:output_type -> agent.GetLoadedSIPResponse + 28, // 51: agent.AgentControlService.ApplySIP:output_type -> agent.ApplySIPResponse + 33, // 52: agent.AgentControlService.RequestRecordingUpload:output_type -> agent.RequestRecordingUploadResponse + 35, // 53: agent.AgentControlService.ReportCallEnded:output_type -> agent.ReportCallEndedResponse + 38, // 54: agent.AgentControlService.ReportCallResult:output_type -> agent.ReportCallResultResponse + 46, // [46:55] is the sub-list for method output_type + 37, // [37:46] is the sub-list for method input_type + 37, // [37:37] is the sub-list for extension type_name + 37, // [37:37] is the sub-list for extension extendee + 0, // [0:37] is the sub-list for field type_name } func init() { file_agent_agent_proto_init() } @@ -3176,7 +3288,7 @@ func file_agent_agent_proto_init() { GoPackagePath: reflect.TypeOf(x{}).PkgPath(), RawDescriptor: unsafe.Slice(unsafe.StringData(file_agent_agent_proto_rawDesc), len(file_agent_agent_proto_rawDesc)), NumEnums: 7, - NumMessages: 31, + NumMessages: 34, NumExtensions: 0, NumServices: 1, }, diff --git a/gen/agent/agent_grpc.pb.go b/gen/agent/agent_grpc.pb.go index 272c08c..ee79b95 100644 --- a/gen/agent/agent_grpc.pb.go +++ b/gen/agent/agent_grpc.pb.go @@ -24,6 +24,7 @@ const ( AgentControlService_ExecuteApproved_FullMethodName = "/agent.AgentControlService/ExecuteApproved" AgentControlService_ApplyApprovedTaskControl_FullMethodName = "/agent.AgentControlService/ApplyApprovedTaskControl" AgentControlService_GetLoadedSIP_FullMethodName = "/agent.AgentControlService/GetLoadedSIP" + AgentControlService_ApplySIP_FullMethodName = "/agent.AgentControlService/ApplySIP" AgentControlService_RequestRecordingUpload_FullMethodName = "/agent.AgentControlService/RequestRecordingUpload" AgentControlService_ReportCallEnded_FullMethodName = "/agent.AgentControlService/ReportCallEnded" AgentControlService_ReportCallResult_FullMethodName = "/agent.AgentControlService/ReportCallResult" @@ -42,6 +43,8 @@ type AgentControlServiceClient interface { ExecuteApproved(ctx context.Context, in *ExecuteApprovedRequest, opts ...grpc.CallOption) (*ExecuteApprovedResponse, error) ApplyApprovedTaskControl(ctx context.Context, in *ApplyApprovedTaskControlRequest, opts ...grpc.CallOption) (*ApplyApprovedTaskControlResponse, error) GetLoadedSIP(ctx context.Context, in *GetLoadedSIPRequest, opts ...grpc.CallOption) (*GetLoadedSIPResponse, error) + // Only a pinned, activated Dispatcher may apply the complete approved SIP snapshot. + ApplySIP(ctx context.Context, in *ApplySIPRequest, opts ...grpc.CallOption) (*ApplySIPResponse, error) // Agent requests a fresh bounded token explicitly; the Dispatcher owns the original target. RequestRecordingUpload(ctx context.Context, in *RequestRecordingUploadRequest, opts ...grpc.CallOption) (*RequestRecordingUploadResponse, error) ReportCallEnded(ctx context.Context, in *ReportCallEndedRequest, opts ...grpc.CallOption) (*ReportCallEndedResponse, error) @@ -106,6 +109,16 @@ func (c *agentControlServiceClient) GetLoadedSIP(ctx context.Context, in *GetLoa return out, nil } +func (c *agentControlServiceClient) ApplySIP(ctx context.Context, in *ApplySIPRequest, opts ...grpc.CallOption) (*ApplySIPResponse, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(ApplySIPResponse) + err := c.cc.Invoke(ctx, AgentControlService_ApplySIP_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + func (c *agentControlServiceClient) RequestRecordingUpload(ctx context.Context, in *RequestRecordingUploadRequest, opts ...grpc.CallOption) (*RequestRecordingUploadResponse, error) { cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) out := new(RequestRecordingUploadResponse) @@ -149,6 +162,8 @@ type AgentControlServiceServer interface { ExecuteApproved(context.Context, *ExecuteApprovedRequest) (*ExecuteApprovedResponse, error) ApplyApprovedTaskControl(context.Context, *ApplyApprovedTaskControlRequest) (*ApplyApprovedTaskControlResponse, error) GetLoadedSIP(context.Context, *GetLoadedSIPRequest) (*GetLoadedSIPResponse, error) + // Only a pinned, activated Dispatcher may apply the complete approved SIP snapshot. + ApplySIP(context.Context, *ApplySIPRequest) (*ApplySIPResponse, error) // Agent requests a fresh bounded token explicitly; the Dispatcher owns the original target. RequestRecordingUpload(context.Context, *RequestRecordingUploadRequest) (*RequestRecordingUploadResponse, error) ReportCallEnded(context.Context, *ReportCallEndedRequest) (*ReportCallEndedResponse, error) @@ -178,6 +193,9 @@ func (UnimplementedAgentControlServiceServer) ApplyApprovedTaskControl(context.C func (UnimplementedAgentControlServiceServer) GetLoadedSIP(context.Context, *GetLoadedSIPRequest) (*GetLoadedSIPResponse, error) { return nil, status.Errorf(codes.Unimplemented, "method GetLoadedSIP not implemented") } +func (UnimplementedAgentControlServiceServer) ApplySIP(context.Context, *ApplySIPRequest) (*ApplySIPResponse, error) { + return nil, status.Errorf(codes.Unimplemented, "method ApplySIP not implemented") +} func (UnimplementedAgentControlServiceServer) RequestRecordingUpload(context.Context, *RequestRecordingUploadRequest) (*RequestRecordingUploadResponse, error) { return nil, status.Errorf(codes.Unimplemented, "method RequestRecordingUpload not implemented") } @@ -298,6 +316,24 @@ func _AgentControlService_GetLoadedSIP_Handler(srv interface{}, ctx context.Cont return interceptor(ctx, in, info, handler) } +func _AgentControlService_ApplySIP_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(ApplySIPRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(AgentControlServiceServer).ApplySIP(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: AgentControlService_ApplySIP_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(AgentControlServiceServer).ApplySIP(ctx, req.(*ApplySIPRequest)) + } + return interceptor(ctx, in, info, handler) +} + func _AgentControlService_RequestRecordingUpload_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { in := new(RequestRecordingUploadRequest) if err := dec(in); err != nil { @@ -379,6 +415,10 @@ var AgentControlService_ServiceDesc = grpc.ServiceDesc{ MethodName: "GetLoadedSIP", Handler: _AgentControlService_GetLoadedSIP_Handler, }, + { + MethodName: "ApplySIP", + Handler: _AgentControlService_ApplySIP_Handler, + }, { MethodName: "RequestRecordingUpload", Handler: _AgentControlService_RequestRecordingUpload_Handler, diff --git a/internal/asterisk/loader.go b/internal/asterisk/loader.go new file mode 100644 index 0000000..8323c9e --- /dev/null +++ b/internal/asterisk/loader.go @@ -0,0 +1,284 @@ +package asterisk + +import ( + "bytes" + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "os" + "os/exec" + "path/filepath" + "regexp" + "sort" + "strings" + "time" + + "git.ipao.vip/rogee/go-sip/internal/configread" +) + +// Loader writes only the Agent-owned endpoint include. A reviewed static +// pjsip.conf owns the transport, whose reload can interrupt live calls. +type Loader struct { + ConfigDir string + Asterisk string + LibraryDir string +} + +type loadedTrunk struct { + ID string `json:"id"` + Host string `json:"host"` + Port int `json:"port"` +} + +type appliedState struct { + Revision int64 `json:"revision"` + Hash string `json:"hash"` + SnapshotHash string `json:"snapshot_hash"` + Trunks []loadedTrunk `json:"trunks"` +} + +var endpointLine = regexp.MustCompile(`(?m)^\s*Endpoint:\s+(\S+)`) + +func (l Loader) paths() (string, string, error) { + if l.ConfigDir == "" || l.Asterisk == "" || l.LibraryDir == "" { + return "", "", errors.New("native Asterisk executable, library and config paths are required") + } + base, err := os.ReadFile(filepath.Join(l.ConfigDir, "pjsip.conf")) + if err != nil { + return "", "", fmt.Errorf("read management-owned PJSIP base: %w", err) + } + if !bytes.Contains(base, []byte("#tryinclude go-sip-managed.conf")) || !bytes.Contains(base, []byte("[go-sip-udp]")) { + return "", "", errors.New("PJSIP base is missing the approved static transport and managed include") + } + return filepath.Join(l.ConfigDir, "go-sip-managed.conf"), filepath.Join(l.ConfigDir, "go-sip-applied.json"), nil +} + +func (l Loader) command(ctx context.Context, name string, args ...string) ([]byte, error) { + cmd := exec.CommandContext(ctx, name, args...) + cmd.Env = append(os.Environ(), "LD_LIBRARY_PATH="+l.LibraryDir) + out, err := cmd.CombinedOutput() + if err != nil { + return nil, fmt.Errorf("native Asterisk command %q failed: %w (output_sha256=%x)", filepath.Base(name), err, sha256.Sum256(out)) + } + return out, nil +} + +func (l Loader) cli(ctx context.Context, command string) ([]byte, error) { + return l.command(ctx, l.Asterisk, "-C", filepath.Join(l.ConfigDir, "asterisk.conf"), "-rx", command) +} + +func (l Loader) observe(ctx context.Context, state appliedState) error { + if state.Revision < 1 || state.SnapshotHash == "" || len(state.Trunks) == 0 { + return errors.New("no verified active SIP trunks") + } + transport, err := l.cli(ctx, "pjsip show transports") + if err != nil { + return fmt.Errorf("inspect native UDP transport: %w", err) + } + approvedTransport := false + for _, line := range strings.Split(string(transport), "\n") { + fields := strings.Fields(line) + if len(fields) >= 6 && fields[0] == "Transport:" && fields[1] == "go-sip-udp" && fields[2] == "udp" && fields[len(fields)-1] == "0.0.0.0:5060" { + approvedTransport = true + } + } + if !approvedTransport { + return errors.New("approved native UDP transport/address is not loaded") + } + all, err := l.cli(ctx, "pjsip show endpoints") + if err != nil { + return err + } + expected := make(map[string]loadedTrunk, len(state.Trunks)) + for _, t := range state.Trunks { + expected[t.ID] = t + } + seen := make(map[string]bool) + for _, match := range endpointLine.FindAllSubmatch(all, -1) { + id := string(match[1]) + if !trunkName.MatchString(id) { // skip Asterisk's heading + continue + } + if _, ok := expected[id]; !ok { + return fmt.Errorf("unapproved Asterisk endpoint %q is loaded", id) + } + seen[id] = true + } + if len(seen) != len(expected) { + return errors.New("approved SIP endpoint set is not loaded") + } + for _, t := range state.Trunks { + endpoint, err := l.cli(ctx, "pjsip show endpoint "+t.ID) + if err != nil || !bytes.Contains(endpoint, []byte(t.ID+"-aor")) || !bytes.Contains(endpoint, []byte("alaw")) || !bytes.Contains(endpoint, []byte("go-sip-udp")) || !bytes.Contains(endpoint, []byte("go-sip-no-inbound")) { + return fmt.Errorf("SIP endpoint %q did not load its approved AOR/codec", t.ID) + } + aor, err := l.cli(ctx, "pjsip show aor "+t.ID+"-aor") + if err != nil || !bytes.Contains(aor, []byte(fmt.Sprintf("sip:%s:%d", t.Host, t.Port))) { + return fmt.Errorf("SIP AOR %q did not load its approved contact", t.ID) + } + } + return nil +} + +// LoadedSIP never trusts a persisted revision alone: it also checks the file +// hash and live PJSIP objects after every Agent boot or Dispatcher query. +func (l Loader) LoadedSIP(ctx context.Context) (map[string]int64, error) { + config, statePath, err := l.paths() + if err != nil { + return nil, err + } + raw, err := os.ReadFile(statePath) + if err != nil { + return nil, fmt.Errorf("read last verified SIP revision: %w", err) + } + var state appliedState + if err := json.Unmarshal(raw, &state); err != nil { + return nil, fmt.Errorf("decode last verified SIP revision: %w", err) + } + data, err := os.ReadFile(config) + if err != nil { + return nil, err + } + hash := sha256.Sum256(data) + if hex.EncodeToString(hash[:]) != state.Hash { + return nil, errors.New("SIP config no longer matches last verified revision") + } + if err := l.observe(ctx, state); err != nil { + return nil, err + } + loaded := make(map[string]int64, len(state.Trunks)) + for _, t := range state.Trunks { + loaded[t.ID] = state.Revision + } + return loaded, nil +} + +// Apply changes no transport, routes or existing call. It records a revision +// only after the native module reload and readback succeed. +func (l Loader) Apply(ctx context.Context, approvedJSON []byte) (map[string]int64, error) { + config, statePath, err := l.paths() + if err != nil { + return nil, err + } + var sip configread.SIP + if err := json.Unmarshal(approvedJSON, &sip); err != nil { + return nil, err + } + // The same renderer is used by RPC validation and disk application. + text, err := Render(sip) + if err != nil { + return nil, err + } + var trunks []struct { + ID string `json:"trunk_id"` + Host string `json:"server_host"` + Port int `json:"server_port"` + Enabled bool `json:"enabled"` + } + if err := json.Unmarshal(sip.Trunks, &trunks); err != nil { + return nil, err + } + canonical, err := json.Marshal(sip) + if err != nil { + return nil, err + } + state := appliedState{Revision: sip.Revision, Hash: fmt.Sprintf("%x", sha256.Sum256([]byte(text))), SnapshotHash: fmt.Sprintf("%x", sha256.Sum256(canonical))} + for _, t := range trunks { + if t.Enabled { + state.Trunks = append(state.Trunks, loadedTrunk{ID: t.ID, Host: t.Host, Port: t.Port}) + } + } + sort.Slice(state.Trunks, func(i, j int) bool { return state.Trunks[i].ID < state.Trunks[j].ID }) + recoverExisting := false + if existing, err := os.ReadFile(statePath); err == nil { + var previous appliedState + if err := json.Unmarshal(existing, &previous); err != nil { + return nil, err + } + if sip.Revision < previous.Revision || (sip.Revision == previous.Revision && (state.Hash != previous.Hash || state.SnapshotHash != previous.SnapshotHash)) { + return nil, errors.New("approved SIP revision regressed or changed content") + } + if sip.Revision == previous.Revision { + return l.LoadedSIP(ctx) + } + } else if !errors.Is(err, os.ErrNotExist) { + return nil, err + } else if current, err := os.ReadFile(config); err == nil { + if sha256.Sum256(current) != sha256.Sum256([]byte(text)) { + return nil, errors.New("unverified managed SIP config does not match the approved snapshot") + } + recoverExisting = true + } else if !errors.Is(err, os.ErrNotExist) { + return nil, err + } + if !recoverExisting { + if err := writeAtomic(config, []byte(text)); err != nil { + return nil, err + } + if _, err := l.command(ctx, "systemctl", "--user", "reload", "go-sip-asterisk.service"); err != nil { + return nil, err + } + } + if err := l.waitObserved(ctx, state); err != nil { + return nil, err + } + encoded, err := json.Marshal(state) + if err != nil { + return nil, err + } + if err := writeAtomic(statePath, encoded); err != nil { + return nil, err + } + return l.LoadedSIP(ctx) +} + +func (l Loader) waitObserved(ctx context.Context, state appliedState) error { + deadline := time.NewTimer(15 * time.Second) + defer deadline.Stop() + var last error + for attempts := 1; ; attempts++ { + if last = l.observe(ctx, state); last == nil { + return nil + } + select { + case <-ctx.Done(): + return fmt.Errorf("native SIP readback cancelled after %d attempts: %w", attempts, ctx.Err()) + case <-deadline.C: + return fmt.Errorf("native SIP readback timed out after %d attempts: %w", attempts, last) + case <-time.After(250 * time.Millisecond): + } + } +} + +func writeAtomic(path string, data []byte) error { + file, err := os.CreateTemp(filepath.Dir(path), ".go-sip-stage-*") + if err != nil { + return err + } + defer os.Remove(file.Name()) + defer file.Close() + if err := file.Chmod(0600); err != nil { + return err + } + if _, err := file.Write(data); err != nil { + return err + } + if err := file.Sync(); err != nil { + return err + } + if err := file.Close(); err != nil { + return err + } + if err := os.Rename(file.Name(), path); err != nil { + return err + } + dir, err := os.Open(filepath.Dir(path)) + if err != nil { + return err + } + defer dir.Close() + return dir.Sync() +} diff --git a/internal/asterisk/loader_test.go b/internal/asterisk/loader_test.go new file mode 100644 index 0000000..cd98d00 --- /dev/null +++ b/internal/asterisk/loader_test.go @@ -0,0 +1,100 @@ +package asterisk + +import ( + "context" + "encoding/json" + "os" + "path/filepath" + "strings" + "testing" +) + +func TestLoaderPersistsOnlyVerifiedNativeReload(t *testing.T) { + dir := t.TempDir() + if err := os.WriteFile(filepath.Join(dir, "pjsip.conf"), []byte("[go-sip-udp]\ntype=transport\nprotocol=udp\nbind=0.0.0.0:5060\n#tryinclude go-sip-managed.conf\n"), 0600); err != nil { + t.Fatal(err) + } + fake := filepath.Join(dir, "asterisk") + script := `#!/bin/sh +case "$*" in + *"pjsip show transports"*) echo 'Transport: go-sip-udp udp 0 0 0.0.0.0:5060' ;; + *"pjsip show endpoints"*) echo 'Endpoint: '; echo 'Endpoint: trunk-shuqi Not in use' ;; + *"pjsip show endpoint trunk-shuqi"*) echo 'Aor: trunk-shuqi-aor allow: alaw transport: go-sip-udp context: go-sip-no-inbound' ;; + *"pjsip show aor trunk-shuqi-aor"*) echo 'Contact: sip:61.132.228.221:5060' ;; + *) exit 9 ;; +esac +` + if err := os.WriteFile(fake, []byte(script), 0700); err != nil { + t.Fatal(err) + } + bin := filepath.Join(dir, "bin") + if err := os.Mkdir(bin, 0700); err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(bin, "systemctl"), []byte("#!/bin/sh\n[ \"$*\" = '--user reload go-sip-asterisk.service' ]\n"), 0700); err != nil { + t.Fatal(err) + } + t.Setenv("PATH", bin+":"+os.Getenv("PATH")) + loader := Loader{ConfigDir: dir, Asterisk: fake, LibraryDir: dir} + sip := testSIP(t) + body, err := json.Marshal(sip) + if err != nil { + t.Fatal(err) + } + // Recover an interrupted first apply only when the approved file and + // native PJSIP objects both agree with the complete snapshot. + rendered, err := Render(sip) + if err != nil { + t.Fatal(err) + } + if err := os.WriteFile(filepath.Join(dir, "go-sip-managed.conf"), []byte(rendered), 0600); err != nil { + t.Fatal(err) + } + loaded, err := loader.Apply(context.Background(), body) + if err != nil || loaded["trunk-shuqi"] != 9 { + t.Fatalf("real command path did not verify SIP revision: %v %v", loaded, err) + } + if _, err := loader.Apply(context.Background(), body); err != nil { + t.Fatalf("same revision should observe without rewriting: %v", err) + } + if err := os.WriteFile(fake, []byte(strings.Replace(script, "Transport: go-sip-udp udp", "Transport: go-sip-udp tcp", 1)), 0700); err != nil { + t.Fatal(err) + } + if _, err := loader.LoadedSIP(context.Background()); err == nil { + t.Fatal("accepted native TCP transport as approved UDP") + } + if err := os.WriteFile(fake, []byte(script), 0700); err != nil { + t.Fatal(err) + } + changed := testSIP(t) + changed.Trunks = []byte(strings.Replace(string(changed.Trunks), "BD93205882", "BD93205883", 1)) + changedBody, err := json.Marshal(changed) + if err != nil { + t.Fatal(err) + } + if _, err := loader.Apply(context.Background(), changedBody); err == nil { + t.Fatal("same revision changed non-rendered caller authorization") + } + p := filepath.Join(dir, "go-sip-managed.conf") + if err := os.WriteFile(p, []byte("[tampered]"), 0600); err != nil { + t.Fatal(err) + } + if _, err := loader.LoadedSIP(context.Background()); err == nil || !strings.Contains(err.Error(), "no longer matches") { + t.Fatalf("tampered config was accepted: %v", err) + } +} + +func TestLoaderFailsClosedWithoutApprovedStaticTransport(t *testing.T) { + dir := t.TempDir() + loader := Loader{ConfigDir: dir, Asterisk: "/bin/true", LibraryDir: dir} + body, err := json.Marshal(testSIP(t)) + if err != nil { + t.Fatal(err) + } + if _, err := loader.Apply(context.Background(), body); err == nil { + t.Fatal("applied SIP without management-owned base transport") + } + if _, err := os.Stat(filepath.Join(dir, "go-sip-managed.conf")); !os.IsNotExist(err) { + t.Fatalf("mutated disk before static transport check: %v", err) + } +} diff --git a/internal/config/agent_runtime.go b/internal/config/agent_runtime.go index f420274..1605bbc 100644 --- a/internal/config/agent_runtime.go +++ b/internal/config/agent_runtime.go @@ -28,11 +28,14 @@ type AgentEnvironment struct { PeerFingerprints map[string]struct{} MockScenarioFile string MockAppliedSIPFile string + AsteriskConfigDir string + AsteriskBin string + AsteriskLibraryDir string } func LoadAgentEnvironment(mode string) (AgentEnvironment, error) { - if mode != "mock" { - return AgentEnvironment{}, errors.New("Agent accepts only isolated Mock mode") + if mode != "mock" && mode != "sip-only" { + return AgentEnvironment{}, errors.New("Agent accepts only isolated Mock or SIP-only mode") } get := func(name string) (string, error) { value := os.Getenv(name) @@ -54,8 +57,6 @@ func LoadAgentEnvironment(mode string) (AgentEnvironment, error) { {"DISPATCHER_GRPC_SERVER_NAME", &settings.DispatcherServerName}, {"MTLS_CA_FILE", &settings.CAFile}, {"MTLS_CERT_FILE", &settings.CertFile}, {"MTLS_KEY_FILE", &settings.KeyFile}, - {"AGENT_MOCK_SCENARIO_FILE", &settings.MockScenarioFile}, - {"AGENT_MOCK_APPLIED_SIP_FILE", &settings.MockAppliedSIPFile}, } { value, err := get(field.name) if err != nil { @@ -63,6 +64,28 @@ func LoadAgentEnvironment(mode string) (AgentEnvironment, error) { } *field.value = value } + var modeFields []struct { + name string + value *string + } + if mode == "mock" { + modeFields = []struct { + name string + value *string + }{{"AGENT_MOCK_SCENARIO_FILE", &settings.MockScenarioFile}, {"AGENT_MOCK_APPLIED_SIP_FILE", &settings.MockAppliedSIPFile}} + } else { + modeFields = []struct { + name string + value *string + }{{"ASTERISK_CONFIG_DIR", &settings.AsteriskConfigDir}, {"ASTERISK_BIN", &settings.AsteriskBin}, {"ASTERISK_LIBRARY_DIR", &settings.AsteriskLibraryDir}} + } + for _, field := range modeFields { + value, err := get(field.name) + if err != nil { + return AgentEnvironment{}, err + } + *field.value = value + } if tenant.ValidateDispatcherID(settings.DispatcherID) != nil { return AgentEnvironment{}, errors.New("DISPATCHER_ID must be a canonical UUID v4") } diff --git a/internal/config/agent_runtime_test.go b/internal/config/agent_runtime_test.go index a39a11f..278c387 100644 --- a/internal/config/agent_runtime_test.go +++ b/internal/config/agent_runtime_test.go @@ -40,6 +40,23 @@ func TestLoadAgentEnvironmentRefusesNonMockBeforeResources(t *testing.T) { } } +func TestLoadSIPOnlyAgentEnvironmentDoesNotRequireMockCalls(t *testing.T) { + setAgentEnvironment(t) + t.Setenv("AGENT_MOCK_SCENARIO_FILE", "") + t.Setenv("AGENT_MOCK_APPLIED_SIP_FILE", "") + for name, value := range map[string]string{ + "ASTERISK_CONFIG_DIR": "/tmp/asterisk-config", + "ASTERISK_BIN": "/tmp/asterisk", + "ASTERISK_LIBRARY_DIR": "/tmp/asterisk-libraries", + } { + t.Setenv(name, value) + } + settings, err := LoadAgentEnvironment("sip-only") + if err != nil || settings.AsteriskConfigDir != "/tmp/asterisk-config" || settings.MockScenarioFile != "" { + t.Fatalf("SIP-only Agent requires real native paths, not Mock fixtures: %+v %v", settings, err) + } +} + func TestLoadAgentEnvironmentRequiresExplicitDeploymentValues(t *testing.T) { for _, name := range []string{ "AGENT_ID", "CELL_ID", "AGENT_GRPC_LISTEN", "AGENT_SESSION_PATH", "AGENT_RECOVERY_ROOT", diff --git a/internal/config/dispatcher_runtime.go b/internal/config/dispatcher_runtime.go index b20d905..c050c9f 100644 --- a/internal/config/dispatcher_runtime.go +++ b/internal/config/dispatcher_runtime.go @@ -37,9 +37,7 @@ func LoadDispatcherRuntimeEnvironment(mode string) (DispatcherRuntimeEnvironment name string value *string }{ - {"DISPATCHER_GRPC_LISTEN", &settings.Listen}, {"DISPATCHER_AGENT_ENDPOINTS_FILE", &settings.AgentEndpointsFile}, - {"DISPATCHER_OSS_CONFIG_FILE", &settings.OSSConfigFile}, {"MTLS_CA_FILE", &settings.CAFile}, {"MTLS_CERT_FILE", &settings.CertFile}, {"MTLS_KEY_FILE", &settings.KeyFile}, @@ -50,6 +48,15 @@ func LoadDispatcherRuntimeEnvironment(mode string) (DispatcherRuntimeEnvironment } *field.value = value } + if mode == "sip-only" { + return settings, nil + } + if settings.Listen, err = get("DISPATCHER_GRPC_LISTEN"); err != nil { + return DispatcherRuntimeEnvironment{}, err + } + if settings.OSSConfigFile, err = get("DISPATCHER_OSS_CONFIG_FILE"); err != nil { + return DispatcherRuntimeEnvironment{}, err + } if !localGRPCAddress(settings.Listen, true) { return DispatcherRuntimeEnvironment{}, errors.New("DISPATCHER_GRPC_LISTEN must be an isolated local Mock address") } diff --git a/internal/config/dispatcher_runtime_test.go b/internal/config/dispatcher_runtime_test.go index 97befea..6388293 100644 --- a/internal/config/dispatcher_runtime_test.go +++ b/internal/config/dispatcher_runtime_test.go @@ -25,6 +25,17 @@ func setCurrentDispatcherRuntimeEnvironment(t *testing.T) string { return database } +func TestLoadSIPOnlyDispatcherRequiresNoRecordingOrCallListener(t *testing.T) { + setCurrentDispatcherRuntimeEnvironment(t) + t.Setenv("DISPATCHER_GRPC_LISTEN", "") + t.Setenv("DISPATCHER_OSS_CONFIG_FILE", "") + t.Setenv("MTLS_PEER_CERT_FINGERPRINTS", "") + settings, err := LoadDispatcherRuntimeEnvironment("sip-only") + if err != nil || settings.AgentEndpointsFile == "" || settings.OSSConfigFile != "" || settings.Listen != "" { + t.Fatalf("SIP-only process inherited Mock business listeners: %+v %v", settings, err) + } +} + func TestLoadDispatcherRuntimeEnvironmentRequiresExplicitMockDeployment(t *testing.T) { database := setCurrentDispatcherRuntimeEnvironment(t) settings, err := LoadDispatcherRuntimeEnvironment("mock") diff --git a/internal/config/runtime.go b/internal/config/runtime.go index f22a71d..39ac4ca 100644 --- a/internal/config/runtime.go +++ b/internal/config/runtime.go @@ -22,10 +22,10 @@ type DispatcherEnvironment struct { } // LoadDispatcherEnvironment is pure inspection: it opens no database, queue or -// network connection. The current executable accepts only isolated Mock mode. +// network connection. Only isolated Mock or SIP-only mode may start. func LoadDispatcherEnvironment(mode string) (DispatcherEnvironment, error) { - if mode != "mock" { - return DispatcherEnvironment{}, errors.New("Dispatcher accepts only isolated Mock mode") + if mode != "mock" && mode != "sip-only" { + return DispatcherEnvironment{}, errors.New("Dispatcher accepts only isolated Mock or SIP-only mode") } get := func(name string) (string, error) { value := os.Getenv(name) diff --git a/internal/dispatcher/approved_originator.go b/internal/dispatcher/approved_originator.go index 0f60d41..9896584 100644 --- a/internal/dispatcher/approved_originator.go +++ b/internal/dispatcher/approved_originator.go @@ -3,6 +3,7 @@ package dispatcher import ( "bytes" "context" + "crypto/sha256" "encoding/json" "errors" "fmt" @@ -18,6 +19,7 @@ import ( // transport or an MQ execution fallback. type ApprovedAgentRPC interface { GetLoadedSIP(context.Context, *agentpb.GetLoadedSIPRequest, ...grpc.CallOption) (*agentpb.GetLoadedSIPResponse, error) + ApplySIP(context.Context, *agentpb.ApplySIPRequest, ...grpc.CallOption) (*agentpb.ApplySIPResponse, error) ExecuteApproved(context.Context, *agentpb.ExecuteApprovedRequest, ...grpc.CallOption) (*agentpb.ExecuteApprovedResponse, error) ApplyApprovedTaskControl(context.Context, *agentpb.ApplyApprovedTaskControlRequest, ...grpc.CallOption) (*agentpb.ApplyApprovedTaskControlResponse, error) } @@ -68,6 +70,49 @@ func (o *ApprovedOriginator) LoadedTrunks(ctx context.Context) (map[string]int64 return loaded, nil } +// ApplySIP sends the complete approved partition through the pinned, active +// Agent session. The response must describe the exact revision and trunk set. +func (o *ApprovedOriginator) ApplySIP(ctx context.Context, sip configread.SIP) error { + if o == nil || sip.DispatcherID != o.DispatcherID || sip.Revision <= 0 { + return errors.New("approved SIP snapshot does not belong to this Dispatcher") + } + meta, err := o.activeMeta(ctx) + if err != nil { + return err + } + body, err := json.Marshal(sip) + if err != nil { + return fmt.Errorf("encode approved SIP snapshot: %w", err) + } + hash := sha256.Sum256(body) + meta.OperationId = fmt.Sprintf("sip-%d", sip.Revision) + meta.IdempotencyKey = fmt.Sprintf("sip-%d-%x", sip.Revision, hash[:8]) + response, err := o.Client.ApplySIP(ctx, &agentpb.ApplySIPRequest{Meta: meta, ApprovedSnapshotJson: body}) + if err != nil { + return fmt.Errorf("Agent SIP apply revision %d is unconfirmed: %w", sip.Revision, err) + } + var trunks []struct { + ID string `json:"trunk_id"` + Enabled bool `json:"enabled"` + } + if err := json.Unmarshal(sip.Trunks, &trunks); err != nil { + return err + } + expected := 0 + for _, trunk := range trunks { + if trunk.Enabled { + expected++ + if response == nil || response.TrunkRevision[trunk.ID] != sip.Revision { + return fmt.Errorf("Agent did not confirm approved SIP trunk %q", trunk.ID) + } + } + } + if response == nil || len(response.TrunkRevision) != expected { + return errors.New("Agent confirmed a different SIP trunk set") + } + return nil +} + // VerifySIP requires an exact match between the approved entire SIP partition // and the Agent's actual loaded set; a single matching trunk is insufficient. func (o *ApprovedOriginator) VerifySIP(ctx context.Context, sip configread.SIP) error { @@ -76,6 +121,7 @@ func (o *ApprovedOriginator) VerifySIP(ctx context.Context, sip configread.SIP) } var trunks []struct { TrunkID string `json:"trunk_id"` + Enabled bool `json:"enabled"` } if err := json.Unmarshal(sip.Trunks, &trunks); err != nil || len(trunks) == 0 { return errors.New("approved SIP trunk list is invalid") @@ -84,14 +130,19 @@ func (o *ApprovedOriginator) VerifySIP(ctx context.Context, sip configread.SIP) if err != nil { return err } - if len(loaded) != len(trunks) { - return fmt.Errorf("Agent loaded %d trunks; approved snapshot has %d", len(loaded), len(trunks)) - } + active := 0 for _, trunk := range trunks { + if !trunk.Enabled { + continue + } + active++ if trunk.TrunkID == "" || loaded[trunk.TrunkID] != sip.Revision { return fmt.Errorf("approved SIP trunk %q revision %d is not loaded", trunk.TrunkID, sip.Revision) } } + if len(loaded) != active { + return fmt.Errorf("Agent loaded %d trunks; approved snapshot enables %d", len(loaded), active) + } return nil } diff --git a/internal/dispatcher/approved_originator_test.go b/internal/dispatcher/approved_originator_test.go index 25eb5e3..a85bbfd 100644 --- a/internal/dispatcher/approved_originator_test.go +++ b/internal/dispatcher/approved_originator_test.go @@ -23,6 +23,9 @@ type fakeApprovedAgentRPC struct { controlResponse *agentpb.ApplyApprovedTaskControlResponse controlErr error controlCalls int + sipRequest *agentpb.ApplySIPRequest + sipResponse *agentpb.ApplySIPResponse + sipErr error } func (f *fakeApprovedAgentRPC) ExecuteApproved(_ context.Context, req *agentpb.ExecuteApprovedRequest, _ ...grpc.CallOption) (*agentpb.ExecuteApprovedResponse, error) { @@ -45,6 +48,17 @@ func (f *fakeApprovedAgentRPC) ApplyApprovedTaskControl(_ context.Context, req * return &agentpb.ApplyApprovedTaskControlResponse{Accepted: true}, nil } +func (f *fakeApprovedAgentRPC) ApplySIP(_ context.Context, req *agentpb.ApplySIPRequest, _ ...grpc.CallOption) (*agentpb.ApplySIPResponse, error) { + f.sipRequest = req + if f.sipErr != nil { + return nil, f.sipErr + } + if f.sipResponse != nil { + return f.sipResponse, nil + } + return &agentpb.ApplySIPResponse{TrunkRevision: f.loaded}, nil +} + func (f *fakeApprovedAgentRPC) GetLoadedSIP(_ context.Context, _ *agentpb.GetLoadedSIPRequest, _ ...grpc.CallOption) (*agentpb.GetLoadedSIPResponse, error) { if f.loadedErr != nil { return nil, f.loadedErr @@ -72,6 +86,49 @@ func approvedOriginatorFixture(t *testing.T, client *fakeApprovedAgentRPC) (*App return orig, spec } +func TestApprovedOriginatorAppliesFullSIPSnapshotBeforeReportingLoaded(t *testing.T) { + fake := &fakeApprovedAgentRPC{loaded: map[string]int64{"trunk-mock": 8}} + orig, spec := approvedOriginatorFixture(t, fake) + if err := orig.ApplySIP(context.Background(), spec.Snapshot.SIP); err != nil { + t.Fatal(err) + } + if fake.sipRequest == nil || fake.sipRequest.Meta.IdempotencyKey == "" { + t.Fatal("SIP apply omitted operation identity") + } + var sent configread.SIP + if err := json.Unmarshal(fake.sipRequest.ApprovedSnapshotJson, &sent); err != nil || sent.Revision != 8 || sent.DispatcherID != orig.DispatcherID { + t.Fatal("SIP apply lost the approved full snapshot") + } + fake.sipResponse = &agentpb.ApplySIPResponse{TrunkRevision: map[string]int64{"trunk-mock": 7}} + if err := orig.ApplySIP(context.Background(), spec.Snapshot.SIP); err == nil { + t.Fatal("accepted a stale Agent revision") + } +} + +func TestApprovedOriginatorIgnoresDisabledTrunksDuringNativeReadback(t *testing.T) { + fake := &fakeApprovedAgentRPC{loaded: map[string]int64{"trunk-mock": 8}} + orig, spec := approvedOriginatorFixture(t, fake) + var trunks []map[string]any + if err := json.Unmarshal(spec.Snapshot.SIP.Trunks, &trunks); err != nil { + t.Fatal(err) + } + disabled := make(map[string]any) + for key, value := range trunks[0] { + disabled[key] = value + } + disabled["trunk_id"] = "disabled-line" + disabled["enabled"] = false + trunks = append(trunks, disabled) + data, err := json.Marshal(trunks) + if err != nil { + t.Fatal(err) + } + spec.Snapshot.SIP.Trunks = data + if err := orig.VerifySIP(context.Background(), spec.Snapshot.SIP); err != nil { + t.Fatal(err) + } +} + func TestApprovedOriginatorSendsExactImmutableAgentInstruction(t *testing.T) { fake := &fakeApprovedAgentRPC{loaded: map[string]int64{"trunk-mock": 8}} orig, spec := approvedOriginatorFixture(t, fake) diff --git a/internal/dispatcher/config.go b/internal/dispatcher/config.go index 0f8ec2b..05bd1a9 100644 --- a/internal/dispatcher/config.go +++ b/internal/dispatcher/config.go @@ -16,6 +16,7 @@ type Bootstrap struct { Client *configread.Client Store *store.Store DispatcherID string + ApplySIP func(context.Context, configread.SIP) error // nil for isolated Mock; real SIP-only must apply before verifying VerifySIP func(context.Context, configread.SIP) error DrainControls func(context.Context) error Cursor *string // memory-only; restart always starts from a full snapshot @@ -46,6 +47,11 @@ func (b Bootstrap) Run(ctx context.Context) error { if sip.DispatcherID != b.DispatcherID { return errors.New("SIP snapshot belongs to another Dispatcher") } + if b.ApplySIP != nil { + if err := b.ApplySIP(ctx, sip); err != nil { + return fmt.Errorf("apply approved SIP revision %d on Agent/Asterisk: %w", sip.Revision, err) + } + } if err := b.VerifySIP(ctx, sip); err != nil { return fmt.Errorf("verify SIP revision %d loaded by Agent/Asterisk: %w", sip.Revision, err) } diff --git a/internal/dispatcher/config_test.go b/internal/dispatcher/config_test.go index cb32f9f..e9bd770 100644 --- a/internal/dispatcher/config_test.go +++ b/internal/dispatcher/config_test.go @@ -61,13 +61,20 @@ func TestBootstrapRequiresSIPLoadingAndControlDrain(t *testing.T) { t.Fatal(err) } defer db.Close() - verifierCalled, drained := false, false + applied, verifierCalled, drained := false, false, false var cursor string var approvedSIP configread.SIP bootstrap := Bootstrap{ Client: client, Store: db, DispatcherID: id, Cursor: &cursor, SIP: &approvedSIP, + ApplySIP: func(_ context.Context, sip configread.SIP) error { + applied = sip.Revision == 8 + return nil + }, VerifySIP: func(_ context.Context, sip configread.SIP) error { verifierCalled = true + if !applied { + t.Fatal("verified Agent load before applying SIP") + } if sip.Revision != 8 { t.Fatalf("verified wrong SIP revision %d", sip.Revision) } @@ -102,7 +109,7 @@ func TestBootstrapRequiresSIPLoadingAndControlDrain(t *testing.T) { func TestBootstrapFailsClosedOnVerifierOrDrainError(t *testing.T) { id := "c046b893-8628-4589-ae50-619d049248a6" - for _, failure := range []string{"verify", "drain"} { + for _, failure := range []string{"apply", "verify", "drain"} { t.Run(failure, func(t *testing.T) { server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { var response []byte @@ -139,6 +146,12 @@ func TestBootstrapFailsClosedOnVerifierOrDrainError(t *testing.T) { defer db.Close() b := Bootstrap{ Client: client, Store: db, DispatcherID: id, + ApplySIP: func(context.Context, configread.SIP) error { + if failure == "apply" { + return errors.New("native SIP application rejected") + } + return nil + }, VerifySIP: func(context.Context, configread.SIP) error { if failure == "verify" { return errors.New("applied SIP revision unavailable") diff --git a/internal/dispatcher/sip_only.go b/internal/dispatcher/sip_only.go new file mode 100644 index 0000000..a57d130 --- /dev/null +++ b/internal/dispatcher/sip_only.go @@ -0,0 +1,83 @@ +package dispatcher + +import ( + "context" + "encoding/json" + "errors" + "fmt" + + "git.ipao.vip/rogee/go-sip/internal/configread" + "git.ipao.vip/rogee/go-sip/internal/contract" + "git.ipao.vip/rogee/go-sip/internal/store" +) + +// SIPOnly is a distinct update lane: it never discovers tasks or enables +// call admission, even when native Asterisk confirms the entire snapshot. +type SIPOnly struct { + DispatcherID string + Store *store.Store + Client *configread.Client + ApplySIP func(context.Context, configread.SIP) error + VerifySIP func(context.Context, configread.SIP) error +} + +// HandleNotification persists the assigned SIP revision before the broker ACK. +// Non-SIP controls must be requeued for a business Dispatcher, never swallowed. +func (s SIPOnly) HandleNotification(_ context.Context, route string, body []byte) error { + if s.Store == nil || s.DispatcherID == "" || route != "d."+s.DispatcherID+".control.in" { + return errors.New("SIP-only notification has no approved queue owner") + } + if err := contract.ValidateCurrent("mq", body); err != nil { + return fmt.Errorf("SIP-only notification violates current contract: %w", err) + } + var notice struct { + EventType string `json:"event_type"` + DispatcherID string `json:"dispatcher_id"` + Payload struct { + Revision int64 `json:"revision"` + } `json:"payload"` + } + if err := json.Unmarshal(body, ¬ice); err != nil { + return err + } + if notice.DispatcherID != s.DispatcherID || notice.EventType != "sip.config" || notice.Payload.Revision <= 0 { + return errors.New("SIP-only cannot consume another Dispatcher or business control") + } + return s.Store.NoteSIPChange(s.DispatcherID, notice.Payload.Revision) +} + +// Sync fetches only the complete SIP snapshot; it never reads task, quota or +// AI config. The Agent revision and actual native load must both agree before +// the durable pending revision is cleared, while admission remains closed. +func (s SIPOnly) Sync(ctx context.Context) error { + if s.Store == nil || s.Client == nil || s.DispatcherID == "" || s.ApplySIP == nil || s.VerifySIP == nil { + return errors.New("SIP-only sync requires a bound client, durable store and native Agent") + } + sip, err := s.Client.ReadSIP(ctx) + if err != nil { + return fmt.Errorf("read approved SIP-only snapshot: %w", err) + } + if sip.DispatcherID != s.DispatcherID || sip.Revision <= 0 { + return errors.New("SIP-only snapshot owner or revision is invalid") + } + applied, pending, err := s.Store.SIPState(s.DispatcherID) + if err != nil { + return err + } + if sip.Revision < pending || sip.Revision < applied { + return fmt.Errorf("approved SIP snapshot revision %d is behind durable applied=%d pending=%d", sip.Revision, applied, pending) + } + if err := s.Store.NoteSIPChange(s.DispatcherID, sip.Revision); err != nil { + return err + } + if err := s.ApplySIP(ctx, sip); err != nil { + return fmt.Errorf("native SIP apply revision %d: %w", sip.Revision, err) + } + if err := s.VerifySIP(ctx, sip); err != nil { + return fmt.Errorf("native SIP load revision %d: %w", sip.Revision, err) + } + if err := s.Store.MarkSIPOnlyVerified(s.DispatcherID, sip.Revision); err != nil { + return err + } + return nil +} diff --git a/internal/dispatcher/sip_only_test.go b/internal/dispatcher/sip_only_test.go new file mode 100644 index 0000000..e547d87 --- /dev/null +++ b/internal/dispatcher/sip_only_test.go @@ -0,0 +1,67 @@ +package dispatcher + +import ( + "context" + "net/http" + "net/http/httptest" + "path/filepath" + "testing" + + "git.ipao.vip/rogee/go-sip/internal/configread" + "git.ipao.vip/rogee/go-sip/internal/store" +) + +func TestSIPOnlyNotificationCheckpointsNativeLoadWithoutOpeningBusiness(t *testing.T) { + const id = "c046b893-8628-4589-ae50-619d049248a6" + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/internal/v1/dispatcher/sip" { + t.Errorf("SIP-only fetched business config %s", r.URL.Path) + w.WriteHeader(404) + return + } + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write(configExample(t, "config-read-sip")) + })) + defer server.Close() + client, err := configread.NewClient(server.URL, id, "test-secret", server.Client()) + if err != nil { + t.Fatal(err) + } + db, err := store.Open(filepath.Join(t.TempDir(), "sip-only.db")) + if err != nil { + t.Fatal(err) + } + defer db.Close() + applied := false + syncer := SIPOnly{DispatcherID: id, Store: db, Client: client, + ApplySIP: func(_ context.Context, sip configread.SIP) error { + applied = true + if sip.Revision != 8 { + t.Fatal("unexpected SIP revision") + } + return nil + }, + VerifySIP: func(_ context.Context, sip configread.SIP) error { + if !applied { + t.Fatal("verified before native apply") + } + return nil + }, + } + body := configExample(t, "mq-sip-change") + if err := syncer.HandleNotification(context.Background(), "d."+id+".control.in", body); err != nil { + t.Fatal(err) + } + if appliedRev, pending, err := db.SIPState(id); err != nil || appliedRev != 0 || pending != 8 { + t.Fatalf("notification was not durably fenced: %d %d %v", appliedRev, pending, err) + } + if err := syncer.Sync(context.Background()); err != nil { + t.Fatal(err) + } + if appliedRev, pending, err := db.SIPState(id); err != nil || appliedRev != 8 || pending != 0 { + t.Fatalf("verified readback was not recorded: %d %d %v", appliedRev, pending, err) + } + if err := syncer.HandleNotification(context.Background(), "d."+id+".control.in", configExample(t, "mq-control")); err == nil { + t.Fatal("SIP-only consumed a business control event") + } +} diff --git a/internal/dispatcher/sip_reload.go b/internal/dispatcher/sip_reload.go index 6e27ad6..ec92b02 100644 --- a/internal/dispatcher/sip_reload.go +++ b/internal/dispatcher/sip_reload.go @@ -54,6 +54,12 @@ func (r *Runtime) refreshSIP(ctx context.Context, approvedSIP *configread.SIP) e r.Logger.Debug("SaaS SIP snapshot is behind durable SIP notification; admission stays closed", "dispatcher_id", r.Bootstrap.DispatcherID, "pending_revision", pending, "available_revision", current.Revision) return nil } + if r.Bootstrap.ApplySIP != nil { + if err := r.Bootstrap.ApplySIP(ctx, current); err != nil { + r.Logger.Warn("approved SIP update remains fenced until native reload confirms", "dispatcher_id", r.Bootstrap.DispatcherID, "revision", current.Revision, "error", err) + return nil + } + } if err := r.Bootstrap.VerifySIP(ctx, current); err != nil { r.Logger.Debug("SIP reload waits for actual Agent/Asterisk revision", "dispatcher_id", r.Bootstrap.DispatcherID, "pending_revision", pending, "error", err) return nil diff --git a/internal/dispatcher/sip_runtime_test.go b/internal/dispatcher/sip_runtime_test.go index b13f636..c726570 100644 --- a/internal/dispatcher/sip_runtime_test.go +++ b/internal/dispatcher/sip_runtime_test.go @@ -42,13 +42,27 @@ func TestSIPNotificationPersistsBarrierAndWaitsForLoadedFullSnapshot(t *testing. t.Fatal(err) } var loaded atomic.Bool + var applied atomic.Int32 verify := func(_ context.Context, sip configread.SIP) error { + if applied.Load() == 0 { + return errors.New("SIP load checked before native apply") + } if sip.Revision != 9 || !loaded.Load() { return errors.New("Agent and Asterisk have not applied SIP revision 9") } return nil } - runtime := &Runtime{Bootstrap: Bootstrap{DispatcherID: executor.DispatcherID, Client: client, Store: s, VerifySIP: verify}, Logger: slog.Default()} + runtime := &Runtime{Bootstrap: Bootstrap{DispatcherID: executor.DispatcherID, Client: client, Store: s, + ApplySIP: func(_ context.Context, sip configread.SIP) error { + if sip.Revision != 9 { + return errors.New("wrong SIP revision sent to Agent") + } + applied.Add(1) + if !loaded.Load() { + return errors.New("native Asterisk SIP reload is still unavailable") + } + return nil + }, VerifySIP: verify}, Logger: slog.Default()} body := []byte(strings.Replace(string(configExample(t, "mq-sip-change")), `"revision":8`, `"revision":9`, 1)) if err := runtime.handleControl(context.Background(), "", body); err != nil { @@ -66,6 +80,9 @@ func TestSIPNotificationPersistsBarrierAndWaitsForLoadedFullSnapshot(t *testing. if admitted, err := s.CanAdmit(executor.DispatcherID, 1001, "task-asr"); err != nil || admitted { t.Fatalf("unloaded SIP reopened admission: %v %v", admitted, err) } + if applied.Load() == 0 { + t.Fatal("notification never sent to Agent") + } loaded.Store(true) if err := runtime.refreshSIP(context.Background(), &approved); err != nil { t.Fatal(err) diff --git a/internal/rpc/server.go b/internal/rpc/server.go index 1febc85..4392d57 100644 --- a/internal/rpc/server.go +++ b/internal/rpc/server.go @@ -32,8 +32,10 @@ type ServerOptions struct { PeerAgentIDs map[string]string PeerCertificateFingerprints map[string]struct{} StatePath string - // LoadedSIP reports the revision the mock Agent actually loaded; nil fails closed. + // LoadedSIP reports the revision the Agent actually observed in Asterisk; nil fails closed. LoadedSIP func(context.Context) (map[string]int64, error) + // ApplySIP writes and reloads only a validated, Dispatcher-approved full snapshot. + ApplySIP func(context.Context, []byte) (map[string]int64, error) // The mock may have issued a call even if its outcome is unknown. MockApprovedOriginate func(context.Context, ApprovedExecution) error // ApprovedTaskCalls is shared with the approved call runner; nil rejects task controls. @@ -53,6 +55,8 @@ type Server struct { peerAgentIDs map[string]string peerCertificateFingerprints map[string]struct{} loadedSIP func(context.Context) (map[string]int64, error) + applySIP func(context.Context, []byte) (map[string]int64, error) + sipMu sync.Mutex mockApprovedOriginate func(context.Context, ApprovedExecution) error approvedTaskCalls *agent.TaskCalls sessions *SessionRegistry @@ -87,6 +91,7 @@ func NewServer(options ServerOptions) *Server { peerAgentIDs: cloneStringMap(options.PeerAgentIDs), peerCertificateFingerprints: cloneSet(options.PeerCertificateFingerprints), loadedSIP: options.LoadedSIP, + applySIP: options.ApplySIP, mockApprovedOriginate: options.MockApprovedOriginate, approvedTaskCalls: options.ApprovedTaskCalls, sessions: NewSessionRegistry(options.StatePath), diff --git a/internal/rpc/service_test.go b/internal/rpc/service_test.go index ac32d07..32370cd 100644 --- a/internal/rpc/service_test.go +++ b/internal/rpc/service_test.go @@ -15,7 +15,7 @@ func TestAgentControlServiceOnlyExposesApprovedMethods(t *testing.T) { } want := []string{ "GetAgentStatus", "ActivateAgent", - "ExecuteApproved", "ApplyApprovedTaskControl", "GetLoadedSIP", + "ExecuteApproved", "ApplyApprovedTaskControl", "GetLoadedSIP", "ApplySIP", "RequestRecordingUpload", "ReportCallEnded", "ReportCallResult", } if !reflect.DeepEqual(got, want) { diff --git a/internal/rpc/sip_apply.go b/internal/rpc/sip_apply.go new file mode 100644 index 0000000..e05ab8e --- /dev/null +++ b/internal/rpc/sip_apply.go @@ -0,0 +1,70 @@ +package rpc + +import ( + "context" + "encoding/json" + "log/slog" + + agentpb "git.ipao.vip/rogee/go-sip/gen/agent" + "git.ipao.vip/rogee/go-sip/internal/asterisk" + "git.ipao.vip/rogee/go-sip/internal/configread" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" +) + +// ApplySIP is deliberately unavailable on the Mock call server. The real +// SIP-only Agent accepts a complete Dispatcher-owned snapshot, then reports +// only revisions it actually applied and inspected in native Asterisk. +func (s *Server) ApplySIP(ctx context.Context, req *agentpb.ApplySIPRequest) (*agentpb.ApplySIPResponse, error) { + if req == nil || req.Meta == nil || len(req.ApprovedSnapshotJson) == 0 || len(req.ApprovedSnapshotJson) > 1<<20 { + return nil, status.Error(codes.InvalidArgument, "approved SIP request or bounded snapshot is required") + } + if err := s.authorize(ctx, req.Meta); err != nil { + return nil, err + } + dispatcherID, err := s.sessions.ApprovedDispatcher(req.Meta, s.now()) + if err != nil { + return nil, err + } + if s.mode != "sip-only" || s.applySIP == nil { + return nil, status.Error(codes.FailedPrecondition, "real SIP apply is unavailable on this Agent") + } + var approved configread.SIP + if err := json.Unmarshal(req.ApprovedSnapshotJson, &approved); err != nil || approved.DispatcherID != dispatcherID || approved.Revision <= 0 { + return nil, status.Error(codes.InvalidArgument, "SIP snapshot owner or content is invalid") + } + if _, err := asterisk.Render(approved); err != nil { + return nil, status.Error(codes.InvalidArgument, "SIP snapshot is not approved for native reload") + } + var trunks []struct { + ID string `json:"trunk_id"` + Enabled bool `json:"enabled"` + } + if err := json.Unmarshal(approved.Trunks, &trunks); err != nil { + return nil, status.Error(codes.InvalidArgument, "SIP trunk list is invalid") + } + expected := make(map[string]bool) + for _, trunk := range trunks { + if trunk.Enabled { + expected[trunk.ID] = true + } + } + s.sipMu.Lock() + defer s.sipMu.Unlock() + loaded, err := s.applySIP(ctx, req.ApprovedSnapshotJson) + if err != nil { + slog.Error("native SIP apply failed", "dispatcher_id", dispatcherID, "revision", approved.Revision, "error", err) + return nil, status.Error(codes.Unavailable, "native Asterisk SIP reload failed") + } + if len(loaded) != len(expected) { + slog.Error("native SIP loaded-set mismatch", "dispatcher_id", dispatcherID, "revision", approved.Revision, "expected_count", len(expected), "loaded_count", len(loaded)) + return nil, status.Error(codes.Unavailable, "native Asterisk SIP loaded-set mismatch") + } + for trunk, revision := range loaded { + if !expected[trunk] || revision != approved.Revision { + slog.Error("native SIP loaded-revision mismatch", "dispatcher_id", dispatcherID, "revision", approved.Revision, "trunk_id", trunk, "observed_revision", revision) + return nil, status.Error(codes.Unavailable, "native Asterisk SIP loaded-revision mismatch") + } + } + return &agentpb.ApplySIPResponse{TrunkRevision: loaded}, nil +} diff --git a/internal/rpc/sip_apply_test.go b/internal/rpc/sip_apply_test.go new file mode 100644 index 0000000..9bf2a01 --- /dev/null +++ b/internal/rpc/sip_apply_test.go @@ -0,0 +1,55 @@ +package rpc + +import ( + "context" + "os" + "strings" + "testing" + "time" + + agentpb "git.ipao.vip/rogee/go-sip/gen/agent" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/status" +) + +func TestApplySIPOnlyAllowsActivatedSIPService(t *testing.T) { + now := time.Date(2026, 10, 3, 9, 0, 0, 0, time.UTC) + const id = "c046b893-8628-4589-ae50-619d049248a6" + attempts := 0 + apply := func(_ context.Context, _ []byte) (map[string]int64, error) { + attempts++ + return map[string]int64{"trunk-mock": 8}, nil + } + server := NewServer(ServerOptions{ + Mode: "sip-only", Now: func() time.Time { return now }, + Status: &agentpb.AgentStatus{AgentId: "agent-1", CellId: "cell-1", BootId: "boot-1"}, + ApplySIP: apply, + }) + _, err := server.ActivateAgent(context.Background(), &agentpb.ActivateAgentRequest{ + Meta: testMeta("activate-sip", "", 0), + Binding: &agentpb.AgentBinding{AgentId: "agent-1", CellId: "cell-1", ExpectedBootId: "boot-1", DispatcherEpoch: "epoch-1", SessionGeneration: 1, DispatcherId: id}, + ActivationOperationId: "activate-sip", + }) + if err != nil { + t.Fatal(err) + } + body, err := os.ReadFile("../../contracts/local/examples/config-read-sip.json") + if err != nil { + t.Fatal(err) + } + body = []byte(strings.ReplaceAll(string(body), `"transport":null`, `"transport":"udp"`)) + body = []byte(strings.ReplaceAll(string(body), `"auth_mode":null`, `"auth_mode":"ip"`)) + body = []byte(strings.ReplaceAll(string(body), `"registration_required":null`, `"registration_required":false`)) + body = []byte(strings.ReplaceAll(string(body), `"server_host":"sip.example.invalid"`, `"server_host":"127.0.0.1"`)) + req := &agentpb.ApplySIPRequest{Meta: testMeta("apply-sip", "apply-sip", 1), ApprovedSnapshotJson: body} + if _, err := server.ApplySIP(context.Background(), req); err != nil { + t.Fatal(err) + } + if attempts != 1 { + t.Fatalf("apply attempts = %d", attempts) + } + server.mode = "mock" + if _, err := server.ApplySIP(context.Background(), req); status.Code(err) != codes.FailedPrecondition || attempts != 1 { + t.Fatalf("mock boundary allowed real SIP apply: attempts=%d err=%v", attempts, err) + } +} diff --git a/internal/store/sip.go b/internal/store/sip.go index 4336ced..e34fd0d 100644 --- a/internal/store/sip.go +++ b/internal/store/sip.go @@ -45,6 +45,28 @@ func (s *Store) NoteSIPChange(dispatcherID string, revision int64) error { return nil } +// MarkSIPOnlyVerified checkpoints native Asterisk readback without ever +// enabling business call admission. A newer durable notification wins races. +func (s *Store) MarkSIPOnlyVerified(dispatcherID string, revision int64) error { + if dispatcherID == "" || revision <= 0 { + return errors.New("SIP-only checkpoint requires identity and positive verified revision") + } + result, err := s.db.Exec(`UPDATE dispatcher_state + SET applied_sip_revision=?,pending_sip_revision=0,discovery_ready=0 + WHERE dispatcher_id=? AND applied_sip_revision<=? AND pending_sip_revision<=?`, revision, dispatcherID, revision, revision) + if err != nil { + return fmt.Errorf("checkpoint verified native SIP: %w", err) + } + count, err := result.RowsAffected() + if err != nil { + return err + } + if count != 1 { + return errors.New("SIP-only checkpoint conflicts with a newer durable notification or loaded revision") + } + return nil +} + func (s *Store) SIPState(dispatcherID string) (applied, pending int64, err error) { err = s.db.QueryRow(`SELECT applied_sip_revision,pending_sip_revision FROM dispatcher_state WHERE dispatcher_id=?`, dispatcherID).Scan(&applied, &pending) if err != nil { diff --git a/internal/store/sip_change_test.go b/internal/store/sip_change_test.go index 4e666de..8adca99 100644 --- a/internal/store/sip_change_test.go +++ b/internal/store/sip_change_test.go @@ -6,6 +6,33 @@ import ( "testing" ) +func TestSIPOnlyVerifiedNeverOpensCallAdmission(t *testing.T) { + s, err := Open(filepath.Join(t.TempDir(), "sip-only.db")) + if err != nil { + t.Fatal(err) + } + defer s.Close() + if err := s.NoteSIPChange(currentDispatcherID, 9); err != nil { + t.Fatal(err) + } + if err := s.MarkSIPOnlyVerified(currentDispatcherID, 9); err != nil { + t.Fatal(err) + } + if applied, pending, err := s.SIPState(currentDispatcherID); err != nil || applied != 9 || pending != 0 { + t.Fatalf("SIP-only checkpoint: %d %d %v", applied, pending, err) + } + var admitted int + if err := s.db.QueryRow(`SELECT discovery_ready FROM dispatcher_state WHERE dispatcher_id=?`, currentDispatcherID).Scan(&admitted); err != nil || admitted != 0 { + t.Fatalf("SIP-only opened call admission: %d %v", admitted, err) + } + if err := s.NoteSIPChange(currentDispatcherID, 10); err != nil { + t.Fatal(err) + } + if err := s.MarkSIPOnlyVerified(currentDispatcherID, 9); err == nil { + t.Fatal("stale native SIP confirmation discarded newer notification") + } +} + func TestSIPNotificationRequiresFullDrainAndExactLoadedRevision(t *testing.T) { s := preparedCurrentCallStore(t) if err := s.MarkReadyForSIP(currentDispatcherID, 8); err != nil { diff --git a/proto/agent/agent.proto b/proto/agent/agent.proto index edf25cd..fc4e8fa 100644 --- a/proto/agent/agent.proto +++ b/proto/agent/agent.proto @@ -13,6 +13,8 @@ service AgentControlService { rpc ExecuteApproved(ExecuteApprovedRequest) returns (ExecuteApprovedResponse); rpc ApplyApprovedTaskControl(ApplyApprovedTaskControlRequest) returns (ApplyApprovedTaskControlResponse); rpc GetLoadedSIP(GetLoadedSIPRequest) returns (GetLoadedSIPResponse); + // Only a pinned, activated Dispatcher may apply the complete approved SIP snapshot. + rpc ApplySIP(ApplySIPRequest) returns (ApplySIPResponse); // Agent requests a fresh bounded token explicitly; the Dispatcher owns the original target. rpc RequestRecordingUpload(RequestRecordingUploadRequest) returns (RequestRecordingUploadResponse); rpc ReportCallEnded(ReportCallEndedRequest) returns (ReportCallEndedResponse); @@ -265,6 +267,15 @@ message GetLoadedSIPResponse { map trunk_revision = 1; } +message ApplySIPRequest { + RequestMeta meta = 1; + bytes approved_snapshot_json = 2; +} + +message ApplySIPResponse { + map trunk_revision = 1; +} + // Current task-level control has no external command ID, revision CAS or // execution binding. The authenticated session must match dispatcher_id. message ApplyApprovedTaskControlRequest { diff --git a/proto/manifest.json b/proto/manifest.json index 36e4c8f..5472c1d 100644 --- a/proto/manifest.json +++ b/proto/manifest.json @@ -34,18 +34,18 @@ }, { "path": "proto/agent/agent.proto", - "bytes": 8626, - "sha256": "9a76ca785c6466b9d5eb217a9c28cda80bb738f308210df66ad6818c2e74805c" + "bytes": 8933, + "sha256": "f63f3477bceaea4c41ed47d86191f30a4dce5c0e2767b90834fbf93b557fb29b" }, { "path": "gen/agent/agent.pb.go", - "bytes": 104492, - "sha256": "e1deb02528c51837cdc0c66c4344b4e09704e7db096de8b3fe527654f98e4647" + "bytes": 108405, + "sha256": "f8d735ab7c29975bce8207b507167c31e10053733de86cae47f917472ab4b480" }, { "path": "gen/agent/agent_grpc.pb.go", - "bytes": 18310, - "sha256": "8aba17a4840f41d6766e21da75f2cb994613c596014e346a75e46006d15115d5" + "bytes": 20100, + "sha256": "4e3960698ab2d361753f38a7d06c7971fe8d4d043fb8a4806a324c85dc09c4f1" } ] }