split daemon request and runtime setup
This commit is contained in:
@@ -186,3 +186,10 @@
|
|||||||
- 将容器镜像解析的平台校验、tag 校验、仓库 digest 收集拆分,保留摘要唯一性要求。
|
- 将容器镜像解析的平台校验、tag 校验、仓库 digest 收集拆分,保留摘要唯一性要求。
|
||||||
|
|
||||||
验证:`GOCACHE=/tmp/yms-go-cache go test ./...` 通过。
|
验证:`GOCACHE=/tmp/yms-go-cache go test ./...` 通过。
|
||||||
|
|
||||||
|
### 本轮继续拆分
|
||||||
|
|
||||||
|
- daemon server 将请求读取、进度回调、更新/重启分派、诊断分派拆开,保留协议错误文本与响应行为。
|
||||||
|
- `runServe` 将 native/container 运行时构建与资源关闭拆出,避免入口函数承担全部分支。
|
||||||
|
|
||||||
|
验证:`GOCACHE=/tmp/yms-go-cache go test ./...`、`git diff --check` 通过。
|
||||||
|
|||||||
@@ -130,84 +130,23 @@ func (s *Server) Serve(ctx context.Context) error {
|
|||||||
// 最终无论成功失败都会写入一条 ResponseResult 响应并记录日志。
|
// 最终无论成功失败都会写入一条 ResponseResult 响应并记录日志。
|
||||||
func (s *Server) handle(ctx context.Context, connection net.Conn) {
|
func (s *Server) handle(ctx context.Context, connection net.Conn) {
|
||||||
defer connection.Close()
|
defer connection.Close()
|
||||||
_ = connection.SetReadDeadline(time.Now().Add(10 * time.Second))
|
request, err := s.readRequest(connection)
|
||||||
request, err := decodeRequest(connection)
|
|
||||||
_ = connection.SetReadDeadline(time.Time{})
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
_ = s.writeResponse(connection, daemonapi.Response{Kind: daemonapi.ResponseResult, Error: err.Error()})
|
s.writeError(connection, err.Error())
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
if request.Service != "backend" {
|
if request.Service != "backend" {
|
||||||
_ = s.writeResponse(connection, daemonapi.Response{Kind: daemonapi.ResponseResult, Error: "service must be backend"})
|
s.writeError(connection, "service must be backend")
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
progressWritable := true
|
progress := s.progressReporter(ctx, connection)
|
||||||
report := func(progress backendupdate.Progress) {
|
if request.Operation == daemonapi.OperationStatus || request.Operation == daemonapi.OperationDoctor || request.Operation == daemonapi.OperationReconcile {
|
||||||
if !progressWritable {
|
s.handleDiagnosis(ctx, connection, request)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
err := s.writeResponse(connection, daemonapi.Response{
|
record, updateErr := s.handleBackendOperation(ctx, request, progress)
|
||||||
Kind: daemonapi.ResponseProgress,
|
if updateErr != nil && record.ID == "" {
|
||||||
TransactionID: progress.TransactionID,
|
s.writeError(connection, updateErr.Error())
|
||||||
State: string(progress.State),
|
|
||||||
Message: progress.Message,
|
|
||||||
})
|
|
||||||
if err != nil {
|
|
||||||
progressWritable = false
|
|
||||||
s.logger.WarnContext(ctx, "write daemon update progress", "error", err)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
var record transaction.Transaction
|
|
||||||
var updateErr error
|
|
||||||
switch request.Operation {
|
|
||||||
case daemonapi.OperationUpdate:
|
|
||||||
switch request.InputType {
|
|
||||||
case daemonapi.InputTypeRepackZIP:
|
|
||||||
if request.File == "" || request.ImageReference != "" {
|
|
||||||
_ = s.writeResponse(connection, daemonapi.Response{Kind: daemonapi.ResponseResult, Error: "repack-zip requires file and does not accept imageReference"})
|
|
||||||
return
|
|
||||||
}
|
|
||||||
record, updateErr = s.updater.UpdateRepack(ctx, request.File, report)
|
|
||||||
case daemonapi.InputTypeNativeJAR:
|
|
||||||
if request.File == "" || request.ImageReference != "" {
|
|
||||||
_ = s.writeResponse(connection, daemonapi.Response{Kind: daemonapi.ResponseResult, Error: "native-jar requires file and does not accept imageReference"})
|
|
||||||
return
|
|
||||||
}
|
|
||||||
record, updateErr = s.updater.UpdateNativeJAR(ctx, request.File, report)
|
|
||||||
case daemonapi.InputTypeContainerImage:
|
|
||||||
if request.File != "" || request.ImageReference == "" {
|
|
||||||
_ = s.writeResponse(connection, daemonapi.Response{Kind: daemonapi.ResponseResult, Error: "container-image requires imageReference and does not accept file"})
|
|
||||||
return
|
|
||||||
}
|
|
||||||
record, updateErr = s.updater.UpdateContainerImage(ctx, request.ImageReference, request.StartLog, report)
|
|
||||||
default:
|
|
||||||
_ = s.writeResponse(connection, daemonapi.Response{Kind: daemonapi.ResponseResult, Error: "inputType must be repack-zip, native-jar, or container-image"})
|
|
||||||
return
|
|
||||||
}
|
|
||||||
case daemonapi.OperationRestart:
|
|
||||||
if request.InputType != "" || request.File != "" || request.ImageReference != "" {
|
|
||||||
_ = s.writeResponse(connection, daemonapi.Response{Kind: daemonapi.ResponseResult, Error: "restart does not accept inputType or file"})
|
|
||||||
return
|
|
||||||
}
|
|
||||||
record, updateErr = s.updater.Restart(ctx, report)
|
|
||||||
case daemonapi.OperationStatus, daemonapi.OperationDoctor:
|
|
||||||
if request.InputType != "" || request.File != "" || request.ImageReference != "" || request.Apply {
|
|
||||||
_ = s.writeResponse(connection, daemonapi.Response{Kind: daemonapi.ResponseResult, Error: request.Operation + " does not accept inputType, file, imageReference, or apply"})
|
|
||||||
return
|
|
||||||
}
|
|
||||||
diagnosis, diagnosisErr := s.diagnoser.Diagnose(ctx)
|
|
||||||
s.writeDiagnosis(ctx, connection, request.Operation, diagnosis, diagnosisErr)
|
|
||||||
return
|
|
||||||
case daemonapi.OperationReconcile:
|
|
||||||
if request.InputType != "" || request.File != "" || request.ImageReference != "" {
|
|
||||||
_ = s.writeResponse(connection, daemonapi.Response{Kind: daemonapi.ResponseResult, Error: "reconcile does not accept inputType, file, or imageReference"})
|
|
||||||
return
|
|
||||||
}
|
|
||||||
diagnosis, diagnosisErr := s.diagnoser.Reconcile(ctx, request.Apply)
|
|
||||||
s.writeDiagnosis(ctx, connection, request.Operation, diagnosis, diagnosisErr)
|
|
||||||
return
|
|
||||||
default:
|
|
||||||
_ = s.writeResponse(connection, daemonapi.Response{Kind: daemonapi.ResponseResult, Error: "operation must be update, restart, status, doctor, or reconcile"})
|
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
response := daemonapi.Response{Kind: daemonapi.ResponseResult, TransactionID: record.ID, State: string(record.State)}
|
response := daemonapi.Response{Kind: daemonapi.ResponseResult, TransactionID: record.ID, State: string(record.State)}
|
||||||
@@ -222,6 +161,82 @@ func (s *Server) handle(ctx context.Context, connection net.Conn) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (s *Server) readRequest(connection net.Conn) (daemonapi.Request, error) {
|
||||||
|
_ = connection.SetReadDeadline(time.Now().Add(10 * time.Second))
|
||||||
|
request, err := decodeRequest(connection)
|
||||||
|
_ = connection.SetReadDeadline(time.Time{})
|
||||||
|
return request, err
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *Server) writeError(connection net.Conn, message string) {
|
||||||
|
_ = s.writeResponse(connection, daemonapi.Response{Kind: daemonapi.ResponseResult, Error: message})
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *Server) progressReporter(ctx context.Context, connection net.Conn) backendupdate.ProgressReporter {
|
||||||
|
progressWritable := true
|
||||||
|
return func(progress backendupdate.Progress) {
|
||||||
|
if !progressWritable {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
err := s.writeResponse(connection, daemonapi.Response{Kind: daemonapi.ResponseProgress, TransactionID: progress.TransactionID, State: string(progress.State), Message: progress.Message})
|
||||||
|
if err != nil {
|
||||||
|
progressWritable = false
|
||||||
|
s.logger.WarnContext(ctx, "write daemon update progress", "error", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *Server) handleBackendOperation(ctx context.Context, request daemonapi.Request, report backendupdate.ProgressReporter) (transaction.Transaction, error) {
|
||||||
|
switch request.Operation {
|
||||||
|
case daemonapi.OperationUpdate:
|
||||||
|
return s.handleUpdate(ctx, request, report)
|
||||||
|
case daemonapi.OperationRestart:
|
||||||
|
if request.InputType != "" || request.File != "" || request.ImageReference != "" {
|
||||||
|
return transaction.Transaction{}, errors.New("restart does not accept inputType or file")
|
||||||
|
}
|
||||||
|
return s.updater.Restart(ctx, report)
|
||||||
|
default:
|
||||||
|
return transaction.Transaction{}, errors.New("operation must be update, restart, status, doctor, or reconcile")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *Server) handleUpdate(ctx context.Context, request daemonapi.Request, report backendupdate.ProgressReporter) (transaction.Transaction, error) {
|
||||||
|
switch request.InputType {
|
||||||
|
case daemonapi.InputTypeRepackZIP:
|
||||||
|
if request.File == "" || request.ImageReference != "" {
|
||||||
|
return transaction.Transaction{}, errors.New("repack-zip requires file and does not accept imageReference")
|
||||||
|
}
|
||||||
|
return s.updater.UpdateRepack(ctx, request.File, report)
|
||||||
|
case daemonapi.InputTypeNativeJAR:
|
||||||
|
if request.File == "" || request.ImageReference != "" {
|
||||||
|
return transaction.Transaction{}, errors.New("native-jar requires file and does not accept imageReference")
|
||||||
|
}
|
||||||
|
return s.updater.UpdateNativeJAR(ctx, request.File, report)
|
||||||
|
case daemonapi.InputTypeContainerImage:
|
||||||
|
if request.File != "" || request.ImageReference == "" {
|
||||||
|
return transaction.Transaction{}, errors.New("container-image requires imageReference and does not accept file")
|
||||||
|
}
|
||||||
|
return s.updater.UpdateContainerImage(ctx, request.ImageReference, request.StartLog, report)
|
||||||
|
default:
|
||||||
|
return transaction.Transaction{}, errors.New("inputType must be repack-zip, native-jar, or container-image")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (s *Server) handleDiagnosis(ctx context.Context, connection net.Conn, request daemonapi.Request) {
|
||||||
|
if request.InputType != "" || request.File != "" || request.ImageReference != "" || (request.Operation != daemonapi.OperationReconcile && request.Apply) {
|
||||||
|
s.writeError(connection, request.Operation+" does not accept inputType, file, imageReference, or apply")
|
||||||
|
return
|
||||||
|
}
|
||||||
|
var diagnosis daemonapi.Diagnosis
|
||||||
|
var err error
|
||||||
|
if request.Operation == daemonapi.OperationReconcile {
|
||||||
|
diagnosis, err = s.diagnoser.Reconcile(ctx, request.Apply)
|
||||||
|
} else {
|
||||||
|
diagnosis, err = s.diagnoser.Diagnose(ctx)
|
||||||
|
}
|
||||||
|
s.writeDiagnosis(ctx, connection, request.Operation, diagnosis, err)
|
||||||
|
}
|
||||||
|
|
||||||
// writeResponse 把一条响应以换行分隔的 JSON 编码写入连接,失败时返回底层写入错误。
|
// writeResponse 把一条响应以换行分隔的 JSON 编码写入连接,失败时返回底层写入错误。
|
||||||
func (s *Server) writeResponse(connection net.Conn, response daemonapi.Response) error {
|
func (s *Server) writeResponse(connection net.Conn, response daemonapi.Response) error {
|
||||||
return json.NewEncoder(connection).Encode(response)
|
return json.NewEncoder(connection).Encode(response)
|
||||||
|
|||||||
@@ -9,6 +9,7 @@ import (
|
|||||||
"flag"
|
"flag"
|
||||||
"fmt"
|
"fmt"
|
||||||
"io"
|
"io"
|
||||||
|
"log/slog"
|
||||||
"net/http"
|
"net/http"
|
||||||
"os"
|
"os"
|
||||||
"os/signal"
|
"os/signal"
|
||||||
@@ -342,34 +343,9 @@ func runServe(ctx context.Context) (result error) {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
httpClient := &http.Client{Timeout: 5 * time.Second}
|
httpClient := &http.Client{Timeout: 5 * time.Second}
|
||||||
var updater *backendupdate.Updater
|
updater, diagnoser, closeRuntime, err := buildRuntime(config, store, coordinator, gateway, httpClient, logger)
|
||||||
var diagnoser *backendstatus.Diagnoser
|
if closeRuntime != nil {
|
||||||
switch config.Backend.Type {
|
defer func() { result = errors.Join(result, closeRuntime()) }()
|
||||||
case deploymentconfig.BackendTypeNative:
|
|
||||||
releaseStore, err := filestore.New(config.Backend.ReleaseDir)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
units, err := systemd.NewSystemctl(config.Backend.SystemctlPath)
|
|
||||||
if err != nil {
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
updater, err = backendupdate.New(config, runtimepaths.WorkRoot, store, coordinator, releaseStore, units, gateway, httpClient, logger)
|
|
||||||
if err == nil {
|
|
||||||
diagnoser, err = backendstatus.New(config, store, coordinator, nil, units, gateway)
|
|
||||||
}
|
|
||||||
case deploymentconfig.BackendTypeContainer:
|
|
||||||
var engine *containerengine.MobyEngine
|
|
||||||
engine, err = containerengine.NewMobyEngine()
|
|
||||||
if err == nil {
|
|
||||||
defer func() { result = errors.Join(result, engine.Close()) }()
|
|
||||||
updater, err = backendupdate.NewContainer(config, runtimepaths.WorkRoot, store, coordinator, engine, gateway, httpClient, logger)
|
|
||||||
if err == nil {
|
|
||||||
diagnoser, err = backendstatus.New(config, store, coordinator, engine, nil, gateway)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
default:
|
|
||||||
err = fmt.Errorf("unsupported backend.type %q", config.Backend.Type)
|
|
||||||
}
|
}
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
@@ -381,6 +357,50 @@ func runServe(ctx context.Context) (result error) {
|
|||||||
return server.Serve(ctx)
|
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 写出命令行客户端各子命令的用法说明。
|
// writeUsage 向 output 写出命令行客户端各子命令的用法说明。
|
||||||
func writeUsage(output io.Writer) {
|
func writeUsage(output io.Writer) {
|
||||||
fmt.Fprintln(output, "usage:")
|
fmt.Fprintln(output, "usage:")
|
||||||
|
|||||||
Reference in New Issue
Block a user