// Package daemonclient 提供了命令行客户端向正在运行的守护进程服务提交请求的能力。 // 它负责把调用方参数封装成 daemonapi.Request、通过 Unix Socket 发送请求并解析流式响应, // 将进度事件通过回调逐条透传给调用方,最终返回携带终态的最终结果。 package daemonclient import ( "context" "encoding/json" "errors" "fmt" "net" "path/filepath" "yms-daemon/internal/daemonapi" ) // Update 向守护进程提交一次后端更新请求,输入类型由 inputType 指定。 // socketPath 守护进程 Unix Socket 的绝对路径,service 是目标服务名, // file 本地制品文件的绝对路径,progress 在非空时会在每收到一条进度事件时被调用一次。 // 返回的 daemonapi.Response 为最终结果;当路径校验失败、连接失败或守护进程返回错误时返回 error。 // 该函数不修改传入文件,可并发调用,每次调用都会建立独立的连接。 func Update(ctx context.Context, socketPath string, service string, inputType string, file string, progress func(daemonapi.Response)) (daemonapi.Response, error) { if !filepath.IsAbs(socketPath) || !filepath.IsAbs(file) { return daemonapi.Response{}, errors.New("daemon socket and update file paths must be absolute") } request := daemonapi.Request{Operation: daemonapi.OperationUpdate, Service: service, InputType: inputType, File: file} return submit(ctx, socketPath, request, progress) } // UpdateContainerImage 向守护进程提交一次基于容器镜像的后端更新请求。 // socketPath 守护进程 Unix Socket 的绝对路径,service 是目标服务名, // imageReference 必填的容器镜像引用,startLog 控制更新完成后是否输出容器启动日志, // progress 在非空时会在每收到一条进度事件时被调用一次。 // 返回的 daemonapi.Response 为最终结果;当路径非法、镜像引用为空、连接失败或守护进程返回错误时返回 error。 func UpdateContainerImage(ctx context.Context, socketPath string, service string, imageReference string, startLog bool, progress func(daemonapi.Response)) (daemonapi.Response, error) { if !filepath.IsAbs(socketPath) { return daemonapi.Response{}, errors.New("daemon socket path must be absolute") } if imageReference == "" { return daemonapi.Response{}, errors.New("container image reference is required") } request := daemonapi.Request{Operation: daemonapi.OperationUpdate, Service: service, InputType: daemonapi.InputTypeContainerImage, ImageReference: imageReference, StartLog: startLog} return submit(ctx, socketPath, request, progress) } // Restart 向守护进程提交一次后端重启请求,不更换任何制品。 // socketPath 守护进程 Unix Socket 的绝对路径,service 是目标服务名, // progress 在非空时会在每收到一条进度事件时被调用一次。 // 返回的 daemonapi.Response 为最终结果;当路径非法、连接失败或守护进程返回错误时返回 error。 func Restart(ctx context.Context, socketPath string, service string, progress func(daemonapi.Response)) (daemonapi.Response, error) { if !filepath.IsAbs(socketPath) { return daemonapi.Response{}, errors.New("daemon socket path must be absolute") } request := daemonapi.Request{Operation: daemonapi.OperationRestart, Service: service} return submit(ctx, socketPath, request, progress) } // Status 向守护进程查询后端当前状态,返回结构化诊断结果。 // socketPath 守护进程 Unix Socket 的绝对路径,service 是目标服务名。 // 当路径非法、连接失败或守护进程返回错误时返回 error。 func Status(ctx context.Context, socketPath string, service string) (daemonapi.Diagnosis, error) { if !filepath.IsAbs(socketPath) { return daemonapi.Diagnosis{}, errors.New("daemon socket path must be absolute") } return submitDiagnosis(ctx, socketPath, daemonapi.Request{Operation: daemonapi.OperationStatus, Service: service}) } // Doctor 向守护进程提交一次只读诊断请求,返回全部漂移项与建议动作。 // socketPath 守护进程 Unix Socket 的绝对路径,service 是目标服务名。 // 当路径非法、连接失败或守护进程返回错误时返回 error。 func Doctor(ctx context.Context, socketPath string, service string) (daemonapi.Diagnosis, error) { if !filepath.IsAbs(socketPath) { return daemonapi.Diagnosis{}, errors.New("daemon socket path must be absolute") } return submitDiagnosis(ctx, socketPath, daemonapi.Request{Operation: daemonapi.OperationDoctor, Service: service}) } // Reconcile 向守护进程提交对账请求,apply 为 true 时执行自动修复,否则只生成修复计划。 // socketPath 守护进程 Unix Socket 的绝对路径,service 是目标服务名。 // 当路径非法、连接失败或守护进程返回错误时返回 error。 func Reconcile(ctx context.Context, socketPath string, service string, apply bool) (daemonapi.Diagnosis, error) { if !filepath.IsAbs(socketPath) { return daemonapi.Diagnosis{}, errors.New("daemon socket path must be absolute") } return submitDiagnosis(ctx, socketPath, daemonapi.Request{Operation: daemonapi.OperationReconcile, Service: service, Apply: apply}) } // submitDiagnosis 提交一次诊断类请求并解析结果中的诊断结构,供 status/doctor/reconcile 复用。 func submitDiagnosis(ctx context.Context, socketPath string, request daemonapi.Request) (daemonapi.Diagnosis, error) { response, err := submit(ctx, socketPath, request, nil) if err != nil { return daemonapi.Diagnosis{}, err } if response.Diagnosis == nil { return daemonapi.Diagnosis{}, errors.New("daemon returned no diagnosis result") } return *response.Diagnosis, nil } // submit 通过 Unix Socket 把 request 发送给守护进程并读取响应流,供本包各公开函数复用。 // 它会以换行分隔的 JSON 编码请求,写完后关闭写端以向对端发送结束信号, // 随后循环解码响应:进度事件通过 progress 回调透传,最终结果则返回给调用方。 // 当连接失败、编解码失败、对端返回未知响应类型或最终结果携带错误时返回 error。 func submit(ctx context.Context, socketPath string, request daemonapi.Request, progress func(daemonapi.Response)) (daemonapi.Response, error) { dialer := net.Dialer{} connection, err := dialer.DialContext(ctx, "unix", socketPath) if err != nil { return daemonapi.Response{}, fmt.Errorf("connect to daemon service at %s: %w", socketPath, err) } defer connection.Close() if err := json.NewEncoder(connection).Encode(request); err != nil { return daemonapi.Response{}, fmt.Errorf("submit daemon request: %w", err) } if writer, ok := connection.(interface{ CloseWrite() error }); ok { if err := writer.CloseWrite(); err != nil { return daemonapi.Response{}, fmt.Errorf("finish daemon request: %w", err) } } decoder := json.NewDecoder(connection) decoder.DisallowUnknownFields() for { var response daemonapi.Response if err := decoder.Decode(&response); err != nil { return daemonapi.Response{}, fmt.Errorf("read daemon response: %w", err) } switch response.Kind { case daemonapi.ResponseProgress: if progress != nil { progress(response) } case daemonapi.ResponseResult: if response.Error != "" { return response, errors.New(response.Error) } return response, nil default: return daemonapi.Response{}, fmt.Errorf("daemon returned unsupported response kind %q", response.Kind) } } }