2026-08-16 01:27:30 +08:00
|
|
|
package daemonserver
|
|
|
|
|
|
|
|
|
|
import (
|
|
|
|
|
"context"
|
|
|
|
|
"errors"
|
|
|
|
|
"io"
|
|
|
|
|
"log/slog"
|
|
|
|
|
"os"
|
|
|
|
|
"path/filepath"
|
|
|
|
|
"syscall"
|
|
|
|
|
"testing"
|
|
|
|
|
"time"
|
|
|
|
|
|
|
|
|
|
"yms-daemon/internal/backendupdate"
|
|
|
|
|
"yms-daemon/internal/daemonapi"
|
|
|
|
|
"yms-daemon/internal/daemonclient"
|
|
|
|
|
"yms-daemon/internal/transaction"
|
|
|
|
|
)
|
|
|
|
|
|
2026-08-17 10:10:14 +08:00
|
|
|
// TestServerAcceptsBackendUpdateThroughUnixSocket 验证服务端能通过 Unix Socket 接收后端 ZIP 更新请求,
|
|
|
|
|
// 正确分派到更新器,并向客户端返回事务终态以及流式的更新进度事件。
|
2026-08-16 01:27:30 +08:00
|
|
|
func TestServerAcceptsBackendUpdateThroughUnixSocket(t *testing.T) {
|
|
|
|
|
socketPath := shortSocketPath(t)
|
|
|
|
|
updater := &fakeUpdater{record: transaction.Transaction{ID: "transaction-01", State: transaction.StateCommitted}}
|
2026-08-22 15:34:09 +08:00
|
|
|
server, err := New(socketPath, updater, &fakeDiagnoser{}, slog.New(slog.NewTextHandler(io.Discard, nil)))
|
2026-08-16 01:27:30 +08:00
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("create daemon server: %v", err)
|
|
|
|
|
}
|
|
|
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
|
|
|
serveResult := make(chan error, 1)
|
|
|
|
|
go func() { serveResult <- server.Serve(ctx) }()
|
|
|
|
|
waitForSocket(t, socketPath, serveResult)
|
|
|
|
|
|
|
|
|
|
packagePath := filepath.Join(t.TempDir(), "package.zip")
|
|
|
|
|
if err := os.WriteFile(packagePath, []byte("zip"), 0o600); err != nil {
|
|
|
|
|
t.Fatalf("write update package: %v", err)
|
|
|
|
|
}
|
|
|
|
|
var progress []daemonapi.Response
|
|
|
|
|
response, err := daemonclient.Update(context.Background(), socketPath, "backend", daemonapi.InputTypeRepackZIP, packagePath, func(event daemonapi.Response) {
|
|
|
|
|
progress = append(progress, event)
|
|
|
|
|
})
|
|
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("submit backend update: %v", err)
|
|
|
|
|
}
|
|
|
|
|
if response.TransactionID != updater.record.ID || response.State != string(transaction.StateCommitted) || updater.file != packagePath || updater.inputType != daemonapi.InputTypeRepackZIP {
|
|
|
|
|
t.Fatalf("unexpected update response or dispatch: response=%+v file=%s", response, updater.file)
|
|
|
|
|
}
|
|
|
|
|
if len(progress) != 1 || progress[0].Message != "test update progress" {
|
|
|
|
|
t.Fatalf("unexpected streamed update progress: %+v", progress)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
cancel()
|
|
|
|
|
select {
|
|
|
|
|
case err := <-serveResult:
|
|
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("stop daemon server: %v", err)
|
|
|
|
|
}
|
|
|
|
|
case <-time.After(3 * time.Second):
|
|
|
|
|
t.Fatal("daemon server did not stop")
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-08-17 10:10:14 +08:00
|
|
|
// TestServerAcceptsDirectNativeJARThroughUnixSocket 验证服务端能接收原生后端 JAR 直接更新请求,
|
|
|
|
|
// 并将制品路径与输入类型正确分派到更新器。
|
2026-08-16 01:27:30 +08:00
|
|
|
func TestServerAcceptsDirectNativeJARThroughUnixSocket(t *testing.T) {
|
|
|
|
|
socketPath := shortSocketPath(t)
|
|
|
|
|
updater := &fakeUpdater{record: transaction.Transaction{ID: "transaction-direct-01", State: transaction.StateCommitted}}
|
2026-08-22 15:34:09 +08:00
|
|
|
server, err := New(socketPath, updater, &fakeDiagnoser{}, slog.New(slog.NewTextHandler(io.Discard, nil)))
|
2026-08-16 01:27:30 +08:00
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("create daemon server: %v", err)
|
|
|
|
|
}
|
|
|
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
|
|
|
defer cancel()
|
|
|
|
|
serveResult := make(chan error, 1)
|
|
|
|
|
go func() { serveResult <- server.Serve(ctx) }()
|
|
|
|
|
waitForSocket(t, socketPath, serveResult)
|
|
|
|
|
|
|
|
|
|
jarPath := filepath.Join(t.TempDir(), "glory-soft-yms.jar")
|
|
|
|
|
if err := os.WriteFile(jarPath, []byte("jar"), 0o600); err != nil {
|
|
|
|
|
t.Fatalf("write direct native backend JAR: %v", err)
|
|
|
|
|
}
|
|
|
|
|
response, err := daemonclient.Update(context.Background(), socketPath, "backend", daemonapi.InputTypeNativeJAR, jarPath, nil)
|
|
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("submit direct native backend update: %v", err)
|
|
|
|
|
}
|
|
|
|
|
if response.TransactionID != updater.record.ID || updater.file != jarPath || updater.inputType != daemonapi.InputTypeNativeJAR {
|
|
|
|
|
t.Fatalf("unexpected direct update response or dispatch: response=%+v file=%s inputType=%s", response, updater.file, updater.inputType)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-08-17 10:10:14 +08:00
|
|
|
// TestServerAcceptsContainerImageThroughUnixSocket 验证服务端能接收容器镜像更新请求,
|
|
|
|
|
// 并把镜像引用作为制品路径分派到更新器。
|
2026-08-16 17:12:06 +08:00
|
|
|
func TestServerAcceptsContainerImageThroughUnixSocket(t *testing.T) {
|
|
|
|
|
socketPath := shortSocketPath(t)
|
|
|
|
|
updater := &fakeUpdater{record: transaction.Transaction{ID: "transaction-container-01", State: transaction.StateCommitted}}
|
2026-08-22 15:34:09 +08:00
|
|
|
server, err := New(socketPath, updater, &fakeDiagnoser{}, slog.New(slog.NewTextHandler(io.Discard, nil)))
|
2026-08-16 17:12:06 +08:00
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("create daemon server: %v", err)
|
|
|
|
|
}
|
|
|
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
|
|
|
defer cancel()
|
|
|
|
|
serveResult := make(chan error, 1)
|
|
|
|
|
go func() { serveResult <- server.Serve(ctx) }()
|
|
|
|
|
waitForSocket(t, socketPath, serveResult)
|
|
|
|
|
|
|
|
|
|
imageReference := "harbor.ymswell.asia/ymswell/glory-ymswell:20260813-184902-a37bf50d-v1.1.8.1"
|
2026-08-17 02:10:10 +08:00
|
|
|
response, err := daemonclient.UpdateContainerImage(context.Background(), socketPath, "backend", imageReference, true, nil)
|
2026-08-16 17:12:06 +08:00
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("submit container backend image: %v", err)
|
|
|
|
|
}
|
|
|
|
|
if response.TransactionID != updater.record.ID || updater.file != imageReference || updater.inputType != daemonapi.InputTypeContainerImage {
|
|
|
|
|
t.Fatalf("unexpected container update dispatch: response=%+v image=%s inputType=%s", response, updater.file, updater.inputType)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-08-17 10:10:14 +08:00
|
|
|
// TestServerReturnsTransactionFailure 验证更新器返回错误时,服务端会在最终结果中携带事务标识、
|
|
|
|
|
// 回滚终态以及错误信息,客户端据此返回错误。
|
2026-08-16 01:27:30 +08:00
|
|
|
func TestServerReturnsTransactionFailure(t *testing.T) {
|
|
|
|
|
socketPath := shortSocketPath(t)
|
|
|
|
|
updater := &fakeUpdater{
|
|
|
|
|
record: transaction.Transaction{ID: "transaction-02", State: transaction.StateRolledBack},
|
|
|
|
|
err: errors.New("health check failed"),
|
|
|
|
|
}
|
2026-08-22 15:34:09 +08:00
|
|
|
server, err := New(socketPath, updater, &fakeDiagnoser{}, slog.New(slog.NewTextHandler(io.Discard, nil)))
|
2026-08-16 01:27:30 +08:00
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("create daemon server: %v", err)
|
|
|
|
|
}
|
|
|
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
|
|
|
defer cancel()
|
|
|
|
|
serveResult := make(chan error, 1)
|
|
|
|
|
go func() { serveResult <- server.Serve(ctx) }()
|
|
|
|
|
waitForSocket(t, socketPath, serveResult)
|
|
|
|
|
|
|
|
|
|
packagePath := filepath.Join(t.TempDir(), "package.zip")
|
|
|
|
|
response, err := daemonclient.Update(context.Background(), socketPath, "backend", daemonapi.InputTypeRepackZIP, packagePath, nil)
|
|
|
|
|
if err == nil || response.TransactionID != updater.record.ID || response.State != string(transaction.StateRolledBack) {
|
|
|
|
|
t.Fatalf("unexpected failed update response: response=%+v err=%v", response, err)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-08-17 10:10:14 +08:00
|
|
|
// TestServerAcceptsBackendRestartThroughUnixSocket 验证服务端能接收后端重启请求,
|
|
|
|
|
// 并返回事务终态与流式的重启进度事件。
|
2026-08-16 01:27:30 +08:00
|
|
|
func TestServerAcceptsBackendRestartThroughUnixSocket(t *testing.T) {
|
|
|
|
|
socketPath := shortSocketPath(t)
|
|
|
|
|
updater := &fakeUpdater{record: transaction.Transaction{ID: "transaction-restart-01", State: transaction.StateCommitted}}
|
2026-08-22 15:34:09 +08:00
|
|
|
server, err := New(socketPath, updater, &fakeDiagnoser{}, slog.New(slog.NewTextHandler(io.Discard, nil)))
|
2026-08-16 01:27:30 +08:00
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("create daemon server: %v", err)
|
|
|
|
|
}
|
|
|
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
|
|
|
defer cancel()
|
|
|
|
|
serveResult := make(chan error, 1)
|
|
|
|
|
go func() { serveResult <- server.Serve(ctx) }()
|
|
|
|
|
waitForSocket(t, socketPath, serveResult)
|
|
|
|
|
|
|
|
|
|
var progress []daemonapi.Response
|
|
|
|
|
response, err := daemonclient.Restart(context.Background(), socketPath, "backend", func(event daemonapi.Response) {
|
|
|
|
|
progress = append(progress, event)
|
|
|
|
|
})
|
|
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("submit backend restart: %v", err)
|
|
|
|
|
}
|
|
|
|
|
if response.TransactionID != updater.record.ID || updater.operation != daemonapi.OperationRestart {
|
|
|
|
|
t.Fatalf("unexpected restart response or dispatch: response=%+v operation=%s", response, updater.operation)
|
|
|
|
|
}
|
|
|
|
|
if len(progress) != 1 || progress[0].Message != "test restart progress" {
|
|
|
|
|
t.Fatalf("unexpected streamed restart progress: %+v", progress)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-08-22 15:34:09 +08:00
|
|
|
// TestServerReturnsStatusThroughUnixSocket 验证服务端能通过 Unix Socket 处理 status 请求,
|
|
|
|
|
// 并把诊断结果返回给客户端。
|
|
|
|
|
func TestServerReturnsStatusThroughUnixSocket(t *testing.T) {
|
|
|
|
|
socketPath := shortSocketPath(t)
|
|
|
|
|
updater := &fakeUpdater{}
|
|
|
|
|
diagnosis := daemonapi.Diagnosis{Service: "backend", Type: "container", Healthy: true, Items: []daemonapi.DiagnosisItem{{Level: daemonapi.DiagnosisLevelOK, Code: "ok", Message: "ok"}}}
|
|
|
|
|
diagnoser := &fakeDiagnoser{diagnosis: diagnosis}
|
|
|
|
|
server, err := New(socketPath, updater, diagnoser, slog.New(slog.NewTextHandler(io.Discard, nil)))
|
|
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("create daemon server: %v", err)
|
|
|
|
|
}
|
|
|
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
|
|
|
defer cancel()
|
|
|
|
|
serveResult := make(chan error, 1)
|
|
|
|
|
go func() { serveResult <- server.Serve(ctx) }()
|
|
|
|
|
waitForSocket(t, socketPath, serveResult)
|
|
|
|
|
|
|
|
|
|
result, err := daemonclient.Status(context.Background(), socketPath, "backend")
|
|
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("submit status: %v", err)
|
|
|
|
|
}
|
|
|
|
|
if result.Service != "backend" || !result.Healthy {
|
|
|
|
|
t.Fatalf("unexpected status result: %+v", result)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// TestServerReconcileForwardsApply 验证 reconcile --apply 会把 apply 标志正确传递给诊断器。
|
|
|
|
|
func TestServerReconcileForwardsApply(t *testing.T) {
|
|
|
|
|
socketPath := shortSocketPath(t)
|
|
|
|
|
updater := &fakeUpdater{}
|
|
|
|
|
diagnoser := &fakeDiagnoser{diagnosis: daemonapi.Diagnosis{Service: "backend", Type: "container", Healthy: true}}
|
|
|
|
|
server, err := New(socketPath, updater, diagnoser, slog.New(slog.NewTextHandler(io.Discard, nil)))
|
|
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("create daemon server: %v", err)
|
|
|
|
|
}
|
|
|
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
|
|
|
defer cancel()
|
|
|
|
|
serveResult := make(chan error, 1)
|
|
|
|
|
go func() { serveResult <- server.Serve(ctx) }()
|
|
|
|
|
waitForSocket(t, socketPath, serveResult)
|
|
|
|
|
|
|
|
|
|
if _, err := daemonclient.Reconcile(context.Background(), socketPath, "backend", true); err != nil {
|
|
|
|
|
t.Fatalf("submit reconcile apply: %v", err)
|
|
|
|
|
}
|
|
|
|
|
if !diagnoser.applied {
|
|
|
|
|
t.Fatal("reconcile --apply was not forwarded to diagnoser")
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-08-17 10:10:14 +08:00
|
|
|
// shortSocketPath 创建一个短目录并返回其下的 daemon.sock 路径,用于规避 Unix Socket 路径长度上限。
|
|
|
|
|
// 目录会在测试结束时自动清理。
|
2026-08-16 01:27:30 +08:00
|
|
|
func shortSocketPath(t *testing.T) string {
|
|
|
|
|
t.Helper()
|
|
|
|
|
directory, err := os.MkdirTemp("", "yd-")
|
|
|
|
|
if err != nil {
|
|
|
|
|
t.Fatalf("create short Unix Socket directory: %v", err)
|
|
|
|
|
}
|
|
|
|
|
t.Cleanup(func() { _ = os.RemoveAll(directory) })
|
|
|
|
|
return filepath.Join(directory, "daemon.sock")
|
|
|
|
|
}
|
|
|
|
|
|
2026-08-17 10:10:14 +08:00
|
|
|
// waitForSocket 轮询等待 Unix Socket 文件出现,用于在测试中同步服务端就绪状态。
|
|
|
|
|
// 若服务端在 Socket 创建前就退出且原因为权限不足则跳过测试,否则视为失败;
|
|
|
|
|
// 超过 3 秒仍未见 Socket 也视为失败。
|
2026-08-16 01:27:30 +08:00
|
|
|
func waitForSocket(t *testing.T, socketPath string, serveResult <-chan error) {
|
|
|
|
|
t.Helper()
|
|
|
|
|
deadline := time.Now().Add(3 * time.Second)
|
|
|
|
|
for time.Now().Before(deadline) {
|
|
|
|
|
select {
|
|
|
|
|
case err := <-serveResult:
|
|
|
|
|
if errors.Is(err, syscall.EPERM) {
|
|
|
|
|
t.Skip("Unix Socket creation is not permitted by the test sandbox")
|
|
|
|
|
}
|
|
|
|
|
t.Fatalf("daemon server stopped before creating Unix Socket: %v", err)
|
|
|
|
|
default:
|
|
|
|
|
}
|
|
|
|
|
info, err := os.Lstat(socketPath)
|
|
|
|
|
if err == nil && info.Mode()&os.ModeSocket != 0 {
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
time.Sleep(5 * time.Millisecond)
|
|
|
|
|
}
|
|
|
|
|
t.Fatalf("daemon Unix Socket was not created: %s", socketPath)
|
|
|
|
|
}
|
|
|
|
|
|
2026-08-17 10:10:14 +08:00
|
|
|
// fakeUpdater backendUpdater 接口的测试替身,记录被分派的参数并返回预设结果。
|
2026-08-16 01:27:30 +08:00
|
|
|
type fakeUpdater struct {
|
2026-08-17 10:10:14 +08:00
|
|
|
// record 各更新方法返回的预设事务。
|
|
|
|
|
record transaction.Transaction
|
|
|
|
|
// err 各更新方法返回的预设错误。
|
|
|
|
|
err error
|
|
|
|
|
// file 记录最近一次分派收到的制品路径或镜像引用。
|
|
|
|
|
file string
|
|
|
|
|
// inputType 记录最近一次分派收到的输入类型。
|
2026-08-16 01:27:30 +08:00
|
|
|
inputType string
|
2026-08-17 10:10:14 +08:00
|
|
|
// operation 记录最近一次分派收到的操作类型。
|
2026-08-16 01:27:30 +08:00
|
|
|
operation string
|
|
|
|
|
}
|
|
|
|
|
|
2026-08-17 10:10:14 +08:00
|
|
|
// UpdateRepack 记录分派参数并上报一条更新进度,返回预设事务与错误。
|
2026-08-16 01:27:30 +08:00
|
|
|
func (u *fakeUpdater) UpdateRepack(_ context.Context, file string, report backendupdate.ProgressReporter) (transaction.Transaction, error) {
|
|
|
|
|
u.file = file
|
|
|
|
|
u.inputType = daemonapi.InputTypeRepackZIP
|
|
|
|
|
u.operation = daemonapi.OperationUpdate
|
|
|
|
|
if report != nil {
|
|
|
|
|
report(backendupdate.Progress{TransactionID: u.record.ID, State: transaction.StateStarting, Message: "test update progress"})
|
|
|
|
|
}
|
|
|
|
|
return u.record, u.err
|
|
|
|
|
}
|
|
|
|
|
|
2026-08-17 10:10:14 +08:00
|
|
|
// UpdateNativeJAR 记录分派参数并上报一条更新进度,返回预设事务与错误。
|
2026-08-16 01:27:30 +08:00
|
|
|
func (u *fakeUpdater) UpdateNativeJAR(_ context.Context, file string, report backendupdate.ProgressReporter) (transaction.Transaction, error) {
|
|
|
|
|
u.file = file
|
|
|
|
|
u.inputType = daemonapi.InputTypeNativeJAR
|
|
|
|
|
u.operation = daemonapi.OperationUpdate
|
|
|
|
|
if report != nil {
|
|
|
|
|
report(backendupdate.Progress{TransactionID: u.record.ID, State: transaction.StateStarting, Message: "test update progress"})
|
|
|
|
|
}
|
|
|
|
|
return u.record, u.err
|
|
|
|
|
}
|
|
|
|
|
|
2026-08-17 10:10:14 +08:00
|
|
|
// UpdateContainerImage 记录分派参数并上报一条更新进度,返回预设事务与错误。
|
2026-08-17 02:10:10 +08:00
|
|
|
func (u *fakeUpdater) UpdateContainerImage(_ context.Context, imageReference string, _ bool, report backendupdate.ProgressReporter) (transaction.Transaction, error) {
|
2026-08-16 17:12:06 +08:00
|
|
|
u.file = imageReference
|
|
|
|
|
u.inputType = daemonapi.InputTypeContainerImage
|
|
|
|
|
u.operation = daemonapi.OperationUpdate
|
|
|
|
|
if report != nil {
|
|
|
|
|
report(backendupdate.Progress{TransactionID: u.record.ID, State: transaction.StateStarting, Message: "test update progress"})
|
|
|
|
|
}
|
|
|
|
|
return u.record, u.err
|
|
|
|
|
}
|
|
|
|
|
|
2026-08-17 10:10:14 +08:00
|
|
|
// Restart 记录重启操作并上报一条重启进度,返回预设事务与错误。
|
2026-08-16 01:27:30 +08:00
|
|
|
func (u *fakeUpdater) Restart(_ context.Context, report backendupdate.ProgressReporter) (transaction.Transaction, error) {
|
|
|
|
|
u.operation = daemonapi.OperationRestart
|
|
|
|
|
if report != nil {
|
|
|
|
|
report(backendupdate.Progress{TransactionID: u.record.ID, State: transaction.StateStarting, Message: "test restart progress"})
|
|
|
|
|
}
|
|
|
|
|
return u.record, u.err
|
|
|
|
|
}
|
2026-08-22 15:34:09 +08:00
|
|
|
|
|
|
|
|
// fakeDiagnoser backendDiagnoser 接口的测试替身,返回预设的空诊断结果。
|
|
|
|
|
type fakeDiagnoser struct {
|
|
|
|
|
diagnosis daemonapi.Diagnosis
|
|
|
|
|
err error
|
|
|
|
|
applied bool
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Diagnose 返回预设的诊断结果。
|
|
|
|
|
func (d *fakeDiagnoser) Diagnose(context.Context) (daemonapi.Diagnosis, error) {
|
|
|
|
|
return d.diagnosis, d.err
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Reconcile 返回预设的诊断结果,并记录 apply 标志。
|
|
|
|
|
func (d *fakeDiagnoser) Reconcile(_ context.Context, apply bool) (daemonapi.Diagnosis, error) {
|
|
|
|
|
d.applied = apply
|
|
|
|
|
return d.diagnosis, d.err
|
|
|
|
|
}
|