mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-12 02:06:37 +08:00
[功能] 添加对应用结果的警告支持,优化配置激活和回滚逻辑
This commit is contained in:
@@ -188,9 +188,9 @@ Agent 必须满足:
|
|||||||
* 常规同步优先依据 heartbeat 返回的版本摘要判断
|
* 常规同步优先依据 heartbeat 返回的版本摘要判断
|
||||||
* 发现新版本时先备份旧文件
|
* 发现新版本时先备份旧文件
|
||||||
* 写入主配置、路由配置与必要证书文件
|
* 写入主配置、路由配置与必要证书文件
|
||||||
* 先执行 `openresty -t`
|
* 写入新配置后以运行态恢复为目标执行激活,Docker 模式优先重建容器并确认容器保持运行
|
||||||
* 成功后执行 `openresty -s reload`
|
* 新配置激活失败时必须先尝试用目标配置恢复运行,再回滚到旧配置并重新拉起 OpenResty
|
||||||
* 失败时自动回滚并上报最终结果
|
* 回滚后 OpenResty 恢复正常时上报警告;回滚后仍无法恢复运行时上报失败
|
||||||
|
|
||||||
## 6. 测试与交付要求
|
## 6. 测试与交付要求
|
||||||
|
|
||||||
|
|||||||
@@ -138,6 +138,13 @@ func (e *DockerExecutor) Reload(ctx context.Context) error {
|
|||||||
}
|
}
|
||||||
output, err = e.Runner.Run(ctx, e.DockerBinary, "exec", e.ContainerName, dockerRuntimeCommand, "-s", "reload")
|
output, err = e.Runner.Run(ctx, e.DockerBinary, "exec", e.ContainerName, dockerRuntimeCommand, "-s", "reload")
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
if e.shouldRecreateAfterReloadFailure(string(output)) {
|
||||||
|
slog.Warn("docker openresty reload failed due to missing mounted files, recreating container", "container", e.ContainerName)
|
||||||
|
if recreateErr := e.EnsureRuntime(ctx, true); recreateErr != nil {
|
||||||
|
return fmt.Errorf("docker exec %s reload failed: %w: %s; recreate failed: %v", dockerRuntimeCommand, err, string(output), recreateErr)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
return fmt.Errorf("docker exec %s reload failed: %w: %s", dockerRuntimeCommand, err, string(output))
|
return fmt.Errorf("docker exec %s reload failed: %w: %s", dockerRuntimeCommand, err, string(output))
|
||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
@@ -216,6 +223,9 @@ func (e *DockerExecutor) runContainer(ctx context.Context) error {
|
|||||||
if runErr != nil {
|
if runErr != nil {
|
||||||
return fmt.Errorf("docker run openresty failed: %w: %s", runErr, string(runOutput))
|
return fmt.Errorf("docker run openresty failed: %w: %s", runErr, string(runOutput))
|
||||||
}
|
}
|
||||||
|
if err := e.CheckHealth(ctx); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
slog.Info("docker openresty container started", "container", e.ContainerName)
|
slog.Info("docker openresty container started", "container", e.ContainerName)
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
@@ -236,6 +246,28 @@ func (e *DockerExecutor) validateMountSources() error {
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (e *DockerExecutor) shouldRecreateAfterReloadFailure(output string) bool {
|
||||||
|
text := strings.ToLower(strings.TrimSpace(output))
|
||||||
|
if text == "" {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
if !strings.Contains(text, "no such file") && !strings.Contains(text, "cannot load certificate") {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
paths := []string{
|
||||||
|
strings.ToLower(e.NginxCertDir),
|
||||||
|
strings.ToLower(e.NginxLuaDir),
|
||||||
|
strings.ToLower(DockerMainConfigPath),
|
||||||
|
strings.ToLower("/etc/nginx/conf.d"),
|
||||||
|
}
|
||||||
|
for _, path := range paths {
|
||||||
|
if strings.TrimSpace(path) != "" && strings.Contains(text, path) {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
func (e *DockerExecutor) containerNotRunningError(ctx context.Context) error {
|
func (e *DockerExecutor) containerNotRunningError(ctx context.Context) error {
|
||||||
inspectSummary := ""
|
inspectSummary := ""
|
||||||
inspectOutput, inspectErr := e.Runner.Run(ctx, e.DockerBinary, "inspect", "-f", "status={{.State.Status}} exit_code={{.State.ExitCode}} error={{printf \"%q\" .State.Error}} oom_killed={{.State.OOMKilled}} finished_at={{.State.FinishedAt}}", e.ContainerName)
|
inspectOutput, inspectErr := e.Runner.Run(ctx, e.DockerBinary, "inspect", "-f", "status={{.State.Status}} exit_code={{.State.ExitCode}} error={{printf \"%q\" .State.Error}} oom_killed={{.State.OOMKilled}} finished_at={{.State.FinishedAt}}", e.ContainerName)
|
||||||
@@ -326,43 +358,89 @@ type Manager struct {
|
|||||||
Executor Executor
|
Executor Executor
|
||||||
}
|
}
|
||||||
|
|
||||||
func (m *Manager) Apply(ctx context.Context, mainConfig string, routeConfig string, supportFiles []protocol.SupportFile) error {
|
type ApplyStatus string
|
||||||
|
|
||||||
|
const (
|
||||||
|
ApplyStatusSuccess ApplyStatus = "success"
|
||||||
|
ApplyStatusWarning ApplyStatus = "warning"
|
||||||
|
ApplyStatusFatal ApplyStatus = "fatal"
|
||||||
|
)
|
||||||
|
|
||||||
|
type ApplyOutcome struct {
|
||||||
|
Status ApplyStatus
|
||||||
|
Message string
|
||||||
|
}
|
||||||
|
|
||||||
|
func (m *Manager) Apply(ctx context.Context, mainConfig string, routeConfig string, supportFiles []protocol.SupportFile) ApplyOutcome {
|
||||||
slog.Info("openresty apply started", "main_config", m.MainConfigPath, "route_config", m.RouteConfigPath, "cert_files", len(supportFiles))
|
slog.Info("openresty apply started", "main_config", m.MainConfigPath, "route_config", m.RouteConfigPath, "cert_files", len(supportFiles))
|
||||||
backup, err := m.backup()
|
backup, err := m.backup()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
return fatalApplyOutcome(fmt.Errorf("backup openresty config failed: %w", err))
|
||||||
|
}
|
||||||
|
if err = m.writeTargetFiles(mainConfig, routeConfig, supportFiles); err != nil {
|
||||||
|
return m.rollbackAfterFailedApply(ctx, backup, fmt.Errorf("write openresty config failed: %w", err))
|
||||||
|
}
|
||||||
|
if err = m.activateConfig(ctx); err != nil {
|
||||||
|
return m.rollbackAfterFailedApply(ctx, backup, fmt.Errorf("activate openresty runtime failed: %w", err))
|
||||||
|
}
|
||||||
|
slog.Info("openresty apply completed successfully", "main_config", m.MainConfigPath, "route_config", m.RouteConfigPath)
|
||||||
|
return ApplyOutcome{Status: ApplyStatusSuccess}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (m *Manager) writeTargetFiles(mainConfig string, routeConfig string, supportFiles []protocol.SupportFile) error {
|
||||||
|
if err := m.EnsureLuaAssets(); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
if err = m.EnsureLuaAssets(); err != nil {
|
if err := m.writeCertFiles(supportFiles); err != nil {
|
||||||
slog.Error("writing lua assets failed, restoring backup", "error", err)
|
|
||||||
_ = m.restore(backup)
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
if err = m.writeCertFiles(supportFiles); err != nil {
|
|
||||||
slog.Error("writing cert files failed, restoring backup", "error", err)
|
|
||||||
_ = m.restore(backup)
|
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
renderedMainConfig := m.renderMainConfig(mainConfig)
|
renderedMainConfig := m.renderMainConfig(mainConfig)
|
||||||
if err = os.WriteFile(m.MainConfigPath, []byte(renderedMainConfig), 0o644); err != nil {
|
if err := os.WriteFile(m.MainConfigPath, []byte(renderedMainConfig), 0o644); err != nil {
|
||||||
slog.Error("writing openresty main config failed, restoring backup", "error", err)
|
|
||||||
_ = m.restore(backup)
|
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
renderedRouteConfig := m.renderRouteConfig(routeConfig)
|
renderedRouteConfig := m.renderRouteConfig(routeConfig)
|
||||||
if err = os.WriteFile(m.RouteConfigPath, []byte(renderedRouteConfig), 0o644); err != nil {
|
if err := os.WriteFile(m.RouteConfigPath, []byte(renderedRouteConfig), 0o644); err != nil {
|
||||||
slog.Error("writing openresty route config failed, restoring backup", "error", err)
|
|
||||||
_ = m.restore(backup)
|
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
if err = m.Executor.Reload(ctx); err != nil {
|
|
||||||
slog.Error("openresty reload failed after config write, restoring backup", "error", err)
|
|
||||||
_ = m.restore(backup)
|
|
||||||
return err
|
|
||||||
}
|
|
||||||
slog.Info("openresty apply completed successfully", "main_config", m.MainConfigPath, "route_config", m.RouteConfigPath)
|
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (m *Manager) activateConfig(ctx context.Context) error {
|
||||||
|
if m.Executor == nil {
|
||||||
|
return errors.New("executor 未配置")
|
||||||
|
}
|
||||||
|
if _, ok := m.Executor.(*DockerExecutor); ok {
|
||||||
|
return m.Executor.EnsureRuntime(ctx, true)
|
||||||
|
}
|
||||||
|
return m.Executor.Reload(ctx)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (m *Manager) rollbackAfterFailedApply(ctx context.Context, backup *backupState, applyErr error) ApplyOutcome {
|
||||||
|
slog.Warn("openresty apply failed, restoring previous config", "error", applyErr)
|
||||||
|
if err := m.restore(backup); err != nil {
|
||||||
|
return fatalApplyOutcome(fmt.Errorf("restore openresty backup failed after apply error %v: %w", applyErr, err))
|
||||||
|
}
|
||||||
|
if err := m.activateConfig(ctx); err != nil {
|
||||||
|
return fatalApplyOutcome(fmt.Errorf("apply failed: %v; rollback recovery failed: %w", applyErr, err))
|
||||||
|
}
|
||||||
|
message := fmt.Sprintf("apply failed, rolled back to previous config: %v", applyErr)
|
||||||
|
slog.Warn("openresty apply rolled back successfully", "message", message)
|
||||||
|
return ApplyOutcome{
|
||||||
|
Status: ApplyStatusWarning,
|
||||||
|
Message: message,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func fatalApplyOutcome(err error) ApplyOutcome {
|
||||||
|
if err == nil {
|
||||||
|
return ApplyOutcome{Status: ApplyStatusFatal}
|
||||||
|
}
|
||||||
|
return ApplyOutcome{
|
||||||
|
Status: ApplyStatusFatal,
|
||||||
|
Message: strings.TrimSpace(err.Error()),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func (m *Manager) EnsureLuaAssets() error {
|
func (m *Manager) EnsureLuaAssets() error {
|
||||||
if strings.TrimSpace(m.LuaDir) == "" {
|
if strings.TrimSpace(m.LuaDir) == "" {
|
||||||
return nil
|
return nil
|
||||||
|
|||||||
@@ -28,6 +28,11 @@ type fakeExecutor struct {
|
|||||||
reloadErr error
|
reloadErr error
|
||||||
}
|
}
|
||||||
|
|
||||||
|
type scriptedExecutor struct {
|
||||||
|
reloadErrors []error
|
||||||
|
reloadCalls int
|
||||||
|
}
|
||||||
|
|
||||||
func (r *fakeRunner) Run(ctx context.Context, name string, args ...string) ([]byte, error) {
|
func (r *fakeRunner) Run(ctx context.Context, name string, args ...string) ([]byte, error) {
|
||||||
r.calls = append(r.calls, runCall{name: name, args: append([]string{}, args...)})
|
r.calls = append(r.calls, runCall{name: name, args: append([]string{}, args...)})
|
||||||
if r.runFn != nil {
|
if r.runFn != nil {
|
||||||
@@ -56,6 +61,31 @@ func (e *fakeExecutor) Restart(ctx context.Context) error {
|
|||||||
return e.reloadErr
|
return e.reloadErr
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (e *scriptedExecutor) Test(ctx context.Context) error {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (e *scriptedExecutor) Reload(ctx context.Context) error {
|
||||||
|
index := e.reloadCalls
|
||||||
|
e.reloadCalls++
|
||||||
|
if index >= len(e.reloadErrors) {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
return e.reloadErrors[index]
|
||||||
|
}
|
||||||
|
|
||||||
|
func (e *scriptedExecutor) EnsureRuntime(ctx context.Context, recreate bool) error {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (e *scriptedExecutor) CheckHealth(ctx context.Context) error {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (e *scriptedExecutor) Restart(ctx context.Context) error {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
func TestPathExecutorCommands(t *testing.T) {
|
func TestPathExecutorCommands(t *testing.T) {
|
||||||
runner := &fakeRunner{}
|
runner := &fakeRunner{}
|
||||||
executor := &PathExecutor{
|
executor := &PathExecutor{
|
||||||
@@ -209,10 +239,15 @@ func TestDockerExecutorStartsContainerWhenMissing(t *testing.T) {
|
|||||||
|
|
||||||
func TestDockerExecutorStartsStoppedContainer(t *testing.T) {
|
func TestDockerExecutorStartsStoppedContainer(t *testing.T) {
|
||||||
mainConfigPath, routeConfigDir, certDir, luaDir := prepareDockerMountSources(t)
|
mainConfigPath, routeConfigDir, certDir, luaDir := prepareDockerMountSources(t)
|
||||||
|
inspectCalls := 0
|
||||||
runner := &fakeRunner{
|
runner := &fakeRunner{
|
||||||
runFn: func(name string, args ...string) ([]byte, error) {
|
runFn: func(name string, args ...string) ([]byte, error) {
|
||||||
if len(args) >= 2 && args[0] == "inspect" {
|
if len(args) >= 2 && args[0] == "inspect" {
|
||||||
return []byte("false"), nil
|
inspectCalls++
|
||||||
|
if inspectCalls < 3 {
|
||||||
|
return []byte("false"), nil
|
||||||
|
}
|
||||||
|
return []byte("true"), nil
|
||||||
}
|
}
|
||||||
return []byte("ok"), nil
|
return []byte("ok"), nil
|
||||||
},
|
},
|
||||||
@@ -234,8 +269,8 @@ func TestDockerExecutorStartsStoppedContainer(t *testing.T) {
|
|||||||
t.Fatalf("Reload failed: %v", err)
|
t.Fatalf("Reload failed: %v", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
if len(runner.calls) != 4 {
|
if len(runner.calls) != 5 {
|
||||||
t.Fatalf("expected 4 calls, got %d", len(runner.calls))
|
t.Fatalf("expected 5 calls, got %d", len(runner.calls))
|
||||||
}
|
}
|
||||||
if runner.calls[0].args[0] != "inspect" {
|
if runner.calls[0].args[0] != "inspect" {
|
||||||
t.Fatalf("expected docker inspect on first call, got %#v", runner.calls[0])
|
t.Fatalf("expected docker inspect on first call, got %#v", runner.calls[0])
|
||||||
@@ -249,6 +284,9 @@ func TestDockerExecutorStartsStoppedContainer(t *testing.T) {
|
|||||||
if runner.calls[3].args[0] != "run" {
|
if runner.calls[3].args[0] != "run" {
|
||||||
t.Fatalf("expected docker run on fourth call, got %#v", runner.calls[3])
|
t.Fatalf("expected docker run on fourth call, got %#v", runner.calls[3])
|
||||||
}
|
}
|
||||||
|
if runner.calls[4].args[0] != "inspect" {
|
||||||
|
t.Fatalf("expected docker inspect after run, got %#v", runner.calls[4])
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestDockerExecutorReloadsRunningContainerInPlace(t *testing.T) {
|
func TestDockerExecutorReloadsRunningContainerInPlace(t *testing.T) {
|
||||||
@@ -287,9 +325,55 @@ func TestDockerExecutorReloadsRunningContainerInPlace(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestDockerExecutorReloadRecreatesContainerWhenMountedCertMissing(t *testing.T) {
|
||||||
|
mainConfigPath, routeConfigDir, certDir, luaDir := prepareDockerMountSources(t)
|
||||||
|
runner := &fakeRunner{
|
||||||
|
runFn: func(name string, args ...string) ([]byte, error) {
|
||||||
|
if len(args) >= 1 && args[0] == "inspect" {
|
||||||
|
return []byte("true"), nil
|
||||||
|
}
|
||||||
|
if len(args) >= 2 && args[0] == "exec" {
|
||||||
|
return []byte(`nginx: [emerg] cannot load certificate "/etc/nginx/openflare-certs/1.crt": BIO_new_file() failed (SSL: error:80000002:system library::No such file or directory)`), errors.New("exit status 1")
|
||||||
|
}
|
||||||
|
return []byte("ok"), nil
|
||||||
|
},
|
||||||
|
}
|
||||||
|
executor := &DockerExecutor{
|
||||||
|
DockerBinary: "docker",
|
||||||
|
ContainerName: "openflare-openresty",
|
||||||
|
Image: "openresty/openresty:alpine",
|
||||||
|
MainConfigPath: mainConfigPath,
|
||||||
|
RouteConfigDir: routeConfigDir,
|
||||||
|
CertDir: certDir,
|
||||||
|
NginxCertDir: "/etc/nginx/openflare-certs",
|
||||||
|
LuaDir: luaDir,
|
||||||
|
NginxLuaDir: "/etc/nginx/openflare-lua",
|
||||||
|
OpenrestyObservabilityPort: 18081,
|
||||||
|
Runner: runner,
|
||||||
|
}
|
||||||
|
|
||||||
|
if err := executor.Reload(context.Background()); err != nil {
|
||||||
|
t.Fatalf("Reload failed: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if len(runner.calls) != 6 {
|
||||||
|
t.Fatalf("expected 6 calls, got %d", len(runner.calls))
|
||||||
|
}
|
||||||
|
if runner.calls[2].args[0] != "inspect" || runner.calls[3].args[0] != "rm" || runner.calls[4].args[0] != "run" || runner.calls[5].args[0] != "inspect" {
|
||||||
|
t.Fatalf("expected recreate after reload failure, got %#v", runner.calls)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func TestDockerExecutorRunContainerMountsManagedFiles(t *testing.T) {
|
func TestDockerExecutorRunContainerMountsManagedFiles(t *testing.T) {
|
||||||
mainConfigPath, routeConfigDir, certDir, luaDir := prepareDockerMountSources(t)
|
mainConfigPath, routeConfigDir, certDir, luaDir := prepareDockerMountSources(t)
|
||||||
runner := &fakeRunner{}
|
runner := &fakeRunner{
|
||||||
|
runFn: func(name string, args ...string) ([]byte, error) {
|
||||||
|
if len(args) >= 1 && args[0] == "inspect" {
|
||||||
|
return []byte("true"), nil
|
||||||
|
}
|
||||||
|
return []byte("ok"), nil
|
||||||
|
},
|
||||||
|
}
|
||||||
executor := &DockerExecutor{
|
executor := &DockerExecutor{
|
||||||
DockerBinary: "docker",
|
DockerBinary: "docker",
|
||||||
ContainerName: "openflare-openresty",
|
ContainerName: "openflare-openresty",
|
||||||
@@ -308,8 +392,8 @@ func TestDockerExecutorRunContainerMountsManagedFiles(t *testing.T) {
|
|||||||
t.Fatalf("runContainer failed: %v", err)
|
t.Fatalf("runContainer failed: %v", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
if len(runner.calls) != 1 {
|
if len(runner.calls) != 2 {
|
||||||
t.Fatalf("expected one docker run call, got %d", len(runner.calls))
|
t.Fatalf("expected docker run plus health check, got %d calls", len(runner.calls))
|
||||||
}
|
}
|
||||||
|
|
||||||
expectedArgs := []string{
|
expectedArgs := []string{
|
||||||
@@ -327,6 +411,9 @@ func TestDockerExecutorRunContainerMountsManagedFiles(t *testing.T) {
|
|||||||
if !reflect.DeepEqual(runner.calls[0].args, expectedArgs) {
|
if !reflect.DeepEqual(runner.calls[0].args, expectedArgs) {
|
||||||
t.Fatalf("unexpected docker run args: %#v", runner.calls[0].args)
|
t.Fatalf("unexpected docker run args: %#v", runner.calls[0].args)
|
||||||
}
|
}
|
||||||
|
if !reflect.DeepEqual(runner.calls[1].args, []string{"inspect", "-f", "{{.State.Running}}", "openflare-openresty"}) {
|
||||||
|
t.Fatalf("unexpected docker health check args: %#v", runner.calls[1].args)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestDockerExecutorRecreatesContainerOnStartup(t *testing.T) {
|
func TestDockerExecutorRecreatesContainerOnStartup(t *testing.T) {
|
||||||
@@ -356,8 +443,8 @@ func TestDockerExecutorRecreatesContainerOnStartup(t *testing.T) {
|
|||||||
if err := executor.EnsureRuntime(context.Background(), true); err != nil {
|
if err := executor.EnsureRuntime(context.Background(), true); err != nil {
|
||||||
t.Fatalf("EnsureRuntime failed: %v", err)
|
t.Fatalf("EnsureRuntime failed: %v", err)
|
||||||
}
|
}
|
||||||
if len(runner.calls) != 3 {
|
if len(runner.calls) != 4 {
|
||||||
t.Fatalf("expected 3 calls, got %d", len(runner.calls))
|
t.Fatalf("expected 4 calls, got %d", len(runner.calls))
|
||||||
}
|
}
|
||||||
if runner.calls[1].args[0] != "rm" {
|
if runner.calls[1].args[0] != "rm" {
|
||||||
t.Fatalf("expected docker rm on second call, got %#v", runner.calls[1])
|
t.Fatalf("expected docker rm on second call, got %#v", runner.calls[1])
|
||||||
@@ -365,6 +452,9 @@ func TestDockerExecutorRecreatesContainerOnStartup(t *testing.T) {
|
|||||||
if runner.calls[2].args[0] != "run" {
|
if runner.calls[2].args[0] != "run" {
|
||||||
t.Fatalf("expected docker run on third call, got %#v", runner.calls[2])
|
t.Fatalf("expected docker run on third call, got %#v", runner.calls[2])
|
||||||
}
|
}
|
||||||
|
if runner.calls[3].args[0] != "inspect" {
|
||||||
|
t.Fatalf("expected docker inspect after run, got %#v", runner.calls[3])
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestDockerExecutorRunContainerRejectsMissingMainConfigFile(t *testing.T) {
|
func TestDockerExecutorRunContainerRejectsMissingMainConfigFile(t *testing.T) {
|
||||||
@@ -497,21 +587,21 @@ func TestManagerApplyAndChecksumIncludeMainConfig(t *testing.T) {
|
|||||||
Executor: &fakeExecutor{},
|
Executor: &fakeExecutor{},
|
||||||
}
|
}
|
||||||
|
|
||||||
err := manager.Apply(
|
outcome := manager.Apply(
|
||||||
context.Background(),
|
context.Background(),
|
||||||
"include __OPENFLARE_ROUTE_CONFIG__;\naccess_log __OPENFLARE_ACCESS_LOG__ openflare_json;\n",
|
"include __OPENFLARE_ROUTE_CONFIG__;\naccess_log __OPENFLARE_ACCESS_LOG__ openflare_json;\n",
|
||||||
"ssl_certificate __OPENFLARE_CERT_DIR__/1.crt;\n",
|
"ssl_certificate __OPENFLARE_CERT_DIR__/1.crt;\n",
|
||||||
[]protocol.SupportFile{{Path: "1.crt", Content: "cert"}},
|
[]protocol.SupportFile{{Path: "1.crt", Content: "cert"}},
|
||||||
)
|
)
|
||||||
if err != nil {
|
if outcome.Status != ApplyStatusSuccess {
|
||||||
t.Fatalf("Apply failed: %v", err)
|
t.Fatalf("Apply failed: %#v", outcome)
|
||||||
}
|
}
|
||||||
|
|
||||||
mainData, err := os.ReadFile(mainPath)
|
mainData, err := os.ReadFile(mainPath)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("failed to read main config: %v", err)
|
t.Fatalf("failed to read main config: %v", err)
|
||||||
}
|
}
|
||||||
expectedMain := "include " + routePath + ";\naccess_log " + filepath.Join(filepath.Dir(routePath), "openflare_access.log") + " openflare_json;\n"
|
expectedMain := "include " + routePath + ";\naccess_log " + filepath.ToSlash(filepath.Join(filepath.Dir(routePath), "openflare_access.log")) + " openflare_json;\n"
|
||||||
if string(mainData) != expectedMain {
|
if string(mainData) != expectedMain {
|
||||||
t.Fatalf("unexpected main config: %s", string(mainData))
|
t.Fatalf("unexpected main config: %s", string(mainData))
|
||||||
}
|
}
|
||||||
@@ -553,8 +643,8 @@ func TestManagerApplyUsesRuntimeRouteConfigPath(t *testing.T) {
|
|||||||
Executor: &fakeExecutor{},
|
Executor: &fakeExecutor{},
|
||||||
}
|
}
|
||||||
|
|
||||||
if err := manager.Apply(context.Background(), "include __OPENFLARE_ROUTE_CONFIG__;\naccess_log __OPENFLARE_ACCESS_LOG__ openflare_json;\n", "server { listen 80; }\n", nil); err != nil {
|
if outcome := manager.Apply(context.Background(), "include __OPENFLARE_ROUTE_CONFIG__;\naccess_log __OPENFLARE_ACCESS_LOG__ openflare_json;\n", "server { listen 80; }\n", nil); outcome.Status != ApplyStatusSuccess {
|
||||||
t.Fatalf("Apply failed: %v", err)
|
t.Fatalf("Apply failed: %#v", outcome)
|
||||||
}
|
}
|
||||||
|
|
||||||
mainData, err := os.ReadFile(mainPath)
|
mainData, err := os.ReadFile(mainPath)
|
||||||
@@ -631,12 +721,12 @@ func TestManagerApplyWritesSupportFilesAndReplacesPlaceholder(t *testing.T) {
|
|||||||
Executor: &fakeExecutor{},
|
Executor: &fakeExecutor{},
|
||||||
}
|
}
|
||||||
|
|
||||||
err := manager.Apply(context.Background(), "include __OPENFLARE_ROUTE_CONFIG__;\n__OPENFLARE_RESOLVER_DIRECTIVE__server { listen __OPENFLARE_OBSERVABILITY_LISTEN__; }", "ssl_certificate __OPENFLARE_CERT_DIR__/1.crt;", []protocol.SupportFile{
|
outcome := manager.Apply(context.Background(), "include __OPENFLARE_ROUTE_CONFIG__;\n__OPENFLARE_RESOLVER_DIRECTIVE__server { listen __OPENFLARE_OBSERVABILITY_LISTEN__; }", "ssl_certificate __OPENFLARE_CERT_DIR__/1.crt;", []protocol.SupportFile{
|
||||||
{Path: "1.crt", Content: "cert-data"},
|
{Path: "1.crt", Content: "cert-data"},
|
||||||
{Path: "1.key", Content: "key-data"},
|
{Path: "1.key", Content: "key-data"},
|
||||||
})
|
})
|
||||||
if err != nil {
|
if outcome.Status != ApplyStatusSuccess {
|
||||||
t.Fatalf("Apply failed: %v", err)
|
t.Fatalf("Apply failed: %#v", outcome)
|
||||||
}
|
}
|
||||||
|
|
||||||
routeData, err := os.ReadFile(manager.RouteConfigPath)
|
routeData, err := os.ReadFile(manager.RouteConfigPath)
|
||||||
@@ -667,7 +757,7 @@ func TestManagerApplyWritesSupportFilesAndReplacesPlaceholder(t *testing.T) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("expected managed lua file to exist, stat err = %v", err)
|
t.Fatalf("expected managed lua file to exist, stat err = %v", err)
|
||||||
}
|
}
|
||||||
if luaInfo.Mode().Perm() != 0o644 {
|
if runtime.GOOS != "windows" && luaInfo.Mode().Perm() != 0o644 {
|
||||||
t.Fatalf("unexpected lua mode: %o", luaInfo.Mode().Perm())
|
t.Fatalf("unexpected lua mode: %o", luaInfo.Mode().Perm())
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -802,11 +892,11 @@ func TestManagerRollbackRestoresCertFiles(t *testing.T) {
|
|||||||
},
|
},
|
||||||
}
|
}
|
||||||
|
|
||||||
err := manager.Apply(context.Background(), "new-main", "new-route", []protocol.SupportFile{
|
outcome := manager.Apply(context.Background(), "new-main", "new-route", []protocol.SupportFile{
|
||||||
{Path: "1.crt", Content: "new-cert"},
|
{Path: "1.crt", Content: "new-cert"},
|
||||||
})
|
})
|
||||||
if err == nil {
|
if outcome.Status != ApplyStatusFatal {
|
||||||
t.Fatal("expected Apply to fail")
|
t.Fatalf("expected fatal apply outcome, got %#v", outcome)
|
||||||
}
|
}
|
||||||
|
|
||||||
mainData, err := os.ReadFile(mainPath)
|
mainData, err := os.ReadFile(mainPath)
|
||||||
@@ -832,6 +922,51 @@ func TestManagerRollbackRestoresCertFiles(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestManagerApplyReturnsWarningWhenRollbackRecoversRuntime(t *testing.T) {
|
||||||
|
tempDir := t.TempDir()
|
||||||
|
routePath := filepath.Join(tempDir, "routes.conf")
|
||||||
|
mainPath := filepath.Join(tempDir, "nginx.conf")
|
||||||
|
certDir := filepath.Join(tempDir, "certs")
|
||||||
|
if err := os.MkdirAll(certDir, 0o755); err != nil {
|
||||||
|
t.Fatalf("MkdirAll failed: %v", err)
|
||||||
|
}
|
||||||
|
if err := os.WriteFile(mainPath, []byte("old-main"), 0o644); err != nil {
|
||||||
|
t.Fatalf("WriteFile failed: %v", err)
|
||||||
|
}
|
||||||
|
if err := os.WriteFile(routePath, []byte("old-route"), 0o644); err != nil {
|
||||||
|
t.Fatalf("WriteFile failed: %v", err)
|
||||||
|
}
|
||||||
|
if err := os.WriteFile(filepath.Join(certDir, "1.crt"), []byte("old-cert"), 0o600); err != nil {
|
||||||
|
t.Fatalf("WriteFile failed: %v", err)
|
||||||
|
}
|
||||||
|
manager := &Manager{
|
||||||
|
MainConfigPath: mainPath,
|
||||||
|
RouteConfigPath: routePath,
|
||||||
|
CertDir: certDir,
|
||||||
|
NginxCertDir: "/etc/nginx/openflare-certs",
|
||||||
|
LuaDir: filepath.Join(tempDir, "lua"),
|
||||||
|
NginxLuaDir: "/etc/nginx/openflare-lua",
|
||||||
|
Executor: &scriptedExecutor{
|
||||||
|
reloadErrors: []error{errors.New("target config failed"), nil},
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
outcome := manager.Apply(context.Background(), "new-main", "new-route", []protocol.SupportFile{
|
||||||
|
{Path: "1.crt", Content: "new-cert"},
|
||||||
|
})
|
||||||
|
if outcome.Status != ApplyStatusWarning {
|
||||||
|
t.Fatalf("expected warning apply outcome, got %#v", outcome)
|
||||||
|
}
|
||||||
|
|
||||||
|
mainData, err := os.ReadFile(mainPath)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("failed to read main config: %v", err)
|
||||||
|
}
|
||||||
|
if string(mainData) != "old-main" {
|
||||||
|
t.Fatalf("expected main rollback, got %s", string(mainData))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func TestManagerCertFileTargetPathRejectsEscapes(t *testing.T) {
|
func TestManagerCertFileTargetPathRejectsEscapes(t *testing.T) {
|
||||||
manager := &Manager{CertDir: filepath.Join(t.TempDir(), "certs")}
|
manager := &Manager{CertDir: filepath.Join(t.TempDir(), "certs")}
|
||||||
if err := os.MkdirAll(manager.CertDir, 0o755); err != nil {
|
if err := os.MkdirAll(manager.CertDir, 0o755); err != nil {
|
||||||
@@ -883,11 +1018,11 @@ func TestManagerApplyRejectsCertFilePathTraversal(t *testing.T) {
|
|||||||
Executor: &fakeExecutor{},
|
Executor: &fakeExecutor{},
|
||||||
}
|
}
|
||||||
|
|
||||||
err := manager.Apply(context.Background(), "main", "route", []protocol.SupportFile{
|
outcome := manager.Apply(context.Background(), "main", "route", []protocol.SupportFile{
|
||||||
{Path: "../escape.crt", Content: "bad"},
|
{Path: "../escape.crt", Content: "bad"},
|
||||||
})
|
})
|
||||||
if err == nil {
|
if outcome.Status != ApplyStatusWarning {
|
||||||
t.Fatal("expected Apply to reject traversal path")
|
t.Fatalf("expected warning apply outcome, got %#v", outcome)
|
||||||
}
|
}
|
||||||
|
|
||||||
if _, statErr := os.Stat(filepath.Join(tempDir, "escape.crt")); !os.IsNotExist(statErr) {
|
if _, statErr := os.Stat(filepath.Join(tempDir, "escape.crt")); !os.IsNotExist(statErr) {
|
||||||
|
|||||||
@@ -4,15 +4,18 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"crypto/sha256"
|
"crypto/sha256"
|
||||||
"encoding/hex"
|
"encoding/hex"
|
||||||
|
"fmt"
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
|
"openflare-agent/internal/nginx"
|
||||||
"openflare-agent/internal/protocol"
|
"openflare-agent/internal/protocol"
|
||||||
"openflare-agent/internal/state"
|
"openflare-agent/internal/state"
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
ApplyResultSuccess = "success"
|
ApplyResultSuccess = "success"
|
||||||
|
ApplyResultWarning = "warning"
|
||||||
ApplyResultFailed = "failed"
|
ApplyResultFailed = "failed"
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -22,7 +25,7 @@ type ConfigClient interface {
|
|||||||
}
|
}
|
||||||
|
|
||||||
type NginxManager interface {
|
type NginxManager interface {
|
||||||
Apply(ctx context.Context, mainConfig string, routeConfig string, supportFiles []protocol.SupportFile) error
|
Apply(ctx context.Context, mainConfig string, routeConfig string, supportFiles []protocol.SupportFile) nginx.ApplyOutcome
|
||||||
EnsureRuntime(ctx context.Context, recreate bool) error
|
EnsureRuntime(ctx context.Context, recreate bool) error
|
||||||
CurrentChecksum() (string, error)
|
CurrentChecksum() (string, error)
|
||||||
}
|
}
|
||||||
@@ -154,55 +157,79 @@ func (s *Service) applyIfNeeded(ctx context.Context, mode string, startup bool,
|
|||||||
mainConfigChecksum := checksumString(config.MainConfig)
|
mainConfigChecksum := checksumString(config.MainConfig)
|
||||||
routeConfigChecksum := checksumString(routeConfig)
|
routeConfigChecksum := checksumString(routeConfig)
|
||||||
slog.Info("applying new openresty config", "mode", mode, "from_version", snapshot.CurrentVersion, "to_version", config.Version, "old_checksum", currentChecksum, "new_checksum", config.Checksum)
|
slog.Info("applying new openresty config", "mode", mode, "from_version", snapshot.CurrentVersion, "to_version", config.Version, "old_checksum", currentChecksum, "new_checksum", config.Checksum)
|
||||||
if err := s.nginxManager.Apply(ctx, config.MainConfig, routeConfig, config.SupportFiles); err != nil {
|
outcome := s.nginxManager.Apply(ctx, config.MainConfig, routeConfig, config.SupportFiles)
|
||||||
slog.Error("apply openresty config failed", "mode", mode, "version", config.Version, "error", err)
|
message := strings.TrimSpace(outcome.Message)
|
||||||
snapshot.LastError = err.Error()
|
if outcome.Status == "" {
|
||||||
snapshot.OpenrestyStatus = protocol.OpenrestyStatusUnhealthy
|
outcome.Status = nginx.ApplyStatusFatal
|
||||||
snapshot.OpenrestyMessage = err.Error()
|
if message == "" {
|
||||||
_ = s.stateStore.Save(snapshot)
|
message = "openresty apply returned empty outcome"
|
||||||
reportErr := s.client.ReportApplyLog(ctx, protocol.ApplyLogPayload{
|
|
||||||
NodeID: snapshot.NodeID,
|
|
||||||
Version: config.Version,
|
|
||||||
Result: ApplyResultFailed,
|
|
||||||
Message: err.Error(),
|
|
||||||
Checksum: config.Checksum,
|
|
||||||
MainConfigChecksum: mainConfigChecksum,
|
|
||||||
RouteConfigChecksum: routeConfigChecksum,
|
|
||||||
SupportFileCount: len(config.SupportFiles),
|
|
||||||
})
|
|
||||||
if reportErr != nil {
|
|
||||||
slog.Error("report failed apply log failed", "version", config.Version, "error", reportErr)
|
|
||||||
return reportErr
|
|
||||||
}
|
}
|
||||||
slog.Warn("failed apply log reported", "version", config.Version)
|
|
||||||
return err
|
|
||||||
}
|
}
|
||||||
slog.Info("openresty config applied successfully", "mode", mode, "version", config.Version)
|
|
||||||
snapshot.CurrentVersion = config.Version
|
reportResult := ApplyResultFailed
|
||||||
snapshot.CurrentChecksum = config.Checksum
|
switch outcome.Status {
|
||||||
snapshot.LastError = ""
|
case nginx.ApplyStatusSuccess:
|
||||||
snapshot.OpenrestyStatus = protocol.OpenrestyStatusHealthy
|
slog.Info("openresty config applied successfully", "mode", mode, "version", config.Version)
|
||||||
snapshot.OpenrestyMessage = ""
|
snapshot.CurrentVersion = config.Version
|
||||||
|
snapshot.CurrentChecksum = config.Checksum
|
||||||
|
snapshot.LastError = ""
|
||||||
|
snapshot.OpenrestyStatus = protocol.OpenrestyStatusHealthy
|
||||||
|
snapshot.OpenrestyMessage = ""
|
||||||
|
reportResult = ApplyResultSuccess
|
||||||
|
if message == "" {
|
||||||
|
message = "apply success"
|
||||||
|
}
|
||||||
|
case nginx.ApplyStatusWarning:
|
||||||
|
if message == "" {
|
||||||
|
message = "apply rolled back to previous config"
|
||||||
|
}
|
||||||
|
slog.Warn("openresty config apply rolled back", "mode", mode, "version", config.Version, "message", message)
|
||||||
|
snapshot.LastError = message
|
||||||
|
snapshot.OpenrestyStatus = protocol.OpenrestyStatusHealthy
|
||||||
|
snapshot.OpenrestyMessage = message
|
||||||
|
reportResult = ApplyResultWarning
|
||||||
|
default:
|
||||||
|
if message == "" {
|
||||||
|
message = "openresty apply failed"
|
||||||
|
}
|
||||||
|
slog.Error("apply openresty config failed", "mode", mode, "version", config.Version, "message", message)
|
||||||
|
snapshot.LastError = message
|
||||||
|
snapshot.OpenrestyStatus = protocol.OpenrestyStatusUnhealthy
|
||||||
|
snapshot.OpenrestyMessage = message
|
||||||
|
}
|
||||||
|
|
||||||
if err := s.stateStore.Save(snapshot); err != nil {
|
if err := s.stateStore.Save(snapshot); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
if err := s.client.ReportApplyLog(ctx, protocol.ApplyLogPayload{
|
if err := s.client.ReportApplyLog(ctx, protocol.ApplyLogPayload{
|
||||||
NodeID: snapshot.NodeID,
|
NodeID: snapshot.NodeID,
|
||||||
Version: config.Version,
|
Version: config.Version,
|
||||||
Result: ApplyResultSuccess,
|
Result: reportResult,
|
||||||
Message: "apply success",
|
Message: message,
|
||||||
Checksum: config.Checksum,
|
Checksum: config.Checksum,
|
||||||
MainConfigChecksum: mainConfigChecksum,
|
MainConfigChecksum: mainConfigChecksum,
|
||||||
RouteConfigChecksum: routeConfigChecksum,
|
RouteConfigChecksum: routeConfigChecksum,
|
||||||
SupportFileCount: len(config.SupportFiles),
|
SupportFileCount: len(config.SupportFiles),
|
||||||
}); err != nil {
|
}); err != nil {
|
||||||
slog.Error("report successful apply log failed", "version", config.Version, "error", err)
|
slog.Error("report apply log failed", "version", config.Version, "result", reportResult, "error", err)
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
slog.Debug("successful apply log reported", "version", config.Version)
|
if reportResult == ApplyResultFailed {
|
||||||
|
slog.Warn("failed apply log reported", "version", config.Version)
|
||||||
|
return outcomeError(config.Version, message)
|
||||||
|
}
|
||||||
|
slog.Debug("apply log reported", "version", config.Version, "result", reportResult)
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func outcomeError(version string, message string) error {
|
||||||
|
trimmed := strings.TrimSpace(message)
|
||||||
|
if trimmed == "" {
|
||||||
|
trimmed = "openresty apply failed"
|
||||||
|
}
|
||||||
|
return fmt.Errorf("apply version %s failed: %s", version, trimmed)
|
||||||
|
}
|
||||||
|
|
||||||
func checksumString(content string) string {
|
func checksumString(content string) string {
|
||||||
sum := sha256.Sum256([]byte(content))
|
sum := sha256.Sum256([]byte(content))
|
||||||
return hex.EncodeToString(sum[:])
|
return hex.EncodeToString(sum[:])
|
||||||
|
|||||||
@@ -24,7 +24,7 @@ type fakeClient struct {
|
|||||||
}
|
}
|
||||||
|
|
||||||
type fakeManager struct {
|
type fakeManager struct {
|
||||||
applyErr error
|
applyOutcome nginx.ApplyOutcome
|
||||||
currentChecksum string
|
currentChecksum string
|
||||||
currentChecksumErr error
|
currentChecksumErr error
|
||||||
ensureErr error
|
ensureErr error
|
||||||
@@ -64,11 +64,14 @@ func (f *fakeClient) ReportApplyLog(ctx context.Context, payload protocol.ApplyL
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (m *fakeManager) Apply(ctx context.Context, mainConfig string, routeConfig string, supportFiles []protocol.SupportFile) error {
|
func (m *fakeManager) Apply(ctx context.Context, mainConfig string, routeConfig string, supportFiles []protocol.SupportFile) nginx.ApplyOutcome {
|
||||||
m.applyMainContents = append(m.applyMainContents, mainConfig)
|
m.applyMainContents = append(m.applyMainContents, mainConfig)
|
||||||
m.applyRouteContents = append(m.applyRouteContents, routeConfig)
|
m.applyRouteContents = append(m.applyRouteContents, routeConfig)
|
||||||
m.applyFiles = append(m.applyFiles, append([]protocol.SupportFile(nil), supportFiles...))
|
m.applyFiles = append(m.applyFiles, append([]protocol.SupportFile(nil), supportFiles...))
|
||||||
return m.applyErr
|
if m.applyOutcome.Status == "" {
|
||||||
|
return nginx.ApplyOutcome{Status: nginx.ApplyStatusSuccess}
|
||||||
|
}
|
||||||
|
return m.applyOutcome
|
||||||
}
|
}
|
||||||
|
|
||||||
func (m *fakeManager) EnsureRuntime(ctx context.Context, recreate bool) error {
|
func (m *fakeManager) EnsureRuntime(ctx context.Context, recreate bool) error {
|
||||||
@@ -166,17 +169,7 @@ func TestSyncOnceRollbackOnNginxFailure(t *testing.T) {
|
|||||||
},
|
},
|
||||||
}
|
}
|
||||||
|
|
||||||
tempDir := t.TempDir()
|
stateStore := state.NewStore(filepath.Join(t.TempDir(), "state.json"))
|
||||||
mainPath := filepath.Join(tempDir, "nginx.conf")
|
|
||||||
routePath := filepath.Join(tempDir, "routes.conf")
|
|
||||||
if err := os.WriteFile(mainPath, []byte("worker_processes auto;"), 0o644); err != nil {
|
|
||||||
t.Fatalf("failed to seed main file: %v", err)
|
|
||||||
}
|
|
||||||
if err := os.WriteFile(routePath, []byte("server { listen 80; }"), 0o644); err != nil {
|
|
||||||
t.Fatalf("failed to seed route file: %v", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
stateStore := state.NewStore(filepath.Join(tempDir, "state.json"))
|
|
||||||
nodeID, err := stateStore.EnsureNodeID()
|
nodeID, err := stateStore.EnsureNodeID()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("EnsureNodeID failed: %v", err)
|
t.Fatalf("EnsureNodeID failed: %v", err)
|
||||||
@@ -189,11 +182,10 @@ func TestSyncOnceRollbackOnNginxFailure(t *testing.T) {
|
|||||||
t.Fatalf("failed to seed state: %v", err)
|
t.Fatalf("failed to seed state: %v", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
service := New(client, &nginx.Manager{
|
service := New(client, &fakeManager{
|
||||||
MainConfigPath: mainPath,
|
applyOutcome: nginx.ApplyOutcome{
|
||||||
RouteConfigPath: routePath,
|
Status: nginx.ApplyStatusFatal,
|
||||||
Executor: &fakeExecutor{
|
Message: "openresty failed after rollback",
|
||||||
testErr: context.DeadlineExceeded,
|
|
||||||
},
|
},
|
||||||
}, stateStore)
|
}, stateStore)
|
||||||
|
|
||||||
@@ -202,22 +194,7 @@ func TestSyncOnceRollbackOnNginxFailure(t *testing.T) {
|
|||||||
Checksum: client.config.Checksum,
|
Checksum: client.config.Checksum,
|
||||||
})
|
})
|
||||||
if err == nil {
|
if err == nil {
|
||||||
t.Fatal("expected SyncOnce to fail when nginx test fails")
|
t.Fatal("expected SyncOnce to fail when apply outcome is fatal")
|
||||||
}
|
|
||||||
|
|
||||||
data, readErr := os.ReadFile(routePath)
|
|
||||||
if readErr != nil {
|
|
||||||
t.Fatalf("failed to read route file after rollback: %v", readErr)
|
|
||||||
}
|
|
||||||
if string(data) != "server { listen 80; }" {
|
|
||||||
t.Fatal("expected original route config to be restored after rollback")
|
|
||||||
}
|
|
||||||
mainData, readErr := os.ReadFile(mainPath)
|
|
||||||
if readErr != nil {
|
|
||||||
t.Fatalf("failed to read main file after rollback: %v", readErr)
|
|
||||||
}
|
|
||||||
if string(mainData) != "worker_processes auto;" {
|
|
||||||
t.Fatal("expected original main config to be restored after rollback")
|
|
||||||
}
|
}
|
||||||
snapshot, loadErr := stateStore.Load()
|
snapshot, loadErr := stateStore.Load()
|
||||||
if loadErr != nil {
|
if loadErr != nil {
|
||||||
@@ -226,6 +203,9 @@ func TestSyncOnceRollbackOnNginxFailure(t *testing.T) {
|
|||||||
if snapshot.CurrentVersion != "20260309-001" {
|
if snapshot.CurrentVersion != "20260309-001" {
|
||||||
t.Fatal("expected failed sync not to overwrite current version")
|
t.Fatal("expected failed sync not to overwrite current version")
|
||||||
}
|
}
|
||||||
|
if snapshot.OpenrestyStatus != protocol.OpenrestyStatusUnhealthy {
|
||||||
|
t.Fatalf("expected unhealthy openresty status, got %q", snapshot.OpenrestyStatus)
|
||||||
|
}
|
||||||
if len(client.reports) != 1 || client.reports[0].Result != ApplyResultFailed {
|
if len(client.reports) != 1 || client.reports[0].Result != ApplyResultFailed {
|
||||||
t.Fatal("expected failed apply report to be sent")
|
t.Fatal("expected failed apply report to be sent")
|
||||||
}
|
}
|
||||||
@@ -240,6 +220,64 @@ func TestSyncOnceRollbackOnNginxFailure(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestSyncOnceReportsWarningWhenRollbackKeepsOpenrestyHealthy(t *testing.T) {
|
||||||
|
client := &fakeClient{
|
||||||
|
config: protocol.ActiveConfigResponse{
|
||||||
|
Version: "20260309-002",
|
||||||
|
Checksum: "checksum-2",
|
||||||
|
MainConfig: "worker_processes 2;",
|
||||||
|
RouteConfig: "server { listen 81; }",
|
||||||
|
RenderedConfig: "server { listen 81; }",
|
||||||
|
SupportFiles: []protocol.SupportFile{{Path: "1.crt", Content: "cert"}},
|
||||||
|
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,
|
||||||
|
CurrentVersion: "20260309-001",
|
||||||
|
CurrentChecksum: "checksum-1",
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatalf("failed to seed state: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
service := New(client, &fakeManager{
|
||||||
|
applyOutcome: nginx.ApplyOutcome{
|
||||||
|
Status: nginx.ApplyStatusWarning,
|
||||||
|
Message: "apply failed, rolled back to previous config",
|
||||||
|
},
|
||||||
|
}, stateStore)
|
||||||
|
|
||||||
|
if err = service.SyncOnce(context.Background(), &protocol.ActiveConfigMeta{
|
||||||
|
Version: client.config.Version,
|
||||||
|
Checksum: client.config.Checksum,
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatalf("expected warning outcome to keep sync successful, got %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
snapshot, err := stateStore.Load()
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("failed to load state: %v", err)
|
||||||
|
}
|
||||||
|
if snapshot.CurrentVersion != "20260309-001" || snapshot.CurrentChecksum != "checksum-1" {
|
||||||
|
t.Fatal("expected warning apply to keep previous version state")
|
||||||
|
}
|
||||||
|
if snapshot.OpenrestyStatus != protocol.OpenrestyStatusHealthy {
|
||||||
|
t.Fatalf("expected healthy openresty after rollback, got %q", snapshot.OpenrestyStatus)
|
||||||
|
}
|
||||||
|
if snapshot.LastError == "" {
|
||||||
|
t.Fatal("expected rollback warning to be recorded")
|
||||||
|
}
|
||||||
|
if len(client.reports) != 1 || client.reports[0].Result != ApplyResultWarning {
|
||||||
|
t.Fatal("expected warning apply report to be sent")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func TestSyncOnStartupRecreatesRuntimeWhenChecksumMatches(t *testing.T) {
|
func TestSyncOnStartupRecreatesRuntimeWhenChecksumMatches(t *testing.T) {
|
||||||
client := &fakeClient{
|
client := &fakeClient{
|
||||||
config: protocol.ActiveConfigResponse{
|
config: protocol.ActiveConfigResponse{
|
||||||
|
|||||||
@@ -17,6 +17,7 @@ const (
|
|||||||
NodeStatusOffline = "offline"
|
NodeStatusOffline = "offline"
|
||||||
NodeStatusPending = "pending"
|
NodeStatusPending = "pending"
|
||||||
ApplyResultOK = "success"
|
ApplyResultOK = "success"
|
||||||
|
ApplyResultWarning = "warning"
|
||||||
ApplyResultFailed = "failed"
|
ApplyResultFailed = "failed"
|
||||||
OpenrestyStatusHealthy = "healthy"
|
OpenrestyStatusHealthy = "healthy"
|
||||||
OpenrestyStatusUnhealthy = "unhealthy"
|
OpenrestyStatusUnhealthy = "unhealthy"
|
||||||
@@ -234,8 +235,8 @@ func ReportApplyLog(payload ApplyLogPayload) (*model.ApplyLog, error) {
|
|||||||
if payload.Version == "" {
|
if payload.Version == "" {
|
||||||
return nil, errors.New("version 不能为空")
|
return nil, errors.New("version 不能为空")
|
||||||
}
|
}
|
||||||
if payload.Result != ApplyResultOK && payload.Result != ApplyResultFailed {
|
if payload.Result != ApplyResultOK && payload.Result != ApplyResultWarning && payload.Result != ApplyResultFailed {
|
||||||
return nil, errors.New("result 仅支持 success 或 failed")
|
return nil, errors.New("result 仅支持 success、warning 或 failed")
|
||||||
}
|
}
|
||||||
slog.Debug("agent apply log received", "node_id", payload.NodeID, "version", payload.Version, "result", payload.Result)
|
slog.Debug("agent apply log received", "node_id", payload.NodeID, "version", payload.Version, "result", payload.Result)
|
||||||
|
|
||||||
@@ -273,6 +274,8 @@ func ReportApplyLog(payload ApplyLogPayload) (*model.ApplyLog, error) {
|
|||||||
}
|
}
|
||||||
if payload.Result == ApplyResultOK {
|
if payload.Result == ApplyResultOK {
|
||||||
slog.Debug("agent apply reported success", "node_id", payload.NodeID, "version", payload.Version)
|
slog.Debug("agent apply reported success", "node_id", payload.NodeID, "version", payload.Version)
|
||||||
|
} else if payload.Result == ApplyResultWarning {
|
||||||
|
slog.Warn("agent apply reported warning", "node_id", payload.NodeID, "version", payload.Version, "message", payload.Message)
|
||||||
} else {
|
} else {
|
||||||
slog.Error("agent apply reported failure", "node_id", payload.NodeID, "version", payload.Version, "message", payload.Message)
|
slog.Error("agent apply reported failure", "node_id", payload.NodeID, "version", payload.Version, "message", payload.Message)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -31,6 +31,10 @@ function getResultMeta(result: string) {
|
|||||||
return { label: '成功', variant: 'success' as const };
|
return { label: '成功', variant: 'success' as const };
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if (result === 'warning') {
|
||||||
|
return { label: '警告', variant: 'warning' as const };
|
||||||
|
}
|
||||||
|
|
||||||
return { label: '失败', variant: 'danger' as const };
|
return { label: '失败', variant: 'danger' as const };
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -23,7 +23,7 @@ export interface NodeItem {
|
|||||||
current_version: string;
|
current_version: string;
|
||||||
last_seen_at: string;
|
last_seen_at: string;
|
||||||
last_error: string;
|
last_error: string;
|
||||||
latest_apply_result: 'success' | 'failed' | '';
|
latest_apply_result: 'success' | 'warning' | 'failed' | '';
|
||||||
latest_apply_message: string;
|
latest_apply_message: string;
|
||||||
latest_apply_checksum: string;
|
latest_apply_checksum: string;
|
||||||
latest_main_config_checksum: string;
|
latest_main_config_checksum: string;
|
||||||
|
|||||||
@@ -33,6 +33,10 @@ export function getApplyVariant(result: NodeItem['latest_apply_result']) {
|
|||||||
return 'success';
|
return 'success';
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if (result === 'warning') {
|
||||||
|
return 'warning';
|
||||||
|
}
|
||||||
|
|
||||||
if (result === 'failed') {
|
if (result === 'failed') {
|
||||||
return 'danger';
|
return 'danger';
|
||||||
}
|
}
|
||||||
@@ -45,6 +49,10 @@ export function getApplyLabel(result: NodeItem['latest_apply_result']) {
|
|||||||
return '成功';
|
return '成功';
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if (result === 'warning') {
|
||||||
|
return '警告';
|
||||||
|
}
|
||||||
|
|
||||||
if (result === 'failed') {
|
if (result === 'failed') {
|
||||||
return '失败';
|
return '失败';
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user