diff --git a/AGENTS.md b/AGENTS.md index 5c869ce..ed32e6d 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -80,8 +80,8 @@ - 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 启动即拒绝。另有严格隔离的 `--mode sip-only`:只允许 Dispatcher 读完整 SIP、持久接纳归属 `sip.config`、经已激活双向 TLS Agent 会话将完整快照应用到原生 Asterisk,并从运行态核对版本;不发现任务、不启动业务呼叫、不开放准入、不处理其他业务控制。测试机已凭使用者批准,将三条历史登记地址以**测试快照显式声明的** UDP/IP 鉴权、无需 REGISTER 配置写入并核对 Asterisk 运行态;供应商尚未确认这些实际线路是否满足上述参数,无拨号,不证明线路可用或生产签收。没有真实 SaaS、management、真实通话或生产签收;真实百炼 LLM 与 OSS 已各做一次**不拨号、独立的最小连接/写入诊断**(见 [`docs/evidence/real-ai-oss-one-shot-20261003.md`](docs/evidence/real-ai-oss-one-shot-20261003.md)),这不是获批任务的 ASR/LLM/TTS、录音和上传全链路签收。生产发布包仍为 `production_approval=false`。任何本机 Mock 或 SIP-only 核验均不授权真实呼叫。 -- 使用者批准独立的测试专用 `deploys/test/saas-mock/` 仅模拟缺失的 SaaS:从 0600 的静态快照提供五类正式 HTTP 配置,在专用 RabbitMQ vhost 中由模拟 SaaS 预建现行拓扑;不在 Dispatcher 内注入快照,不发布外呼消息,不替代 Agent/Asterisk/AI/OSS,也不解除当前 Mock-only 呼叫屏障。只有完成真实链路、线路鉴权和拨号前诊断,才可另行安排逐次试拨。此模拟不构成真实 SaaS 签收。 +- 当前业务范围仍仅**单节点、单 Dispatcher、单 Agent、单 Cell、单租户**。根命令只接受显式 `agent`/`dispatcher`;`mixed` 和裸 `real` 仍拒绝;隔离 `mock`、只读 `sip-only` 和需完整非生产证据/正式获批指令的 `nonprod-real` 为彼此独立的入口。新增非生产真实入口尚未在测试机完成全链路核验或拨号;不得以编译/本机 Mock 通过宣称可用。另有严格隔离的 `--mode sip-only`:只允许 Dispatcher 读完整 SIP、持久接纳归属 `sip.config`、经已激活双向 TLS Agent 会话将完整快照应用到原生 Asterisk,并从运行态核对版本;不发现任务、不启动业务呼叫、不开放准入、不处理其他业务控制。测试机已凭使用者批准,将三条历史登记地址以**测试快照显式声明的** UDP/IP 鉴权、无需 REGISTER 配置写入并核对 Asterisk 运行态;供应商尚未确认这些实际线路是否满足上述参数,无拨号,不证明线路可用或生产签收。没有真实 SaaS、management、真实通话或生产签收;真实百炼 LLM 与 OSS 已各做一次**不拨号、独立的最小连接/写入诊断**(见 [`docs/evidence/real-ai-oss-one-shot-20261003.md`](docs/evidence/real-ai-oss-one-shot-20261003.md)),这不是获批任务的 ASR/LLM/TTS、录音和上传全链路签收。生产发布包仍为 `production_approval=false`。任何本机 Mock 或 SIP-only 核验均不授权真实呼叫。 +- 使用者批准独立的测试专用 `deploys/test/saas-mock/` 仅模拟缺失的 SaaS:从 0600 的静态快照提供五类正式 HTTP 配置,在专用 RabbitMQ vhost 中由模拟 SaaS 预建现行拓扑;不在 Dispatcher 内注入快照;目前不发布外呼消息,也不替代 Agent/Asterisk/AI/OSS。测试用正式入口必须同时具备只读配置、专用 MQ、真实 Agent/Asterisk/AI/录音/OSS、逐通确认和拨号前活跃抓包证据,缺一项不得试拨。此模拟不构成真实 SaaS 签收。 - 开发按 TDD 分批,小步提交;不得覆盖使用者现存修改/未跟踪文件,不自动清理、迁移或覆盖任何现存 SQLite、spool、outbox 和 Agent 恢复文件。旧 `.executions` 及恢复根目录中旧 `.uploads`、`.upload-locks`、逐执行 `state.json` 的发现须只读失败关闭,现存未交付事实由使用者确认处置。真实云账号、EIP、线路、拨号、生产部署和共享数据操作分别需要明确授权。 ## SaaS、Dispatcher 与 Agent 的现行边界 @@ -105,7 +105,7 @@ | 中鼎 | `60.171.24.90:5060` | `mbkq` | 无 | | 百应 | `160.202.254.79:5060` | `KQ91526` | `mka755` | -- 全线路**只允许原始号码** `15003164745`、`15830461047`,但白名单和上述登记绝非真实拨号授权。已授权的真实/旧路径仅在 Asia/Shanghai 每日 `09:00`(含)至 `20:00`(不含)放行,每条 trunk 对每个原始号码每天最多 3 次,窗口外直接拒绝,不等候/自动延迟/自动重试/静默换线。每次真实试拨仍须使用者明确安排并受现行程序 Mock-only 启动屏障约束;不能拿 Mock 时段测试宣称 real 放行。 +- 全线路**只允许原始号码** `15003164745`、`15830461047`,但白名单和上述登记绝非真实拨号授权。已授权的真实/旧路径仅在 Asia/Shanghai 每日 `09:00`(含)至 `20:00`(不含)放行,每条 trunk 对每个原始号码每天最多 3 次,窗口外直接拒绝,不等候/自动延迟/自动重试/静默换线。每次真实试拨仍须使用者明确安排,并由专用主机脚本在拨号前启用抓包和 PJSIP logger、签发与该通 `event_id`/trunk/原始号码绑定的短时有效活跃抓包凭证;Agent 拒绝缺失/失效/不匹配的凭证。不能拿 Mock 时段测试宣称真实放行。 - 任务按周一至周日多个时段与指定排除日期配置,线路只有每周允许时段(**没有线路排除日期**),Asia/Shanghai 左闭右开、跨日拆分;缺失或不确定 fail-closed,不自动重拨。由 Dispatcher 在持久接纳与实际发出指令前判定,并取任务/获批 AI 较小通话时限;Agent 仅校验会话和签发期限,不重算外呼策略。本地策略 Mock 与固定真实门禁必须分别报告。 - `BD` 等主叫原值不得清洗或当作 Digest 用户名;业务原始被叫号码不变,仅被选定数企 trunk 按规则构造 `7089<原号>`(其它线路使用自己的前缀),不重复加前缀。三条 trunk 独立,不能把同地址伪造为备用线路或换线重拨;服务商反馈 PCMA,对应 Asterisk `allow=alaw`,传输/注册/鉴权/并发仍待真实签收。不以 sipgo/diago 另造 Asterisk 替代架构。 diff --git a/cmd/sip-go-agent/agent_command.go b/cmd/sip-go-agent/agent_command.go index 0afabde..0d01c03 100644 --- a/cmd/sip-go-agent/agent_command.go +++ b/cmd/sip-go-agent/agent_command.go @@ -22,7 +22,7 @@ import ( func newAgentCommand() *cobra.Command { var mode string command := &cobra.Command{ - Use: "agent", Short: "Run the isolated Agent", Args: cobra.NoArgs, + Use: "agent", Short: "Run the bound Agent", Args: cobra.NoArgs, SilenceUsage: true, RunE: func(cmd *cobra.Command, _ []string) error { settings, err := config.LoadAgentEnvironment(mode) @@ -50,14 +50,16 @@ func newAgentCommand() *cobra.Command { 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 + if mode == "mock" { + scenario, err = loadApprovedMockScenario(settings.MockScenarioFile) + if err != nil { + return err + } + applied, err = loadMockAppliedSIP(settings.MockAppliedSIPFile) + if err != nil { + return err + } } var clientTLS *tls.Config clientTLS, err = rpc.NewClientTLSConfig(ca, cert, key, settings.DispatcherServerName) @@ -69,7 +71,12 @@ func newAgentCommand() *cobra.Command { return errors.New("Agent cannot create pinned Dispatcher connection") } defer connection.Close() - handler, err = newAgentServer(cmd.Context(), settings, scenario, applied, agentpb.NewAgentControlServiceClient(connection)) + client := agentpb.NewAgentControlServiceClient(connection) + if mode == "mock" { + handler, err = newAgentServer(cmd.Context(), settings, scenario, applied, client) + } else { + handler, err = newRealAgentServer(cmd.Context(), settings, client) + } } if err != nil { return err @@ -100,12 +107,12 @@ func newAgentCommand() *cobra.Command { return nil }, } - command.Flags().StringVar(&mode, "mode", "mock", "isolated Mock or SIP-only mode") + command.Flags().StringVar(&mode, "mode", "mock", "isolated Mock, SIP-only or explicit nonprod-real mode") return command } // This is an isolated deployment fixture, not SaaS SIP configuration or proof -// that Asterisk actually loaded the revisions. Real/mixed startup is rejected. +// that Asterisk actually loaded the revisions. Nonprod-real cannot load it. func loadMockAppliedSIP(path string) (map[string]int64, error) { if strings.TrimSpace(path) == "" { return nil, errors.New("AGENT_MOCK_APPLIED_SIP_FILE is required") diff --git a/cmd/sip-go-agent/agent_real.go b/cmd/sip-go-agent/agent_real.go new file mode 100644 index 0000000..784cc5f --- /dev/null +++ b/cmd/sip-go-agent/agent_real.go @@ -0,0 +1,89 @@ +package main + +import ( + "context" + "errors" + "log" + "os" + "path/filepath" + "reflect" + "strings" + "time" + + agentpb "git.ipao.vip/rogee/go-sip/gen/agent" + "git.ipao.vip/rogee/go-sip/internal/agent" + "git.ipao.vip/rogee/go-sip/internal/asterisk" + "git.ipao.vip/rogee/go-sip/internal/config" + "git.ipao.vip/rogee/go-sip/internal/rpc" + "git.ipao.vip/rogee/go-sip/internal/tenant" + "github.com/google/uuid" +) + +// newRealAgentServer only exposes the signed Dispatcher execution path. There +// is no scenario, synthetic recording, local dial entry, or Mock upload host. +func newRealAgentServer(ctx context.Context, settings config.AgentEnvironment, dispatcher agentpb.AgentControlServiceClient) (*rpc.Server, error) { + if ctx == nil || ctx.Err() != nil || settings.Mode != "nonprod-real" || settings.AgentID == "" || settings.CellID == "" || + settings.SessionPath == "" || settings.RecoveryRoot == "" || len(settings.PeerFingerprints) == 0 || + settings.OSSAllowedHost == "" || dispatcher == nil || tenant.ValidateDispatcherID(settings.DispatcherID) != nil { + return nil, errors.New("real Agent requires approved nonproduction identity, pinned Dispatcher and private recovery") + } + value := reflect.ValueOf(dispatcher) + if value.Kind() == reflect.Pointer && value.IsNil() { + return nil, errors.New("real Agent requires a live pinned Dispatcher transport") + } + for _, path := range []string{settings.AsteriskConfigDir, settings.AsteriskBin, settings.AsteriskLibraryDir, settings.EvidenceRoot, settings.RecoveryRoot} { + if !filepath.IsAbs(path) { + return nil, errors.New("real Agent paths must be absolute") + } + } + if strings.ContainsAny(settings.OSSAllowedHost, "/@ ") || settings.MockScenarioFile != "" || settings.MockAppliedSIPFile != "" { + return nil, errors.New("real Agent must not use Mock assets or an invalid OSS host") + } + root, err := os.Stat(settings.RecoveryRoot) + if err != nil || !root.IsDir() || root.Mode().Perm() != 0700 { + return nil, errors.New("real Agent recovery directory must exist with mode 0700") + } + 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} + var handler *rpc.Server + worker := &rpc.ApprovedCallWorker{ + Lifecycle: ctx, Calls: &agent.TaskCalls{}, + Prepare: func(execution rpc.ApprovedExecution) (func(context.Context) error, error) { + if handler == nil { + return nil, errors.New("Agent session is unavailable") + } + delivery := &agent.RecordingDelivery{ + Call: agent.RecordingClient{ + Client: dispatcher, DispatcherID: execution.DispatcherID, TenantID: execution.TenantID, + SourceEventID: execution.SourceEventID, Session: func(context.Context) (*agentpb.RequestMeta, error) { return handler.ActiveSessionMeta() }, + }, + Recovery: &agent.RecordingRecovery{ + Root: settings.RecoveryRoot, + Upload: agent.UploadClient{AllowedHosts: map[string]struct{}{strings.ToLower(settings.OSSAllowedHost): {}}}, + }, + } + return (&rpc.ApprovedRecordedRealCall{ + Loader: loader, MediaPayloadType: 118, EvidenceRoot: settings.EvidenceRoot, + MaxWAVBytes: 64 << 20, ReportTimeout: 15 * time.Minute, Delivery: delivery, + }).Prepare(execution) + }, + OnFailure: func(execution rpc.ApprovedExecution, cause error) error { + log.Printf("Agent real call requires inspection: event_id=%q task_id=%q cause_type=%T", execution.SourceEventID, execution.TaskID, cause) + return nil + }, + } + handler, err = rpc.NewApprovedAgentServer(rpc.ServerOptions{ + Mode: settings.Mode, 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, + }, worker) + return handler, err +} diff --git a/cmd/sip-go-agent/agent_real_test.go b/cmd/sip-go-agent/agent_real_test.go new file mode 100644 index 0000000..fed75a8 --- /dev/null +++ b/cmd/sip-go-agent/agent_real_test.go @@ -0,0 +1,116 @@ +package main + +import ( + "context" + "io" + "net" + "os" + "path/filepath" + "testing" + "time" + + agentpb "git.ipao.vip/rogee/go-sip/gen/agent" + "git.ipao.vip/rogee/go-sip/internal/rpc" + "google.golang.org/grpc" + "google.golang.org/grpc/credentials" +) + +func TestRealAgentCommandStartsPinnedServerWithoutMockFixturesOrDial(t *testing.T) { + ca, agentCert, agentKey, dispatcherCert, dispatcherKey, dispatcherLeaf := localCommandCertificates(t) + settings, _ := currentAgentSetupFixture(t) + for path, data := range map[string][]byte{settings.CAFile: ca, settings.CertFile: agentCert, settings.KeyFile: agentKey} { + if err := os.WriteFile(path, data, 0600); err != nil { + t.Fatal(err) + } + } + listener, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + listen := listener.Addr().String() + listener.Close() + for name, value := range map[string]string{ + "AGENT_ID": settings.AgentID, "CELL_ID": settings.CellID, + "AGENT_GRPC_LISTEN": listen, "AGENT_SESSION_PATH": settings.SessionPath, + "AGENT_RECOVERY_ROOT": settings.RecoveryRoot, "DISPATCHER_ID": settings.DispatcherID, + "DISPATCHER_GRPC_ENDPOINT": settings.DispatcherEndpoint, "DISPATCHER_GRPC_SERVER_NAME": settings.DispatcherServerName, + "MTLS_CA_FILE": settings.CAFile, "MTLS_CERT_FILE": settings.CertFile, "MTLS_KEY_FILE": settings.KeyFile, + "MTLS_PEER_CERT_FINGERPRINTS": rpc.CertificateFingerprint(dispatcherLeaf), + "ASTERISK_CONFIG_DIR": filepath.Join(settings.RecoveryRoot, "asterisk"), + "ASTERISK_BIN": "/usr/sbin/asterisk", "ASTERISK_LIBRARY_DIR": "/usr/lib/asterisk/modules", + "AGENT_EVIDENCE_ROOT": filepath.Join(settings.RecoveryRoot, "evidence"), + "AGENT_OSS_ALLOWED_HOST": "bucket.oss-cn-beijing.aliyuncs.com", + "AGENT_MOCK_SCENARIO_FILE": "", "AGENT_MOCK_APPLIED_SIP_FILE": "", + } { + t.Setenv(name, value) + } + ctx, cancel := context.WithCancel(context.Background()) + cmd := newAgentCommand() + cmd.SetOut(io.Discard) + cmd.SetErr(io.Discard) + cmd.SetContext(ctx) + cmd.SetArgs([]string{"--mode", "nonprod-real"}) + finished := make(chan error, 1) + go func() { finished <- cmd.Execute() }() + t.Cleanup(func() { + cancel() + select { + case err := <-finished: + if err != nil { + t.Errorf("real Agent shutdown failed: %v", err) + } + case <-time.After(5 * time.Second): + t.Error("real Agent did not stop after cancellation") + } + }) + tlsConfig, err := rpc.NewClientTLSConfig(ca, dispatcherCert, dispatcherKey, "agent.local") + if err != nil { + t.Fatal(err) + } + dialCtx, stop := context.WithTimeout(context.Background(), 6*time.Second) + defer stop() + conn, err := grpc.DialContext(dialCtx, listen, grpc.WithTransportCredentials(credentials.NewTLS(tlsConfig)), grpc.WithBlock()) + if err != nil { + select { + case startupErr := <-finished: + t.Fatalf("nonprod-real startup failed: %v (dial: %v)", startupErr, err) + default: + t.Fatalf("nonprod-real pinned server unavailable: %v", err) + } + } + defer conn.Close() + status, err := agentpb.NewAgentControlServiceClient(conn).GetAgentStatus(dialCtx, &agentpb.GetAgentStatusRequest{ + Meta: &agentpb.RequestMeta{ProtocolVersion: "agent.v1", RequestId: "status", TraceId: "status", OperationId: "status", AgentId: settings.AgentID, CellId: settings.CellID}, + Target: &agentpb.AgentBinding{AgentId: settings.AgentID, CellId: settings.CellID}, + }) + if err != nil || status.GetStatus().GetAgentId() != settings.AgentID { + t.Fatalf("nonprod-real Agent did not expose its pinned status: status=%v err=%v", status, err) + } +} + +func TestRealAgentServerRequiresBoundNativeSIPAndNoMock(t *testing.T) { + settings, _ := currentAgentSetupFixture(t) + settings.Mode = "nonprod-real" + settings.AsteriskConfigDir = filepath.Join(settings.RecoveryRoot, "asterisk") + settings.AsteriskBin = "/usr/sbin/asterisk" + settings.AsteriskLibraryDir = "/usr/lib/asterisk/modules" + settings.EvidenceRoot = filepath.Join(settings.RecoveryRoot, "evidence") + settings.OSSAllowedHost = "bucket.oss-cn-beijing.aliyuncs.com" + client := &isolatedAgentRecordingClient{} + server, err := newRealAgentServer(context.Background(), settings, client) + if err != nil || server == nil { + t.Fatalf("real Agent must use the native SIP and non-Mock delivery boundary: %v", err) + } + if _, err := os.Lstat(settings.SessionPath); !os.IsNotExist(err) { + t.Fatalf("server assembly must not touch the existing session: %v", err) + } + var typedNil *isolatedAgentRecordingClient + var nilClient agentpb.AgentControlServiceClient = typedNil + if server, err := newRealAgentServer(context.Background(), settings, nilClient); err == nil || server != nil { + t.Fatal("typed-nil Dispatcher must block real Agent startup") + } + settings.MockScenarioFile = "old-mock.json" + if server, err := newRealAgentServer(context.Background(), settings, client); err == nil || server != nil { + t.Fatal("real Agent must not accept isolated Mock fixtures") + } +} diff --git a/cmd/sip-go-agent/dispatcher_command.go b/cmd/sip-go-agent/dispatcher_command.go index 0fd7bb7..a3eb076 100644 --- a/cmd/sip-go-agent/dispatcher_command.go +++ b/cmd/sip-go-agent/dispatcher_command.go @@ -24,13 +24,13 @@ import ( func newDispatcherCommand() *cobra.Command { var mode string command := &cobra.Command{ - Use: "dispatcher", Short: "Run the isolated Dispatcher", Args: cobra.NoArgs, + Use: "dispatcher", Short: "Run the bound Dispatcher", Args: cobra.NoArgs, SilenceUsage: true, RunE: func(cmd *cobra.Command, _ []string) error { return runDispatcher(cmd.Context(), mode) }, } - command.Flags().StringVar(&mode, "mode", "mock", "isolated Mock or SIP-only mode") + command.Flags().StringVar(&mode, "mode", "mock", "isolated Mock, SIP-only or explicit nonprod-real mode") return command } @@ -48,7 +48,11 @@ func runDispatcher(ctx context.Context, mode string) (result error) { if err != nil { return err } - ossConfiguration, err := config.LoadOSSConfig(settings.OSSConfigFile, settings.DispatcherID) + loadOSS := config.LoadOSSConfig + if mode == "nonprod-real" { + loadOSS = config.LoadNonprodRealOSSConfig + } + ossConfiguration, err := loadOSS(settings.OSSConfigFile, settings.DispatcherID) if err != nil { return err } @@ -147,7 +151,7 @@ func runDispatcher(ctx context.Context, mode string) (result error) { } listener, err := net.Listen("tcp", settings.Listen) if err != nil { - return fmt.Errorf("Dispatcher cannot listen on configured Mock address: %w", err) + return fmt.Errorf("Dispatcher cannot listen on configured local address: %w", err) } server := grpc.NewServer(grpc.Creds(credentials.NewTLS(serverTLS))) agentpb.RegisterAgentControlServiceServer(server, &rpc.RecordingServer{ diff --git a/deploys/test/README.md b/deploys/test/README.md index 40a041e..f3e8635 100644 --- a/deploys/test/README.md +++ b/deploys/test/README.md @@ -17,7 +17,7 @@ systemd service or add it to a production host. ## Native Asterisk validation `nonprod-call-evidence.sh` is the mandatory capture-first wrapper for -non-production `mock`, `mixed` and `real` call checks. It runs on a validation +non-production `mock`, `mixed` and explicit `nonprod-real` call checks. It runs on a validation host with native Asterisk and required diagnostics; it is not an Asterisk or Agent replacement and is not containerized. Run it explicitly with the current call authorization and the approved target/trunk. It refuses production mode @@ -29,7 +29,15 @@ of the installed binary and configuration. Missing facts fail the validation; `--preflight-only` never authorizes a call. After capture, missing capture or recording SHA-256, Asterisk journal, SIP summary, logger shutdown, timestamp, or restricted evidence ownership is recorded in `diagnostic-errors.txt` and fails -the check; an already failed call keeps its nonzero result. Isolated tests +the check; an already failed call keeps its nonzero result. Once live tcpdump +and the PJSIP logger are running, the script creates a root-owned, group-readable +`.active` capture arm under `--proof-root` (default +`/run/sip-go-agent/nonprod-armed`). The real Agent must set `AGENT_EVIDENCE_ROOT` +to this directory; it checks the exact approved event ID, trunk, raw callee, +recent arm and live capture PID before origination. The script removes the arm +before stopping capture. Run one approved call per invocation, with +`--call-id` equal to that call's MQ `event_id`; never reuse a stale arm. +Isolated tests replace host tools with fakes: they do not prove a real host or supplier is ready. For the **nonproduction user-level Asterisk service only**, pass diff --git a/deploys/test/nonprod-call-evidence.sh b/deploys/test/nonprod-call-evidence.sh index 08fa4b3..5948028 100755 --- a/deploys/test/nonprod-call-evidence.sh +++ b/deploys/test/nonprod-call-evidence.sh @@ -39,6 +39,9 @@ target="" call_command=() preflight_only=0 attempt_ledger="/var/lib/sip-go-agent/state/real-call-attempts.tsv" +proof_root="/run/sip-go-agent/nonprod-armed" +proof_file="" +proof_created=0 attempt_number=0 while (($#)); do @@ -55,6 +58,7 @@ while (($#)); do --rtp-end) [[ $# -ge 2 ]] || usage; rtp_end=$2; shift 2 ;; --preflight-only) preflight_only=1; shift ;; --attempt-ledger) [[ $# -ge 2 ]] || usage; attempt_ledger=$2; shift 2 ;; + --proof-root) [[ $# -ge 2 ]] || usage; proof_root=$2; shift 2 ;; --trunk) [[ $# -ge 2 ]] || usage; trunk=$2; shift 2 ;; --target) [[ $# -ge 2 ]] || usage; target=$2; shift 2 ;; --) shift; call_command=("$@"); break ;; @@ -73,6 +77,7 @@ esac [[ "$trunk" =~ ^(provider-primary|provider-second|provider-third|trunk-[A-Za-z0-9._-]+)$ ]] || { echo 'trunk is not an approved non-production trunk id' >&2; exit 1; } [[ "$target" =~ ^(15003164745|15830461047)$ ]] || { echo 'target is outside the approved outbound whitelist' >&2; exit 1; } [[ "$attempt_ledger" =~ ^/[A-Za-z0-9._/-]+$ ]] || { echo 'invalid attempt ledger path' >&2; exit 1; } +[[ "$proof_root" =~ ^/[A-Za-z0-9._/-]+$ ]] || { echo 'invalid proof root path' >&2; exit 1; } [[ ${#call_command[@]} -gt 0 ]] || { echo 'call command is required after --' >&2; exit 1; } [[ "$interface" =~ ^[A-Za-z0-9_.:-]+$ ]] || { echo 'invalid capture interface' >&2; exit 1; } [[ "$sip_port" =~ ^[0-9]+$ && "$rtp_start" =~ ^[0-9]+$ && "$rtp_end" =~ ^[0-9]+$ ]] || { echo 'invalid port' >&2; exit 1; } @@ -294,6 +299,10 @@ cleanup() { local call_exit=$? evidence_failed=0 trap - EXIT set +e + if ((proof_created)) && ! rm -f -- "$proof_file"; then + printf 'capture arm removal failed\n' >>"$evidence_dir/diagnostic-errors.txt" + evidence_failed=1 + fi stop_capture if ((logger_enabled)) && ! asterisk_cli "pjsip set logger off" >"$evidence_dir/pjsip-logger-off.txt" 2>&1; then printf 'PJSIP logger stop failed\n' >>"$evidence_dir/diagnostic-errors.txt" @@ -394,6 +403,18 @@ if ((preflight_only)); then fi require_call_window +# The Agent checks this root-owned, call-specific live capture arm before any +# originate. A stale arm is never overwritten; the trap removes it first. +install -d -o root -g "$run_as" -m 0750 -- "$proof_root" +proof_file="$proof_root/$call_id.active" +[[ ! -e "$proof_file" ]] || { echo 'stale or concurrent capture arm; fail-closed' >&2; exit 1; } +proof_tmp="$(mktemp --tmpdir="$proof_root" ".$call_id.XXXXXX")" +printf '%s\t%s\t%s\n' "$capture_pid" "$trunk" "$target" >"$proof_tmp" +chown "root:$run_as" "$proof_tmp" +chmod 0640 "$proof_tmp" +ln -- "$proof_tmp" "$proof_file" || { rm -f -- "$proof_tmp"; echo 'capture arm already exists; fail-closed' >&2; exit 1; } +proof_created=1 +rm -f -- "$proof_tmp" call_status=0 set +e runuser -u "$run_as" -- "${call_command[@]}" >"$evidence_dir/call-output.private" 2>&1 diff --git a/internal/config/agent_runtime.go b/internal/config/agent_runtime.go index 883a53b..d248d4f 100644 --- a/internal/config/agent_runtime.go +++ b/internal/config/agent_runtime.go @@ -5,6 +5,7 @@ import ( "fmt" "net" "os" + "path/filepath" "strconv" "strings" @@ -14,6 +15,7 @@ import ( // AgentEnvironment is deployment-owned and is inspected without opening any // durable file, media adapter, credential or network connection. type AgentEnvironment struct { + Mode string AgentID string CellID string DispatcherID string @@ -31,11 +33,13 @@ type AgentEnvironment struct { AsteriskConfigDir string AsteriskBin string AsteriskLibraryDir string + EvidenceRoot string + OSSAllowedHost string } func LoadAgentEnvironment(mode string) (AgentEnvironment, error) { - if mode != "mock" && mode != "sip-only" { - return AgentEnvironment{}, errors.New("Agent accepts only isolated Mock or SIP-only mode") + if mode != "mock" && mode != "sip-only" && mode != "nonprod-real" { + return AgentEnvironment{}, errors.New("Agent accepts only isolated Mock, SIP-only or explicit nonprod-real mode") } get := func(name string) (string, error) { value := os.Getenv(name) @@ -44,7 +48,7 @@ func LoadAgentEnvironment(mode string) (AgentEnvironment, error) { } return value, nil } - var settings AgentEnvironment + settings := AgentEnvironment{Mode: mode} for _, field := range []struct { name string value *string @@ -67,13 +71,15 @@ func LoadAgentEnvironment(mode string) (AgentEnvironment, error) { value *string } var err error - if mode == "mock" { + if mode == "mock" || mode == "nonprod-real" { if settings.DispatcherEndpoint, err = get("DISPATCHER_GRPC_ENDPOINT"); err != nil { return AgentEnvironment{}, err } if settings.DispatcherServerName, err = get("DISPATCHER_GRPC_SERVER_NAME"); err != nil { return AgentEnvironment{}, err } + } + if mode == "mock" { modeFields = []struct { name string value *string @@ -91,15 +97,31 @@ func LoadAgentEnvironment(mode string) (AgentEnvironment, error) { } *field.value = value } + if mode == "nonprod-real" { + if settings.EvidenceRoot, err = get("AGENT_EVIDENCE_ROOT"); err != nil { + return AgentEnvironment{}, err + } + if settings.OSSAllowedHost, err = get("AGENT_OSS_ALLOWED_HOST"); err != nil { + return AgentEnvironment{}, err + } + } if tenant.ValidateDispatcherID(settings.DispatcherID) != nil { return AgentEnvironment{}, errors.New("DISPATCHER_ID must be a canonical UUID v4") } if !localGRPCAddress(settings.Listen, true) { return AgentEnvironment{}, errors.New("AGENT_GRPC_LISTEN must be an isolated local Mock address") } - if mode == "mock" && !localGRPCAddress(settings.DispatcherEndpoint, false) { + if mode != "sip-only" && !localGRPCAddress(settings.DispatcherEndpoint, false) { return AgentEnvironment{}, errors.New("DISPATCHER_GRPC_ENDPOINT must be an isolated local Mock address") } + if mode == "nonprod-real" { + if !filepath.IsAbs(settings.EvidenceRoot) { + return AgentEnvironment{}, errors.New("AGENT_EVIDENCE_ROOT must be absolute") + } + if os.Getenv("AGENT_MOCK_SCENARIO_FILE") != "" || os.Getenv("AGENT_MOCK_APPLIED_SIP_FILE") != "" { + return AgentEnvironment{}, errors.New("nonprod-real Agent cannot load isolated Mock fixtures") + } + } rawFingerprints, err := get("MTLS_PEER_CERT_FINGERPRINTS") if err != nil { return AgentEnvironment{}, err diff --git a/internal/config/agent_runtime_test.go b/internal/config/agent_runtime_test.go index 5b12737..bf0ddd8 100644 --- a/internal/config/agent_runtime_test.go +++ b/internal/config/agent_runtime_test.go @@ -59,6 +59,41 @@ func TestLoadSIPOnlyAgentEnvironmentDoesNotRequireMockCalls(t *testing.T) { } } +func TestLoadNonprodRealAgentRequiresNativeCallAndEvidenceSettings(t *testing.T) { + state := 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", + "AGENT_EVIDENCE_ROOT": "/tmp/nonprod-call-evidence", + "AGENT_OSS_ALLOWED_HOST": "test.oss-cn-beijing.aliyuncs.com", + } { + t.Setenv(name, value) + } + settings, err := LoadAgentEnvironment("nonprod-real") + if err != nil || settings.AsteriskConfigDir == "" || settings.DispatcherEndpoint == "" || settings.EvidenceRoot != "/tmp/nonprod-call-evidence" || settings.OSSAllowedHost != "test.oss-cn-beijing.aliyuncs.com" || settings.MockScenarioFile != "" { + t.Fatalf("nonproduction real mode must use explicit native, capture and OSS settings: %+v %v", settings, err) + } + if _, err := os.Stat(state); !os.IsNotExist(err) { + t.Fatalf("configuration inspection opened durable session: %v", err) + } + for _, name := range []string{"ASTERISK_CONFIG_DIR", "ASTERISK_BIN", "ASTERISK_LIBRARY_DIR", "AGENT_EVIDENCE_ROOT", "AGENT_OSS_ALLOWED_HOST", "DISPATCHER_GRPC_ENDPOINT", "DISPATCHER_GRPC_SERVER_NAME"} { + t.Run(name, func(t *testing.T) { + t.Setenv(name, "") + if _, err := LoadAgentEnvironment("nonprod-real"); err == nil || !strings.Contains(err.Error(), name) { + t.Fatalf("missing real-call setting must block startup: %v", err) + } + t.Setenv(name, map[string]string{ + "ASTERISK_CONFIG_DIR": "/tmp/asterisk-config", "ASTERISK_BIN": "/tmp/asterisk", + "ASTERISK_LIBRARY_DIR": "/tmp/asterisk-libraries", "AGENT_EVIDENCE_ROOT": "/tmp/nonprod-call-evidence", + "AGENT_OSS_ALLOWED_HOST": "test.oss-cn-beijing.aliyuncs.com", "DISPATCHER_GRPC_ENDPOINT": "127.0.0.1:39443", + "DISPATCHER_GRPC_SERVER_NAME": "dispatcher.local", + }[name]) + }) + } +} + 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/oss_config.go b/internal/config/oss_config.go index 70ed220..d815818 100644 --- a/internal/config/oss_config.go +++ b/internal/config/oss_config.go @@ -20,6 +20,16 @@ import ( // Dispatcher. Credentials are resolved from explicit environment references; // the file, its contents and secrets are never echoed in errors. func LoadOSSConfig(filename, dispatcherID string) (oss.Config, error) { + return loadOSSConfig(filename, dispatcherID, false) +} + +// LoadNonprodRealOSSConfig accepts only the approved Beijing HTTPS OSS endpoint. +// The old isolated Mock loader remains loopback-only. +func LoadNonprodRealOSSConfig(filename, dispatcherID string) (oss.Config, error) { + return loadOSSConfig(filename, dispatcherID, true) +} + +func loadOSSConfig(filename, dispatcherID string, nonprodReal bool) (oss.Config, error) { if strings.TrimSpace(filename) == "" || tenant.ValidateDispatcherID(dispatcherID) != nil { return oss.Config{}, errors.New("DISPATCHER_OSS_CONFIG_FILE and canonical Dispatcher ID are required") } @@ -60,9 +70,11 @@ func LoadOSSConfig(filename, dispatcherID string) (oss.Config, error) { return oss.Config{}, errors.New("OSS configuration is not assigned to this Dispatcher") } endpoint, err := url.Parse(input.OSS.Endpoint) - if err != nil || !localEndpoint(input.OSS.Endpoint, "https") || endpoint.User != nil || endpoint.RawQuery != "" || endpoint.Fragment != "" || - (endpoint.Path != "" && endpoint.Path != "/") || endpoint.RawPath != "" || !localGRPCAddress(endpoint.Host, false) { - return oss.Config{}, errors.New("OSS Mock endpoint must be a local HTTPS service with an explicit port") + if err != nil || endpoint.Scheme != "https" || endpoint.User != nil || endpoint.RawQuery != "" || endpoint.Fragment != "" || + (endpoint.Path != "" && endpoint.Path != "/") || endpoint.RawPath != "" || + (!nonprodReal && (!localEndpoint(input.OSS.Endpoint, "https") || !localGRPCAddress(endpoint.Host, false))) || + (nonprodReal && (endpoint.Host != "oss-cn-beijing.aliyuncs.com" || input.OSS.Region != "cn-beijing")) { + return oss.Config{}, errors.New("OSS endpoint does not match the selected isolated Mock or approved Beijing real deployment") } prefix := input.OSS.ObjectPrefix if prefix == "" || strings.TrimSpace(prefix) != prefix || path.IsAbs(prefix) || path.Clean(prefix) != prefix || diff --git a/internal/config/oss_config_test.go b/internal/config/oss_config_test.go index 7a93f32..21d7043 100644 --- a/internal/config/oss_config_test.go +++ b/internal/config/oss_config_test.go @@ -33,6 +33,28 @@ func TestLoadOSSConfigRequiresBoundLocalHTTPSAndExplicitCredentials(t *testing.T } } +func TestNonprodRealOSSOnlyAcceptsApprovedBeijingEndpoint(t *testing.T) { + t.Setenv("MOCK_OSS_KEY_ID", "isolated-id") + t.Setenv("MOCK_OSS_KEY_SECRET", "isolated-secret") + realBody := strings.Replace(strings.Replace(mockOSSDeployment, "https://127.0.0.1:19445", "https://oss-cn-beijing.aliyuncs.com", 1), `"region":"cn-mock"`, `"region":"cn-beijing"`, 1) + cfg, err := LoadNonprodRealOSSConfig(writeMockOSSDeployment(t, realBody), dispatcherFixtureID) + if err != nil || cfg.Endpoint != "https://oss-cn-beijing.aliyuncs.com" || cfg.Region != "cn-beijing" { + t.Fatalf("explicit Beijing OSS deployment must be accepted: err=%v", err) + } + if _, err := LoadOSSConfig(writeMockOSSDeployment(t, realBody), dispatcherFixtureID); err == nil { + t.Fatal("isolated Mock must not reach real OSS") + } + for _, body := range []string{ + mockOSSDeployment, + strings.Replace(realBody, "oss-cn-beijing.aliyuncs.com", "oss-cn-shanghai.aliyuncs.com", 1), + strings.Replace(realBody, "https://oss-cn-beijing.aliyuncs.com", "http://oss-cn-beijing.aliyuncs.com", 1), + } { + if _, err := LoadNonprodRealOSSConfig(writeMockOSSDeployment(t, body), dispatcherFixtureID); err == nil { + t.Fatal("nonproduction real OSS admitted a Mock, insecure or wrong-region endpoint") + } + } +} + func TestLoadOSSConfigIgnoresLegacyOSSEnvironment(t *testing.T) { for _, key := range []string{"DISPATCHER_OSS_REGION", "DISPATCHER_OSS_ENDPOINT", "DISPATCHER_OSS_BUCKET", "DISPATCHER_OSS_ACCESS_KEY_ID", "DISPATCHER_OSS_ACCESS_KEY_SECRET", "DISPATCHER_OSS_KEY_PREFIX", "DISPATCHER_OSS_GRANT_TTL_SECONDS"} { t.Setenv(key, "legacy-value") diff --git a/internal/config/runtime.go b/internal/config/runtime.go index 39ac4ca..93ae1f8 100644 --- a/internal/config/runtime.go +++ b/internal/config/runtime.go @@ -22,10 +22,11 @@ type DispatcherEnvironment struct { } // LoadDispatcherEnvironment is pure inspection: it opens no database, queue or -// network connection. Only isolated Mock or SIP-only mode may start. +// network connection. Nonproduction real calls still use only the local SaaS +// simulator and predeclared local AMQP test queues. func LoadDispatcherEnvironment(mode string) (DispatcherEnvironment, error) { - if mode != "mock" && mode != "sip-only" { - return DispatcherEnvironment{}, errors.New("Dispatcher accepts only isolated Mock or SIP-only mode") + if mode != "mock" && mode != "sip-only" && mode != "nonprod-real" { + return DispatcherEnvironment{}, errors.New("Dispatcher accepts only isolated Mock, SIP-only or explicit nonprod-real mode") } get := func(name string) (string, error) { value := os.Getenv(name) diff --git a/internal/config/runtime_test.go b/internal/config/runtime_test.go index 171a837..10332e8 100644 --- a/internal/config/runtime_test.go +++ b/internal/config/runtime_test.go @@ -34,6 +34,21 @@ func TestLoadDispatcherEnvironmentRefusesNonMockBeforeResources(t *testing.T) { } } +func TestLoadDispatcherEnvironmentNonprodRealUsesOnlyLocalSaaSMockAndQueue(t *testing.T) { + path := setDispatcherEnvironment(t) + settings, err := LoadDispatcherEnvironment("nonprod-real") + if err != nil || settings.SQLitePath != path { + t.Fatalf("explicit nonprod-real should use the same bound SaaS HTTP and AMQP inputs: %v", err) + } + if _, err := os.Stat(path); !os.IsNotExist(err) { + t.Fatalf("real configuration inspection opened SQLite: %v", err) + } + t.Setenv("SAAS_BASE_URL", "https://saas.example.invalid") + if _, err := LoadDispatcherEnvironment("nonprod-real"); err == nil { + t.Fatal("external SaaS is not authorized for this nonproduction test") + } +} + func TestLoadDispatcherEnvironmentRequiresExplicitIdentityAndNoFallback(t *testing.T) { for _, name := range []string{"DISPATCHER_ID", "DISPATCHER_SECRET_KEY", "SAAS_BASE_URL", "RABBITMQ_URL", "DISPATCHER_SQLITE_PATH"} { t.Run(name, func(t *testing.T) { diff --git a/internal/rpc/approved_control.go b/internal/rpc/approved_control.go index 884648b..7259142 100644 --- a/internal/rpc/approved_control.go +++ b/internal/rpc/approved_control.go @@ -28,8 +28,8 @@ func (s *Server) ApplyApprovedTaskControl(ctx context.Context, req *agentpb.Appl if req.DispatcherId != dispatcherID { return nil, status.Error(codes.PermissionDenied, "Dispatcher does not own the active Agent session") } - if s.mode != "mock" { - return nil, status.Error(codes.FailedPrecondition, "approved task control requires isolated Mock mode") + if s.mode != "mock" && s.mode != "nonprod-real" { + return nil, status.Error(codes.FailedPrecondition, "approved task control requires an approved call mode") } if req.TenantId <= 0 || strings.TrimSpace(req.TaskId) == "" { return nil, status.Error(codes.InvalidArgument, "approved task control identity is incomplete") diff --git a/internal/rpc/approved_control_test.go b/internal/rpc/approved_control_test.go index 6f6cd26..d8e2b82 100644 --- a/internal/rpc/approved_control_test.go +++ b/internal/rpc/approved_control_test.go @@ -53,6 +53,12 @@ func TestApprovedTaskControlRequiresCurrentSessionAndValidTaskPolicy(t *testing. if response, err := server.ApplyApprovedTaskControl(context.Background(), request); err != nil || response == nil || !response.Accepted { t.Fatalf("approved resume was not applied: response=%v err=%v", response, err) } + server.mode = "nonprod-real" + request.Action = agentpb.ControlAction_CONTROL_ACTION_PAUSE + request.ActiveCallPolicy = agentpb.ActiveCallPolicy_ACTIVE_CALL_POLICY_HANGUP + if response, err := server.ApplyApprovedTaskControl(context.Background(), request); err != nil || response == nil || !response.Accepted { + t.Fatalf("nonprod-real task control must use the same authenticated entry: response=%v err=%v", response, err) + } server.approvedTaskCalls = nil if _, err := server.ApplyApprovedTaskControl(context.Background(), request); status.Code(err) != codes.FailedPrecondition { t.Fatalf("missing Agent control adapter reported success: %v", err) diff --git a/internal/rpc/approved_recorded_real.go b/internal/rpc/approved_recorded_real.go index fcecdf7..0c63844 100644 --- a/internal/rpc/approved_recorded_real.go +++ b/internal/rpc/approved_recorded_real.go @@ -4,8 +4,11 @@ import ( "context" "errors" "os" + "path/filepath" + "strconv" "strings" "sync/atomic" + "syscall" "time" "git.ipao.vip/rogee/go-sip/internal/agent" @@ -20,6 +23,7 @@ import ( type ApprovedRecordedRealCall struct { Loader asterisk.Loader MediaPayloadType uint8 + EvidenceRoot string MaxWAVBytes int64 ReportTimeout time.Duration Delivery *agent.RecordingDelivery @@ -52,8 +56,13 @@ func (r *ApprovedRecordedRealCall) Prepare(approved ApprovedExecution) (func(con if err != nil || !info.IsDir() || info.Mode().Perm() != 0700 { return nil, errors.New("real recording recovery directory must exist with mode 0700") } - if r.originator == nil && (r.Loader.ConfigDir == "" || r.Loader.Asterisk == "" || r.Loader.LibraryDir == "") { - return nil, errors.New("native Asterisk configuration required for real call") + if r.originator == nil { + if r.Loader.ConfigDir == "" || r.Loader.Asterisk == "" || r.Loader.LibraryDir == "" { + return nil, errors.New("native Asterisk configuration required for real call") + } + if err := verifyCallEvidence(r.EvidenceRoot, approved.SourceEventID, approved.SelectedTrunkID, approved.Callee); err != nil { + return nil, err + } } if _, err := ai.NewCall(approved.AI, func(context.Context) error { return nil }); err != nil { return nil, err @@ -81,6 +90,11 @@ func (r *ApprovedRecordedRealCall) Prepare(approved ApprovedExecution) (func(con if ctx == nil || ctx.Err() != nil || !time.Now().Before(approved.DialBefore) { return errors.New("real call instruction expired or cancelled before origination") } + if r.originator == nil { + if err := verifyCallEvidence(r.EvidenceRoot, approved.SourceEventID, approved.SelectedTrunkID, approved.Callee); err != nil { + return err + } + } call, err := originate(ctx, dial) if err != nil { return err // unknown origination is never retried or reported as a completed call @@ -131,6 +145,43 @@ func (r *ApprovedRecordedRealCall) Prepare(approved ApprovedExecution) (func(con }, nil } +// verifyCallEvidence accepts only a current, live capture armed by the root +// nonproduction host script for this exact accepted call and trunk. +func verifyCallEvidence(root, id, trunk, target string) error { + if !filepath.IsAbs(root) || id == "" || id == "." || id == ".." || filepath.Base(id) != id || strings.ContainsAny(id, "/\\") { + return errors.New("real call evidence path or identity is invalid") + } + marker := filepath.Join(root, id+".active") + info, err := os.Lstat(marker) + if err != nil || !info.Mode().IsRegular() || info.Mode().Perm() != 0640 || time.Since(info.ModTime()) > 10*time.Minute { + return errors.New("real call capture is not armed") + } + proof, err := os.ReadFile(marker) + if err != nil || len(proof) > 160 { + return errors.New("real call capture proof is unreadable") + } + fields := strings.Split(strings.TrimSuffix(string(proof), "\n"), "\t") + if len(fields) != 3 || fields[1] != trunk || fields[2] != target { + return errors.New("real call capture does not match signed call identity") + } + pid, err := strconv.Atoi(fields[0]) + if err != nil || pid <= 1 { + return errors.New("real call capture PID invalid") + } + if err := syscall.Kill(pid, 0); err != nil && !errors.Is(err, syscall.EPERM) { + return errors.New("real call capture process is not running") + } + state, err := os.ReadFile(filepath.Join("/proc", fields[0], "stat")) + if err != nil { + return errors.New("real call capture process state is unavailable") + } + end := strings.LastIndex(string(state), ") ") + if end == -1 || len(state) <= end+2 || state[end+2] == 'Z' || state[end+2] == 'X' { + return errors.New("real call capture process has ended") + } + return nil +} + func (r *ApprovedRecordedRealCall) Run(ctx context.Context, approved ApprovedExecution) error { if ctx == nil { return errors.New("real call requires a context") diff --git a/internal/rpc/approved_recorded_real_test.go b/internal/rpc/approved_recorded_real_test.go index 3dd7e67..0a9d469 100644 --- a/internal/rpc/approved_recorded_real_test.go +++ b/internal/rpc/approved_recorded_real_test.go @@ -4,7 +4,10 @@ import ( "context" "encoding/json" "errors" + "fmt" "io" + "os" + "path/filepath" "strings" "testing" "time" @@ -20,6 +23,33 @@ func (endedRealMedia) ReadPayload(context.Context) ([]byte, error) { return nil func (endedRealMedia) SendPCM16(context.Context, []byte, int) error { return nil } func (endedRealMedia) Stats() media.RTPStats { return media.RTPStats{} } +func TestRealCallRequiresLiveCaptureBoundToSignedIdentity(t *testing.T) { + root := t.TempDir() + id, trunk, target := "event-1", "trunk-1", "15003164745" + marker := filepath.Join(root, id+".active") + if err := verifyCallEvidence(root, id, trunk, target); err == nil { + t.Fatal("missing SIP/RTP evidence arm must block the call") + } + if err := os.WriteFile(marker, []byte(fmt.Sprintf("%d\t%s\t%s\n", os.Getpid(), trunk, target)), 0640); err != nil { + t.Fatal(err) + } + if err := verifyCallEvidence(root, id, trunk, target); err != nil { + t.Fatal(err) + } + if err := verifyCallEvidence(root, id, trunk, "15830461047"); err == nil { + t.Fatal("capture bound to a different callee cannot authorize a real call") + } + if err := verifyCallEvidence(root, "../event-1", trunk, target); err == nil { + t.Fatal("call ID cannot escape the evidence directory") + } + if err := os.WriteFile(marker, []byte(fmt.Sprintf("%d\t%s\t%s\n", 2147483647, trunk, target)), 0640); err != nil { + t.Fatal(err) + } + if err := verifyCallEvidence(root, id, trunk, target); err == nil { + t.Fatal("terminated tcpdump must block the call") + } +} + func TestApprovedRecordedRealCallReportsOnlyEndedObservedCall(t *testing.T) { fixture, stub, puts, _, approved := recordedMockFixture(t, nil, time.Second) approved.CallerID, approved.DialedCallee, approved.RingTimeout = "BD93205882", "7089"+approved.Callee, time.Second diff --git a/internal/rpc/approved_server.go b/internal/rpc/approved_server.go index 1edad14..e17792f 100644 --- a/internal/rpc/approved_server.go +++ b/internal/rpc/approved_server.go @@ -11,14 +11,14 @@ var ErrApprovedWorkerRequired = errors.New("approved call worker is required") // NewApprovedAgentServer binds the control RPC and one-shot originate callback // to the same task-call registry. The separate legacy hooks are not accepted -// on this current Mock-only entry: otherwise a control could confirm while a -// call started by a different adapter remains active. +// on either explicitly selected entry: otherwise a control could confirm while +// a call started by a different adapter remains active. func NewApprovedAgentServer(options ServerOptions, worker *ApprovedCallWorker) (*Server, error) { if worker == nil { return nil, ErrApprovedWorkerRequired } - if options.Mode != "mock" || options.StatePath == "" || options.Status == nil || options.Status.AgentId == "" || options.Status.CellId == "" || options.Status.BootId == "" || strings.TrimSpace(options.ApprovedDispatcherID) == "" || options.LoadedSIP == nil { - return nil, errors.New("approved Agent requires explicit Mock mode, durable state, Agent and Dispatcher identity, and observed SIP") + if (options.Mode != "mock" && options.Mode != "nonprod-real") || options.StatePath == "" || options.Status == nil || options.Status.AgentId == "" || options.Status.CellId == "" || options.Status.BootId == "" || strings.TrimSpace(options.ApprovedDispatcherID) == "" || options.LoadedSIP == nil { + return nil, errors.New("approved Agent requires an explicit supported mode, durable state, Agent and Dispatcher identity, and observed SIP") } if worker.Lifecycle == nil || worker.Calls == nil || worker.Prepare == nil || worker.OnFailure == nil { return nil, errors.New("approved Agent requires a process lifecycle, task calls, runner and failure reporting") @@ -28,13 +28,17 @@ func NewApprovedAgentServer(options ServerOptions, worker *ApprovedCallWorker) ( } else if !errors.Is(err, os.ErrNotExist) { return nil, fmt.Errorf("inspect legacy Agent execution state: %w", err) } - if options.MockApprovedOriginate != nil || options.ApprovedTaskCalls != nil { + if options.MockApprovedOriginate != nil || options.NonprodRealOriginate != nil || options.ApprovedTaskCalls != nil { return nil, errors.New("approved Agent cannot use competing call or control adapters") } // Copy the configuration so later mutation of worker fields cannot cause // ExecuteApproved and task controls to observe different registries. frozen := *worker - options.MockApprovedOriginate = frozen.Originate + if options.Mode == "mock" { + options.MockApprovedOriginate = frozen.Originate + } else { + options.NonprodRealOriginate = frozen.Originate + } options.ApprovedTaskCalls = frozen.Calls return NewServer(options), nil } diff --git a/internal/rpc/approved_server_test.go b/internal/rpc/approved_server_test.go index b06a995..8b77712 100644 --- a/internal/rpc/approved_server_test.go +++ b/internal/rpc/approved_server_test.go @@ -175,6 +175,21 @@ func TestApprovedAgentServerRefusesMissingOrCompetingAdapters(t *testing.T) { } } +func TestApprovedAgentServerBindsOnlyRealWorkerInNonprodMode(t *testing.T) { + worker := &ApprovedCallWorker{Lifecycle: context.Background(), Calls: &agent.TaskCalls{}, + Prepare: prepareWorkerRun(func(context.Context, ApprovedExecution) error { return nil }), + OnFailure: func(ApprovedExecution, error) error { return nil }, + } + server, err := NewApprovedAgentServer(ServerOptions{ + Mode: "nonprod-real", StatePath: filepath.Join(t.TempDir(), "agent-session.json"), + ApprovedDispatcherID: "dispatcher-1", Status: &agentpb.AgentStatus{AgentId: "agent-1", CellId: "cell-1", BootId: "boot-1"}, + LoadedSIP: func(context.Context) (map[string]int64, error) { return map[string]int64{"trunk-1": 8}, nil }, + }, worker) + if err != nil || server == nil || server.nonprodRealOriginate == nil || server.mockApprovedOriginate != nil { + t.Fatalf("nonproduction real mode must bind only its real worker: server=%v err=%v", server, err) + } +} + func TestApprovedAgentServerRefusesLegacyExecutionJournalWithoutMutation(t *testing.T) { state := filepath.Join(t.TempDir(), "agent-session.json") legacy := state + ".executions" diff --git a/internal/rpc/sip_apply.go b/internal/rpc/sip_apply.go index 2d81e1d..0cb121a 100644 --- a/internal/rpc/sip_apply.go +++ b/internal/rpc/sip_apply.go @@ -12,8 +12,8 @@ import ( "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 +// ApplySIP is deliberately unavailable on the Mock call server. A native +// 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 { @@ -26,7 +26,7 @@ func (s *Server) ApplySIP(ctx context.Context, req *agentpb.ApplySIPRequest) (*a if err != nil { return nil, err } - if s.mode != "sip-only" || s.applySIP == nil { + if (s.mode != "sip-only" && s.mode != "nonprod-real") || s.applySIP == nil { return nil, status.Error(codes.FailedPrecondition, "real SIP apply is unavailable on this Agent") } var approved configread.SIP diff --git a/internal/rpc/sip_apply_test.go b/internal/rpc/sip_apply_test.go index b9afa83..caafc9f 100644 --- a/internal/rpc/sip_apply_test.go +++ b/internal/rpc/sip_apply_test.go @@ -68,4 +68,8 @@ func TestApplySIPOnlyAllowsActivatedSIPService(t *testing.T) { 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) } + server.mode, server.applySIP = "nonprod-real", apply + if _, err := server.ApplySIP(context.Background(), req); err != nil || attempts != 2 { + t.Fatalf("nonprod-real must observe native SIP through the same approved path: attempts=%d err=%v", attempts, err) + } } diff --git a/scripts/check-release-gates-local.sh b/scripts/check-release-gates-local.sh index b9fe9ea..5095d04 100755 --- a/scripts/check-release-gates-local.sh +++ b/scripts/check-release-gates-local.sh @@ -88,11 +88,18 @@ for mode in mixed real; do echo "release accepted unapproved $mode execution" >&2 exit 1 fi - if ! grep -Fq 'Dispatcher accepts only isolated Mock or SIP-only mode' "$WORK/blocked-$mode.log" || test -e "$blocked_db"; then + if ! grep -Fq 'Dispatcher accepts only isolated Mock, SIP-only or explicit nonprod-real mode' "$WORK/blocked-$mode.log" || test -e "$blocked_db"; then echo "release reached resources before rejecting $mode execution" >&2 exit 1 fi done +# The explicit nonproduction real path cannot open durable state unless every +# bound deployment input is present; this package test supplies none of them. +if DISPATCHER_SQLITE_PATH="$WORK/unconfigured-real.sqlite" "$WORK/release/sip-go-agent" dispatcher --mode nonprod-real > "$WORK/unconfigured-real.log" 2>&1 || + test -e "$WORK/unconfigured-real.sqlite"; then + echo 'unconfigured real Dispatcher opened state or started' >&2 + exit 1 +fi rm -- "$PACKAGE_RELEASE/sentinel" rmdir -- "$PACKAGE_RELEASE" "$ROOT/deploys/build-package.sh" "$PACKAGE_VERSION" > "$WORK/package.log" 2>&1 || {