Files

514 lines
20 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// 本包是 yms-daemon 的入口,根据可执行文件名决定运行角色:
// 当二进制名为 ymsd 时作为守护进程服务端长期运行,否则作为 ymsctl 命令行客户端处理用户命令。
package main
import (
"context"
"encoding/json"
"errors"
"flag"
"fmt"
"io"
"log/slog"
"net/http"
"os"
"os/signal"
"path/filepath"
"strings"
"syscall"
"time"
"yms-daemon/internal/backendstatus"
"yms-daemon/internal/backendupdate"
"yms-daemon/internal/containerengine"
"yms-daemon/internal/daemonapi"
"yms-daemon/internal/daemonclient"
"yms-daemon/internal/daemonserver"
"yms-daemon/internal/deploymentconfig"
"yms-daemon/internal/filestore"
"yms-daemon/internal/hostnginx"
"yms-daemon/internal/logging"
"yms-daemon/internal/processlock"
"yms-daemon/internal/runtimepaths"
"yms-daemon/internal/systemd"
"yms-daemon/internal/transaction"
)
// serviceBackend 当前命令行客户端与守护进程唯一支持的目标服务名。
const serviceBackend = "backend"
// main 进程入口,负责根据可执行文件名分流到守护进程服务端或命令行客户端。
// 它先建立可被 os.Interrupt 与 SIGTERM 中断的上下文,再判断二进制名是否为 ymsd:
// 是则校验不接受任何参数并运行 runServe,否则把剩余参数交给 run 处理并以返回码退出。
func main() {
ctx, cancel := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer cancel()
if filepath.Base(os.Args[0]) == "ymsd" {
if len(os.Args) != 1 {
fmt.Fprintln(os.Stderr, "ymsd does not accept arguments")
os.Exit(2)
}
if err := runServe(ctx); err != nil {
fmt.Fprintln(os.Stderr, "ymsd failed:", err)
os.Exit(1)
}
return
}
os.Exit(run(ctx, os.Args[1:], os.Stdout, os.Stderr))
}
// run 解析并执行命令行客户端的一个子命令,返回进程退出码。
// arguments 去除程序名后的命令行参数,stdout 与 stderr 分别接收正常输出与错误输出。
// 它支持 update、list、restart、help 子命令;无参数或未知命令时向 stderr 输出用法并返回 2,
// 命令执行失败返回 1,成功返回 0。
func run(ctx context.Context, arguments []string, stdout io.Writer, stderr io.Writer) int {
if len(arguments) == 0 {
writeUsage(stderr)
return 2
}
switch arguments[0] {
case "update":
return runUpdate(ctx, arguments[1:], stdout, stderr)
case "list":
return runListCommand(ctx, arguments[1:], stdout, stderr)
case "restart":
return runRestart(ctx, arguments[1:], stdout, stderr)
case "status", "doctor", "reconcile":
return runDiagnosis(ctx, arguments[0], arguments[1:], stdout, stderr)
case "help", "-h", "--help":
writeUsage(stdout)
return 0
default:
fmt.Fprintf(stderr, "unknown command: %s\n", arguments[0])
writeUsage(stderr)
return 2
}
}
func runUpdate(ctx context.Context, arguments []string, stdout io.Writer, stderr io.Writer) int {
request, err := parseUpdateArgs(arguments, stderr)
if err != nil {
fmt.Fprintln(stderr, err)
return 2
}
var progress func(daemonapi.Response)
if !request.quite {
progress = func(event daemonapi.Response) {
writeUpdateProgress(stdout, event)
}
}
var response daemonapi.Response
if request.inputType == daemonapi.InputTypeContainerImage {
response, err = daemonclient.UpdateContainerImage(ctx, runtimepaths.Socket, request.service, request.imageReference, request.startLog, progress)
} else {
response, err = daemonclient.Update(ctx, runtimepaths.Socket, request.service, request.inputType, request.file, progress)
}
if err != nil {
if response.TransactionID != "" {
fmt.Fprintf(stderr, "transaction=%s state=%s error=%v\n", response.TransactionID, response.State, err)
} else {
fmt.Fprintln(stderr, "ymsctl update failed:", err)
}
return 1
}
if !request.quite {
fmt.Fprintf(stdout, "transaction=%s state=%s\n", response.TransactionID, response.State)
}
return 0
}
func runListCommand(ctx context.Context, arguments []string, stdout io.Writer, stderr io.Writer) int {
if err := runList(ctx, arguments, stdout, stderr); err != nil {
fmt.Fprintln(stderr, "ymsctl list failed:", err)
return 1
}
return 0
}
func runRestart(ctx context.Context, arguments []string, stdout io.Writer, stderr io.Writer) int {
request, err := parseRestartArgs(arguments, stderr)
if err != nil {
fmt.Fprintln(stderr, err)
return 2
}
var progress func(daemonapi.Response)
if !request.quite {
progress = func(event daemonapi.Response) {
writeUpdateProgress(stdout, event)
}
}
response, err := daemonclient.Restart(ctx, runtimepaths.Socket, request.service, progress)
if err != nil {
if response.TransactionID != "" {
fmt.Fprintf(stderr, "transaction=%s state=%s error=%v\n", response.TransactionID, response.State, err)
} else {
fmt.Fprintln(stderr, "ymsctl restart failed:", err)
}
return 1
}
if !request.quite {
fmt.Fprintf(stdout, "transaction=%s state=%s\n", response.TransactionID, response.State)
}
return 0
}
func runDiagnosis(ctx context.Context, command string, arguments []string, stdout io.Writer, stderr io.Writer) int {
request, err := parseDiagnosisArgs(command, arguments, stderr)
if err != nil {
fmt.Fprintln(stderr, err)
return 2
}
var diagnosis daemonapi.Diagnosis
switch command {
case "status":
diagnosis, err = daemonclient.Status(ctx, runtimepaths.Socket, request.service)
case "doctor":
diagnosis, err = daemonclient.Doctor(ctx, runtimepaths.Socket, request.service)
case "reconcile":
diagnosis, err = daemonclient.Reconcile(ctx, runtimepaths.Socket, request.service, request.apply)
}
if err != nil {
fmt.Fprintf(stderr, "ymsctl %s failed: %v\n", command, err)
return 1
}
writeDiagnosis(stdout, command, diagnosis)
return 0
}
// updateArguments 保存 update 子命令解析后的参数。
type updateArguments struct {
// service 目标服务名,当前仅接受 backend。
service string
// inputType 更新输入类型,由 -f 或 --native-jar 推导得到。
inputType string
// file 本地制品的绝对路径,容器镜像输入时为空。
file string
// imageReference 容器镜像引用,仅在 --container-image 输入时非空。
imageReference string
// quite 为真时抑制进度与成功结果的输出。
quite bool
// startLog 为真时在容器更新完成后输出启动日志。
startLog bool
}
// restartArguments 保存 restart 子命令解析后的参数。
type restartArguments struct {
// service 目标服务名,当前仅接受 backend。
service string
// quite 为真时抑制进度与成功结果的输出。
quite bool
}
// parseUpdateArgs 解析 update 子命令的 flag 参数并做业务校验。
// 它要求 --service 必须为 backend-f、--native-jar、--container-image 三者必须且只能提供一个;
// 对文件类输入会解析为绝对路径并校验其为非符号链接的普通文件,容器镜像输入不允许附带 --no-start-log。
// 校验失败时向 output 写入原因并返回 error。
func parseUpdateArgs(arguments []string, output io.Writer) (updateArguments, error) {
flags := flag.NewFlagSet("update", flag.ContinueOnError)
flags.SetOutput(output)
service := flags.String("service", "", "service to update")
file := flags.String("f", "", "repack ZIP path")
nativeJAR := flags.String("native-jar", "", "direct native backend JAR path")
containerImage := flags.String("container-image", "", "development backend container image reference")
quite := flags.Bool("quite", false, "suppress progress and successful result output")
noStartLog := flags.Bool("no-start-log", false, "do not print container startup logs")
if err := flags.Parse(arguments); err != nil {
return updateArguments{}, err
}
if flags.NArg() != 0 {
return updateArguments{}, errors.New("update does not accept positional arguments")
}
if *service != serviceBackend {
return updateArguments{}, errors.New("--service currently accepts only backend")
}
inputCount := 0
for _, value := range []string{*file, *nativeJAR, *containerImage} {
if value != "" {
inputCount++
}
}
if inputCount != 1 {
return updateArguments{}, errors.New("exactly one of -f, --native-jar, and --container-image is required; update inputs are mutually exclusive")
}
if *containerImage != "" {
return updateArguments{service: *service, inputType: daemonapi.InputTypeContainerImage, imageReference: *containerImage, quite: *quite, startLog: !*noStartLog}, nil
}
if *noStartLog {
return updateArguments{}, errors.New("--no-start-log requires --container-image")
}
inputType := daemonapi.InputTypeRepackZIP
inputFile := *file
if *nativeJAR != "" {
inputType = daemonapi.InputTypeNativeJAR
inputFile = *nativeJAR
}
absoluteFile, err := filepath.Abs(inputFile)
if err != nil {
return updateArguments{}, fmt.Errorf("resolve update package path: %w", err)
}
info, err := os.Lstat(absoluteFile)
if err != nil {
return updateArguments{}, fmt.Errorf("inspect update package %s: %w", absoluteFile, err)
}
if !info.Mode().IsRegular() || info.Mode()&os.ModeSymlink != 0 {
return updateArguments{}, fmt.Errorf("update package is not a direct regular file: %s", absoluteFile)
}
return updateArguments{service: *service, inputType: inputType, file: absoluteFile, quite: *quite}, nil
}
// parseRestartArgs 解析 restart 子命令的 flag 参数并做业务校验。
// 它要求 --service 必须为 backend,且不接受任何位置参数;校验失败时向 output 写入原因并返回 error。
func parseRestartArgs(arguments []string, output io.Writer) (restartArguments, error) {
flags := flag.NewFlagSet("restart", flag.ContinueOnError)
flags.SetOutput(output)
service := flags.String("service", "", "service to restart")
quite := flags.Bool("quite", false, "suppress progress and successful result output")
if err := flags.Parse(arguments); err != nil {
return restartArguments{}, err
}
if flags.NArg() != 0 {
return restartArguments{}, errors.New("restart does not accept positional arguments")
}
if *service != serviceBackend {
return restartArguments{}, errors.New("--service currently accepts only backend")
}
return restartArguments{service: *service, quite: *quite}, nil
}
// diagnosisArguments 保存 status/doctor/reconcile 子命令解析后的参数。
type diagnosisArguments struct {
// service 目标服务名,当前仅接受 backend。
service string
// apply 为真时表示 reconcile 执行自动修复,其余命令不接受该选项。
apply bool
}
// parseDiagnosisArgs 解析 status/doctor/reconcile 子命令的 flag 参数并做业务校验。
// 它要求 --service 必须为 backend,且不接受任何位置参数;--apply 仅允许 reconcile 使用。
// 校验失败时向 output 写入原因并返回 error。
func parseDiagnosisArgs(command string, arguments []string, output io.Writer) (diagnosisArguments, error) {
flags := flag.NewFlagSet(command, flag.ContinueOnError)
flags.SetOutput(output)
service := flags.String("service", "", "service to diagnose")
apply := flags.Bool("apply", false, "apply automatic fixes (reconcile only)")
if err := flags.Parse(arguments); err != nil {
return diagnosisArguments{}, err
}
if flags.NArg() != 0 {
return diagnosisArguments{}, errors.New(command + " does not accept positional arguments")
}
if *service != serviceBackend {
return diagnosisArguments{}, errors.New("--service currently accepts only backend")
}
if *apply && command != "reconcile" {
return diagnosisArguments{}, errors.New("--apply requires reconcile")
}
return diagnosisArguments{service: *service, apply: *apply}, nil
}
// runServe 以守护进程服务端角色运行:初始化日志、获取进程锁、加载部署配置、
// 打开事务存储与协调器、构建 Nginx 控制器与后端更新器,最后通过 daemonserver 开始服务。
// 它使用具名返回 result 以便在各资源清理阶段合并所有 Close 错误,最终返回服务端的运行错误。
func runServe(ctx context.Context) (result error) {
logger, logCloser, err := logging.New(runtimepaths.Log)
if err != nil {
return err
}
defer func() { result = errors.Join(result, logCloser.Close()) }()
lock, err := processlock.Acquire(runtimepaths.Lock)
if err != nil {
return err
}
defer func() { result = errors.Join(result, lock.Close()) }()
config, err := deploymentconfig.Load(deploymentconfig.DefaultPath)
if err != nil {
return err
}
store, err := transaction.OpenStore(ctx, runtimepaths.Database)
if err != nil {
return err
}
defer func() { result = errors.Join(result, store.Close()) }()
coordinator, err := transaction.NewCoordinator(store, logger)
if err != nil {
return err
}
gateway, err := hostnginx.NewController(
runtimepaths.HostNginxConfig,
runtimepaths.HostNginxExecutable,
)
if err != nil {
return err
}
httpClient := &http.Client{Timeout: 5 * time.Second}
updater, diagnoser, closeRuntime, err := buildRuntime(config, store, coordinator, gateway, httpClient, logger)
if closeRuntime != nil {
defer func() { result = errors.Join(result, closeRuntime()) }()
}
if err != nil {
return err
}
server, err := daemonserver.New(runtimepaths.Socket, updater, diagnoser, logger)
if err != nil {
return err
}
return server.Serve(ctx)
}
type runtimeCloser func() error
func buildRuntime(config deploymentconfig.Config, store *transaction.Store, coordinator *transaction.Coordinator, gateway *hostnginx.Controller, httpClient *http.Client, logger *slog.Logger) (*backendupdate.Updater, *backendstatus.Diagnoser, runtimeCloser, error) {
switch config.Backend.Type {
case deploymentconfig.BackendTypeNative:
return buildNativeRuntime(config, store, coordinator, gateway, httpClient, logger)
case deploymentconfig.BackendTypeContainer:
return buildContainerRuntime(config, store, coordinator, gateway, httpClient, logger)
default:
return nil, nil, nil, fmt.Errorf("unsupported backend.type %q", config.Backend.Type)
}
}
func buildNativeRuntime(config deploymentconfig.Config, store *transaction.Store, coordinator *transaction.Coordinator, gateway *hostnginx.Controller, httpClient *http.Client, logger *slog.Logger) (*backendupdate.Updater, *backendstatus.Diagnoser, runtimeCloser, error) {
releaseStore, err := filestore.New(config.Backend.ReleaseDir)
if err != nil {
return nil, nil, nil, err
}
units, err := systemd.NewSystemctl(config.Backend.SystemctlPath)
if err != nil {
return nil, nil, nil, err
}
updater, err := backendupdate.New(config, runtimepaths.WorkRoot, store, coordinator, releaseStore, units, gateway, httpClient, logger)
if err != nil {
return nil, nil, nil, err
}
diagnoser, err := backendstatus.New(config, store, coordinator, nil, units, gateway)
return updater, diagnoser, nil, err
}
func buildContainerRuntime(config deploymentconfig.Config, store *transaction.Store, coordinator *transaction.Coordinator, gateway *hostnginx.Controller, httpClient *http.Client, logger *slog.Logger) (*backendupdate.Updater, *backendstatus.Diagnoser, runtimeCloser, error) {
engine, err := containerengine.NewMobyEngine()
if err != nil {
return nil, nil, nil, err
}
closeEngine := runtimeCloser(engine.Close)
updater, err := backendupdate.NewContainer(config, runtimepaths.WorkRoot, store, coordinator, engine, gateway, httpClient, logger)
if err != nil {
return nil, nil, closeEngine, err
}
diagnoser, err := backendstatus.New(config, store, coordinator, engine, nil, gateway)
return updater, diagnoser, closeEngine, err
}
// writeUsage 向 output 写出命令行客户端各子命令的用法说明。
func writeUsage(output io.Writer) {
fmt.Fprintln(output, "usage:")
fmt.Fprintln(output, " ymsctl update --service backend -f <repack.zip> [--quite]")
fmt.Fprintln(output, " ymsctl update --service backend --native-jar <backend.jar> [--quite]")
fmt.Fprintln(output, " ymsctl update --service backend --container-image <image-ref> [--no-start-log] [--quite]")
fmt.Fprintln(output, " ymsctl restart --service backend [--quite]")
fmt.Fprintln(output, " ymsctl list [--limit <n>] [--service backend] [--state <state>] [--json]")
fmt.Fprintln(output, " ymsctl status --service backend")
fmt.Fprintln(output, " ymsctl doctor --service backend")
fmt.Fprintln(output, " ymsctl reconcile --service backend [--apply]")
}
// writeUpdateProgress 把一条进度事件格式化写入 output:状态字段左对齐占 14 列,
// 状态为空时显示为 INFO,随后输出事件消息。
func writeUpdateProgress(output io.Writer, event daemonapi.Response) {
state := event.State
if state == "" {
state = "INFO"
}
fmt.Fprintf(output, "%-14s %s\n", state, event.Message)
}
// writeDiagnosis 把诊断结果格式化写入 output。
// status 只展示需要关注的漂移与可修复项,doctor 展示全部诊断项,
// reconcile 只展示可自动修复项作为修复计划。
func writeDiagnosis(output io.Writer, command string, diagnosis daemonapi.Diagnosis) {
fmt.Fprintf(output, "service=%s type=%s status=%s\n", diagnosis.Service, diagnosis.Type, diagnosisState(diagnosis))
if diagnosis.RepairApplied {
fmt.Fprintf(output, "APPLIED 对账修复已执行 transaction=%s\n", diagnosis.RepairTransactionID)
fmt.Fprintln(output, "RESULT 现场已重新核对")
}
shown := 0
for _, item := range diagnosis.Items {
if !showDiagnosisItem(command, item.Level) {
continue
}
shown++
fmt.Fprintf(output, "%-8s %s\n", strings.ToUpper(item.Level), item.Message)
if item.Action != "" {
fmt.Fprintf(output, "%-8s %s\n", "action", item.Action)
}
}
if shown == 0 {
fmt.Fprintln(output, "(no items)")
} else if command == "reconcile" && !diagnosis.RepairApplied {
fmt.Fprintln(output, "PLAN 仅生成修复计划,未执行;如需执行请追加 --apply")
}
}
func diagnosisState(diagnosis daemonapi.Diagnosis) string {
if !diagnosis.Healthy {
return "DRIFT"
}
for _, item := range diagnosis.Items {
if item.Level == daemonapi.DiagnosisLevelFixable {
return "REPAIRABLE"
}
}
return "HEALTHY"
}
func showDiagnosisItem(command string, level string) bool {
switch command {
case "status":
return level != daemonapi.DiagnosisLevelOK
case "reconcile":
return level == daemonapi.DiagnosisLevelFixable
default:
return true
}
}
// runList 执行 list 子命令:解析过滤参数,打开事务存储读取最近事务,
// 按是否指定 --json 决定以稳定 JSON 或制表符分隔的表格形式输出到 stdout。
// 它支持 --limit、--service、--state、--json 四个选项,返回查询或输出阶段的错误。
func runList(ctx context.Context, arguments []string, stdout, stderr io.Writer) error {
flags := flag.NewFlagSet("list", flag.ContinueOnError)
flags.SetOutput(stderr)
limit := flags.Int("limit", 20, "maximum number of operations to show")
service := flags.String("service", "", "exact service filter")
state := flags.String("state", "", "exact transaction state filter")
jsonOutput := flags.Bool("json", false, "write stable JSON")
if err := flags.Parse(arguments); err != nil {
return err
}
if flags.NArg() != 0 {
return errors.New("list does not accept positional arguments")
}
if *service != "" && *service != serviceBackend {
return errors.New("--service currently accepts only backend")
}
store, err := transaction.OpenStore(ctx, runtimepaths.Database)
if err != nil {
return err
}
defer store.Close()
items, err := store.ListRecent(ctx, transaction.ListFilter{Limit: *limit, Service: *service, State: transaction.State(*state)})
if err != nil {
return err
}
if *jsonOutput {
return json.NewEncoder(stdout).Encode(items)
}
fmt.Fprintln(stdout, "TIME\tTRANSACTION\tSERVICE\tSTATE\tSOURCE")
for _, item := range items {
fmt.Fprintf(stdout, "%s\t%s\t%s\t%s\t%s\n", item.CreatedAt.Local().Format("2006-01-02 15:04:05"), item.ID, item.Service, item.State, item.Source)
}
return nil
}