mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-06 15:46:37 +08:00
feat: openresty状态上报
This commit is contained in:
@@ -41,25 +41,27 @@ func main() {
|
||||
|
||||
client := httpclient.New(cfg.ServerURL, cfg.InitialAuthToken(), cfg.RequestTimeout.Duration())
|
||||
stateStore := state.NewStore(cfg.StatePath)
|
||||
runtimeManager := &nginx.Manager{
|
||||
RouteConfigPath: cfg.RouteConfigPath,
|
||||
CertDir: cfg.CertDir,
|
||||
NginxCertDir: cfg.OpenrestyCertDir,
|
||||
Executor: nginx.NewExecutor(nginx.ExecutorOptions{
|
||||
NginxPath: cfg.OpenrestyPath,
|
||||
DockerBinary: cfg.DockerBinary,
|
||||
ContainerName: cfg.OpenrestyContainerName,
|
||||
Image: cfg.OpenrestyDockerImage,
|
||||
RouteConfigPath: cfg.RouteConfigPath,
|
||||
CertDir: cfg.CertDir,
|
||||
NginxCertDir: cfg.OpenrestyCertDir,
|
||||
}),
|
||||
}
|
||||
runner := &agent.Runner{
|
||||
Config: cfg,
|
||||
StateStore: stateStore,
|
||||
HeartbeatService: heartbeat.New(client),
|
||||
SyncService: syncservice.New(client, &nginx.Manager{
|
||||
RouteConfigPath: cfg.RouteConfigPath,
|
||||
CertDir: cfg.CertDir,
|
||||
NginxCertDir: cfg.OpenrestyCertDir,
|
||||
Executor: nginx.NewExecutor(nginx.ExecutorOptions{
|
||||
NginxPath: cfg.OpenrestyPath,
|
||||
DockerBinary: cfg.DockerBinary,
|
||||
ContainerName: cfg.OpenrestyContainerName,
|
||||
Image: cfg.OpenrestyDockerImage,
|
||||
RouteConfigPath: cfg.RouteConfigPath,
|
||||
CertDir: cfg.CertDir,
|
||||
NginxCertDir: cfg.OpenrestyCertDir,
|
||||
}),
|
||||
}, stateStore),
|
||||
Updater: updater.New(),
|
||||
SyncService: syncservice.New(client, runtimeManager, stateStore),
|
||||
Updater: updater.New(),
|
||||
RuntimeManager: runtimeManager,
|
||||
}
|
||||
|
||||
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
|
||||
|
||||
@@ -27,6 +27,11 @@ type Updater interface {
|
||||
CheckAndUpdate(ctx context.Context, repo string, options UpdateOptions) error
|
||||
}
|
||||
|
||||
type RuntimeManager interface {
|
||||
CheckHealth(ctx context.Context) error
|
||||
Restart(ctx context.Context) error
|
||||
}
|
||||
|
||||
type UpdateOptions struct {
|
||||
Channel string
|
||||
TagName string
|
||||
@@ -39,12 +44,14 @@ type Runner struct {
|
||||
HeartbeatService HeartbeatService
|
||||
SyncService SyncService
|
||||
Updater Updater
|
||||
RuntimeManager RuntimeManager
|
||||
|
||||
autoUpdate bool
|
||||
updateNow bool
|
||||
updateRepo string
|
||||
updateChan string
|
||||
updateTag string
|
||||
autoUpdate bool
|
||||
updateNow bool
|
||||
updateRepo string
|
||||
updateChan string
|
||||
updateTag string
|
||||
restartOpenrestyNow bool
|
||||
}
|
||||
|
||||
func (r *Runner) Run(ctx context.Context) error {
|
||||
@@ -60,12 +67,14 @@ func (r *Runner) Run(ctx context.Context) error {
|
||||
} else {
|
||||
log.Printf("agent startup sync completed")
|
||||
}
|
||||
r.refreshOpenrestyHealth(ctx)
|
||||
settings, hbErr := r.HeartbeatService.Heartbeat(ctx, r.nodePayload(nodeID))
|
||||
if hbErr != nil {
|
||||
log.Printf("agent startup heartbeat failed: %v", hbErr)
|
||||
} else {
|
||||
log.Printf("agent startup heartbeat succeeded: node_id=%s", nodeID)
|
||||
r.applySettings(settings)
|
||||
r.tryRestartOpenresty(ctx)
|
||||
}
|
||||
} else if err = r.tryRegister(ctx, &nodeID); err != nil {
|
||||
log.Printf("agent initial discovery register failed: %v", err)
|
||||
@@ -88,6 +97,7 @@ func (r *Runner) Run(ctx context.Context) error {
|
||||
}
|
||||
continue
|
||||
}
|
||||
r.refreshOpenrestyHealth(ctx)
|
||||
settings, hbErr := r.HeartbeatService.Heartbeat(ctx, r.nodePayload(nodeID))
|
||||
if hbErr != nil {
|
||||
log.Printf("agent heartbeat failed: %v", hbErr)
|
||||
@@ -96,6 +106,7 @@ func (r *Runner) Run(ctx context.Context) error {
|
||||
heartbeatTicker.Reset(r.Config.HeartbeatInterval.Duration())
|
||||
syncTicker.Reset(r.Config.SyncInterval.Duration())
|
||||
}
|
||||
r.tryRestartOpenresty(ctx)
|
||||
r.tryAutoUpdate(ctx)
|
||||
}
|
||||
case <-syncTicker.C:
|
||||
@@ -143,9 +154,28 @@ func (r *Runner) applySettings(settings *protocol.AgentSettings) bool {
|
||||
r.updateRepo = strings.TrimSpace(settings.UpdateRepo)
|
||||
r.updateChan = strings.TrimSpace(settings.UpdateChannel)
|
||||
r.updateTag = strings.TrimSpace(settings.UpdateTag)
|
||||
r.restartOpenrestyNow = settings.RestartOpenrestyNow
|
||||
return changed
|
||||
}
|
||||
|
||||
func (r *Runner) tryRestartOpenresty(ctx context.Context) {
|
||||
if !r.restartOpenrestyNow {
|
||||
return
|
||||
}
|
||||
r.restartOpenrestyNow = false
|
||||
if r.RuntimeManager == nil {
|
||||
return
|
||||
}
|
||||
log.Printf("agent openresty restart requested by server")
|
||||
if err := r.RuntimeManager.Restart(ctx); err != nil {
|
||||
log.Printf("agent openresty restart failed: %v", err)
|
||||
r.recordOpenrestyUnhealthy(err, false)
|
||||
return
|
||||
}
|
||||
log.Printf("agent openresty restart succeeded")
|
||||
r.recordOpenrestyHealthy()
|
||||
}
|
||||
|
||||
func (r *Runner) tryAutoUpdate(ctx context.Context) {
|
||||
force := r.updateNow
|
||||
shouldCheck := r.autoUpdate || force
|
||||
@@ -224,15 +254,70 @@ func (r *Runner) recordSyncError(err error) {
|
||||
}
|
||||
}
|
||||
|
||||
func (r *Runner) nodePayload(nodeID string) protocol.NodePayload {
|
||||
snapshot, _ := r.StateStore.Load()
|
||||
return protocol.NodePayload{
|
||||
NodeID: nodeID,
|
||||
Name: r.Config.NodeName,
|
||||
IP: r.Config.NodeIP,
|
||||
AgentVersion: r.Config.AgentVersion,
|
||||
NginxVersion: r.Config.NginxVersion,
|
||||
CurrentVersion: snapshot.CurrentVersion,
|
||||
LastError: snapshot.LastError,
|
||||
func (r *Runner) refreshOpenrestyHealth(ctx context.Context) {
|
||||
if r.RuntimeManager == nil || r.StateStore == nil {
|
||||
return
|
||||
}
|
||||
if err := r.RuntimeManager.CheckHealth(ctx); err != nil {
|
||||
r.recordOpenrestyUnhealthy(err, true)
|
||||
return
|
||||
}
|
||||
r.recordOpenrestyHealthy()
|
||||
}
|
||||
|
||||
func (r *Runner) recordOpenrestyHealthy() {
|
||||
if r.StateStore == nil {
|
||||
return
|
||||
}
|
||||
snapshot, err := r.StateStore.Load()
|
||||
if err != nil {
|
||||
log.Printf("load state before recording openresty health failed: %v", err)
|
||||
return
|
||||
}
|
||||
if snapshot.OpenrestyStatus == protocol.OpenrestyStatusHealthy && strings.TrimSpace(snapshot.OpenrestyMessage) == "" {
|
||||
return
|
||||
}
|
||||
snapshot.OpenrestyStatus = protocol.OpenrestyStatusHealthy
|
||||
snapshot.OpenrestyMessage = ""
|
||||
if err = r.StateStore.Save(snapshot); err != nil {
|
||||
log.Printf("save state after recording openresty health failed: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func (r *Runner) recordOpenrestyUnhealthy(err error, fallbackOnly bool) {
|
||||
if err == nil || r.StateStore == nil {
|
||||
return
|
||||
}
|
||||
snapshot, loadErr := r.StateStore.Load()
|
||||
if loadErr != nil {
|
||||
log.Printf("load state before recording openresty error failed: %v", loadErr)
|
||||
return
|
||||
}
|
||||
message := strings.TrimSpace(err.Error())
|
||||
if !fallbackOnly || strings.TrimSpace(snapshot.OpenrestyMessage) == "" {
|
||||
snapshot.OpenrestyMessage = message
|
||||
}
|
||||
snapshot.OpenrestyStatus = protocol.OpenrestyStatusUnhealthy
|
||||
if saveErr := r.StateStore.Save(snapshot); saveErr != nil {
|
||||
log.Printf("save state after recording openresty error failed: %v", saveErr)
|
||||
}
|
||||
}
|
||||
|
||||
func (r *Runner) nodePayload(nodeID string) protocol.NodePayload {
|
||||
snapshot, _ := r.StateStore.Load()
|
||||
openrestyStatus := strings.TrimSpace(snapshot.OpenrestyStatus)
|
||||
if openrestyStatus == "" {
|
||||
openrestyStatus = protocol.OpenrestyStatusUnknown
|
||||
}
|
||||
return protocol.NodePayload{
|
||||
NodeID: nodeID,
|
||||
Name: r.Config.NodeName,
|
||||
IP: r.Config.NodeIP,
|
||||
AgentVersion: r.Config.AgentVersion,
|
||||
NginxVersion: r.Config.NginxVersion,
|
||||
CurrentVersion: snapshot.CurrentVersion,
|
||||
LastError: snapshot.LastError,
|
||||
OpenrestyStatus: openrestyStatus,
|
||||
OpenrestyMessage: snapshot.OpenrestyMessage,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -15,14 +15,16 @@ import (
|
||||
)
|
||||
|
||||
type fakeHeartbeatService struct {
|
||||
mu sync.Mutex
|
||||
registerCalls int
|
||||
heartbeatCalls int
|
||||
registerErr error
|
||||
registerResp *protocol.RegisterNodeResponse
|
||||
heartbeatErrs []error
|
||||
onHeartbeat func(int)
|
||||
lastToken string
|
||||
mu sync.Mutex
|
||||
registerCalls int
|
||||
heartbeatCalls int
|
||||
registerErr error
|
||||
registerResp *protocol.RegisterNodeResponse
|
||||
heartbeatErrs []error
|
||||
heartbeatSettings []*protocol.AgentSettings
|
||||
heartbeatPayloads []protocol.NodePayload
|
||||
onHeartbeat func(int)
|
||||
lastToken string
|
||||
}
|
||||
|
||||
func (f *fakeHeartbeatService) Register(ctx context.Context, payload protocol.NodePayload) (*protocol.RegisterNodeResponse, error) {
|
||||
@@ -36,16 +38,21 @@ func (f *fakeHeartbeatService) Heartbeat(ctx context.Context, payload protocol.N
|
||||
f.mu.Lock()
|
||||
f.heartbeatCalls++
|
||||
callIndex := f.heartbeatCalls
|
||||
f.heartbeatPayloads = append(f.heartbeatPayloads, payload)
|
||||
var err error
|
||||
if len(f.heartbeatErrs) >= callIndex {
|
||||
err = f.heartbeatErrs[callIndex-1]
|
||||
}
|
||||
var settings *protocol.AgentSettings
|
||||
if len(f.heartbeatSettings) >= callIndex {
|
||||
settings = f.heartbeatSettings[callIndex-1]
|
||||
}
|
||||
onHeartbeat := f.onHeartbeat
|
||||
f.mu.Unlock()
|
||||
if onHeartbeat != nil {
|
||||
onHeartbeat(callIndex)
|
||||
}
|
||||
return nil, err
|
||||
return settings, err
|
||||
}
|
||||
|
||||
func (f *fakeHeartbeatService) SetToken(token string) {
|
||||
@@ -63,6 +70,30 @@ type fakeSyncService struct {
|
||||
onSyncOnceCall func(int)
|
||||
}
|
||||
|
||||
type fakeRuntimeManager struct {
|
||||
mu sync.Mutex
|
||||
healthErr error
|
||||
restartErr error
|
||||
restartCalls int
|
||||
clearHealthOnRestart bool
|
||||
}
|
||||
|
||||
func (f *fakeRuntimeManager) CheckHealth(ctx context.Context) error {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
return f.healthErr
|
||||
}
|
||||
|
||||
func (f *fakeRuntimeManager) Restart(ctx context.Context) error {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
f.restartCalls++
|
||||
if f.clearHealthOnRestart && f.restartErr == nil {
|
||||
f.healthErr = nil
|
||||
}
|
||||
return f.restartErr
|
||||
}
|
||||
|
||||
func (f *fakeSyncService) SyncOnStartup(ctx context.Context) error {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
@@ -182,6 +213,71 @@ func TestRunnerDoesNotExitOnHeartbeatOrSyncError(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestRunnerReportsOpenrestyHealthAndExecutesRestart(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
|
||||
stateStore := state.NewStore(filepath.Join(t.TempDir(), "state.json"))
|
||||
if err := stateStore.Save(&state.Snapshot{
|
||||
OpenrestyStatus: protocol.OpenrestyStatusUnhealthy,
|
||||
OpenrestyMessage: "docker run openresty failed: bind 80 already allocated",
|
||||
}); err != nil {
|
||||
t.Fatalf("failed to seed state: %v", err)
|
||||
}
|
||||
heartbeatService := &fakeHeartbeatService{
|
||||
heartbeatSettings: []*protocol.AgentSettings{{RestartOpenrestyNow: true}},
|
||||
onHeartbeat: func(callCount int) {
|
||||
if callCount >= 1 {
|
||||
cancel()
|
||||
}
|
||||
},
|
||||
}
|
||||
runtimeManager := &fakeRuntimeManager{
|
||||
healthErr: errors.New("docker openresty container is not running"),
|
||||
clearHealthOnRestart: true,
|
||||
}
|
||||
runner := &Runner{
|
||||
Config: &config.Config{
|
||||
AgentToken: "agent-token",
|
||||
NodeName: "edge-01",
|
||||
NodeIP: "10.0.0.8",
|
||||
AgentVersion: config.AgentVersion,
|
||||
NginxVersion: "1.27.1.2",
|
||||
HeartbeatInterval: config.MillisecondDuration(10 * time.Millisecond),
|
||||
SyncInterval: config.MillisecondDuration(100 * time.Millisecond),
|
||||
},
|
||||
StateStore: stateStore,
|
||||
HeartbeatService: heartbeatService,
|
||||
SyncService: &fakeSyncService{},
|
||||
RuntimeManager: runtimeManager,
|
||||
}
|
||||
|
||||
err := runner.Run(ctx)
|
||||
if !errors.Is(err, context.Canceled) {
|
||||
t.Fatalf("expected context cancellation, got %v", err)
|
||||
}
|
||||
if len(heartbeatService.heartbeatPayloads) == 0 {
|
||||
t.Fatal("expected at least one heartbeat payload")
|
||||
}
|
||||
payload := heartbeatService.heartbeatPayloads[0]
|
||||
if payload.OpenrestyStatus != protocol.OpenrestyStatusUnhealthy {
|
||||
t.Fatalf("expected unhealthy openresty status in heartbeat payload, got %q", payload.OpenrestyStatus)
|
||||
}
|
||||
if payload.OpenrestyMessage != "docker run openresty failed: bind 80 already allocated" {
|
||||
t.Fatalf("unexpected openresty message: %q", payload.OpenrestyMessage)
|
||||
}
|
||||
if runtimeManager.restartCalls != 1 {
|
||||
t.Fatalf("expected one openresty restart attempt, got %d", runtimeManager.restartCalls)
|
||||
}
|
||||
snapshot, loadErr := stateStore.Load()
|
||||
if loadErr != nil {
|
||||
t.Fatalf("failed to load state: %v", loadErr)
|
||||
}
|
||||
if snapshot.OpenrestyStatus != protocol.OpenrestyStatusHealthy || snapshot.OpenrestyMessage != "" {
|
||||
t.Fatal("expected restart success to mark openresty healthy")
|
||||
}
|
||||
}
|
||||
|
||||
func TestRunnerDiscoveryRegisterUpdatesTokenAndNodeID(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
|
||||
@@ -25,6 +25,8 @@ type Executor interface {
|
||||
Test(ctx context.Context) error
|
||||
Reload(ctx context.Context) error
|
||||
EnsureRuntime(ctx context.Context, recreate bool) error
|
||||
CheckHealth(ctx context.Context) error
|
||||
Restart(ctx context.Context) error
|
||||
}
|
||||
|
||||
type CommandRunner interface {
|
||||
@@ -68,6 +70,27 @@ func (e *PathExecutor) EnsureRuntime(ctx context.Context, recreate bool) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (e *PathExecutor) CheckHealth(ctx context.Context) error {
|
||||
return e.Test(ctx)
|
||||
}
|
||||
|
||||
func (e *PathExecutor) Restart(ctx context.Context) error {
|
||||
log.Printf("restarting openresty with binary: %s", e.Path)
|
||||
output, err := e.Runner.Run(ctx, e.Path, "-s", "quit")
|
||||
if err != nil {
|
||||
text := string(output)
|
||||
if !isIgnorableOpenrestyStopError(text) {
|
||||
return fmt.Errorf("openresty stop failed: %w: %s", err, text)
|
||||
}
|
||||
}
|
||||
output, err = e.Runner.Run(ctx, e.Path)
|
||||
if err != nil {
|
||||
return fmt.Errorf("openresty start failed: %w: %s", err, string(output))
|
||||
}
|
||||
log.Printf("openresty restart succeeded with binary: %s", e.Path)
|
||||
return nil
|
||||
}
|
||||
|
||||
type DockerExecutor struct {
|
||||
DockerBinary string
|
||||
ContainerName string
|
||||
@@ -114,6 +137,22 @@ func (e *DockerExecutor) EnsureRuntime(ctx context.Context, recreate bool) error
|
||||
return e.runContainer(ctx)
|
||||
}
|
||||
|
||||
func (e *DockerExecutor) CheckHealth(ctx context.Context) error {
|
||||
log.Printf("checking docker openresty runtime health: container=%s", e.ContainerName)
|
||||
output, err := e.Runner.Run(ctx, e.DockerBinary, "inspect", "-f", "{{.State.Running}}", e.ContainerName)
|
||||
if err != nil {
|
||||
return fmt.Errorf("docker inspect openresty failed: %w: %s", err, string(output))
|
||||
}
|
||||
if strings.TrimSpace(string(output)) != "true" {
|
||||
return errors.New("docker openresty container is not running")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (e *DockerExecutor) Restart(ctx context.Context) error {
|
||||
return e.EnsureRuntime(ctx, true)
|
||||
}
|
||||
|
||||
func (e *DockerExecutor) removeContainer(ctx context.Context) error {
|
||||
log.Printf("removing docker openresty container: container=%s", e.ContainerName)
|
||||
output, err := e.Runner.Run(ctx, e.DockerBinary, "rm", "-f", e.ContainerName)
|
||||
@@ -193,6 +232,21 @@ func (m *Manager) EnsureRuntime(ctx context.Context, recreate bool) error {
|
||||
return m.Executor.EnsureRuntime(ctx, recreate)
|
||||
}
|
||||
|
||||
func (m *Manager) CheckHealth(ctx context.Context) error {
|
||||
if m.Executor == nil {
|
||||
return errors.New("executor 未配置")
|
||||
}
|
||||
return m.Executor.CheckHealth(ctx)
|
||||
}
|
||||
|
||||
func (m *Manager) Restart(ctx context.Context) error {
|
||||
if m.Executor == nil {
|
||||
return errors.New("executor 未配置")
|
||||
}
|
||||
log.Printf("openresty restart requested")
|
||||
return m.Executor.Restart(ctx)
|
||||
}
|
||||
|
||||
func (m *Manager) CurrentChecksum() (string, error) {
|
||||
if m.RouteConfigPath == "" {
|
||||
return "", errors.New("route config path 不能为空")
|
||||
@@ -300,6 +354,14 @@ func parseNginxVersion(output string) string {
|
||||
|
||||
var nginxVersionPattern = regexp.MustCompile(`(?im)(?:nginx|openresty) version:\s*(?:nginx|openresty)/([^\s]+)`)
|
||||
|
||||
func isIgnorableOpenrestyStopError(output string) bool {
|
||||
text := strings.ToLower(strings.TrimSpace(output))
|
||||
if text == "" {
|
||||
return false
|
||||
}
|
||||
return strings.Contains(text, "invalid pid") || strings.Contains(text, "no such process")
|
||||
}
|
||||
|
||||
func (e *DockerExecutor) runEphemeralRuntimeCommand(ctx context.Context, args ...string) ([]byte, error) {
|
||||
return e.runEphemeralRuntimeCommandWithBinary(ctx, dockerRuntimeCommand, args...)
|
||||
}
|
||||
|
||||
@@ -47,6 +47,14 @@ func (e *fakeExecutor) EnsureRuntime(ctx context.Context, recreate bool) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (e *fakeExecutor) CheckHealth(ctx context.Context) error {
|
||||
return e.testErr
|
||||
}
|
||||
|
||||
func (e *fakeExecutor) Restart(ctx context.Context) error {
|
||||
return e.reloadErr
|
||||
}
|
||||
|
||||
func TestPathExecutorCommands(t *testing.T) {
|
||||
runner := &fakeRunner{}
|
||||
executor := &PathExecutor{
|
||||
@@ -80,6 +88,47 @@ func TestPathExecutorEnsureRuntimeNoop(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestPathExecutorRestartIgnoresMissingPID(t *testing.T) {
|
||||
runner := &fakeRunner{
|
||||
runFn: func(name string, args ...string) ([]byte, error) {
|
||||
if len(args) == 2 && args[0] == "-s" && args[1] == "quit" {
|
||||
return []byte("openresty: [error] invalid PID number \"\" in \"/usr/local/openresty/nginx/logs/nginx.pid\""), errors.New("exit status 1")
|
||||
}
|
||||
return []byte(""), nil
|
||||
},
|
||||
}
|
||||
executor := &PathExecutor{
|
||||
Path: "/usr/local/openresty/nginx/sbin/openresty",
|
||||
Runner: runner,
|
||||
}
|
||||
if err := executor.Restart(context.Background()); err != nil {
|
||||
t.Fatalf("Restart failed: %v", err)
|
||||
}
|
||||
if len(runner.calls) != 2 {
|
||||
t.Fatalf("expected 2 restart calls, got %d", len(runner.calls))
|
||||
}
|
||||
}
|
||||
|
||||
func TestDockerExecutorCheckHealthFailsWhenContainerStopped(t *testing.T) {
|
||||
runner := &fakeRunner{
|
||||
runFn: func(name string, args ...string) ([]byte, error) {
|
||||
return []byte("false"), nil
|
||||
},
|
||||
}
|
||||
executor := &DockerExecutor{
|
||||
DockerBinary: "docker",
|
||||
ContainerName: "atsflare-openresty",
|
||||
Image: "openresty/openresty:alpine",
|
||||
RouteConfigDir: filepath.Clean("/tmp/routes"),
|
||||
CertDir: filepath.Clean("/tmp/certs"),
|
||||
NginxCertDir: "/etc/nginx/atsflare-certs",
|
||||
Runner: runner,
|
||||
}
|
||||
if err := executor.CheckHealth(context.Background()); err == nil {
|
||||
t.Fatal("expected CheckHealth to fail when container is not running")
|
||||
}
|
||||
}
|
||||
|
||||
func TestDockerExecutorStartsContainerWhenMissing(t *testing.T) {
|
||||
runner := &fakeRunner{
|
||||
runFn: func(name string, args ...string) ([]byte, error) {
|
||||
|
||||
@@ -14,23 +14,32 @@ type HeartbeatAPIResponse struct {
|
||||
}
|
||||
|
||||
type AgentSettings struct {
|
||||
HeartbeatInterval int `json:"heartbeat_interval"`
|
||||
SyncInterval int `json:"sync_interval"`
|
||||
AutoUpdate bool `json:"auto_update"`
|
||||
UpdateRepo string `json:"update_repo"`
|
||||
UpdateNow bool `json:"update_now"`
|
||||
UpdateChannel string `json:"update_channel"`
|
||||
UpdateTag string `json:"update_tag"`
|
||||
HeartbeatInterval int `json:"heartbeat_interval"`
|
||||
SyncInterval int `json:"sync_interval"`
|
||||
AutoUpdate bool `json:"auto_update"`
|
||||
UpdateRepo string `json:"update_repo"`
|
||||
UpdateNow bool `json:"update_now"`
|
||||
UpdateChannel string `json:"update_channel"`
|
||||
UpdateTag string `json:"update_tag"`
|
||||
RestartOpenrestyNow bool `json:"restart_openresty_now"`
|
||||
}
|
||||
|
||||
const (
|
||||
OpenrestyStatusHealthy = "healthy"
|
||||
OpenrestyStatusUnhealthy = "unhealthy"
|
||||
OpenrestyStatusUnknown = "unknown"
|
||||
)
|
||||
|
||||
type NodePayload struct {
|
||||
NodeID string `json:"node_id"`
|
||||
Name string `json:"name"`
|
||||
IP string `json:"ip"`
|
||||
AgentVersion string `json:"agent_version"`
|
||||
NginxVersion string `json:"nginx_version"`
|
||||
CurrentVersion string `json:"current_version"`
|
||||
LastError string `json:"last_error"`
|
||||
NodeID string `json:"node_id"`
|
||||
Name string `json:"name"`
|
||||
IP string `json:"ip"`
|
||||
AgentVersion string `json:"agent_version"`
|
||||
NginxVersion string `json:"nginx_version"`
|
||||
CurrentVersion string `json:"current_version"`
|
||||
LastError string `json:"last_error"`
|
||||
OpenrestyStatus string `json:"openresty_status"`
|
||||
OpenrestyMessage string `json:"openresty_message"`
|
||||
}
|
||||
|
||||
type RegisterNodeResponse struct {
|
||||
|
||||
@@ -10,10 +10,12 @@ import (
|
||||
)
|
||||
|
||||
type Snapshot struct {
|
||||
NodeID string `json:"node_id"`
|
||||
CurrentVersion string `json:"current_version"`
|
||||
CurrentChecksum string `json:"current_checksum"`
|
||||
LastError string `json:"last_error"`
|
||||
NodeID string `json:"node_id"`
|
||||
CurrentVersion string `json:"current_version"`
|
||||
CurrentChecksum string `json:"current_checksum"`
|
||||
LastError string `json:"last_error"`
|
||||
OpenrestyStatus string `json:"openresty_status"`
|
||||
OpenrestyMessage string `json:"openresty_message"`
|
||||
}
|
||||
|
||||
type Store struct {
|
||||
|
||||
@@ -72,9 +72,14 @@ func (s *Service) sync(ctx context.Context, startup bool) error {
|
||||
if startup {
|
||||
log.Printf("ensuring openresty runtime on startup: version=%s", config.Version)
|
||||
if err = s.nginxManager.EnsureRuntime(ctx, true); err != nil {
|
||||
snapshot.OpenrestyStatus = protocol.OpenrestyStatusUnhealthy
|
||||
snapshot.OpenrestyMessage = err.Error()
|
||||
_ = s.stateStore.Save(snapshot)
|
||||
return err
|
||||
}
|
||||
log.Printf("openresty runtime ensured on startup: version=%s", config.Version)
|
||||
snapshot.OpenrestyStatus = protocol.OpenrestyStatusHealthy
|
||||
snapshot.OpenrestyMessage = ""
|
||||
}
|
||||
snapshot.CurrentVersion = config.Version
|
||||
snapshot.CurrentChecksum = config.Checksum
|
||||
@@ -90,6 +95,8 @@ func (s *Service) sync(ctx context.Context, startup bool) error {
|
||||
if err = s.nginxManager.Apply(ctx, config.RenderedConfig, config.SupportFiles); err != nil {
|
||||
log.Printf("apply openresty config failed: mode=%s version=%s error=%v", mode, config.Version, err)
|
||||
snapshot.LastError = err.Error()
|
||||
snapshot.OpenrestyStatus = protocol.OpenrestyStatusUnhealthy
|
||||
snapshot.OpenrestyMessage = err.Error()
|
||||
_ = s.stateStore.Save(snapshot)
|
||||
reportErr := s.client.ReportApplyLog(ctx, protocol.ApplyLogPayload{
|
||||
NodeID: snapshot.NodeID,
|
||||
@@ -108,6 +115,8 @@ func (s *Service) sync(ctx context.Context, startup bool) error {
|
||||
snapshot.CurrentVersion = config.Version
|
||||
snapshot.CurrentChecksum = config.Checksum
|
||||
snapshot.LastError = ""
|
||||
snapshot.OpenrestyStatus = protocol.OpenrestyStatusHealthy
|
||||
snapshot.OpenrestyMessage = ""
|
||||
if err = s.stateStore.Save(snapshot); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -26,6 +26,7 @@ type fakeManager struct {
|
||||
applyErr error
|
||||
currentChecksum string
|
||||
currentChecksumErr error
|
||||
ensureErr error
|
||||
ensureCalls []bool
|
||||
applyContents []string
|
||||
applyFiles [][]protocol.SupportFile
|
||||
@@ -43,6 +44,14 @@ func (f *fakeExecutor) EnsureRuntime(ctx context.Context, recreate bool) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (f *fakeExecutor) CheckHealth(ctx context.Context) error {
|
||||
return f.testErr
|
||||
}
|
||||
|
||||
func (f *fakeExecutor) Restart(ctx context.Context) error {
|
||||
return f.reloadErr
|
||||
}
|
||||
|
||||
func (f *fakeClient) GetActiveConfig(ctx context.Context) (*protocol.ActiveConfigResponse, error) {
|
||||
return &f.config, nil
|
||||
}
|
||||
@@ -60,7 +69,7 @@ func (m *fakeManager) Apply(ctx context.Context, content string, supportFiles []
|
||||
|
||||
func (m *fakeManager) EnsureRuntime(ctx context.Context, recreate bool) error {
|
||||
m.ensureCalls = append(m.ensureCalls, recreate)
|
||||
return nil
|
||||
return m.ensureErr
|
||||
}
|
||||
|
||||
func (m *fakeManager) CurrentChecksum() (string, error) {
|
||||
@@ -216,4 +225,45 @@ func TestSyncOnStartupRecreatesRuntimeWhenChecksumMatches(t *testing.T) {
|
||||
if snapshot.CurrentChecksum != "checksum-3" || snapshot.CurrentVersion != "20260309-003" {
|
||||
t.Fatal("expected snapshot to be refreshed from active config")
|
||||
}
|
||||
if snapshot.OpenrestyStatus != protocol.OpenrestyStatusHealthy || snapshot.OpenrestyMessage != "" {
|
||||
t.Fatal("expected startup sync to mark openresty healthy")
|
||||
}
|
||||
}
|
||||
|
||||
func TestSyncOnStartupRecordsRuntimeFailure(t *testing.T) {
|
||||
client := &fakeClient{
|
||||
config: protocol.ActiveConfigResponse{
|
||||
Version: "20260309-004",
|
||||
Checksum: "checksum-4",
|
||||
RenderedConfig: "server { listen 83; }",
|
||||
CreatedAt: time.Now().Format(time.RFC3339),
|
||||
},
|
||||
}
|
||||
stateStore := state.NewStore(filepath.Join(t.TempDir(), "state.json"))
|
||||
nodeID, err := stateStore.EnsureNodeID()
|
||||
if err != nil {
|
||||
t.Fatalf("EnsureNodeID failed: %v", err)
|
||||
}
|
||||
if err = stateStore.Save(&state.Snapshot{NodeID: nodeID}); err != nil {
|
||||
t.Fatalf("failed to seed state: %v", err)
|
||||
}
|
||||
|
||||
manager := &fakeManager{
|
||||
currentChecksum: "checksum-4",
|
||||
ensureErr: context.DeadlineExceeded,
|
||||
}
|
||||
service := New(client, manager, stateStore)
|
||||
if err = service.SyncOnStartup(context.Background()); err == nil {
|
||||
t.Fatal("expected SyncOnStartup to fail when runtime recreation fails")
|
||||
}
|
||||
snapshot, err := stateStore.Load()
|
||||
if err != nil {
|
||||
t.Fatalf("failed to load state: %v", err)
|
||||
}
|
||||
if snapshot.OpenrestyStatus != protocol.OpenrestyStatusUnhealthy {
|
||||
t.Fatalf("expected unhealthy openresty status, got %q", snapshot.OpenrestyStatus)
|
||||
}
|
||||
if snapshot.OpenrestyMessage == "" {
|
||||
t.Fatal("expected runtime error message to be recorded")
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user