feat: fixed systemd preproduction deployment with safe state reset

This commit is contained in:
2026-10-09 01:58:52 +08:00
parent 5b2cb11c2f
commit e5f46ba016
11 changed files with 1270 additions and 2 deletions
+1
View File
@@ -0,0 +1 @@
__pycache__/
+60
View File
@@ -0,0 +1,60 @@
# 固定预生产部署
目标是登记测试机 `server.sip`。服务使用 `nonprod-real`、真实 SaaS/RabbitMQ/AI/OSS 和原生 Asterisk;不安装 SaaS Mock 或本机 RabbitMQ,不以一次性运行目录代替正式运行。
## 固定位置
| 内容 | 位置 |
| --- | --- |
| 配置及 TLS | `~/.config/go-sip/` |
| Dispatcher SQLite | `~/.local/share/go-sip/dispatcher.sqlite` |
| Agent 会话 | `~/.local/share/go-sip/agent-session.json` |
| Agent 结果、录音及上传恢复 | `~/.local/share/go-sip/agent-recovery/` |
| 版本制品 | `~/.local/opt/go-sip/releases/<main提交>/` |
| 当前制品指针 | `~/.local/opt/go-sip/current` |
| 操作备份及私密证据 | `~/.local/share/go-sip-backups/` |
| 原生 Asterisk 配置 | `~/.config/go-sip-asterisk/`(沿用已安装实例) |
版本目录只保存制品,不包含业务数据库。正常重部署切换制品,不创建另一套运行数据,也不清除现有状态。
服务名固定为 `go-sip-asterisk.service`、`go-sip-agent.service`、`go-sip-dispatcher.service`,由 `rogee` 的 `systemd --user` 管理。旧日期命名的测试单元退役,不并行启动多个 Agent/Dispatcher。
## 当前项目配置来源
- `saas.env` 为 `KEY=VALUE`:`DispatcherUUID` → `DISPATCHER_ID`,`DispatcherKEY` → `DISPATCHER_SECRET_KEY`,`SaaSBaseURL` → `SAAS_BASE_URL`。
- `rabbitmq.env` 为 `KEY=VALUE`:`RABBITMQ_HOST/PORT/USER/PASS/VHOST` 构造正式 `RABBITMQ_URL`。用户名、密码及 vhost 按 URI 规则编码;不改用旧本机地址,不创建外部队列或交换机。
- `aliyun-oss.env` 不是 shell env:读取准确大小写的 `bucket`、`Endpoint`、`Region` 和 `RAM.username/accessKeyId/accessKeySecret`,转换为 Dispatcher 私有 OSS JSON,并以环境变量名称引用访问凭据。部署的录音前缀明确为 `call-recordings`,最大录音 64 MiB,与现有 Agent 上限一致;不写旧快照的桶或访问凭据。
- SIP、任务、额度、AI 连接和 AI 参数只从真实 SaaS 取得。`.local/provider-ai.env` 不覆盖 SaaS 返回的型号、音色、参数或连接。
三份源文件先核对所有者及 `0600`,只在受限进程中解析,不执行 `source`、不打印字段值。服务配置和密钥为 `0600`,目录为 `0700`;私有文件、完整连接串及原始证据不进入 Git。
## 部署及重置边界
一键脚本须核验登记 SSH 指纹、目标主机、正式 main 制品来源/hash、Debian 13/架构、原生 Asterisk及诊断工具、权限和必要配置。缺项立即失败,不跳过诊断、不安装 Mock 作为回退。
普通部署不得改写数据库、幂等记录、录音、结果、恢复文件或会话。首次安装明确使用上述固定路径;发现不兼容数据时停止,不偷偷切换另一套空目录。
只有显式 `--reset-state` 才能重置本项目非生产数据:停止相关服务、核实 Asterisk 没有活动通话、完成私密备份和清单记录后再处理。当前确认的首次部署可使用;之后每次必须重新确认,不根据报错自动重置。外部 SaaS、RabbitMQ、OSS 数据不属于重置范围,不以重置伪造旧通话已结束或已提交结果。
退役资源只限本项目旧测试部署及自有 Mock 设施。不得删除整个 `~/.local`、原生 Asterisk、其他项目服务或未确认归属的容器;必要凭据和原始证据保留,旧运行数据按已授权范围备份后处理。
## 启动及验证
用户已授权本次目标中的真实 SaaS 全部获批任务,不再逐通询问。程序仍核对任务归属、时段、线路、额度、并发、签发期限和消息身份;不自行生成外呼指令,不盲重拨或换线重拨。
验收包括三项服务安装、启用并启动,Asterisk/Agent `enabled+active`,外部 RabbitMQ 登录/通道、Agent监听、Asterisk基础运行及版本/hash。SaaS 未启动不阻断部署:Dispatcher可明确等待或自动重试,新执行准入保持关闭;SaaS读取、其业务队列供给、业务会话及新SIP加载延期,不把此状态写成业务已就绪。证书、配置、MQ登录及Asterisk故障不得归为SaaS不可用。无获批任务、未接通或AI未调用时如实记录未验证。
## 一键入口
```bash
python3 deploys/preprod/preprod.py deploy
python3 deploys/preprod/preprod.py status
python3 deploys/preprod/preprod.py stop
python3 deploys/preprod/preprod.py start
# 仅本次首次部署已获授权;之后先取得新的重置确认:
python3 deploys/preprod/preprod.py deploy --reset-state --confirm-reset
```
部署要求本地main干净且已推送,源码固定为远程main,文件上传后再次校验SHA-256。`bootstrap` 仅导入原生Asterisk/TLS等主机基础元数据,并按根目录env覆盖业务来源;`retire-mocks` 在停服务、零活动通话和备份后退役已确认的本项目旧设施。普通重部署不重置数据。
自动化测试可以使用短时夹具和资源,结束后清理;它们不是预生产的服务依赖。生产发布门禁仍独立,不能把此次预生产部署写成生产签收。
+169
View File
@@ -0,0 +1,169 @@
#!/usr/bin/env python3
"""Deploy only the registered nonproduction node; never execute source env files."""
import argparse
import hashlib
import json
import os
from pathlib import Path
import shlex
import stat
import subprocess
import tempfile
import urllib.parse
import uuid
HERE = Path(__file__).resolve().parent
REPO = HERE.parent.parent
def private_text(path):
path = Path(path)
s = path.stat()
if path.is_symlink() or s.st_uid != os.getuid() or stat.S_IMODE(s.st_mode) != 0o600:
raise ValueError('private source ownership or permissions invalid')
return path.read_text()
def literal(value):
value = value.strip()
if value.startswith(('"', "'")):
try:
parts = shlex.split(value)
except ValueError:
raise ValueError('invalid quoted configuration') from None
if len(parts) != 1:
raise ValueError('ambiguous quoted configuration')
return parts[0]
return value
def read_env(path):
result = {}
for line in private_text(path).splitlines():
line = line.strip()
if not line or line.startswith(('#', '//', ';')):
continue
if '=' not in line:
raise ValueError('unsupported env syntax')
key, value = line.split('=', 1)
key = key.strip()
if key in result:
raise ValueError('duplicate env field')
result[key] = literal(value)
return result
def read_oss(path):
result, section = {}, ''
for raw in private_text(path).splitlines():
if not raw.strip() or raw.lstrip().startswith('#'):
continue
if ':' not in raw:
raise ValueError('unsupported OSS source syntax')
key, value = raw.strip().split(':', 1)
if key == 'RAM' and not value.strip():
section = 'RAM'
continue
# This historical file is colon-delimited, not YAML: RAM children
# follow the RAM heading even when they are not indented.
name = section + '.' + key if section == 'RAM' and key in ('username', 'accessKeyId', 'accessKeySecret') else key
if name in result:
raise ValueError('duplicate OSS field')
result[name] = literal(value)
needed = {'bucket', 'Endpoint', 'Region', 'RAM.username', 'RAM.accessKeyId', 'RAM.accessKeySecret'}
if not needed.issubset(result) or any(not result[k] for k in needed):
raise ValueError('missing required OSS field')
return result
def sources(root):
root = Path(root)
s, q, o = read_env(root / 'saas.env'), read_env(root / 'rabbitmq.env'), read_oss(root / 'aliyun-oss.env')
for values, needed in ((s, {'DispatcherUUID', 'DispatcherKEY', 'SaaSBaseURL'}), (q, {'RABBITMQ_HOST', 'RABBITMQ_PORT', 'RABBITMQ_USER', 'RABBITMQ_PASS', 'RABBITMQ_VHOST'})):
if set(values) != needed or any(not values[k] for k in needed):
raise ValueError('missing or unsupported source field')
uuid.UUID(s['DispatcherUUID'])
parsed = urllib.parse.urlsplit(s['SaaSBaseURL'])
if parsed.scheme not in ('http', 'https') or not parsed.hostname or parsed.username or parsed.query or parsed.fragment:
raise ValueError('invalid SaaS base URL')
port = int(q['RABBITMQ_PORT'])
if not 1 <= port <= 65535 or any(c in q['RABBITMQ_HOST'] for c in '/@\r\n '):
raise ValueError('invalid RabbitMQ endpoint')
quote = lambda v: urllib.parse.quote(v, safe='')
host = q['RABBITMQ_HOST']
if ':' in host:
host = '[' + host + ']'
mq_url = 'amqp://' + quote(q['RABBITMQ_USER']) + ':' + quote(q['RABBITMQ_PASS']) + '@' + host + ':' + str(port) + '/' + quote(q['RABBITMQ_VHOST'])
return {'dispatcher_id': s['DispatcherUUID'], 'secret': s['DispatcherKEY'], 'saas_url': s['SaaSBaseURL'], 'mq_url': mq_url, 'oss': o, 'probe': {'saas_url': s['SaaSBaseURL'], 'dispatcher_id': s['DispatcherUUID'], 'secret': s['DispatcherKEY'], 'mq_host': q['RABBITMQ_HOST'], 'mq_port': port, 'mq_user': q['RABBITMQ_USER'], 'mq_password': q['RABBITMQ_PASS'], 'mq_vhost': q['RABBITMQ_VHOST']}}
def checked(args, **kwargs):
r = subprocess.run(args, capture_output=True, text=True, **kwargs)
if r.returncode:
if args[0] == 'ssh':
try:
report = json.loads(r.stdout)
if report.get('ok') is False:
return r.stdout.strip()
except (ValueError, AttributeError):
pass
raise ValueError('command failed: ' + Path(args[0]).name)
return r.stdout.strip()
def main():
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument('action', choices=['deploy', 'start', 'stop', 'status', 'retire-mocks', 'bootstrap'])
parser.add_argument('--host', default='server.sip')
parser.add_argument('--known-hosts', type=Path, default=REPO / '.local/agent-call-known_hosts')
parser.add_argument('--source-dir', type=Path, default=REPO)
parser.add_argument('--reset-state', action='store_true')
parser.add_argument('--confirm-reset', action='store_true', help='explicit current operator confirmation; never automatic')
args = parser.parse_args()
if args.host != 'server.sip':
parser.error('this deployment is authorized only for server.sip')
if args.reset_state and (args.action != 'deploy' or not args.confirm_reset):
parser.error('--reset-state requires deploy and --confirm-reset')
if not args.known_hosts.is_file():
parser.error('registered known_hosts is required')
opts = ['-o', 'BatchMode=yes', '-o', 'StrictHostKeyChecking=yes', '-o', 'UserKnownHostsFile=' + str(args.known_hosts.resolve()), '-o', 'GlobalKnownHostsFile=/dev/null', '-o', 'HostKeyAlgorithms=ssh-ed25519', '-o', 'UpdateHostKeys=no', '-o', 'ConnectTimeout=15']
ssh = lambda cmd, **kw: checked(['ssh', *opts, args.host, cmd], timeout=180, **kw)
tool_dir = '/home/rogee/.local/share/go-sip-tools'
ssh("python3 -c \"from pathlib import Path;import os;p=Path('" + tool_dir + "');p.mkdir(mode=0o700,exist_ok=True);os.chmod(p,0o700)\"")
checked(['scp', '-q', *opts, str(HERE / 'remote.py'), args.host + ':' + tool_dir + '/preprod-remote.py'], timeout=30)
payload = {'action': args.action, 'reset_state': args.reset_state, 'confirm_reset': args.confirm_reset}
if args.action in ('deploy', 'bootstrap', 'retire-mocks'):
payload['source'] = sources(args.source_dir)
if args.action == 'deploy':
branch = checked(['git', 'branch', '--show-current'], cwd=REPO)
if branch != 'main' or checked(['git', 'status', '--porcelain'], cwd=REPO):
raise ValueError('deploy requires a clean main checkout')
commit = checked(['git', 'rev-parse', 'HEAD'], cwd=REPO)
remote_commit = checked(['git', 'ls-remote', 'origin', 'refs/heads/main'], cwd=REPO).split()[0]
if commit != remote_commit:
raise ValueError('main is not pushed')
build_root = Path(tempfile.mkdtemp(prefix='go-sip-preprod-build-'))
try:
for name, source in (('sip-go-agent', './cmd/sip-go-agent'), ('preprod-probe', './deploys/preprod')):
env = dict(os.environ, CGO_ENABLED='0', GOOS='linux', GOARCH='amd64')
checked(['go', 'build', '-trimpath', '-o', str(build_root / name), source], cwd=REPO, env=env, timeout=180)
staging = tool_dir + '/' + commit
ssh('mkdir -m 700 -p ' + shlex.quote(staging))
for name in ('sip-go-agent', 'preprod-probe'):
checked(['scp', '-q', *opts, str(build_root / name), args.host + ':' + staging + '/' + name], timeout=60)
payload.update({'commit': commit, 'staging': staging, 'hashes': {name: hashlib.sha256((build_root / name).read_bytes()).hexdigest() for name in ('sip-go-agent', 'preprod-probe')}})
finally:
__import__('shutil').rmtree(build_root)
raw = ssh('python3 ' + shlex.quote(tool_dir + '/preprod-remote.py'), input=json.dumps(payload))
result = json.loads(raw)
print(json.dumps(result, indent=2))
if not result.get('ok'):
raise SystemExit(1)
if __name__ == '__main__':
try:
main()
except (ValueError, OSError, subprocess.TimeoutExpired) as e:
print(json.dumps({'ok': False, 'error_class': type(e).__name__, 'reason': str(e) if isinstance(e, ValueError) else 'deployment operation failed'}))
raise SystemExit(1)
+260
View File
@@ -0,0 +1,260 @@
package main
import (
"context"
"encoding/json"
"errors"
"flag"
"fmt"
"net"
"net/http"
"os"
"time"
"git.ipao.vip/rogee/go-sip/internal/config"
"git.ipao.vip/rogee/go-sip/internal/configread"
"git.ipao.vip/rogee/go-sip/internal/rpc"
amqp "github.com/rabbitmq/amqp091-go"
"github.com/santhosh-tekuri/jsonschema/v6"
)
type input struct {
SaaSURL string `json:"saas_url"`
DispatcherID string `json:"dispatcher_id"`
Secret string `json:"secret"`
MQHost string `json:"mq_host"`
MQPort int `json:"mq_port"`
MQUser string `json:"mq_user"`
MQPassword string `json:"mq_password"`
MQVhost string `json:"mq_vhost"`
ConnectionName string `json:"connection_name"`
HoldMQMS int `json:"hold_mq_ms"`
}
type report struct {
Success bool `json:"success"`
Phase string `json:"phase"`
HTTPStatus int `json:"http_status"`
ErrorClass string `json:"error_class,omitempty"`
SchemaLocations []string `json:"schema_locations,omitempty"`
Trunks int `json:"trunks"`
Providers int `json:"providers"`
Tasks int `json:"tasks"`
ValidatedTasks int `json:"validated_tasks"`
MQConnected bool `json:"mq_connected"`
MQChannelClosed bool `json:"mq_channel_closed"`
MQConnectionClosed bool `json:"mq_connection_closed"`
}
type transport struct{ report *report }
func (t transport) RoundTrip(req *http.Request) (*http.Response, error) {
response, err := http.DefaultTransport.RoundTrip(req)
if response != nil {
t.report.HTTPStatus = response.StatusCode
}
return response, err
}
func rejected(r *report, err error) {
r.Success = false
r.ErrorClass = fmt.Sprintf("%T", err)
var validation *jsonschema.ValidationError
if errors.As(err, &validation) {
var visit func(*jsonschema.ValidationError)
visit = func(e *jsonschema.ValidationError) {
if len(e.Causes) == 0 {
r.SchemaLocations = append(r.SchemaLocations, fmt.Sprint(e.InstanceLocation))
}
for _, cause := range e.Causes {
visit(cause)
}
}
visit(validation)
}
}
func check(in input) (r report) {
ctx, cancel := context.WithTimeout(context.Background(), 90*time.Second)
defer cancel()
r.Phase = "mq_connect"
uri := amqp.URI{Scheme: "amqp", Host: in.MQHost, Port: in.MQPort, Username: in.MQUser, Password: in.MQPassword, Vhost: in.MQVhost}
conn, err := amqp.DialConfig(uri.String(), amqp.Config{Heartbeat: 5 * time.Second, Locale: "en_US", Properties: amqp.Table{"connection_name": in.ConnectionName}, Dial: func(network, addr string) (net.Conn, error) {
c, err := (&net.Dialer{Timeout: 8 * time.Second}).DialContext(ctx, network, addr)
if err != nil {
return nil, err
}
if err := c.SetDeadline(time.Now().Add(15 * time.Second)); err != nil {
return nil, errors.Join(err, c.Close())
}
return c, nil
}})
if err != nil {
rejected(&r, err)
return
}
r.MQConnected = true
if in.HoldMQMS > 0 && in.HoldMQMS <= 10000 {
time.Sleep(time.Duration(in.HoldMQMS) * time.Millisecond)
}
r.Phase = "mq_channel"
channel, err := conn.Channel()
if err != nil {
rejected(&r, errors.Join(err, conn.Close()))
return
}
if err := channel.Close(); err != nil {
rejected(&r, errors.Join(err, conn.Close()))
return
}
r.MQChannelClosed = true
if err := conn.Close(); err != nil {
rejected(&r, err)
return
}
r.MQConnectionClosed = true
r.Phase = "saas_client"
client, err := configread.NewClient(in.SaaSURL, in.DispatcherID, in.Secret, &http.Client{Timeout: 15 * time.Second, Transport: transport{&r}})
if err != nil {
rejected(&r, err)
return
}
r.Phase = "sip"
sip, err := client.ReadSIP(ctx)
if err != nil {
rejected(&r, err)
return
}
var trunks []json.RawMessage
if err := json.Unmarshal(sip.Trunks, &trunks); err != nil {
rejected(&r, err)
return
}
r.Trunks = len(trunks)
r.Phase = "providers"
providers, err := client.ReadProviders(ctx)
if err != nil {
rejected(&r, err)
return
}
r.Providers = len(providers)
var tasks []configread.DiscoveredTask
cursor := ""
seen := map[string]bool{}
for {
r.Phase = "tasks"
page, err := client.ReadTasks(ctx, cursor)
if err != nil {
rejected(&r, err)
return
}
for _, task := range page.Tasks {
if seen[task.TaskID] {
rejected(&r, errors.New("duplicate discovery task"))
return
}
seen[task.TaskID] = true
tasks = append(tasks, task)
}
if page.Cursor == "" {
break
}
cursor = page.Cursor
}
r.Tasks = len(tasks)
for _, task := range tasks {
r.Phase = "task_and_quota"
if _, err := client.ReadTask(ctx, task.TaskID, task.TenantID, sip, providers); err != nil {
rejected(&r, err)
return
}
r.ValidatedTasks++
}
r.Phase = "complete"
r.Success = true
return
}
// environmentCheck opens only explicit local configuration and certificate
// files. It never opens SQLite, network connections, or cloud requests.
func environmentCheck(role string) (r report) {
var ca, cert, key string
r.Phase = role + "_environment"
if role == "agent" {
s, err := config.LoadAgentEnvironment("nonprod-real")
if err != nil {
rejected(&r, err)
return
}
ca, cert, key = s.CAFile, s.CertFile, s.KeyFile
} else if role == "dispatcher" {
s, err := config.LoadDispatcherRuntimeEnvironment("nonprod-real")
if err != nil {
rejected(&r, err)
return
}
r.Phase = "agent_inventory"
if _, err := config.LoadMockAgentEndpoint(s.AgentEndpointsFile); err != nil {
rejected(&r, err)
return
}
r.Phase = "oss_configuration"
if _, err := config.LoadNonprodRealOSSConfig(s.OSSConfigFile, s.DispatcherID); err != nil {
rejected(&r, err)
return
}
ca, cert, key = s.CAFile, s.CertFile, s.KeyFile
} else {
rejected(&r, errors.New("unknown environment role"))
return
}
r.Phase = role + "_tls"
caBytes, err := os.ReadFile(ca)
if err != nil {
rejected(&r, err)
return
}
certBytes, err := os.ReadFile(cert)
if err != nil {
rejected(&r, err)
return
}
keyBytes, err := os.ReadFile(key)
if err != nil {
rejected(&r, err)
return
}
if _, err := rpc.NewServerTLSConfig(caBytes, certBytes, keyBytes); err != nil {
rejected(&r, err)
return
}
r.Phase = "environment_complete"
r.Success = true
return
}
func main() {
role := flag.String("environment", "", "pure agent or dispatcher deployment configuration inspection")
flag.Parse()
if *role != "" {
r := environmentCheck(*role)
if err := json.NewEncoder(os.Stdout).Encode(r); err != nil {
os.Exit(2)
}
if !r.Success {
os.Exit(1)
}
return
}
var in input
d := json.NewDecoder(os.Stdin)
d.DisallowUnknownFields()
if err := d.Decode(&in); err != nil {
_ = json.NewEncoder(os.Stdout).Encode(report{Phase: "input", ErrorClass: fmt.Sprintf("%T", err)})
os.Exit(1)
}
r := check(in)
if err := json.NewEncoder(os.Stdout).Encode(r); err != nil {
os.Exit(2)
}
if !r.Success {
os.Exit(1)
}
}
+30
View File
@@ -0,0 +1,30 @@
package main
import (
"encoding/json"
"errors"
"strings"
"testing"
)
func TestEnvironmentRejectsMissingConfigurationWithoutNetworking(t *testing.T) {
t.Setenv("DISPATCHER_ID", "")
t.Setenv("AGENT_ID", "")
for _, role := range []string{"agent", "dispatcher", "unknown"} {
got := environmentCheck(role)
if got.Success || got.MQConnected || got.ErrorClass == "" {
t.Fatalf("unexpected result: %+v", got)
}
}
}
func TestDiagnosticNeverSerializesErrorContents(t *testing.T) {
r := report{Phase: "sip"}
rejected(&r, errors.New("synthetic-sensitive-value"))
encoded, err := json.Marshal(r)
if err != nil {
t.Fatal(err)
}
if strings.Contains(string(encoded), "synthetic-sensitive-value") {
t.Fatal("diagnostic leaked error content")
}
}
+488
View File
@@ -0,0 +1,488 @@
#!/usr/bin/env python3
"""Private-input remote operations for the single approved preproduction node."""
import base64
import configparser
import datetime
import hashlib
import json
import os
from pathlib import Path
import re
import shlex
import shutil
import socket
import stat
import subprocess
import sys
import time
import urllib.parse
import urllib.request
import uuid
UNITS = ('go-sip-asterisk.service', 'go-sip-agent.service', 'go-sip-dispatcher.service')
def paths(home):
return {'home': home, 'config': home / '.config/go-sip', 'state': home / '.local/share/go-sip', 'releases': home / '.local/opt/go-sip/releases', 'current': home / '.local/opt/go-sip/current', 'backups': home / '.local/share/go-sip-backups', 'units': home / '.config/systemd/user'}
def run(args, **kwargs):
r = subprocess.run(args, capture_output=True, text=True, **kwargs)
if r.returncode:
raise ValueError('command failed: ' + Path(args[0]).name)
return r.stdout.strip()
def mkdir(p):
if p.is_symlink():
raise ValueError('symlink directory refused')
p.mkdir(parents=True, mode=0o700, exist_ok=True)
if p.stat().st_uid != os.getuid():
raise ValueError('directory owner mismatch')
p.chmod(0o700)
def write_private(p, data):
if p.is_symlink():
raise ValueError('symlink file refused')
tmp = p.with_name(p.name + '.new')
fd = os.open(tmp, os.O_WRONLY | os.O_CREAT | os.O_TRUNC, 0o600)
with os.fdopen(fd, 'w') as stream:
stream.write(data)
stream.flush()
os.fsync(stream.fileno())
os.replace(tmp, p)
p.chmod(0o600)
def private_text(p):
s = p.stat()
if p.is_symlink() or s.st_uid != os.getuid() or stat.S_IMODE(s.st_mode) != 0o600:
raise ValueError('private file ownership or permissions invalid')
return p.read_text()
def read_env(p):
result = {}
for raw in private_text(p).splitlines():
raw = raw.strip()
if not raw or raw.startswith('#'):
continue
if '=' not in raw:
raise ValueError('invalid runtime env')
k, v = raw.split('=', 1)
if k in result:
raise ValueError('duplicate runtime env key')
parts = shlex.split(v)
if len(parts) != 1:
raise ValueError('ambiguous runtime env value')
result[k] = parts[0]
return result
def unit_properties(unit):
r = subprocess.run(['systemctl', '--user', 'show', unit, '--property=LoadState,ActiveState,SubState,UnitFileState,FragmentPath,EnvironmentFiles,ExecMainStatus'], capture_output=True, text=True)
if r.returncode and 'LoadState=not-found' not in r.stdout.splitlines():
raise ValueError('systemd status inspection failed')
return dict(line.split('=', 1) for line in r.stdout.splitlines() if '=' in line)
def legacy_env(unit, home):
meta = unit_properties(unit)
fragment = Path(meta['FragmentPath'])
if not fragment.is_file() or fragment.parent != home / '.config/systemd/user':
raise ValueError('legacy unit ownership ambiguous')
references = re.findall(r'(\S+) \(ignore_errors=(?:yes|no)\)', meta['EnvironmentFiles'])
if len(references) != 1:
raise ValueError('legacy environment references ambiguous')
return read_env(Path(references[0]))
def render_roles(agent, dispatcher, source, p):
common = {'MTLS_CA_FILE', 'MTLS_CERT_FILE', 'MTLS_KEY_FILE', 'MTLS_PEER_CERT_FINGERPRINTS'}
agent_keys = common | {'AGENT_ID', 'CELL_ID', 'DISPATCHER_GRPC_SERVER_NAME', 'ASTERISK_CONFIG_DIR', 'ASTERISK_BIN', 'ASTERISK_LIBRARY_DIR', 'AGENT_HEP_LISTEN_ADDR'}
roles = {'agent': {k: v for k, v in agent.items() if k in agent_keys}, 'dispatcher': {k: v for k, v in dispatcher.items() if k in common}}
roles['agent'].update({'AGENT_SESSION_PATH': str(p['state'] / 'agent-session.json'), 'AGENT_RECOVERY_ROOT': str(p['state'] / 'agent-recovery'), 'DISPATCHER_ID': source['dispatcher_id'], 'DISPATCHER_GRPC_ENDPOINT': '127.0.0.1:19444', 'AGENT_GRPC_LISTEN': '127.0.0.1:19443', 'AGENT_EVIDENCE_ROOT': str(p['backups'] / 'debug-calls')})
roles['dispatcher'].update({'DISPATCHER_ID': source['dispatcher_id'], 'SAAS_BASE_URL': source['saas_url'], 'DISPATCHER_SECRET_KEY': source['secret'], 'RABBITMQ_URL': source['mq_url'], 'DISPATCHER_SQLITE_PATH': str(p['state'] / 'dispatcher.sqlite'), 'DISPATCHER_GRPC_LISTEN': '127.0.0.1:19444', 'DISPATCHER_AGENT_ENDPOINTS_FILE': str(p['config'] / 'agent-endpoints.json'), 'DISPATCHER_OSS_CONFIG_FILE': str(p['config'] / 'oss.json')})
return roles
def reset_requested(payload):
return payload.get('reset_state') is True
def require_reset(requested, confirmed, channels):
if not requested or not confirmed or channels != 0:
raise ValueError('reset requires explicit authorization and zero active channels')
def dependency_state(report):
if report.get('success'):
return 'verified'
if report.get('phase') in ('sip', 'providers', 'tasks', 'task_and_quota') and report.get('http_status', 0) in (0, 404, 502, 503, 504):
return 'saas_pending'
if report.get('phase', '').startswith('mq_'):
return 'mq_failed'
return 'dependency_failed'
def ari(home, endpoint='channels'):
root = home / '.config/go-sip-asterisk'
a = configparser.ConfigParser(interpolation=None)
h = configparser.ConfigParser(interpolation=None)
a.read_string(private_text(root / 'ari.conf'))
h.read_string(private_text(root / 'http.conf'))
secret = private_text(root / 'ari-secret').strip()
users = [s for s in a.sections() if s != 'general' and a.get(s, 'password', fallback=None) == secret]
if len(users) != 1 or h.get('general', 'enabled').lower() != 'yes' or h.get('general', 'bindaddr') not in ('127.0.0.1', '0.0.0.0'):
raise ValueError('native ARI configuration uncertain')
port = int(h.get('general', 'bindport'))
authorization = base64.b64encode((users[0] + ':' + secret).encode()).decode()
req = urllib.request.Request('http://127.0.0.1:' + str(port) + '/ari/' + endpoint, headers={'Authorization': 'Basic ' + authorization})
opener = urllib.request.build_opener(urllib.request.ProxyHandler({}))
with opener.open(req, timeout=10) as response:
data = json.load(response)
if not isinstance(data, list):
raise ValueError('unexpected native ARI response')
return data
def stop_for_change(p, legacy=False):
if legacy:
text = run(['systemctl', '--user', 'list-unit-files', '--type=service', '--no-legend'])
units = [line.split()[0] for line in text.splitlines() if line.split()[0].startswith('go-sip-nonprod-') and any(s in line.split()[0] for s in ('dispatcher', 'agent', 'saas'))]
else:
units = ['go-sip-dispatcher.service', 'go-sip-agent.service']
dispatchers = [u for u in units if 'dispatcher' in u or 'saas' in u]
for unit in dispatchers:
if unit_properties(unit).get('FragmentPath'):
run(['systemctl', '--user', 'stop', unit])
if ari(p['home']):
raise ValueError('active calls: Agent remains running; no data changes allowed')
for unit in units:
if 'agent' in unit and unit_properties(unit).get('FragmentPath'):
run(['systemctl', '--user', 'stop', unit])
return units
def backup_dir(p, source, label):
mkdir(p['backups'])
tag = datetime.datetime.now(datetime.UTC).strftime('%Y%m%dT%H%M%S%fZ')
target = p['backups'] / (label + '-' + tag + '.tar')
if source.is_symlink() or source == p['home'] or not source.is_dir():
raise ValueError('unsafe archive source')
run(['sudo', '-n', 'tar', '--one-file-system', '-C', str(source.parent), '-cf', str(target), source.name], timeout=180)
run(['sudo', '-n', 'chown', str(os.getuid()) + ':' + str(os.getgid()), str(target)])
target.chmod(0o600)
if target.stat().st_size < 512:
raise ValueError('archive is incomplete')
# tar's own reread validates the archive before any source removal.
run(['tar', '-tf', str(target)], timeout=120)
with target.open('rb') as stream:
digest = hashlib.file_digest(stream, 'sha256').hexdigest()
write_private(target.with_suffix('.json'), json.dumps({'source': str(source), 'archive': str(target), 'sha256': digest, 'created_at': tag}, indent=2))
return digest
def bootstrap(p, source):
for key in ('config', 'state', 'releases', 'backups', 'units'):
mkdir(p[key])
baseline = p['config'] / 'host-baseline.json'
if baseline.exists():
old = json.loads(private_text(baseline))
agent, dispatcher = old['agent'], old['dispatcher']
else:
agent = legacy_env('go-sip-nonprod-agent-20261004.service', p['home'])
dispatcher = legacy_env('go-sip-nonprod-dispatcher-20261004.service', p['home'])
# This only imports installed Asterisk/TLS/platform metadata. Business
# addresses and credentials are replaced from current root sources.
for d in (agent, dispatcher):
for k in list(d):
if k.startswith(('SAAS_', 'RABBITMQ_', 'OSS_')) or 'MOCK' in k or 'FIXTURE' in k:
del d[k]
tls = p['config'] / 'tls'
mkdir(tls)
for role, d in (('agent', agent), ('dispatcher', dispatcher)):
for key in ('MTLS_CA_FILE', 'MTLS_CERT_FILE', 'MTLS_KEY_FILE'):
old_file = Path(d[key])
data = private_text(old_file)
new_file = tls / (role + '-' + old_file.name)
write_private(new_file, data)
d[key] = str(new_file)
write_private(baseline, json.dumps({'agent': agent, 'dispatcher': dispatcher}, indent=2))
endpoint_file = p['config'] / 'agent-endpoints.json'
if endpoint_file.exists():
inventory = json.loads(private_text(endpoint_file))
else:
inventory = json.loads(private_text(Path(dispatcher['DISPATCHER_AGENT_ENDPOINTS_FILE'])))
if len(inventory) != 1 or inventory[0]['agent_id'] != agent['AGENT_ID'] or inventory[0]['cell_id'] != agent['CELL_ID']:
raise ValueError('single Agent endpoint identity mismatch')
inventory[0]['address'] = '127.0.0.1:19443'
write_private(endpoint_file, json.dumps(inventory, indent=2))
roles = render_roles(agent, dispatcher, source, p)
# Read native mirror destination, rather than guessing which old Agent port
# currently receives HEP. No static transport or trunk changes are made.
hep = p['home'] / '.config/go-sip-asterisk/hep.conf'
if hep.exists():
cfg = configparser.ConfigParser(interpolation=None)
cfg.read(hep)
if cfg.has_option('general', 'capture_address'):
roles['agent']['AGENT_HEP_LISTEN_ADDR'] = cfg.get('general', 'capture_address')
o = source['oss']
endpoint = o['Endpoint']
if '://' not in endpoint:
endpoint = 'https://' + endpoint
oss = {'dispatcher_id': source['dispatcher_id'], 'oss': {'bucket': o['bucket'], 'endpoint': endpoint, 'region': o['Region'], 'object_prefix': 'call-recordings', 'access_key_id_env': 'OSS_ACCESS_KEY_ID', 'access_key_secret_env': 'OSS_ACCESS_KEY_SECRET', 'max_asset_bytes': 64 * 1024 * 1024}}
roles['agent']['AGENT_OSS_ALLOWED_HOST'] = o['bucket'] + '.' + urllib.parse.urlsplit(endpoint).hostname
roles['dispatcher']['OSS_ACCESS_KEY_ID'] = o['RAM.accessKeyId']
roles['dispatcher']['OSS_ACCESS_KEY_SECRET'] = o['RAM.accessKeySecret']
for role, d in roles.items():
for v in d.values():
if any(c in str(v) for c in '\x00\r\n'):
raise ValueError('invalid multiline runtime value')
write_private(p['config'] / (role + '.env'), ''.join(k + '=' + json.dumps(str(v), ensure_ascii=False) + '\n' for k, v in sorted(d.items())))
write_private(p['config'] / 'oss.json', json.dumps(oss, indent=2))
return roles
def host_check(p):
os_release = dict(line.split('=', 1) for line in Path('/etc/os-release').read_text().splitlines() if '=' in line)
if os_release.get('ID', '').strip('"') != 'debian' or os_release.get('VERSION_ID', '').strip('"') != '13' or run(['uname', '-m']) != 'x86_64':
raise ValueError('host is not Debian 13 amd64')
if os.getuid() == 0 or p['home'].name != 'rogee':
raise ValueError('runtime user must be rogee')
for tool in ('python3', 'openssl', 'tcpdump', 'tshark'):
if not shutil.which(tool):
raise ValueError('required host tool missing: ' + tool)
run(['sudo', '-n', 'true'])
channels = len(ari(p['home']))
if channels:
raise ValueError('active native calls prevent deployment')
native = unit_properties('go-sip-asterisk.service')
if native.get('ActiveState') != 'active' or native.get('UnitFileState') != 'enabled':
raise ValueError('native Asterisk is not enabled and active')
return {'os': 'debian13', 'architecture': 'amd64', 'runtime_user': 'rogee', 'active_channels': channels}
def retire(p, payload):
host_check(p)
bootstrap(p, payload['source'])
units = stop_for_change(p, legacy=True)
archived = []
for unit in units:
meta = unit_properties(unit)
fragment = Path(meta['FragmentPath'])
if fragment.parent != p['units'] or fragment.is_symlink() or 'sip' not in fragment.read_text():
raise ValueError('legacy unit ownership unclear')
run(['systemctl', '--user', 'disable', unit])
obsolete_binaries = set()
for unit in units:
detail = run(['systemctl', '--user', 'show', unit, '--property=ExecStart', '--value'])
match = re.search(r'path=(\S+) ;', detail)
if match:
executable = Path(match.group(1))
if executable.is_file() and not executable.is_symlink() and executable.is_relative_to(p['home']) and executable.name in ('sip-go-agent', 'saas-mock'):
obsolete_binaries.add(executable)
artifact_backup = p['backups'] / ('legacy-artifacts-' + datetime.datetime.now(datetime.UTC).strftime('%Y%m%dT%H%M%SZ'))
mkdir(artifact_backup)
for old in obsolete_binaries:
target = artifact_backup / (hashlib.sha256(str(old).encode()).hexdigest()[:12] + '-' + old.name)
shutil.copy2(old, target); target.chmod(0o600)
if hashlib.sha256(old.read_bytes()).digest() != hashlib.sha256(target.read_bytes()).digest():
raise ValueError('legacy artifact backup mismatch')
old.unlink()
# The registered host uses its project-owned system RabbitMQ, not Docker.
local_mq = False
ctl = Path('/usr/sbin/rabbitmqctl')
if ctl.exists() and run(['systemctl', 'show', 'rabbitmq-server.service', '--property=ActiveState', '--value']) == 'active':
legacy = legacy_env('go-sip-nonprod-dispatcher-20261004.service', p['home'])
old_vhost = urllib.parse.unquote(urllib.parse.urlsplit(legacy['RABBITMQ_URL']).path[1:]) or '/'
def query(args):
return json.loads(run(['sudo', '-n', str(ctl), '-q', *args, '--formatter=json'], timeout=30))
vhosts = [v['name'] for v in query(['list_vhosts', 'name'])]
connections = query(['list_connections', 'vhost'])
if any(v not in ('/', old_vhost) for v in vhosts) or query(['list_queues', '-p', '/', 'name']) or any(c['vhost'] != old_vhost for c in connections):
raise ValueError('local RabbitMQ contains unowned resources; retained')
# Correlate a live unique connection, not merely IP addresses: public
# addresses/NAT must not make us uninstall the configured real broker.
marker = 'go-sip-retirement-check-' + uuid.uuid4().hex
config = dict(payload['source']['probe'], connection_name=marker, hold_mq_ms=10000)
probe = p['home'] / '.local/share/go-sip-tools/preprod-probe'
child = subprocess.Popen([str(probe)], stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True)
child.stdin.write(json.dumps(config)); child.stdin.close()
time.sleep(1)
connected = query(['list_connections', 'client_properties'])
target_is_local = any(marker in json.dumps(c) for c in connected)
still_held = child.poll() is None
child.stdout.read(); child.wait(timeout=30)
if not still_held or target_is_local:
raise ValueError('configured RabbitMQ target may be local; no broker removal permitted')
run(['sudo', '-n', 'systemctl', 'stop', 'rabbitmq-server.service'])
backup_dir(p, Path('/var/lib/rabbitmq'), 'mock-rabbitmq-data')
backup_dir(p, Path('/etc/rabbitmq'), 'mock-rabbitmq-config')
run(['sudo', '-n', 'systemctl', 'disable', 'rabbitmq-server.service'])
run(['sudo', '-n', 'apt-get', 'remove', '-y', 'rabbitmq-server'], timeout=180)
run(['sudo', '-n', 'rm', '-rf', '--', '/var/lib/rabbitmq', '/etc/rabbitmq'])
local_mq = True
# Archive legacy project data/config before removal; never touch all ~/.local.
for parent in (p['home'] / '.local/share', p['home'] / '.config'):
for source in sorted(parent.glob('go-sip-nonprod-*')):
if source.is_dir() and not source.is_symlink():
digest = backup_dir(p, source, 'legacy')
archived.append(digest)
run(['sudo', '-n', 'rm', '-rf', '--', str(source)])
unit_backup = p['backups'] / ('legacy-units-' + datetime.datetime.now(datetime.UTC).strftime('%Y%m%dT%H%M%SZ'))
mkdir(unit_backup)
for unit in units:
fragment = Path(unit_properties(unit)['FragmentPath'])
shutil.copy2(fragment, unit_backup / fragment.name)
(unit_backup / fragment.name).chmod(0o600)
fragment.unlink()
run(['systemctl', '--user', 'daemon-reload'])
evidence = {'ok': True, 'legacy_units_removed': len(units), 'legacy_binaries_removed': len(obsolete_binaries), 'legacy_directories_archived': len(archived), 'archive_hashes': archived, 'mock_broker_removed': local_mq, 'active_channels': 0, 'external_services_modified': False}
write_private(p['backups'] / 'retirement.json', json.dumps(evidence, indent=2))
return evidence
def validate_environments(probe, roles):
for role, values in roles.items():
environment = dict(os.environ, **values)
checked = subprocess.run([str(probe), '-environment', role], env=environment, capture_output=True, text=True, timeout=15)
try:
report = json.loads(checked.stdout)
except json.JSONDecodeError:
raise ValueError('environment inspection failed') from None
if checked.returncode or not report.get('success'):
raise ValueError('invalid ' + role + ' configuration at ' + report.get('phase', 'unknown'))
def install(payload, p):
facts = host_check(p)
source = payload['source']
# Changing ownership of existing durable data is never a normal deployment.
old_env = p['config'] / 'dispatcher.env'
if old_env.exists() and (p['state'] / 'dispatcher.sqlite').exists():
previous = read_env(old_env).get('DISPATCHER_ID')
if previous != source['dispatcher_id'] and not reset_requested(payload):
raise ValueError('Dispatcher identity change requires an explicitly confirmed reset')
stop_for_change(p)
if reset_requested(payload):
require_reset(True, payload.get('confirm_reset') is True, len(ari(p['home'])))
if p['state'].exists() and any(p['state'].iterdir()):
backup_dir(p, p['state'], 'reset')
shutil.rmtree(p['state'])
roles = bootstrap(p, source)
mkdir(p['state'] / 'agent-recovery')
release = p['releases'] / payload['commit']
mkdir(release)
staging = Path(payload['staging'])
if staging.parent != p['home'] / '.local/share/go-sip-tools' or staging.is_symlink():
raise ValueError('unexpected uploaded staging path')
for name, expected in payload['hashes'].items():
if name not in ('sip-go-agent', 'preprod-probe'):
raise ValueError('unexpected artifact')
candidate = staging / name
if candidate.is_symlink() or hashlib.sha256(candidate.read_bytes()).hexdigest() != expected:
raise ValueError('uploaded artifact checksum mismatch')
shutil.copy2(candidate, release / name)
(release / name).chmod(0o700)
validate_environments(release / 'preprod-probe', roles)
write_private(release / 'manifest.json', json.dumps({'commit': payload['commit'], 'hashes': payload['hashes'], 'production_approval': False}, indent=2))
current = p['current']
if current.exists() and not current.is_symlink():
raise ValueError('current artifact pointer is not a symlink')
tmp = current.with_name('current.new')
if tmp.exists() or tmp.is_symlink():
raise ValueError('unfinished artifact switch exists')
tmp.symlink_to(release)
os.replace(tmp, current)
for role in ('agent', 'dispatcher'):
unit = '[Unit]\nDescription=go-sip real preproduction ' + role + '\nAfter=go-sip-asterisk.service network-online.target\n\n[Service]\nType=simple\nEnvironmentFile=' + str(p['config'] / (role + '.env')) + '\nExecStart=' + str(current / 'sip-go-agent') + ' ' + role + ' --mode nonprod-real\nRestart=on-failure\nRestartSec=10\nUMask=0077\n\n[Install]\nWantedBy=default.target\n'
write_private(p['units'] / ('go-sip-' + role + '.service'), unit)
run(['systemctl', '--user', 'daemon-reload'])
run(['sudo', '-n', 'loginctl', 'enable-linger', p['home'].name])
run(['systemctl', '--user', 'enable', *UNITS])
# MQ is checked independently even when real SaaS has not started.
probe = subprocess.run([str(release / 'preprod-probe')], input=json.dumps(source['probe']), text=True, capture_output=True, timeout=110)
try:
dependency = json.loads(probe.stdout)
except json.JSONDecodeError:
raise ValueError('redacted external probe failed') from None
state = dependency_state(dependency)
if state not in ('verified', 'saas_pending') or not dependency.get('mq_connected') or not dependency.get('mq_connection_closed'):
raise ValueError('independent dependency failure; not attributable to unavailable SaaS')
write_private(p['config'] / 'dependency-status.json', json.dumps(dependency, indent=2))
run(['systemctl', '--user', 'start', 'go-sip-agent.service'])
time.sleep(2)
if unit_properties('go-sip-agent.service').get('ActiveState') != 'active':
raise ValueError('Agent startup failed independently of SaaS')
run(['systemctl', '--user', 'start', 'go-sip-dispatcher.service'])
time.sleep(2)
result = status(p)
result.update({'host': facts, 'dependency_state': state, 'reset_performed': reset_requested(payload), 'commit': payload['commit']})
write_private(p['backups'] / 'last-deployment.json', json.dumps(result, indent=2))
return result
def status(p):
units = {u: {k: v for k, v in unit_properties(u).items() if k in ('ActiveState', 'SubState', 'UnitFileState', 'ExecMainStatus')} for u in UNITS}
manifest = p['current'] / 'manifest.json'
release = json.loads(private_text(manifest)) if manifest.exists() else None
matches = False
if release:
matches = all(hashlib.sha256((p['current'] / name).read_bytes()).hexdigest() == digest for name, digest in release['hashes'].items())
listen = False
try:
with socket.create_connection(('127.0.0.1', 19443), timeout=3):
listen = True
except OSError:
pass # Report an explicit false, never substitute a successful listener.
dep = p['config'] / 'dependency-status.json'
dependency = dependency_state(json.loads(private_text(dep))) if dep.exists() else 'not_checked'
healthy = matches and listen and all(units[u]['ActiveState'] == 'active' and units[u]['UnitFileState'] == 'enabled' for u in UNITS[:2])
dispatcher = units[UNITS[2]]
waiting_reason = None
if dependency == 'saas_pending':
logs = run(['journalctl', '--user', '-u', 'go-sip-dispatcher.service', '_COMM=sip-go-agent', '-n', '1', '--no-pager', '-o', 'cat'])
if re.search(r'HTTP(?: status)?[ :]+(?:404|502|503|504)', logs):
waiting_reason = 'saas_http_unavailable'
elif 'NOT_FOUND' in logs and ('queue' in logs or 'exchange' in logs):
waiting_reason = 'saas_topology_pending'
elif '/internal/v1/dispatcher/' in logs and any(v in logs.lower() for v in ('timeout', 'connection refused', 'no such host')):
waiting_reason = 'saas_network_unavailable'
acceptable = dispatcher['UnitFileState'] == 'enabled' and (dispatcher['ActiveState'] == 'active' or (waiting_reason is not None and dispatcher['ActiveState'] in ('activating', 'failed')))
return {'ok': healthy and acceptable, 'services': units, 'artifact_hashes_match': matches, 'agent_listener': listen, 'active_channels': len(ari(p['home'])), 'dependency_state': dependency, 'business_ready': dependency == 'verified' and dispatcher['ActiveState'] == 'active', 'commit': release['commit'] if release else None, 'waiting_reason': waiting_reason, 'real_reboot_verified': False}
def main():
payload = json.load(sys.stdin)
p = paths(Path.home())
action = payload['action']
if action == 'bootstrap':
host_check(p)
bootstrap(p, payload['source'])
result = {'ok': True, 'fixed_paths_configured': True}
elif action == 'retire-mocks':
result = retire(p, payload)
elif action == 'deploy':
result = install(payload, p)
elif action == 'status':
result = status(p)
elif action == 'stop':
stop_for_change(p)
result = {'ok': True, 'stopped': True, 'data_modified': False}
elif action == 'start':
run(['systemctl', '--user', 'start', *UNITS])
time.sleep(2)
result = status(p)
else:
raise ValueError('unsupported operation')
print(json.dumps(result, indent=2))
if __name__ == '__main__':
try:
main()
except Exception as e:
print(json.dumps({'ok': False, 'error_class': type(e).__name__, 'reason': str(e) if isinstance(e, ValueError) else 'operation failed; no successful deployment claimed'}))
raise SystemExit(1)
+108
View File
@@ -0,0 +1,108 @@
import hashlib
import importlib.util
import json
from pathlib import Path
import shutil
import subprocess
import tempfile
import unittest
from unittest.mock import patch
HERE = Path(__file__).resolve().parent
spec = importlib.util.spec_from_file_location('remote', HERE / 'remote.py')
remote = importlib.util.module_from_spec(spec)
spec.loader.exec_module(remote)
class LifecycleTest(unittest.TestCase):
def setUp(self):
self.home = Path(tempfile.mkdtemp())
self.addCleanup(lambda: shutil.rmtree(self.home))
self.p = remote.paths(self.home)
for p in self.p.values():
if p != self.p['current']:
remote.mkdir(p)
self.staging = self.home / '.local/share/go-sip-tools' / ('a' * 40)
remote.mkdir(self.staging)
self.hashes = {}
for name in ('sip-go-agent', 'preprod-probe'):
(self.staging / name).write_bytes(b'synthetic executable')
self.hashes[name] = hashlib.sha256((self.staging / name).read_bytes()).hexdigest()
self.payload = {'commit': 'a' * 40, 'staging': str(self.staging), 'hashes': self.hashes, 'source': {'dispatcher_id': 'synthetic', 'probe': {}}}
self.state = self.p['state'] / 'retained-recording'
self.state.write_bytes(b'original recording')
self.calls = []
for name, value in (('host_check', {}), ('ari', []), ('unit_properties', {'ActiveState': 'active'}), ('status', {'ok': True, 'business_ready': False}), ('validate_environments', None)):
p = patch.object(remote, name, return_value=value)
p.start(); self.addCleanup(p.stop)
p = patch.object(remote, 'stop_for_change', side_effect=lambda *a, **k: self.calls.append('stop'))
p.start(); self.addCleanup(p.stop)
p = patch.object(remote, 'bootstrap', side_effect=self.bootstrap)
p.start(); self.addCleanup(p.stop)
p = patch.object(remote, 'run', return_value='')
p.start(); self.addCleanup(p.stop)
p = patch.object(remote.time, 'sleep')
p.start(); self.addCleanup(p.stop)
p = patch.object(remote.subprocess, 'run', return_value=subprocess.CompletedProcess([], 1, json.dumps({'success': False, 'phase': 'sip', 'http_status': 404, 'mq_connected': True, 'mq_connection_closed': True})))
p.start(); self.addCleanup(p.stop)
def bootstrap(self, *args):
remote.mkdir(self.p['state'])
return {'agent': {}, 'dispatcher': {}}
def test_saas_down_installs_and_repeat_preserves_original_state(self):
with patch.object(remote, 'backup_dir') as backup:
first = remote.install(self.payload, self.p)
second = remote.install(self.payload, self.p)
backup.assert_not_called()
self.assertEqual(self.state.read_bytes(), b'original recording')
self.assertEqual(first['dependency_state'], 'saas_pending')
self.assertFalse(second['reset_performed'])
self.assertTrue(self.p['current'].is_symlink())
for role in ('agent', 'dispatcher'):
unit = (self.p['units'] / ('go-sip-' + role + '.service')).read_text()
self.assertIn('Restart=on-failure', unit)
self.assertNotIn('--mode mock', unit)
def test_explicit_reset_stops_and_archives_before_deleting(self):
original = self.state.read_bytes()
def backup(*args):
self.assertIn('stop', self.calls)
self.assertEqual(self.state.read_bytes(), original)
self.calls.append('backup')
with patch.object(remote, 'backup_dir', side_effect=backup):
result = remote.install(dict(self.payload, reset_state=True, confirm_reset=True), self.p)
self.assertTrue(result['reset_performed'])
self.assertEqual(self.calls[:2], ['stop', 'backup'])
self.assertFalse(self.state.exists())
def test_unconfirmed_reset_does_not_remove_original(self):
with patch.object(remote, 'backup_dir') as backup:
with self.assertRaises(ValueError):
remote.install(dict(self.payload, reset_state=True), self.p)
backup.assert_not_called()
self.assertEqual(self.state.read_bytes(), b'original recording')
def test_independent_mq_failure_is_not_saas_waiting(self):
failed = subprocess.CompletedProcess([], 1, json.dumps({'success': False, 'phase': 'mq_connect', 'mq_connected': False}))
with patch.object(remote.subprocess, 'run', return_value=failed):
with self.assertRaisesRegex(ValueError, 'independent dependency failure'):
remote.install(self.payload, self.p)
self.assertEqual(self.state.read_bytes(), b'original recording')
def test_active_calls_prevent_stopping_agent(self):
# Exercise the real stop routine separately from the install mocks.
namespace = {}
exec(compile(HERE.joinpath('remote.py').read_text(), str(HERE / 'remote.py'), 'exec'), namespace)
calls = []
namespace['unit_properties'] = lambda u: {'FragmentPath': '/owned/unit'}
namespace['run'] = lambda args: calls.append(args)
namespace['ari'] = lambda *a: [{'id': 'active'}]
with self.assertRaises(ValueError):
namespace['stop_for_change'](self.p)
self.assertTrue(any('go-sip-dispatcher.service' in a for a in calls))
self.assertFalse(any('go-sip-agent.service' in a for a in calls))
if __name__ == '__main__':
unittest.main()
+116
View File
@@ -0,0 +1,116 @@
import importlib.util
import os
from pathlib import Path
import tempfile
import unittest
HERE = Path(__file__).resolve().parent
def module(name):
spec = importlib.util.spec_from_file_location(name, HERE / (name + '.py'))
result = importlib.util.module_from_spec(spec)
spec.loader.exec_module(result)
return result
class SourcesTest(unittest.TestCase):
def setUp(self):
self.m = module('preprod')
self.root = Path(tempfile.mkdtemp())
self.addCleanup(lambda: __import__('shutil').rmtree(self.root))
def private(self, name, text):
p = self.root / name
p.write_text(text)
p.chmod(0o600)
return p
def test_env_is_literal_and_never_executes(self):
p = self.private('input.env', 'KEY="a b#c"\nOTHER=$(touch sentinel)\n')
self.assertEqual(self.m.read_env(p), {'KEY': 'a b#c', 'OTHER': '$(touch sentinel)'})
self.assertFalse((self.root / 'sentinel').exists())
def test_permissions_duplicates_and_quotes_fail_without_values(self):
for text in ('KEY=private-value\nKEY=another\n', 'KEY="private-value\n'):
with self.assertRaises(ValueError) as ctx:
self.m.read_env(self.private('input.env', text))
self.assertNotIn('private-value', str(ctx.exception))
p = self.private('input.env', 'KEY=private-value\n')
p.chmod(0o664)
with self.assertRaises(ValueError):
self.m.read_env(p)
def test_oss_format_preserves_case_and_required_fields(self):
p = self.private('aliyun-oss.env', 'bucket: test-bucket\nEndpoint: oss-cn-beijing.aliyuncs.com\nRegion: cn-beijing\nRAM:\n username: synthetic\n accessKeyId: test-id\n accessKeySecret: test-secret\n')
d = self.m.read_oss(p)
self.assertEqual(d['bucket'], 'test-bucket')
self.assertEqual(d['RAM.accessKeySecret'], 'test-secret')
self.assertEqual(d['Region'], 'cn-beijing')
# The real source is a colon file, not YAML: RAM children may be flat.
p.write_text(p.read_text().replace(' username:', 'username:').replace(' accessKeyId:', 'accessKeyId:').replace(' accessKeySecret:', 'accessKeySecret:'))
self.assertEqual(self.m.read_oss(p), d)
def test_sources_missing_required_field_do_not_default(self):
self.private('saas.env', 'DispatcherUUID=00000000-0000-4000-8000-000000000001\nSaaSBaseURL=https://saas.example.invalid\n')
self.private('rabbitmq.env', 'RABBITMQ_HOST=host\n')
self.private('aliyun-oss.env', 'bucket: test\n')
with self.assertRaises(ValueError):
self.m.sources(self.root)
class SafetyTest(unittest.TestCase):
def setUp(self):
self.m = module('remote')
def test_fixed_paths_not_per_deployment(self):
p = self.m.paths(Path('/home/synthetic'))
self.assertEqual(str(p['state']), '/home/synthetic/.local/share/go-sip')
self.assertEqual(str(p['config']), '/home/synthetic/.config/go-sip')
self.assertEqual(str(p['releases']), '/home/synthetic/.local/opt/go-sip/releases')
def test_reset_requires_explicit_authorization_and_zero_calls(self):
for requested, confirmed, channels in ((False, True, 0), (True, False, 0), (True, True, 1)):
with self.assertRaises(ValueError):
self.m.require_reset(requested, confirmed, channels)
self.m.require_reset(True, True, 0)
def test_unavailable_saas_is_deferred_not_success_or_abort(self):
status = self.m.dependency_state({'success': False, 'phase': 'sip', 'http_status': 404})
self.assertEqual(status, 'saas_pending')
self.assertEqual(self.m.dependency_state({'success': True, 'phase': 'complete'}), 'verified')
self.assertEqual(self.m.dependency_state({'success': False, 'phase': 'mq_connect'}), 'mq_failed')
def test_role_rendering_overrides_old_mock_business_source(self):
source = {'dispatcher_id': 'new-owner', 'saas_url': 'https://real.example.invalid', 'secret': 'new-secret', 'mq_url': 'amqp://synthetic:placeholder@mq.example.invalid/test'}
p = self.m.paths(Path('/home/synthetic'))
roles = self.m.render_roles({'AGENT_ID': 'agent', 'CELL_ID': 'cell', 'AGENT_SESSION_PATH': '/old/session', 'AGENT_MOCK_SCENARIO_FILE': '/old/scenario', 'AGENT_APPROVED_AI_TASK_MAP_FILE': '/old/ai'}, {'SAAS_BASE_URL': 'http://127.0.0.1:18080', 'DISPATCHER_SECRET_KEY': 'old-secret', 'RABBITMQ_URL': 'amqp://old', 'DISPATCHER_SQLITE_PATH': '/old/db'}, source, p)
self.assertEqual(roles['dispatcher']['SAAS_BASE_URL'], source['saas_url'])
self.assertEqual(roles['dispatcher']['RABBITMQ_URL'], source['mq_url'])
self.assertEqual(roles['agent']['DISPATCHER_ID'], 'new-owner')
self.assertEqual(roles['agent']['AGENT_GRPC_LISTEN'], '127.0.0.1:19443')
self.assertEqual(roles['agent']['AGENT_SESSION_PATH'], str(p['state'] / 'agent-session.json'))
self.assertEqual(roles['dispatcher']['DISPATCHER_SECRET_KEY'], 'new-secret')
self.assertEqual(roles['dispatcher']['DISPATCHER_GRPC_LISTEN'], '127.0.0.1:19444')
self.assertNotIn('SAAS_SECRET_KEY', roles['dispatcher'])
self.assertNotIn('AGENT_APPROVED_AI_TASK_MAP_FILE', roles['agent'])
self.assertTrue(roles['dispatcher']['DISPATCHER_SQLITE_PATH'].endswith('/go-sip/dispatcher.sqlite'))
self.assertFalse(any('MOCK' in k or 'FIXTURE' in k for d in roles.values() for k in d))
def test_absent_unit_is_not_an_inspection_failure(self):
from unittest.mock import patch
import subprocess
with patch.object(self.m.subprocess, 'run', return_value=subprocess.CompletedProcess([], 4, 'LoadState=not-found\nFragmentPath=\n')):
self.assertEqual(self.m.unit_properties('new.service')['LoadState'], 'not-found')
with patch.object(self.m.subprocess, 'run', return_value=subprocess.CompletedProcess([], 1, '')):
with self.assertRaises(ValueError):
self.m.unit_properties('broken.service')
def test_normal_deployment_never_selects_reset(self):
self.assertFalse(self.m.reset_requested({}))
self.assertFalse(self.m.reset_requested({'reset_state': False}))
self.assertTrue(self.m.reset_requested({'reset_state': True}))
if __name__ == '__main__':
unittest.main()