Refactor observability configuration and support files

- Introduced `filterCertificateSupportFiles` function to filter support files for certificates in `agent.go`.
- Updated placeholder constants in `config_version.go` to reflect changes from support directory to certificate directory.
- Modified tests in `https_phase1_test.go` to align with new directory structure and removed obsolete checks for observability Lua scripts.
- Removed observability assets from `openresty_observability_assets.go` and created a new file `observability_assets.go` in `atsf_agent/internal/nginx` to manage observability Lua scripts.
- Updated documentation to reflect changes in configuration paths for certificates and Lua scripts.
- Ensured that the agent writes to the new `cert_dir` and `lua_dir` instead of the old `support_dir`.
This commit is contained in:
ryan
2026-03-15 12:32:24 +08:00
parent e1efbf3868
commit f4d36be2e6
13 changed files with 530 additions and 369 deletions
+2 -1
View File
@@ -215,7 +215,8 @@ Docker 镜像工作流仅构建 `atsf_server`,并产出 `linux/amd64` 与 `lin
| `agent_token` | 节点专属认证 Token | | `agent_token` | 节点专属认证 Token |
| `discovery_token` | 首次自动注册使用的全局 Token | | `discovery_token` | 首次自动注册使用的全局 Token |
| `data_dir` | Agent 托管数据目录 | | `data_dir` | Agent 托管数据目录 |
| `support_dir` | Agent 存放受管附属文件的目录,当前包含证书与 Lua 观测脚本 | | `cert_dir` | Agent 存放受管证书文件的目录 |
| `lua_dir` | Agent 存放受管 Lua 观测脚本的目录;启动时自动覆盖释放 |
| `openresty_observability_port` | Agent 读取 OpenResty Lua 本地观测指标的 loopback 端口 | | `openresty_observability_port` | Agent 读取 OpenResty Lua 本地观测指标的 loopback 端口 |
| `observability_buffer_path` | Agent 本地观测补报缓冲文件路径 | | `observability_buffer_path` | Agent 本地观测补报缓冲文件路径 |
| `observability_replay_minutes` | Agent 恢复 heartbeat 后允许自动补传的最近观测窗口分钟数 | | `observability_replay_minutes` | Agent 恢复 heartbeat 后允许自动补传的最近观测窗口分钟数 |
+18 -7
View File
@@ -39,8 +39,10 @@ func main() {
Image: cfg.OpenrestyDockerImage, Image: cfg.OpenrestyDockerImage,
MainConfigPath: cfg.MainConfigPath, MainConfigPath: cfg.MainConfigPath,
RouteConfigPath: cfg.RouteConfigPath, RouteConfigPath: cfg.RouteConfigPath,
SupportDir: cfg.SupportDir, CertDir: cfg.CertDir,
NginxSupportDir: cfg.OpenrestySupportDir, NginxCertDir: cfg.OpenrestyCertDir,
LuaDir: cfg.LuaDir,
NginxLuaDir: cfg.OpenrestyLuaDir,
OpenrestyObservabilityPort: cfg.OpenrestyObservabilityPort, OpenrestyObservabilityPort: cfg.OpenrestyObservabilityPort,
}, },
) )
@@ -50,7 +52,8 @@ func main() {
"ip", cfg.NodeIP, "ip", cfg.NodeIP,
"heartbeat_interval", cfg.HeartbeatInterval, "heartbeat_interval", cfg.HeartbeatInterval,
"route_config", cfg.RouteConfigPath, "route_config", cfg.RouteConfigPath,
"support_dir", cfg.SupportDir, "cert_dir", cfg.CertDir,
"lua_dir", cfg.LuaDir,
) )
client := httpclient.New(cfg.ServerURL, cfg.InitialAuthToken(), cfg.RequestTimeout.Duration()) client := httpclient.New(cfg.ServerURL, cfg.InitialAuthToken(), cfg.RequestTimeout.Duration())
@@ -64,8 +67,10 @@ func main() {
MainConfigPath: cfg.MainConfigPath, MainConfigPath: cfg.MainConfigPath,
RouteConfigPath: cfg.RouteConfigPath, RouteConfigPath: cfg.RouteConfigPath,
RuntimeRouteConfigPath: runtimeRouteConfigPath, RuntimeRouteConfigPath: runtimeRouteConfigPath,
SupportDir: cfg.SupportDir, CertDir: cfg.CertDir,
NginxSupportDir: cfg.OpenrestySupportDir, NginxCertDir: cfg.OpenrestyCertDir,
LuaDir: cfg.LuaDir,
NginxLuaDir: cfg.OpenrestyLuaDir,
OpenrestyObservabilityListen: nginx.ObservabilityListenAddress(cfg.OpenrestyPath, cfg.OpenrestyObservabilityPort), OpenrestyObservabilityListen: nginx.ObservabilityListenAddress(cfg.OpenrestyPath, cfg.OpenrestyObservabilityPort),
OpenrestyObservabilityPort: cfg.OpenrestyObservabilityPort, OpenrestyObservabilityPort: cfg.OpenrestyObservabilityPort,
Executor: nginx.NewExecutor(nginx.ExecutorOptions{ Executor: nginx.NewExecutor(nginx.ExecutorOptions{
@@ -75,11 +80,17 @@ func main() {
Image: cfg.OpenrestyDockerImage, Image: cfg.OpenrestyDockerImage,
MainConfigPath: cfg.MainConfigPath, MainConfigPath: cfg.MainConfigPath,
RouteConfigPath: cfg.RouteConfigPath, RouteConfigPath: cfg.RouteConfigPath,
SupportDir: cfg.SupportDir, CertDir: cfg.CertDir,
NginxSupportDir: cfg.OpenrestySupportDir, NginxCertDir: cfg.OpenrestyCertDir,
LuaDir: cfg.LuaDir,
NginxLuaDir: cfg.OpenrestyLuaDir,
OpenrestyObservabilityPort: cfg.OpenrestyObservabilityPort, OpenrestyObservabilityPort: cfg.OpenrestyObservabilityPort,
}), }),
} }
if err = runtimeManager.EnsureLuaAssets(); err != nil {
slog.Error("ensure managed lua assets failed", "error", err)
os.Exit(1)
}
runner := &agent.Runner{ runner := &agent.Runner{
Config: cfg, Config: cfg,
StateStore: stateStore, StateStore: stateStore,
+41 -19
View File
@@ -14,10 +14,12 @@ import (
const ( const (
defaultDockerMainConfigRelativePath = "etc/nginx/nginx.conf" defaultDockerMainConfigRelativePath = "etc/nginx/nginx.conf"
defaultDockerRouteConfigRelativePath = "etc/nginx/conf.d/atsflare_routes.conf" defaultDockerRouteConfigRelativePath = "etc/nginx/conf.d/atsflare_routes.conf"
defaultSupportDirRelativePath = "etc/nginx/support" defaultCertDirRelativePath = "etc/nginx/certs"
defaultLuaDirRelativePath = "etc/nginx/lua"
defaultDockerStateRelativePath = "var/lib/atsflare/agent-state.json" defaultDockerStateRelativePath = "var/lib/atsflare/agent-state.json"
defaultObservabilityBufferRelativePath = "var/lib/atsflare/observability-buffer.json" defaultObservabilityBufferRelativePath = "var/lib/atsflare/observability-buffer.json"
defaultDockerOpenRestySupportDir = "/etc/nginx/atsflare-support" defaultDockerOpenRestyCertDir = "/etc/nginx/atsflare-certs"
defaultDockerOpenRestyLuaDir = "/etc/nginx/atsflare-lua"
defaultOpenRestyObservabilityPort = 18081 defaultOpenRestyObservabilityPort = 18081
defaultObservabilityReplayMinutes = 15 defaultObservabilityReplayMinutes = 15
) )
@@ -37,8 +39,10 @@ type Config struct {
DataDir string `json:"data_dir"` DataDir string `json:"data_dir"`
MainConfigPath string `json:"main_config_path"` MainConfigPath string `json:"main_config_path"`
RouteConfigPath string `json:"route_config_path"` RouteConfigPath string `json:"route_config_path"`
SupportDir string `json:"support_dir"` CertDir string `json:"cert_dir"`
OpenrestySupportDir string `json:"openresty_support_dir"` OpenrestyCertDir string `json:"openresty_cert_dir"`
LuaDir string `json:"lua_dir"`
OpenrestyLuaDir string `json:"openresty_lua_dir"`
OpenrestyObservabilityPort int `json:"openresty_observability_port"` OpenrestyObservabilityPort int `json:"openresty_observability_port"`
ObservabilityBufferPath string `json:"observability_buffer_path"` ObservabilityBufferPath string `json:"observability_buffer_path"`
ObservabilityReplayMinutes int `json:"observability_replay_minutes"` ObservabilityReplayMinutes int `json:"observability_replay_minutes"`
@@ -61,10 +65,10 @@ type configFile struct {
DataDir string `json:"data_dir"` DataDir string `json:"data_dir"`
MainConfigPath string `json:"main_config_path"` MainConfigPath string `json:"main_config_path"`
RouteConfigPath string `json:"route_config_path"` RouteConfigPath string `json:"route_config_path"`
SupportDir string `json:"support_dir"` CertDir string `json:"cert_dir"`
OpenrestySupportDir string `json:"openresty_support_dir"` OpenrestyCertDir string `json:"openresty_cert_dir"`
LegacyCertDir string `json:"cert_dir"` LuaDir string `json:"lua_dir"`
LegacyOpenrestyCertDir string `json:"openresty_cert_dir"` OpenrestyLuaDir string `json:"openresty_lua_dir"`
OpenrestyObservabilityPort int `json:"openresty_observability_port"` OpenrestyObservabilityPort int `json:"openresty_observability_port"`
ObservabilityBufferPath string `json:"observability_buffer_path"` ObservabilityBufferPath string `json:"observability_buffer_path"`
ObservabilityReplayMinutes int `json:"observability_replay_minutes"` ObservabilityReplayMinutes int `json:"observability_replay_minutes"`
@@ -95,8 +99,10 @@ func Load(path string) (*Config, error) {
DataDir: file.DataDir, DataDir: file.DataDir,
MainConfigPath: file.MainConfigPath, MainConfigPath: file.MainConfigPath,
RouteConfigPath: file.RouteConfigPath, RouteConfigPath: file.RouteConfigPath,
SupportDir: firstNonEmpty(file.SupportDir, file.LegacyCertDir), CertDir: file.CertDir,
OpenrestySupportDir: firstNonEmpty(file.OpenrestySupportDir, file.LegacyOpenrestyCertDir), OpenrestyCertDir: file.OpenrestyCertDir,
LuaDir: file.LuaDir,
OpenrestyLuaDir: file.OpenrestyLuaDir,
OpenrestyObservabilityPort: file.OpenrestyObservabilityPort, OpenrestyObservabilityPort: file.OpenrestyObservabilityPort,
ObservabilityBufferPath: file.ObservabilityBufferPath, ObservabilityBufferPath: file.ObservabilityBufferPath,
ObservabilityReplayMinutes: file.ObservabilityReplayMinutes, ObservabilityReplayMinutes: file.ObservabilityReplayMinutes,
@@ -148,14 +154,24 @@ func applyDefaults(cfg *Config, baseDir string) {
cfg.StatePath = joinManagedPath(cfg.DataDir, defaultDockerStateRelativePath) cfg.StatePath = joinManagedPath(cfg.DataDir, defaultDockerStateRelativePath)
} }
} }
if cfg.SupportDir == "" { if cfg.CertDir == "" {
cfg.SupportDir = joinManagedPath(cfg.DataDir, defaultSupportDirRelativePath) cfg.CertDir = joinManagedPath(cfg.DataDir, defaultCertDirRelativePath)
} }
if cfg.OpenrestySupportDir == "" { if cfg.OpenrestyCertDir == "" {
if cfg.OpenrestyPath != "" { if cfg.OpenrestyPath != "" {
cfg.OpenrestySupportDir = cfg.SupportDir cfg.OpenrestyCertDir = cfg.CertDir
} else { } else {
cfg.OpenrestySupportDir = defaultDockerOpenRestySupportDir cfg.OpenrestyCertDir = defaultDockerOpenRestyCertDir
}
}
if cfg.LuaDir == "" {
cfg.LuaDir = joinManagedPath(cfg.DataDir, defaultLuaDirRelativePath)
}
if cfg.OpenrestyLuaDir == "" {
if cfg.OpenrestyPath != "" {
cfg.OpenrestyLuaDir = cfg.LuaDir
} else {
cfg.OpenrestyLuaDir = defaultDockerOpenRestyLuaDir
} }
} }
if cfg.OpenrestyObservabilityPort <= 0 { if cfg.OpenrestyObservabilityPort <= 0 {
@@ -189,11 +205,17 @@ func normalizeManagedPaths(cfg *Config) {
if usesSlashPath(cfg.RouteConfigPath) { if usesSlashPath(cfg.RouteConfigPath) {
cfg.RouteConfigPath = filepath.ToSlash(cfg.RouteConfigPath) cfg.RouteConfigPath = filepath.ToSlash(cfg.RouteConfigPath)
} }
if usesSlashPath(cfg.SupportDir) { if usesSlashPath(cfg.CertDir) {
cfg.SupportDir = filepath.ToSlash(cfg.SupportDir) cfg.CertDir = filepath.ToSlash(cfg.CertDir)
} }
if usesSlashPath(cfg.OpenrestySupportDir) { if usesSlashPath(cfg.OpenrestyCertDir) {
cfg.OpenrestySupportDir = filepath.ToSlash(cfg.OpenrestySupportDir) cfg.OpenrestyCertDir = filepath.ToSlash(cfg.OpenrestyCertDir)
}
if usesSlashPath(cfg.LuaDir) {
cfg.LuaDir = filepath.ToSlash(cfg.LuaDir)
}
if usesSlashPath(cfg.OpenrestyLuaDir) {
cfg.OpenrestyLuaDir = filepath.ToSlash(cfg.OpenrestyLuaDir)
} }
if usesSlashPath(cfg.StatePath) { if usesSlashPath(cfg.StatePath) {
cfg.StatePath = filepath.ToSlash(cfg.StatePath) cfg.StatePath = filepath.ToSlash(cfg.StatePath)
+20 -8
View File
@@ -39,8 +39,11 @@ func TestLoadDockerModeUsesManagedPaths(t *testing.T) {
if cfg.RouteConfigPath != filepath.Join(dir, "data", defaultDockerRouteConfigRelativePath) { if cfg.RouteConfigPath != filepath.Join(dir, "data", defaultDockerRouteConfigRelativePath) {
t.Fatalf("unexpected route config path: %s", cfg.RouteConfigPath) t.Fatalf("unexpected route config path: %s", cfg.RouteConfigPath)
} }
if cfg.SupportDir != filepath.Join(dir, "data", defaultSupportDirRelativePath) { if cfg.CertDir != filepath.Join(dir, "data", defaultCertDirRelativePath) {
t.Fatalf("unexpected support dir: %s", cfg.SupportDir) t.Fatalf("unexpected cert dir: %s", cfg.CertDir)
}
if cfg.LuaDir != filepath.Join(dir, "data", defaultLuaDirRelativePath) {
t.Fatalf("unexpected lua dir: %s", cfg.LuaDir)
} }
if cfg.OpenrestyContainerName != "atsflare-openresty" { if cfg.OpenrestyContainerName != "atsflare-openresty" {
t.Fatalf("unexpected openresty container name: %s", cfg.OpenrestyContainerName) t.Fatalf("unexpected openresty container name: %s", cfg.OpenrestyContainerName)
@@ -48,8 +51,11 @@ func TestLoadDockerModeUsesManagedPaths(t *testing.T) {
if cfg.OpenrestyDockerImage != "openresty/openresty:alpine" { if cfg.OpenrestyDockerImage != "openresty/openresty:alpine" {
t.Fatalf("unexpected openresty image: %s", cfg.OpenrestyDockerImage) t.Fatalf("unexpected openresty image: %s", cfg.OpenrestyDockerImage)
} }
if cfg.OpenrestySupportDir != defaultDockerOpenRestySupportDir { if cfg.OpenrestyCertDir != defaultDockerOpenRestyCertDir {
t.Fatalf("unexpected openresty support dir: %s", cfg.OpenrestySupportDir) t.Fatalf("unexpected openresty cert dir: %s", cfg.OpenrestyCertDir)
}
if cfg.OpenrestyLuaDir != defaultDockerOpenRestyLuaDir {
t.Fatalf("unexpected openresty lua dir: %s", cfg.OpenrestyLuaDir)
} }
if cfg.StatePath != filepath.Join(dir, "data", defaultDockerStateRelativePath) { if cfg.StatePath != filepath.Join(dir, "data", defaultDockerStateRelativePath) {
t.Fatalf("unexpected state path: %s", cfg.StatePath) t.Fatalf("unexpected state path: %s", cfg.StatePath)
@@ -103,8 +109,11 @@ func TestLoadPathModeKeepsExplicitPaths(t *testing.T) {
if cfg.ObservabilityBufferPath != filepath.Join(dir, "data", defaultObservabilityBufferRelativePath) { if cfg.ObservabilityBufferPath != filepath.Join(dir, "data", defaultObservabilityBufferRelativePath) {
t.Fatalf("unexpected observability buffer path: %s", cfg.ObservabilityBufferPath) t.Fatalf("unexpected observability buffer path: %s", cfg.ObservabilityBufferPath)
} }
if cfg.OpenrestySupportDir != cfg.SupportDir { if cfg.OpenrestyCertDir != cfg.CertDir {
t.Fatalf("expected path mode openresty support dir to equal support dir, got %s / %s", cfg.OpenrestySupportDir, cfg.SupportDir) t.Fatalf("expected path mode openresty cert dir to equal cert dir, got %s / %s", cfg.OpenrestyCertDir, cfg.CertDir)
}
if cfg.OpenrestyLuaDir != cfg.LuaDir {
t.Fatalf("expected path mode openresty lua dir to equal lua dir, got %s / %s", cfg.OpenrestyLuaDir, cfg.LuaDir)
} }
if cfg.OpenrestyObservabilityPort != defaultOpenRestyObservabilityPort { if cfg.OpenrestyObservabilityPort != defaultOpenRestyObservabilityPort {
t.Fatalf("unexpected path mode openresty observability port: %d", cfg.OpenrestyObservabilityPort) t.Fatalf("unexpected path mode openresty observability port: %d", cfg.OpenrestyObservabilityPort)
@@ -146,8 +155,11 @@ func TestLoadUsesCustomDataDirForGeneratedFiles(t *testing.T) {
if cfg.ObservabilityBufferPath != "/srv/atsflare/"+defaultObservabilityBufferRelativePath { if cfg.ObservabilityBufferPath != "/srv/atsflare/"+defaultObservabilityBufferRelativePath {
t.Fatalf("unexpected observability buffer path: %s", cfg.ObservabilityBufferPath) t.Fatalf("unexpected observability buffer path: %s", cfg.ObservabilityBufferPath)
} }
if cfg.SupportDir != "/srv/atsflare/"+defaultSupportDirRelativePath { if cfg.CertDir != "/srv/atsflare/"+defaultCertDirRelativePath {
t.Fatalf("unexpected support dir: %s", cfg.SupportDir) t.Fatalf("unexpected cert dir: %s", cfg.CertDir)
}
if cfg.LuaDir != "/srv/atsflare/"+defaultLuaDirRelativePath {
t.Fatalf("unexpected lua dir: %s", cfg.LuaDir)
} }
} }
+140 -61
View File
@@ -18,7 +18,7 @@ import (
"atsflare-agent/internal/protocol" "atsflare-agent/internal/protocol"
) )
const SupportDirPlaceholder = "__ATSF_SUPPORT_DIR__" const CertDirPlaceholder = "__ATSF_CERT_DIR__"
const RouteConfigPlaceholder = "__ATSF_ROUTE_CONFIG__" const RouteConfigPlaceholder = "__ATSF_ROUTE_CONFIG__"
const AccessLogPlaceholder = "__ATSF_ACCESS_LOG__" const AccessLogPlaceholder = "__ATSF_ACCESS_LOG__"
const LuaDirPlaceholder = "__ATSF_LUA_DIR__" const LuaDirPlaceholder = "__ATSF_LUA_DIR__"
@@ -106,8 +106,10 @@ type DockerExecutor struct {
Image string Image string
MainConfigPath string MainConfigPath string
RouteConfigDir string RouteConfigDir string
SupportDir string CertDir string
NginxSupportDir string NginxCertDir string
LuaDir string
NginxLuaDir string
OpenrestyObservabilityPort int OpenrestyObservabilityPort int
Runner CommandRunner Runner CommandRunner
} }
@@ -188,7 +190,8 @@ func (e *DockerExecutor) runContainer(ctx context.Context) error {
"-p", fmt.Sprintf("127.0.0.1:%d:%d", e.OpenrestyObservabilityPort, e.OpenrestyObservabilityPort), "-p", fmt.Sprintf("127.0.0.1:%d:%d", e.OpenrestyObservabilityPort, e.OpenrestyObservabilityPort),
"-v", fmt.Sprintf("%s:%s", e.MainConfigPath, DockerMainConfigPath), "-v", fmt.Sprintf("%s:%s", e.MainConfigPath, DockerMainConfigPath),
"-v", fmt.Sprintf("%s:/etc/nginx/conf.d", e.RouteConfigDir), "-v", fmt.Sprintf("%s:/etc/nginx/conf.d", e.RouteConfigDir),
"-v", fmt.Sprintf("%s:%s", e.SupportDir, e.NginxSupportDir), "-v", fmt.Sprintf("%s:%s", e.CertDir, e.NginxCertDir),
"-v", fmt.Sprintf("%s:%s", e.LuaDir, e.NginxLuaDir),
e.Image, e.Image,
} }
runOutput, runErr := e.Runner.Run(ctx, e.DockerBinary, runArgs...) runOutput, runErr := e.Runner.Run(ctx, e.DockerBinary, runArgs...)
@@ -203,21 +206,28 @@ type Manager struct {
MainConfigPath string MainConfigPath string
RouteConfigPath string RouteConfigPath string
RuntimeRouteConfigPath string RuntimeRouteConfigPath string
SupportDir string CertDir string
NginxSupportDir string NginxCertDir string
LuaDir string
NginxLuaDir string
OpenrestyObservabilityListen string OpenrestyObservabilityListen string
OpenrestyObservabilityPort int OpenrestyObservabilityPort int
Executor Executor Executor Executor
} }
func (m *Manager) Apply(ctx context.Context, mainConfig string, routeConfig string, supportFiles []protocol.SupportFile) error { func (m *Manager) Apply(ctx context.Context, mainConfig string, routeConfig string, supportFiles []protocol.SupportFile) error {
slog.Info("openresty apply started", "main_config", m.MainConfigPath, "route_config", m.RouteConfigPath, "support_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 err return err
} }
if err = m.writeSupportFiles(supportFiles); err != nil { if err = m.EnsureLuaAssets(); err != nil {
slog.Error("writing support files failed, restoring backup", "error", err) 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) _ = m.restore(backup)
return err return err
} }
@@ -247,6 +257,31 @@ func (m *Manager) Apply(ctx context.Context, mainConfig string, routeConfig stri
return nil return nil
} }
func (m *Manager) EnsureLuaAssets() error {
if strings.TrimSpace(m.LuaDir) == "" {
return nil
}
if err := os.RemoveAll(m.LuaDir); err != nil && !os.IsNotExist(err) {
return err
}
if err := os.MkdirAll(m.LuaDir, 0o755); err != nil {
return err
}
for _, file := range ManagedObservabilityLuaFiles() {
targetPath, err := luaFileTargetPath(m.LuaDir, file.Path)
if err != nil {
return err
}
if err := os.MkdirAll(filepath.Dir(targetPath), 0o755); err != nil {
return err
}
if err := os.WriteFile(targetPath, []byte(file.Content), 0o644); err != nil {
return err
}
}
return nil
}
func (m *Manager) EnsureRuntime(ctx context.Context, recreate bool) error { func (m *Manager) EnsureRuntime(ctx context.Context, recreate bool) error {
if m.Executor == nil { if m.Executor == nil {
return errors.New("executor 未配置") return errors.New("executor 未配置")
@@ -308,15 +343,15 @@ func (m *Manager) CurrentChecksum() (string, error) {
normalizedMain = strings.ReplaceAll(normalizedMain, fmt.Sprintf("%d", m.OpenrestyObservabilityPort), ObservabilityPortPlaceholder) normalizedMain = strings.ReplaceAll(normalizedMain, fmt.Sprintf("%d", m.OpenrestyObservabilityPort), ObservabilityPortPlaceholder)
} }
normalizedRoute := string(data) normalizedRoute := string(data)
if m.NginxSupportDir != "" { if m.NginxCertDir != "" {
normalizedRoute = strings.ReplaceAll(normalizedRoute, m.NginxSupportDir, SupportDirPlaceholder) normalizedRoute = strings.ReplaceAll(normalizedRoute, m.NginxCertDir, CertDirPlaceholder)
} }
files, err := m.readSupportFiles() files, err := m.readCertFiles()
if err != nil { if err != nil {
return "", err return "", err
} }
result := bundleChecksum(normalizedMain, normalizedRoute, files) result := bundleChecksum(normalizedMain, normalizedRoute, files)
slog.Debug("openresty current checksum calculated", "main_config", m.MainConfigPath, "route_config", m.RouteConfigPath, "checksum", result, "support_files", len(files)) slog.Debug("openresty current checksum calculated", "main_config", m.MainConfigPath, "route_config", m.RouteConfigPath, "checksum", result, "cert_files", len(files))
return result, nil return result, nil
} }
@@ -327,8 +362,10 @@ type ExecutorOptions struct {
Image string Image string
MainConfigPath string MainConfigPath string
RouteConfigPath string RouteConfigPath string
SupportDir string CertDir string
NginxSupportDir string NginxCertDir string
LuaDir string
NginxLuaDir string
OpenrestyObservabilityPort int OpenrestyObservabilityPort int
} }
@@ -341,16 +378,28 @@ func NewExecutor(options ExecutorOptions) Executor {
} }
} }
mainConfigPath := options.MainConfigPath mainConfigPath := options.MainConfigPath
if absPath, err := filepath.Abs(mainConfigPath); err == nil { if mainConfigPath != "" {
mainConfigPath = absPath if absPath, err := filepath.Abs(mainConfigPath); err == nil {
mainConfigPath = absPath
}
} }
routeConfigDir := filepath.Dir(options.RouteConfigPath) routeConfigDir := filepath.Dir(options.RouteConfigPath)
if absDir, err := filepath.Abs(routeConfigDir); err == nil { if options.RouteConfigPath != "" {
routeConfigDir = absDir if absDir, err := filepath.Abs(routeConfigDir); err == nil {
routeConfigDir = absDir
}
} }
supportDir := options.SupportDir certDir := options.CertDir
if absDir, err := filepath.Abs(supportDir); err == nil { if certDir != "" {
supportDir = absDir if absDir, err := filepath.Abs(certDir); err == nil {
certDir = absDir
}
}
luaDir := options.LuaDir
if luaDir != "" {
if absDir, err := filepath.Abs(luaDir); err == nil {
luaDir = absDir
}
} }
return &DockerExecutor{ return &DockerExecutor{
DockerBinary: options.DockerBinary, DockerBinary: options.DockerBinary,
@@ -358,8 +407,10 @@ func NewExecutor(options ExecutorOptions) Executor {
Image: options.Image, Image: options.Image,
MainConfigPath: mainConfigPath, MainConfigPath: mainConfigPath,
RouteConfigDir: routeConfigDir, RouteConfigDir: routeConfigDir,
SupportDir: supportDir, CertDir: certDir,
NginxSupportDir: options.NginxSupportDir, NginxCertDir: options.NginxCertDir,
LuaDir: luaDir,
NginxLuaDir: options.NginxLuaDir,
OpenrestyObservabilityPort: options.OpenrestyObservabilityPort, OpenrestyObservabilityPort: options.OpenrestyObservabilityPort,
Runner: runner, Runner: runner,
} }
@@ -432,7 +483,9 @@ func (e *DockerExecutor) runEphemeralRuntimeCommandWithBinary(ctx context.Contex
"-v", "-v",
fmt.Sprintf("%s:/etc/nginx/conf.d", e.RouteConfigDir), fmt.Sprintf("%s:/etc/nginx/conf.d", e.RouteConfigDir),
"-v", "-v",
fmt.Sprintf("%s:%s", e.SupportDir, e.NginxSupportDir), fmt.Sprintf("%s:%s", e.CertDir, e.NginxCertDir),
"-v",
fmt.Sprintf("%s:%s", e.LuaDir, e.NginxLuaDir),
e.Image, e.Image,
runtimeBinary, runtimeBinary,
} }
@@ -465,8 +518,8 @@ func (m *Manager) backup() (*backupState, error) {
if err := os.MkdirAll(filepath.Dir(m.RouteConfigPath), 0o755); err != nil { if err := os.MkdirAll(filepath.Dir(m.RouteConfigPath), 0o755); err != nil {
return nil, err return nil, err
} }
if m.SupportDir != "" { if m.CertDir != "" {
if err := os.MkdirAll(m.SupportDir, 0o755); err != nil { if err := os.MkdirAll(m.CertDir, 0o755); err != nil {
return nil, err return nil, err
} }
} }
@@ -485,12 +538,12 @@ func (m *Manager) backup() (*backupState, error) {
} else if !os.IsNotExist(err) { } else if !os.IsNotExist(err) {
return nil, err return nil, err
} }
files, err := m.readSupportFiles() files, err := m.readCertFiles()
if err != nil { if err != nil {
return nil, err return nil, err
} }
state.Files = files state.Files = files
slog.Debug("backup captured", "main_exists", state.MainExisted, "route_exists", state.RouteExisted, "support_files", len(state.Files)) slog.Debug("backup captured", "main_exists", state.MainExisted, "route_exists", state.RouteExisted, "cert_files", len(state.Files))
return state, nil return state, nil
} }
@@ -498,7 +551,7 @@ func (m *Manager) restore(state *backupState) error {
if state == nil { if state == nil {
return nil return nil
} }
slog.Warn("restoring nginx backup", "main_existed", state.MainExisted, "route_existed", state.RouteExisted, "support_files", len(state.Files)) slog.Warn("restoring nginx backup", "main_existed", state.MainExisted, "route_existed", state.RouteExisted, "cert_files", len(state.Files))
if state.MainExisted { if state.MainExisted {
if err := os.WriteFile(m.MainConfigPath, state.MainData, 0o644); err != nil { if err := os.WriteFile(m.MainConfigPath, state.MainData, 0o644); err != nil {
return err return err
@@ -513,67 +566,67 @@ func (m *Manager) restore(state *backupState) error {
} else if err := os.Remove(m.RouteConfigPath); err != nil && !os.IsNotExist(err) { } else if err := os.Remove(m.RouteConfigPath); err != nil && !os.IsNotExist(err) {
return err return err
} }
if m.SupportDir == "" { if m.CertDir == "" {
return nil return nil
} }
if err := os.RemoveAll(m.SupportDir); err != nil && !os.IsNotExist(err) { if err := os.RemoveAll(m.CertDir); err != nil && !os.IsNotExist(err) {
return err return err
} }
if err := os.MkdirAll(m.SupportDir, 0o755); err != nil { if err := os.MkdirAll(m.CertDir, 0o755); err != nil {
return err return err
} }
for _, file := range state.Files { for _, file := range state.Files {
targetPath, err := m.supportFileTargetPath(file.Path) targetPath, err := m.certFileTargetPath(file.Path)
if err != nil { if err != nil {
return err return err
} }
if err := os.MkdirAll(filepath.Dir(targetPath), 0o755); err != nil { if err := os.MkdirAll(filepath.Dir(targetPath), 0o755); err != nil {
return err return err
} }
if err := os.WriteFile(targetPath, []byte(file.Content), supportFileMode(file.Path)); err != nil { if err := os.WriteFile(targetPath, []byte(file.Content), certFileMode(file.Path)); err != nil {
return err return err
} }
} }
return nil return nil
} }
func (m *Manager) writeSupportFiles(supportFiles []protocol.SupportFile) error { func (m *Manager) writeCertFiles(certFiles []protocol.SupportFile) error {
if m.SupportDir == "" { if m.CertDir == "" {
return nil return nil
} }
if err := os.RemoveAll(m.SupportDir); err != nil && !os.IsNotExist(err) { if err := os.RemoveAll(m.CertDir); err != nil && !os.IsNotExist(err) {
return err return err
} }
if err := os.MkdirAll(m.SupportDir, 0o755); err != nil { if err := os.MkdirAll(m.CertDir, 0o755); err != nil {
return err return err
} }
for _, file := range supportFiles { for _, file := range certFiles {
targetPath, err := m.supportFileTargetPath(file.Path) targetPath, err := m.certFileTargetPath(file.Path)
if err != nil { if err != nil {
return err return err
} }
if err := os.MkdirAll(filepath.Dir(targetPath), 0o755); err != nil { if err := os.MkdirAll(filepath.Dir(targetPath), 0o755); err != nil {
return err return err
} }
if err := os.WriteFile(targetPath, []byte(file.Content), supportFileMode(file.Path)); err != nil { if err := os.WriteFile(targetPath, []byte(file.Content), certFileMode(file.Path)); err != nil {
return err return err
} }
} }
return nil return nil
} }
func (m *Manager) readSupportFiles() ([]protocol.SupportFile, error) { func (m *Manager) readCertFiles() ([]protocol.SupportFile, error) {
if m.SupportDir == "" { if m.CertDir == "" {
return nil, nil return nil, nil
} }
if _, err := os.Stat(m.SupportDir); err != nil { if _, err := os.Stat(m.CertDir); err != nil {
if os.IsNotExist(err) { if os.IsNotExist(err) {
return nil, nil return nil, nil
} }
return nil, err return nil, err
} }
files := make([]protocol.SupportFile, 0) files := make([]protocol.SupportFile, 0)
err := filepath.Walk(m.SupportDir, func(path string, info os.FileInfo, err error) error { err := filepath.Walk(m.CertDir, func(path string, info os.FileInfo, err error) error {
if err != nil { if err != nil {
return err return err
} }
@@ -584,7 +637,7 @@ func (m *Manager) readSupportFiles() ([]protocol.SupportFile, error) {
if err != nil { if err != nil {
return err return err
} }
relativePath, err := filepath.Rel(m.SupportDir, path) relativePath, err := filepath.Rel(m.CertDir, path)
if err != nil { if err != nil {
return err return err
} }
@@ -603,9 +656,9 @@ func (m *Manager) readSupportFiles() ([]protocol.SupportFile, error) {
return files, nil return files, nil
} }
func (m *Manager) supportFileTargetPath(relativePath string) (string, error) { func (m *Manager) certFileTargetPath(relativePath string) (string, error) {
if strings.TrimSpace(m.SupportDir) == "" { if strings.TrimSpace(m.CertDir) == "" {
return "", errors.New("support dir 不能为空") return "", errors.New("cert dir 不能为空")
} }
candidate := strings.TrimSpace(relativePath) candidate := strings.TrimSpace(relativePath)
if strings.Contains(candidate, `\`) { if strings.Contains(candidate, `\`) {
@@ -613,25 +666,25 @@ func (m *Manager) supportFileTargetPath(relativePath string) (string, error) {
} }
normalizedPath := filepath.Clean(filepath.FromSlash(candidate)) normalizedPath := filepath.Clean(filepath.FromSlash(candidate))
if normalizedPath == "." || normalizedPath == "" { if normalizedPath == "." || normalizedPath == "" {
return "", errors.New("support file path 不能为空") return "", errors.New("cert file path 不能为空")
} }
if filepath.IsAbs(normalizedPath) || filepath.VolumeName(normalizedPath) != "" { if filepath.IsAbs(normalizedPath) || filepath.VolumeName(normalizedPath) != "" {
return "", fmt.Errorf("support file path %q must be relative", relativePath) return "", fmt.Errorf("cert file path %q must be relative", relativePath)
} }
targetPath := filepath.Join(m.SupportDir, normalizedPath) targetPath := filepath.Join(m.CertDir, normalizedPath)
relativeToBase, err := filepath.Rel(m.SupportDir, targetPath) relativeToBase, err := filepath.Rel(m.CertDir, targetPath)
if err != nil { if err != nil {
return "", err return "", err
} }
if relativeToBase == ".." || strings.HasPrefix(relativeToBase, ".."+string(os.PathSeparator)) { if relativeToBase == ".." || strings.HasPrefix(relativeToBase, ".."+string(os.PathSeparator)) {
return "", fmt.Errorf("support file path %q escapes support dir", relativePath) return "", fmt.Errorf("cert file path %q escapes cert dir", relativePath)
} }
return targetPath, nil return targetPath, nil
} }
func supportFileMode(relativePath string) fs.FileMode { func certFileMode(relativePath string) fs.FileMode {
switch strings.ToLower(filepath.Ext(strings.TrimSpace(relativePath))) { switch strings.ToLower(filepath.Ext(strings.TrimSpace(relativePath))) {
case ".lua", ".crt", ".pem": case ".crt", ".pem":
return 0o644 return 0o644
case ".key": case ".key":
return 0o600 return 0o600
@@ -640,11 +693,37 @@ func supportFileMode(relativePath string) fs.FileMode {
} }
} }
func luaFileTargetPath(baseDir string, relativePath string) (string, error) {
if strings.TrimSpace(baseDir) == "" {
return "", errors.New("lua dir 不能为空")
}
candidate := strings.TrimSpace(relativePath)
if strings.Contains(candidate, `\`) {
candidate = strings.ReplaceAll(candidate, `\`, "/")
}
normalizedPath := filepath.Clean(filepath.FromSlash(candidate))
if normalizedPath == "." || normalizedPath == "" {
return "", errors.New("lua file path 不能为空")
}
if filepath.IsAbs(normalizedPath) || filepath.VolumeName(normalizedPath) != "" {
return "", fmt.Errorf("lua file path %q must be relative", relativePath)
}
targetPath := filepath.Join(baseDir, normalizedPath)
relativeToBase, err := filepath.Rel(baseDir, targetPath)
if err != nil {
return "", err
}
if relativeToBase == ".." || strings.HasPrefix(relativeToBase, ".."+string(os.PathSeparator)) {
return "", fmt.Errorf("lua file path %q escapes lua dir", relativePath)
}
return targetPath, nil
}
func (m *Manager) renderRouteConfig(content string) string { func (m *Manager) renderRouteConfig(content string) string {
if m.NginxSupportDir == "" { if m.NginxCertDir == "" {
return content return content
} }
return strings.ReplaceAll(content, SupportDirPlaceholder, m.NginxSupportDir) return strings.ReplaceAll(content, CertDirPlaceholder, m.NginxCertDir)
} }
func (m *Manager) renderMainConfig(content string) string { func (m *Manager) renderMainConfig(content string) string {
@@ -693,10 +772,10 @@ func (m *Manager) accessLogRuntimePath() string {
} }
func (m *Manager) luaRuntimePath() string { func (m *Manager) luaRuntimePath() string {
if strings.TrimSpace(m.NginxSupportDir) == "" { if strings.TrimSpace(m.NginxLuaDir) == "" {
return "" return ""
} }
return filepath.ToSlash(m.NginxSupportDir) return filepath.ToSlash(m.NginxLuaDir)
} }
func checksum(content string) string { func checksum(content string) string {
+99 -87
View File
@@ -117,14 +117,16 @@ func TestDockerExecutorCheckHealthFailsWhenContainerStopped(t *testing.T) {
}, },
} }
executor := &DockerExecutor{ executor := &DockerExecutor{
DockerBinary: "docker", DockerBinary: "docker",
ContainerName: "atsflare-openresty", ContainerName: "atsflare-openresty",
Image: "openresty/openresty:alpine", Image: "openresty/openresty:alpine",
MainConfigPath: filepath.Clean("/tmp/nginx.conf"), MainConfigPath: filepath.Clean("/tmp/nginx.conf"),
RouteConfigDir: filepath.Clean("/tmp/routes"), RouteConfigDir: filepath.Clean("/tmp/routes"),
SupportDir: filepath.Clean("/tmp/support"), CertDir: filepath.Clean("/tmp/certs"),
NginxSupportDir: "/etc/nginx/atsflare-support", NginxCertDir: "/etc/nginx/atsflare-certs",
Runner: runner, LuaDir: filepath.Clean("/tmp/lua"),
NginxLuaDir: "/etc/nginx/atsflare-lua",
Runner: runner,
} }
if err := executor.CheckHealth(context.Background()); err == nil { if err := executor.CheckHealth(context.Background()); err == nil {
t.Fatal("expected CheckHealth to fail when container is not running") t.Fatal("expected CheckHealth to fail when container is not running")
@@ -141,14 +143,16 @@ func TestDockerExecutorStartsContainerWhenMissing(t *testing.T) {
}, },
} }
executor := &DockerExecutor{ executor := &DockerExecutor{
DockerBinary: "docker", DockerBinary: "docker",
ContainerName: "atsflare-openresty", ContainerName: "atsflare-openresty",
Image: "openresty/openresty:alpine", Image: "openresty/openresty:alpine",
MainConfigPath: filepath.Clean("/tmp/nginx.conf"), MainConfigPath: filepath.Clean("/tmp/nginx.conf"),
RouteConfigDir: filepath.Clean("/tmp/routes"), RouteConfigDir: filepath.Clean("/tmp/routes"),
SupportDir: filepath.Clean("/tmp/support"), CertDir: filepath.Clean("/tmp/certs"),
NginxSupportDir: "/etc/nginx/atsflare-support", NginxCertDir: "/etc/nginx/atsflare-certs",
Runner: runner, LuaDir: filepath.Clean("/tmp/lua"),
NginxLuaDir: "/etc/nginx/atsflare-lua",
Runner: runner,
} }
if err := executor.Test(context.Background()); err != nil { if err := executor.Test(context.Background()); err != nil {
@@ -176,14 +180,16 @@ func TestDockerExecutorStartsStoppedContainer(t *testing.T) {
}, },
} }
executor := &DockerExecutor{ executor := &DockerExecutor{
DockerBinary: "docker", DockerBinary: "docker",
ContainerName: "atsflare-openresty", ContainerName: "atsflare-openresty",
Image: "openresty/openresty:alpine", Image: "openresty/openresty:alpine",
MainConfigPath: filepath.Clean("/tmp/nginx.conf"), MainConfigPath: filepath.Clean("/tmp/nginx.conf"),
RouteConfigDir: filepath.Clean("/tmp/routes"), RouteConfigDir: filepath.Clean("/tmp/routes"),
SupportDir: filepath.Clean("/tmp/support"), CertDir: filepath.Clean("/tmp/certs"),
NginxSupportDir: "/etc/nginx/atsflare-support", NginxCertDir: "/etc/nginx/atsflare-certs",
Runner: runner, LuaDir: filepath.Clean("/tmp/lua"),
NginxLuaDir: "/etc/nginx/atsflare-lua",
Runner: runner,
} }
if err := executor.Reload(context.Background()); err != nil { if err := executor.Reload(context.Background()); err != nil {
@@ -207,7 +213,8 @@ func TestDockerExecutorStartsStoppedContainer(t *testing.T) {
func TestDockerExecutorRunContainerMountsManagedFiles(t *testing.T) { func TestDockerExecutorRunContainerMountsManagedFiles(t *testing.T) {
mainConfigPath := filepath.Clean("/tmp/managed/nginx.conf") mainConfigPath := filepath.Clean("/tmp/managed/nginx.conf")
routeConfigDir := filepath.Clean("/tmp/managed/conf.d") routeConfigDir := filepath.Clean("/tmp/managed/conf.d")
supportDir := filepath.Clean("/tmp/managed/support") certDir := filepath.Clean("/tmp/managed/certs")
luaDir := filepath.Clean("/tmp/managed/lua")
runner := &fakeRunner{} runner := &fakeRunner{}
executor := &DockerExecutor{ executor := &DockerExecutor{
DockerBinary: "docker", DockerBinary: "docker",
@@ -215,8 +222,10 @@ func TestDockerExecutorRunContainerMountsManagedFiles(t *testing.T) {
Image: "openresty/openresty:alpine", Image: "openresty/openresty:alpine",
MainConfigPath: mainConfigPath, MainConfigPath: mainConfigPath,
RouteConfigDir: routeConfigDir, RouteConfigDir: routeConfigDir,
SupportDir: supportDir, CertDir: certDir,
NginxSupportDir: "/etc/nginx/atsflare-support", NginxCertDir: "/etc/nginx/atsflare-certs",
LuaDir: luaDir,
NginxLuaDir: "/etc/nginx/atsflare-lua",
OpenrestyObservabilityPort: 18081, OpenrestyObservabilityPort: 18081,
Runner: runner, Runner: runner,
} }
@@ -237,7 +246,8 @@ func TestDockerExecutorRunContainerMountsManagedFiles(t *testing.T) {
"-p", "127.0.0.1:18081:18081", "-p", "127.0.0.1:18081:18081",
"-v", mainConfigPath + ":" + DockerMainConfigPath, "-v", mainConfigPath + ":" + DockerMainConfigPath,
"-v", routeConfigDir + ":/etc/nginx/conf.d", "-v", routeConfigDir + ":/etc/nginx/conf.d",
"-v", supportDir + ":/etc/nginx/atsflare-support", "-v", certDir + ":/etc/nginx/atsflare-certs",
"-v", luaDir + ":/etc/nginx/atsflare-lua",
"openresty/openresty:alpine", "openresty/openresty:alpine",
} }
if !reflect.DeepEqual(runner.calls[0].args, expectedArgs) { if !reflect.DeepEqual(runner.calls[0].args, expectedArgs) {
@@ -260,8 +270,10 @@ func TestDockerExecutorRecreatesContainerOnStartup(t *testing.T) {
Image: "openresty/openresty:alpine", Image: "openresty/openresty:alpine",
MainConfigPath: filepath.Clean("/tmp/nginx.conf"), MainConfigPath: filepath.Clean("/tmp/nginx.conf"),
RouteConfigDir: filepath.Clean("/tmp/routes"), RouteConfigDir: filepath.Clean("/tmp/routes"),
SupportDir: filepath.Clean("/tmp/support"), CertDir: filepath.Clean("/tmp/certs"),
NginxSupportDir: "/etc/nginx/atsflare-support", NginxCertDir: "/etc/nginx/atsflare-certs",
LuaDir: filepath.Clean("/tmp/lua"),
NginxLuaDir: "/etc/nginx/atsflare-lua",
OpenrestyObservabilityPort: 18081, OpenrestyObservabilityPort: 18081,
Runner: runner, Runner: runner,
} }
@@ -287,8 +299,10 @@ func TestNewExecutorUsesAbsoluteDockerMountPath(t *testing.T) {
Image: "openresty/openresty:alpine", Image: "openresty/openresty:alpine",
MainConfigPath: "./data/etc/nginx/nginx.conf", MainConfigPath: "./data/etc/nginx/nginx.conf",
RouteConfigPath: "./data/etc/nginx/conf.d/atsflare_routes.conf", RouteConfigPath: "./data/etc/nginx/conf.d/atsflare_routes.conf",
SupportDir: "./data/etc/nginx/support", CertDir: "./data/etc/nginx/certs",
NginxSupportDir: "/etc/nginx/atsflare-support", NginxCertDir: "/etc/nginx/atsflare-certs",
LuaDir: "./data/etc/nginx/lua",
NginxLuaDir: "/etc/nginx/atsflare-lua",
OpenrestyObservabilityPort: 18081, OpenrestyObservabilityPort: 18081,
}) })
@@ -330,19 +344,21 @@ func TestManagerApplyAndChecksumIncludeMainConfig(t *testing.T) {
tempDir := t.TempDir() tempDir := t.TempDir()
mainPath := filepath.Join(tempDir, "nginx.conf") mainPath := filepath.Join(tempDir, "nginx.conf")
routePath := filepath.Join(tempDir, "conf.d", "atsflare_routes.conf") routePath := filepath.Join(tempDir, "conf.d", "atsflare_routes.conf")
supportDir := filepath.Join(tempDir, "support") certDir := filepath.Join(tempDir, "certs")
manager := &Manager{ manager := &Manager{
MainConfigPath: mainPath, MainConfigPath: mainPath,
RouteConfigPath: routePath, RouteConfigPath: routePath,
SupportDir: supportDir, CertDir: certDir,
NginxSupportDir: "/etc/nginx/atsflare-support", NginxCertDir: "/etc/nginx/atsflare-certs",
LuaDir: filepath.Join(tempDir, "lua"),
NginxLuaDir: "/etc/nginx/atsflare-lua",
Executor: &fakeExecutor{}, Executor: &fakeExecutor{},
} }
err := manager.Apply( err := manager.Apply(
context.Background(), context.Background(),
"include __ATSF_ROUTE_CONFIG__;\naccess_log __ATSF_ACCESS_LOG__ atsflare_json;\n", "include __ATSF_ROUTE_CONFIG__;\naccess_log __ATSF_ACCESS_LOG__ atsflare_json;\n",
"ssl_certificate __ATSF_SUPPORT_DIR__/1.crt;\n", "ssl_certificate __ATSF_CERT_DIR__/1.crt;\n",
[]protocol.SupportFile{{Path: "1.crt", Content: "cert"}}, []protocol.SupportFile{{Path: "1.crt", Content: "cert"}},
) )
if err != nil { if err != nil {
@@ -362,7 +378,7 @@ func TestManagerApplyAndChecksumIncludeMainConfig(t *testing.T) {
if err != nil { if err != nil {
t.Fatalf("failed to read route config: %v", err) t.Fatalf("failed to read route config: %v", err)
} }
if string(routeData) != "ssl_certificate /etc/nginx/atsflare-support/1.crt;\n" { if string(routeData) != "ssl_certificate /etc/nginx/atsflare-certs/1.crt;\n" {
t.Fatalf("unexpected route config: %s", string(routeData)) t.Fatalf("unexpected route config: %s", string(routeData))
} }
@@ -372,7 +388,7 @@ func TestManagerApplyAndChecksumIncludeMainConfig(t *testing.T) {
} }
expected := bundleChecksum( expected := bundleChecksum(
"include __ATSF_ROUTE_CONFIG__;\naccess_log __ATSF_ACCESS_LOG__ atsflare_json;\n", "include __ATSF_ROUTE_CONFIG__;\naccess_log __ATSF_ACCESS_LOG__ atsflare_json;\n",
"ssl_certificate __ATSF_SUPPORT_DIR__/1.crt;\n", "ssl_certificate __ATSF_CERT_DIR__/1.crt;\n",
[]protocol.SupportFile{{Path: "1.crt", Content: "cert"}}, []protocol.SupportFile{{Path: "1.crt", Content: "cert"}},
) )
if value != expected { if value != expected {
@@ -388,8 +404,10 @@ func TestManagerApplyUsesRuntimeRouteConfigPath(t *testing.T) {
MainConfigPath: mainPath, MainConfigPath: mainPath,
RouteConfigPath: routePath, RouteConfigPath: routePath,
RuntimeRouteConfigPath: DockerRouteConfigPath, RuntimeRouteConfigPath: DockerRouteConfigPath,
SupportDir: filepath.Join(tempDir, "support"), CertDir: filepath.Join(tempDir, "certs"),
NginxSupportDir: "/etc/nginx/atsflare-support", NginxCertDir: "/etc/nginx/atsflare-certs",
LuaDir: filepath.Join(tempDir, "lua"),
NginxLuaDir: "/etc/nginx/atsflare-lua",
Executor: &fakeExecutor{}, Executor: &fakeExecutor{},
} }
@@ -462,13 +480,15 @@ func TestManagerApplyWritesSupportFilesAndReplacesPlaceholder(t *testing.T) {
manager := &Manager{ manager := &Manager{
MainConfigPath: filepath.Join(tempDir, "nginx.conf"), MainConfigPath: filepath.Join(tempDir, "nginx.conf"),
RouteConfigPath: filepath.Join(tempDir, "routes.conf"), RouteConfigPath: filepath.Join(tempDir, "routes.conf"),
SupportDir: filepath.Join(tempDir, "support"), CertDir: filepath.Join(tempDir, "certs"),
NginxSupportDir: "/etc/nginx/atsflare-support", NginxCertDir: "/etc/nginx/atsflare-certs",
LuaDir: filepath.Join(tempDir, "lua"),
NginxLuaDir: "/etc/nginx/atsflare-lua",
OpenrestyObservabilityListen: "18081", OpenrestyObservabilityListen: "18081",
Executor: &fakeExecutor{}, Executor: &fakeExecutor{},
} }
err := manager.Apply(context.Background(), "include __ATSF_ROUTE_CONFIG__;\nserver { listen __ATSF_OBSERVABILITY_LISTEN__; }", "ssl_certificate __ATSF_SUPPORT_DIR__/1.crt;", []protocol.SupportFile{ err := manager.Apply(context.Background(), "include __ATSF_ROUTE_CONFIG__;\nserver { listen __ATSF_OBSERVABILITY_LISTEN__; }", "ssl_certificate __ATSF_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"},
}) })
@@ -480,7 +500,7 @@ func TestManagerApplyWritesSupportFilesAndReplacesPlaceholder(t *testing.T) {
if err != nil { if err != nil {
t.Fatalf("failed to read route config: %v", err) t.Fatalf("failed to read route config: %v", err)
} }
if !strings.Contains(string(routeData), "/etc/nginx/atsflare-support/1.crt") { if !strings.Contains(string(routeData), "/etc/nginx/atsflare-certs/1.crt") {
t.Fatalf("expected placeholder replacement in route config, got %s", string(routeData)) t.Fatalf("expected placeholder replacement in route config, got %s", string(routeData))
} }
mainData, err := os.ReadFile(manager.MainConfigPath) mainData, err := os.ReadFile(manager.MainConfigPath)
@@ -490,25 +510,27 @@ func TestManagerApplyWritesSupportFilesAndReplacesPlaceholder(t *testing.T) {
if !strings.Contains(string(mainData), "listen 18081;") { if !strings.Contains(string(mainData), "listen 18081;") {
t.Fatalf("expected observability listen placeholder replacement in main config, got %s", string(mainData)) t.Fatalf("expected observability listen placeholder replacement in main config, got %s", string(mainData))
} }
certData, err := os.ReadFile(filepath.Join(manager.SupportDir, "1.crt")) certData, err := os.ReadFile(filepath.Join(manager.CertDir, "1.crt"))
if err != nil { if err != nil {
t.Fatalf("failed to read cert file: %v", err) t.Fatalf("failed to read cert file: %v", err)
} }
if string(certData) != "cert-data" { if string(certData) != "cert-data" {
t.Fatalf("unexpected cert file content: %s", string(certData)) t.Fatalf("unexpected cert file content: %s", string(certData))
} }
luaInfo, err := os.Stat(filepath.Join(manager.SupportDir, "observability", "log.lua")) luaInfo, err := os.Stat(filepath.Join(manager.LuaDir, "log.lua"))
if err == nil { if err != nil {
t.Fatalf("expected no lua file in this test, got %v", luaInfo) t.Fatalf("expected managed lua file to exist, stat err = %v", err)
}
if luaInfo.Mode().Perm() != 0o644 {
t.Fatalf("unexpected lua mode: %o", luaInfo.Mode().Perm())
} }
} }
func TestSupportFileMode(t *testing.T) { func TestCertFileMode(t *testing.T) {
testCases := []struct { testCases := []struct {
path string path string
want os.FileMode want os.FileMode
}{ }{
{path: "observability/log.lua", want: 0o644},
{path: "1.crt", want: 0o644}, {path: "1.crt", want: 0o644},
{path: "1.pem", want: 0o644}, {path: "1.pem", want: 0o644},
{path: "1.key", want: 0o600}, {path: "1.key", want: 0o600},
@@ -516,53 +538,39 @@ func TestSupportFileMode(t *testing.T) {
} }
for _, testCase := range testCases { for _, testCase := range testCases {
if got := supportFileMode(testCase.path); got != testCase.want { if got := certFileMode(testCase.path); got != testCase.want {
t.Fatalf("unexpected mode for %s: got %o want %o", testCase.path, got, testCase.want) t.Fatalf("unexpected mode for %s: got %o want %o", testCase.path, got, testCase.want)
} }
} }
} }
func TestManagerApplyWritesLuaSupportFilesReadable(t *testing.T) { func TestManagerEnsureLuaAssetsWritesReadableFiles(t *testing.T) {
tempDir := t.TempDir() tempDir := t.TempDir()
manager := &Manager{ manager := &Manager{
MainConfigPath: filepath.Join(tempDir, "nginx.conf"), LuaDir: filepath.Join(tempDir, "lua"),
RouteConfigPath: filepath.Join(tempDir, "routes.conf"), NginxLuaDir: "/etc/nginx/atsflare-lua",
SupportDir: filepath.Join(tempDir, "support"),
NginxSupportDir: "/etc/nginx/atsflare-support",
Executor: &fakeExecutor{},
} }
err := manager.Apply(context.Background(), "main", "route", []protocol.SupportFile{ err := manager.EnsureLuaAssets()
{Path: "observability/log.lua", Content: "return"},
{Path: "1.key", Content: "secret"},
})
if err != nil { if err != nil {
t.Fatalf("Apply failed: %v", err) t.Fatalf("EnsureLuaAssets failed: %v", err)
} }
luaInfo, err := os.Stat(filepath.Join(manager.SupportDir, "observability", "log.lua")) luaInfo, err := os.Stat(filepath.Join(manager.LuaDir, "log.lua"))
if err != nil { if err != nil {
t.Fatalf("failed to stat lua file: %v", err) t.Fatalf("failed to stat lua file: %v", err)
} }
if luaInfo.Mode().Perm() != 0o644 { if luaInfo.Mode().Perm() != 0o644 {
t.Fatalf("unexpected lua mode: %o", luaInfo.Mode().Perm()) t.Fatalf("unexpected lua mode: %o", luaInfo.Mode().Perm())
} }
keyInfo, err := os.Stat(filepath.Join(manager.SupportDir, "1.key"))
if err != nil {
t.Fatalf("failed to stat key file: %v", err)
}
if keyInfo.Mode().Perm() != 0o600 {
t.Fatalf("unexpected key mode: %o", keyInfo.Mode().Perm())
}
} }
func TestManagerRollbackRestoresSupportFiles(t *testing.T) { func TestManagerRollbackRestoresCertFiles(t *testing.T) {
tempDir := t.TempDir() tempDir := t.TempDir()
routePath := filepath.Join(tempDir, "routes.conf") routePath := filepath.Join(tempDir, "routes.conf")
mainPath := filepath.Join(tempDir, "nginx.conf") mainPath := filepath.Join(tempDir, "nginx.conf")
supportDir := filepath.Join(tempDir, "support") certDir := filepath.Join(tempDir, "certs")
if err := os.MkdirAll(supportDir, 0o755); err != nil { if err := os.MkdirAll(certDir, 0o755); err != nil {
t.Fatalf("MkdirAll failed: %v", err) t.Fatalf("MkdirAll failed: %v", err)
} }
if err := os.WriteFile(mainPath, []byte("old-main"), 0o644); err != nil { if err := os.WriteFile(mainPath, []byte("old-main"), 0o644); err != nil {
@@ -571,14 +579,16 @@ func TestManagerRollbackRestoresSupportFiles(t *testing.T) {
if err := os.WriteFile(routePath, []byte("old-route"), 0o644); err != nil { if err := os.WriteFile(routePath, []byte("old-route"), 0o644); err != nil {
t.Fatalf("WriteFile failed: %v", err) t.Fatalf("WriteFile failed: %v", err)
} }
if err := os.WriteFile(filepath.Join(supportDir, "1.crt"), []byte("old-cert"), 0o600); err != nil { if err := os.WriteFile(filepath.Join(certDir, "1.crt"), []byte("old-cert"), 0o600); err != nil {
t.Fatalf("WriteFile failed: %v", err) t.Fatalf("WriteFile failed: %v", err)
} }
manager := &Manager{ manager := &Manager{
MainConfigPath: mainPath, MainConfigPath: mainPath,
RouteConfigPath: routePath, RouteConfigPath: routePath,
SupportDir: supportDir, CertDir: certDir,
NginxSupportDir: "/etc/nginx/atsflare-support", NginxCertDir: "/etc/nginx/atsflare-certs",
LuaDir: filepath.Join(tempDir, "lua"),
NginxLuaDir: "/etc/nginx/atsflare-lua",
Executor: &fakeExecutor{ Executor: &fakeExecutor{
testErr: errors.New("openresty test failed"), testErr: errors.New("openresty test failed"),
}, },
@@ -605,7 +615,7 @@ func TestManagerRollbackRestoresSupportFiles(t *testing.T) {
if string(routeData) != "old-route" { if string(routeData) != "old-route" {
t.Fatalf("expected route rollback, got %s", string(routeData)) t.Fatalf("expected route rollback, got %s", string(routeData))
} }
certData, err := os.ReadFile(filepath.Join(supportDir, "1.crt")) certData, err := os.ReadFile(filepath.Join(certDir, "1.crt"))
if err != nil { if err != nil {
t.Fatalf("failed to read cert file: %v", err) t.Fatalf("failed to read cert file: %v", err)
} }
@@ -614,9 +624,9 @@ func TestManagerRollbackRestoresSupportFiles(t *testing.T) {
} }
} }
func TestManagerSupportFileTargetPathRejectsEscapes(t *testing.T) { func TestManagerCertFileTargetPathRejectsEscapes(t *testing.T) {
manager := &Manager{SupportDir: filepath.Join(t.TempDir(), "support")} manager := &Manager{CertDir: filepath.Join(t.TempDir(), "certs")}
if err := os.MkdirAll(manager.SupportDir, 0o755); err != nil { if err := os.MkdirAll(manager.CertDir, 0o755); err != nil {
t.Fatalf("MkdirAll failed: %v", err) t.Fatalf("MkdirAll failed: %v", err)
} }
@@ -637,7 +647,7 @@ func TestManagerSupportFileTargetPathRejectsEscapes(t *testing.T) {
} }
for _, testCase := range testCases { for _, testCase := range testCases {
targetPath, err := manager.supportFileTargetPath(testCase.path) targetPath, err := manager.certFileTargetPath(testCase.path)
if testCase.shouldErr { if testCase.shouldErr {
if err == nil { if err == nil {
t.Fatalf("expected path %q to be rejected, got target %q", testCase.path, targetPath) t.Fatalf("expected path %q to be rejected, got target %q", testCase.path, targetPath)
@@ -647,19 +657,21 @@ func TestManagerSupportFileTargetPathRejectsEscapes(t *testing.T) {
if err != nil { if err != nil {
t.Fatalf("expected path %q to be accepted: %v", testCase.path, err) t.Fatalf("expected path %q to be accepted: %v", testCase.path, err)
} }
if !strings.HasPrefix(targetPath, manager.SupportDir) { if !strings.HasPrefix(targetPath, manager.CertDir) {
t.Fatalf("expected target path %q to stay under %q", targetPath, manager.SupportDir) t.Fatalf("expected target path %q to stay under %q", targetPath, manager.CertDir)
} }
} }
} }
func TestManagerApplyRejectsSupportFilePathTraversal(t *testing.T) { func TestManagerApplyRejectsCertFilePathTraversal(t *testing.T) {
tempDir := t.TempDir() tempDir := t.TempDir()
manager := &Manager{ manager := &Manager{
MainConfigPath: filepath.Join(tempDir, "nginx.conf"), MainConfigPath: filepath.Join(tempDir, "nginx.conf"),
RouteConfigPath: filepath.Join(tempDir, "routes.conf"), RouteConfigPath: filepath.Join(tempDir, "routes.conf"),
SupportDir: filepath.Join(tempDir, "support"), CertDir: filepath.Join(tempDir, "certs"),
NginxSupportDir: "/etc/nginx/atsflare-support", NginxCertDir: "/etc/nginx/atsflare-certs",
LuaDir: filepath.Join(tempDir, "lua"),
NginxLuaDir: "/etc/nginx/atsflare-lua",
Executor: &fakeExecutor{}, Executor: &fakeExecutor{},
} }
@@ -0,0 +1,158 @@
package nginx
import "atsflare-agent/internal/protocol"
const (
openRestyObservabilityWindowTTL = 7200
openRestyObservabilityWindowSize = 60
)
const openRestyObservabilityInitLua = `local dict = ngx.shared.atsflare_observability
if not dict then
return
end
return
`
const openRestyObservabilityLogLua = `local dict = ngx.shared.atsflare_observability
if not dict then
return
end
local request_uri = tostring(ngx.var.uri or "")
if request_uri == "/atsflare/observability" or request_uri == "/atsflare/stub_status" then
return
end
local ttl = ` + "7200" + `
local now = ngx.time()
local window_size = ` + "60" + `
local window_start = now - (now % window_size)
local function ensure_counter(key)
dict:add(key, 0, ttl)
end
local function incr(key, delta)
ensure_counter(key)
local value, err = dict:incr(key, delta)
if not value and err == "not found" then
dict:set(key, delta, ttl)
end
end
local function remember_value(list_key, marker_key, value)
if value == "" then
return
end
if not dict:add(marker_key, 1, ttl) then
return
end
local existing = dict:get(list_key)
if not existing or existing == "" then
dict:set(list_key, value, ttl)
return
end
dict:set(list_key, existing .. "\n" .. value, ttl)
end
local window_prefix = tostring(window_start)
incr("request_count:" .. window_prefix, 1)
local status = tostring(ngx.status or 0)
if status ~= "0" then
incr("status:" .. window_prefix .. ":" .. status, 1)
remember_value(
"status_keys:" .. window_prefix,
"status_marker:" .. window_prefix .. ":" .. status,
status
)
if tonumber(status) and tonumber(status) >= 500 then
incr("error_count:" .. window_prefix, 1)
end
end
local host = tostring(ngx.var.host or "")
if host ~= "" then
incr("domain:" .. window_prefix .. ":" .. host, 1)
remember_value(
"domain_keys:" .. window_prefix,
"domain_marker:" .. window_prefix .. ":" .. host,
host
)
end
local remote_addr = tostring(ngx.var.binary_remote_addr or ngx.var.remote_addr or "")
if remote_addr ~= "" and dict:add("visitor:" .. window_prefix .. ":" .. remote_addr, 1, ttl) then
incr("unique_visitor_count:" .. window_prefix, 1)
end
local request_length = tonumber(ngx.var.request_length) or 0
if request_length > 0 then
incr("openresty_rx_bytes:" .. window_prefix, request_length)
end
local bytes_sent = tonumber(ngx.var.bytes_sent) or tonumber(ngx.var.body_bytes_sent) or 0
if bytes_sent > 0 then
incr("openresty_tx_bytes:" .. window_prefix, bytes_sent)
end
`
const openRestyObservabilityReadLua = `local cjson = require "cjson.safe"
local dict = ngx.shared.atsflare_observability
if not dict then
ngx.status = ngx.HTTP_SERVICE_UNAVAILABLE
ngx.say(cjson.encode({ message = "shared dict unavailable" }))
return
end
local now = ngx.time()
local window_size = ` + "60" + `
local window_start = now - (now % window_size)
local current_window = tostring(window_start)
local function read_counter(key)
return tonumber(dict:get(key) or 0) or 0
end
local function read_map(window_id, prefix, list_key)
local result = {}
local raw = dict:get(list_key .. ":" .. window_id)
if not raw or raw == "" then
return result
end
for value in string.gmatch(raw, "[^\n]+") do
result[value] = read_counter(prefix .. ":" .. window_id .. ":" .. value)
end
return result
end
local payload = {
window_started_at_unix = window_start,
window_ended_at_unix = now,
request_count = read_counter("request_count:" .. current_window),
error_count = read_counter("error_count:" .. current_window),
unique_visitor_count = read_counter("unique_visitor_count:" .. current_window),
status_codes = read_map(current_window, "status", "status_keys"),
top_domains = read_map(current_window, "domain", "domain_keys"),
source_countries = {},
openresty_rx_bytes = read_counter("openresty_rx_bytes:" .. current_window),
openresty_tx_bytes = read_counter("openresty_tx_bytes:" .. current_window)
}
ngx.header.content_type = "application/json"
ngx.say(cjson.encode(payload))
`
func ManagedObservabilityLuaFiles() []protocol.SupportFile {
return []protocol.SupportFile{
{Path: "init.lua", Content: openRestyObservabilityInitLua},
{Path: "log.lua", Content: openRestyObservabilityLogLua},
{Path: "read.lua", Content: openRestyObservabilityReadLua},
{Path: "observability/init.lua", Content: openRestyObservabilityInitLua},
{Path: "observability/log.lua", Content: openRestyObservabilityLogLua},
{Path: "observability/read.lua", Content: openRestyObservabilityReadLua},
}
}
+16
View File
@@ -191,6 +191,7 @@ func GetActiveConfigForAgent() (*AgentConfigResponse, error) {
return nil, err return nil, err
} }
} }
supportFiles = filterCertificateSupportFiles(supportFiles)
slog.Debug("agent fetched active config", "version", version.Version, "checksum", version.Checksum) slog.Debug("agent fetched active config", "version", version.Version, "checksum", version.Checksum)
return &AgentConfigResponse{ return &AgentConfigResponse{
Version: version.Version, Version: version.Version,
@@ -203,6 +204,21 @@ func GetActiveConfigForAgent() (*AgentConfigResponse, error) {
}, nil }, nil
} }
func filterCertificateSupportFiles(files []SupportFile) []SupportFile {
if len(files) == 0 {
return nil
}
filtered := make([]SupportFile, 0, len(files))
for _, file := range files {
path := strings.ToLower(strings.TrimSpace(file.Path))
switch {
case strings.HasSuffix(path, ".crt"), strings.HasSuffix(path, ".key"), strings.HasSuffix(path, ".pem"):
filtered = append(filtered, file)
}
}
return filtered
}
func ReportApplyLog(payload ApplyLogPayload) (*model.ApplyLog, error) { func ReportApplyLog(payload ApplyLogPayload) (*model.ApplyLog, error) {
now := time.Now() now := time.Now()
payload.NodeID = strings.TrimSpace(payload.NodeID) payload.NodeID = strings.TrimSpace(payload.NodeID)
+3 -4
View File
@@ -115,7 +115,7 @@ type configBundle struct {
} }
const ( const (
nginxSupportDirPlaceholder = "__ATSF_SUPPORT_DIR__" nginxCertDirPlaceholder = "__ATSF_CERT_DIR__"
nginxRouteConfigPlaceholder = "__ATSF_ROUTE_CONFIG__" nginxRouteConfigPlaceholder = "__ATSF_ROUTE_CONFIG__"
nginxAccessLogPlaceholder = "__ATSF_ACCESS_LOG__" nginxAccessLogPlaceholder = "__ATSF_ACCESS_LOG__"
nginxLuaDirPlaceholder = "__ATSF_LUA_DIR__" nginxLuaDirPlaceholder = "__ATSF_LUA_DIR__"
@@ -354,7 +354,6 @@ func buildCurrentConfigBundle(requireRoutes bool) (*configBundle, error) {
if err != nil { if err != nil {
return nil, err return nil, err
} }
supportFiles = append(supportFiles, buildOpenRestyObservabilitySupportFiles()...)
mainConfig := renderMainConfig(openRestyConfig) mainConfig := renderMainConfig(openRestyConfig)
return &configBundle{ return &configBundle{
Routes: routes, Routes: routes,
@@ -752,8 +751,8 @@ func renderHTTPRedirectServer(domain string) string {
} }
func renderHTTPSServer(domain string, originURL string, certificateID uint, customHeaders []ProxyRouteCustomHeaderInput) string { func renderHTTPSServer(domain string, originURL string, certificateID uint, customHeaders []ProxyRouteCustomHeaderInput) string {
certPath := fmt.Sprintf("%s/%s", nginxSupportDirPlaceholder, certificateCertFileName(certificateID)) certPath := fmt.Sprintf("%s/%s", nginxCertDirPlaceholder, certificateCertFileName(certificateID))
keyPath := fmt.Sprintf("%s/%s", nginxSupportDirPlaceholder, certificateKeyFileName(certificateID)) keyPath := fmt.Sprintf("%s/%s", nginxCertDirPlaceholder, certificateKeyFileName(certificateID))
return fmt.Sprintf("server {\n listen 443 ssl;\n server_name %s;\n ssl_certificate %s;\n ssl_certificate_key %s;\n\n location / {\n%s proxy_pass %s;\n }\n}\n\n", domain, certPath, keyPath, renderProxyHeaderBlock(customHeaders), originURL) return fmt.Sprintf("server {\n listen 443 ssl;\n server_name %s;\n ssl_certificate %s;\n ssl_certificate_key %s;\n\n location / {\n%s proxy_pass %s;\n }\n}\n\n", domain, certPath, keyPath, renderProxyHeaderBlock(customHeaders), originURL)
} }
+4 -13
View File
@@ -57,7 +57,7 @@ func TestCreateTLSCertificateAndRenderHTTPSConfig(t *testing.T) {
if !strings.Contains(result.Version.MainConfig, "access_log __ATSF_ACCESS_LOG__ atsflare_json;") { if !strings.Contains(result.Version.MainConfig, "access_log __ATSF_ACCESS_LOG__ atsflare_json;") {
t.Fatal("expected main config to include managed access log placeholder") t.Fatal("expected main config to include managed access log placeholder")
} }
if !strings.Contains(result.Version.MainConfig, "log_by_lua_file __ATSF_LUA_DIR__/observability/log.lua;") { if !strings.Contains(result.Version.MainConfig, "log_by_lua_file __ATSF_LUA_DIR__/log.lua;") {
t.Fatal("expected main config to include managed openresty lua log hook") t.Fatal("expected main config to include managed openresty lua log hook")
} }
if !strings.Contains(result.Version.MainConfig, "listen __ATSF_OBSERVABILITY_LISTEN__;") { if !strings.Contains(result.Version.MainConfig, "listen __ATSF_OBSERVABILITY_LISTEN__;") {
@@ -72,21 +72,12 @@ func TestCreateTLSCertificateAndRenderHTTPSConfig(t *testing.T) {
if !strings.Contains(result.Version.RenderedConfig, "return 301 https://$host$request_uri;") { if !strings.Contains(result.Version.RenderedConfig, "return 301 https://$host$request_uri;") {
t.Fatal("expected rendered config to include http redirect") t.Fatal("expected rendered config to include http redirect")
} }
if !strings.Contains(result.Version.RenderedConfig, "__ATSF_SUPPORT_DIR__/") { if !strings.Contains(result.Version.RenderedConfig, "__ATSF_CERT_DIR__/") {
t.Fatal("expected rendered config to keep support dir placeholder for certificates") t.Fatal("expected rendered config to keep cert dir placeholder for certificates")
} }
if !strings.Contains(result.Version.SupportFilesJSON, ".crt") || !strings.Contains(result.Version.SupportFilesJSON, ".key") { if !strings.Contains(result.Version.SupportFilesJSON, ".crt") || !strings.Contains(result.Version.SupportFilesJSON, ".key") {
t.Fatal("expected support files to contain certificate and key") t.Fatal("expected support files to contain certificate and key")
} }
if !strings.Contains(result.Version.SupportFilesJSON, "observability/log.lua") || !strings.Contains(result.Version.SupportFilesJSON, "observability/read.lua") {
t.Fatal("expected support files to contain managed openresty observability lua scripts")
}
if !strings.Contains(result.Version.SupportFilesJSON, "/atsflare/observability") || !strings.Contains(result.Version.SupportFilesJSON, "/atsflare/stub_status") {
t.Fatal("expected observability log lua to skip self-observability requests")
}
if !strings.Contains(result.Version.SupportFilesJSON, "window_start = now - (now % window_size)") || !strings.Contains(result.Version.SupportFilesJSON, "local window_size = 60") {
t.Fatal("expected support files to use fixed 60-second observability windows")
}
} }
func TestCreateProxyRouteRejectsHTTPSWithoutCertificate(t *testing.T) { func TestCreateProxyRouteRejectsHTTPSWithoutCertificate(t *testing.T) {
@@ -205,7 +196,7 @@ func TestPreviewAndDiffConfigVersion(t *testing.T) {
if !strings.Contains(preview.MainConfig, "include __ATSF_ROUTE_CONFIG__;") { if !strings.Contains(preview.MainConfig, "include __ATSF_ROUTE_CONFIG__;") {
t.Fatal("expected preview main config to include managed route config placeholder") t.Fatal("expected preview main config to include managed route config placeholder")
} }
if !strings.Contains(preview.MainConfig, "log_by_lua_file __ATSF_LUA_DIR__/observability/log.lua;") { if !strings.Contains(preview.MainConfig, "log_by_lua_file __ATSF_LUA_DIR__/log.lua;") {
t.Fatal("expected preview main config to include managed openresty lua log hook") t.Fatal("expected preview main config to include managed openresty lua log hook")
} }
if !strings.Contains(preview.RenderedConfig, `proxy_set_header X-Release "candidate";`) { if !strings.Contains(preview.RenderedConfig, `proxy_set_header X-Release "candidate";`) {
@@ -3,161 +3,11 @@ package service
import "fmt" import "fmt"
const ( const (
openRestyObservabilitySupportDir = "observability" openRestyObservabilityInitLuaPath = "init.lua"
openRestyObservabilityInitLuaPath = openRestyObservabilitySupportDir + "/init.lua" openRestyObservabilityLogLuaPath = "log.lua"
openRestyObservabilityLogLuaPath = openRestyObservabilitySupportDir + "/log.lua" openRestyObservabilityReadLuaPath = "read.lua"
openRestyObservabilityReadLuaPath = openRestyObservabilitySupportDir + "/read.lua"
openRestyObservabilityWindowTTL = 7200
openRestyObservabilityWindowSize = 60
) )
const openRestyObservabilityInitLua = `local dict = ngx.shared.atsflare_observability
if not dict then
return
end
return
`
const openRestyObservabilityLogLua = `local dict = ngx.shared.atsflare_observability
if not dict then
return
end
local request_uri = tostring(ngx.var.uri or "")
if request_uri == "/atsflare/observability" or request_uri == "/atsflare/stub_status" then
return
end
local ttl = ` + "7200" + `
local now = ngx.time()
local window_size = ` + "60" + `
local window_start = now - (now % window_size)
local function ensure_counter(key)
dict:add(key, 0, ttl)
end
local function incr(key, delta)
ensure_counter(key)
local value, err = dict:incr(key, delta)
if not value and err == "not found" then
dict:set(key, delta, ttl)
end
end
local function remember_value(list_key, marker_key, value)
if value == "" then
return
end
if not dict:add(marker_key, 1, ttl) then
return
end
local existing = dict:get(list_key)
if not existing or existing == "" then
dict:set(list_key, value, ttl)
return
end
dict:set(list_key, existing .. "\n" .. value, ttl)
end
local window_prefix = tostring(window_start)
incr("request_count:" .. window_prefix, 1)
local status = tostring(ngx.status or 0)
if status ~= "0" then
incr("status:" .. window_prefix .. ":" .. status, 1)
remember_value(
"status_keys:" .. window_prefix,
"status_marker:" .. window_prefix .. ":" .. status,
status
)
if tonumber(status) and tonumber(status) >= 500 then
incr("error_count:" .. window_prefix, 1)
end
end
local host = tostring(ngx.var.host or "")
if host ~= "" then
incr("domain:" .. window_prefix .. ":" .. host, 1)
remember_value(
"domain_keys:" .. window_prefix,
"domain_marker:" .. window_prefix .. ":" .. host,
host
)
end
local remote_addr = tostring(ngx.var.binary_remote_addr or ngx.var.remote_addr or "")
if remote_addr ~= "" and dict:add("visitor:" .. window_prefix .. ":" .. remote_addr, 1, ttl) then
incr("unique_visitor_count:" .. window_prefix, 1)
end
local request_length = tonumber(ngx.var.request_length) or 0
if request_length > 0 then
incr("openresty_rx_bytes:" .. window_prefix, request_length)
end
local bytes_sent = tonumber(ngx.var.bytes_sent) or tonumber(ngx.var.body_bytes_sent) or 0
if bytes_sent > 0 then
incr("openresty_tx_bytes:" .. window_prefix, bytes_sent)
end
`
const openRestyObservabilityReadLua = `local cjson = require "cjson.safe"
local dict = ngx.shared.atsflare_observability
if not dict then
ngx.status = ngx.HTTP_SERVICE_UNAVAILABLE
ngx.say(cjson.encode({ message = "shared dict unavailable" }))
return
end
local now = ngx.time()
local window_size = ` + "60" + `
local window_start = now - (now % window_size)
local current_window = tostring(window_start)
local function read_counter(key)
return tonumber(dict:get(key) or 0) or 0
end
local function read_map(window_id, prefix, list_key)
local result = {}
local raw = dict:get(list_key .. ":" .. window_id)
if not raw or raw == "" then
return result
end
for value in string.gmatch(raw, "[^\n]+") do
result[value] = read_counter(prefix .. ":" .. window_id .. ":" .. value)
end
return result
end
local payload = {
window_started_at_unix = window_start,
window_ended_at_unix = now,
request_count = read_counter("request_count:" .. current_window),
error_count = read_counter("error_count:" .. current_window),
unique_visitor_count = read_counter("unique_visitor_count:" .. current_window),
status_codes = read_map(current_window, "status", "status_keys"),
top_domains = read_map(current_window, "domain", "domain_keys"),
source_countries = {},
openresty_rx_bytes = read_counter("openresty_rx_bytes:" .. current_window),
openresty_tx_bytes = read_counter("openresty_tx_bytes:" .. current_window)
}
ngx.header.content_type = "application/json"
ngx.say(cjson.encode(payload))
`
func buildOpenRestyObservabilitySupportFiles() []SupportFile {
return []SupportFile{
{Path: openRestyObservabilityInitLuaPath, Content: openRestyObservabilityInitLua},
{Path: openRestyObservabilityLogLuaPath, Content: openRestyObservabilityLogLua},
{Path: openRestyObservabilityReadLuaPath, Content: openRestyObservabilityReadLua},
}
}
func renderOpenRestyObservabilityTemplateBlock() string { func renderOpenRestyObservabilityTemplateBlock() string {
return stringsJoinLines( return stringsJoinLines(
" lua_shared_dict atsflare_observability 10m;", " lua_shared_dict atsflare_observability 10m;",
+19 -11
View File
@@ -260,8 +260,10 @@ go run ./cmd/agent -config ./agent.json
"openresty_path": "/usr/local/openresty/nginx/sbin/openresty", "openresty_path": "/usr/local/openresty/nginx/sbin/openresty",
"main_config_path": "/usr/local/openresty/nginx/conf/nginx.conf", "main_config_path": "/usr/local/openresty/nginx/conf/nginx.conf",
"route_config_path": "/usr/local/openresty/nginx/conf/conf.d/atsflare_routes.conf", "route_config_path": "/usr/local/openresty/nginx/conf/conf.d/atsflare_routes.conf",
"support_dir": "/usr/local/openresty/nginx/conf/support", "cert_dir": "/usr/local/openresty/nginx/conf/certs",
"openresty_support_dir": "/usr/local/openresty/nginx/conf/support", "openresty_cert_dir": "/usr/local/openresty/nginx/conf/certs",
"lua_dir": "/usr/local/openresty/nginx/conf/lua",
"openresty_lua_dir": "/usr/local/openresty/nginx/conf/lua",
"openresty_observability_port": 18081, "openresty_observability_port": 18081,
"observability_buffer_path": "./data/observability-buffer.json", "observability_buffer_path": "./data/observability-buffer.json",
"observability_replay_minutes": 15, "observability_replay_minutes": 15,
@@ -288,8 +290,10 @@ go run ./cmd/agent -config ./agent.json
| `data_dir` | Agent 数据目录,用于存储托管配置、证书和状态文件 | 否 | 配置文件所在目录下的 `data` 子目录 | `./data` | | `data_dir` | Agent 数据目录,用于存储托管配置、证书和状态文件 | 否 | 配置文件所在目录下的 `data` 子目录 | `./data` |
| `main_config_path` | 第五版主配置接管时 OpenResty 主配置文件写入路径 | 第五版本机模式建议必填 | Docker 模式可使用受管默认路径;本机模式建议显式设置 | `/usr/local/openresty/nginx/conf/nginx.conf` | | `main_config_path` | 第五版主配置接管时 OpenResty 主配置文件写入路径 | 第五版本机模式建议必填 | Docker 模式可使用受管默认路径;本机模式建议显式设置 | `/usr/local/openresty/nginx/conf/nginx.conf` |
| `route_config_path` | 路由配置文件写入路径 | 否 | 默认为 `data_dir` 下托管路径 | `/etc/nginx/conf.d/atsflare_routes.conf` | | `route_config_path` | 路由配置文件写入路径 | 否 | 默认为 `data_dir` 下托管路径 | `/etc/nginx/conf.d/atsflare_routes.conf` |
| `support_dir` | Agent 在本机写入受管附属文件的目录,当前包含证书与 Lua 观测脚本 | 否 | 默认为 `data_dir` 下托管 support 目录 | `./data/etc/nginx/support` | | `cert_dir` | Agent 在本机写入受管证书文件的目录 | 否 | 默认为 `data_dir` 下托管 certs 目录 | `./data/etc/nginx/certs` |
| `openresty_support_dir` | OpenResty 实际读取受管附属文件的目录 | 否 | 本机模式默认等于 `support_dir`;Docker 模式默认 `/etc/nginx/atsflare-support` | `/usr/local/openresty/nginx/conf/support` | | `openresty_cert_dir` | OpenResty 实际读取受管证书文件的目录 | 否 | 本机模式默认等于 `cert_dir`;Docker 模式默认 `/etc/nginx/atsflare-certs` | `/usr/local/openresty/nginx/conf/certs` |
| `lua_dir` | Agent 在本机写入受管 Lua 观测脚本的目录;每次启动都会覆盖释放 | 否 | 默认为 `data_dir` 下托管 lua 目录 | `./data/etc/nginx/lua` |
| `openresty_lua_dir` | OpenResty 实际读取受管 Lua 观测脚本的目录 | 否 | 本机模式默认等于 `lua_dir`;Docker 模式默认 `/etc/nginx/atsflare-lua` | `/usr/local/openresty/nginx/conf/lua` |
| `observability_buffer_path` | Agent 本地观测补报缓冲文件路径;用于在 server 短暂离线时按时间窗口落盘待补传数据 | 否 | 默认为 `data_dir` 下托管观测缓冲文件 | `./data/var/lib/atsflare/observability-buffer.json` | | `observability_buffer_path` | Agent 本地观测补报缓冲文件路径;用于在 server 短暂离线时按时间窗口落盘待补传数据 | 否 | 默认为 `data_dir` 下托管观测缓冲文件 | `./data/var/lib/atsflare/observability-buffer.json` |
| `observability_replay_minutes` | Agent 恢复心跳后允许批量补传的最近观测窗口时长(分钟) | 否 | `15` | `30` | | `observability_replay_minutes` | Agent 恢复心跳后允许批量补传的最近观测窗口时长(分钟) | 否 | `15` | `30` |
| `state_path` | Agent 本地状态文件路径 | 否 | 默认为 `data_dir` 下托管状态文件 | `./data/agent-state.json` | | `state_path` | Agent 本地状态文件路径 | 否 | 默认为 `data_dir` 下托管状态文件 | `./data/agent-state.json` |
@@ -306,7 +310,6 @@ go run ./cmd/agent -config ./agent.json
* 未配置 `openresty_path` 时,默认为 Docker OpenResty 模式 * 未配置 `openresty_path` 时,默认为 Docker OpenResty 模式
* `openresty_observability_port` 默认仅绑定本地回环地址;若节点本机已有端口冲突,可改为其他未占用端口 * `openresty_observability_port` 默认仅绑定本地回环地址;若节点本机已有端口冲突,可改为其他未占用端口
* `observability_replay_minutes` 只控制“允许补传最近多少分钟的窗口”;超出该窗口的历史观测会在本地自动裁剪 * `observability_replay_minutes` 只控制“允许补传最近多少分钟的窗口”;超出该窗口的历史观测会在本地自动裁剪
* 为兼容旧节点,Agent 仍可读取历史字段 `cert_dir` / `openresty_cert_dir`,但保存配置时会统一写回 `support_dir` / `openresty_support_dir`
* 配置保存时,`agent_version`、`nginx_version` 由程序运行时维护,不需要写入 JSON * 配置保存时,`agent_version`、`nginx_version` 由程序运行时维护,不需要写入 JSON
* 第五版主配置接管完成后,本机模式下应优先通过 `main_config_path` 由 Agent 写入受管主配置,而不是依赖节点手工维护 include 规则 * 第五版主配置接管完成后,本机模式下应优先通过 `main_config_path` 由 Agent 写入受管主配置,而不是依赖节点手工维护 include 规则
@@ -318,7 +321,8 @@ go run ./cmd/agent -config ./agent.json
| --- | --- | | --- | --- |
| `main_config_path` | 第五版 Docker 模式默认可落在 `data_dir/etc/nginx/nginx.conf`;本机模式建议显式配置 | | `main_config_path` | 第五版 Docker 模式默认可落在 `data_dir/etc/nginx/nginx.conf`;本机模式建议显式配置 |
| `route_config_path` | `data_dir/etc/nginx/conf.d/atsflare_routes.conf` | | `route_config_path` | `data_dir/etc/nginx/conf.d/atsflare_routes.conf` |
| `support_dir` | `data_dir/etc/nginx/support` | | `cert_dir` | `data_dir/etc/nginx/certs` |
| `lua_dir` | `data_dir/etc/nginx/lua` |
| `observability_buffer_path` | `data_dir/var/lib/atsflare/observability-buffer.json` | | `observability_buffer_path` | `data_dir/var/lib/atsflare/observability-buffer.json` |
| `state_path` | `data_dir/var/lib/atsflare/agent-state.json` | | `state_path` | `data_dir/var/lib/atsflare/agent-state.json` |
@@ -326,11 +330,13 @@ Docker OpenResty 模式下:
| 字段 | 默认值 | | 字段 | 默认值 |
| --- | --- | | --- | --- |
| `openresty_support_dir` | `/etc/nginx/atsflare-support` | | `openresty_cert_dir` | `/etc/nginx/atsflare-certs` |
| `openresty_lua_dir` | `/etc/nginx/atsflare-lua` |
补充说明: 补充说明:
* Agent 当前会随受管配置一并向 OpenResty 注入 Lua 观测脚本,并在每次 heartbeat 前通过 `http://127.0.0.1:<openresty_observability_port>/atsflare/observability` 读取最近窗口请求指标 * Agent 会在每次启动时把受管 Lua 观测脚本覆盖释放到 `lua_dir`,不再由 Server 随配置版本下发
* OpenResty 通过 `openresty_lua_dir` 读取这些本地脚本,并在每次 heartbeat 前通过 `http://127.0.0.1:<openresty_observability_port>/atsflare/observability` 读取最近窗口请求指标
* 若 server 短暂离线,Agent 会把最近窗口观测先写入 `observability_buffer_path`,待 heartbeat 恢复后按时间窗口批量补传最近 `observability_replay_minutes` 分钟的数据 * 若 server 短暂离线,Agent 会把最近窗口观测先写入 `observability_buffer_path`,待 heartbeat 恢复后按时间窗口批量补传最近 `observability_replay_minutes` 分钟的数据
* 同一端口还会暴露仅本机可访问的 `stub_status`,用于采集 OpenResty 活动连接数 * 同一端口还会暴露仅本机可访问的 `stub_status`,用于采集 OpenResty 活动连接数
@@ -359,9 +365,11 @@ Docker OpenResty 模式下:
"openresty_path": "/usr/local/openresty/nginx/sbin/openresty", "openresty_path": "/usr/local/openresty/nginx/sbin/openresty",
"main_config_path": "/usr/local/openresty/nginx/conf/nginx.conf", "main_config_path": "/usr/local/openresty/nginx/conf/nginx.conf",
"route_config_path": "/usr/local/openresty/nginx/conf/conf.d/atsflare_routes.conf", "route_config_path": "/usr/local/openresty/nginx/conf/conf.d/atsflare_routes.conf",
"support_dir": "/usr/local/openresty/nginx/conf/support", "cert_dir": "/usr/local/openresty/nginx/conf/certs",
"openresty_support_dir": "/usr/local/openresty/nginx/conf/support" "openresty_cert_dir": "/usr/local/openresty/nginx/conf/certs",
} "lua_dir": "/usr/local/openresty/nginx/conf/lua",
"openresty_lua_dir": "/usr/local/openresty/nginx/conf/lua"
}
``` ```
--- ---
+7 -5
View File
@@ -285,8 +285,8 @@ export LOG_LEVEL='info'
验证点: 验证点:
1. 首次启动后确认 `data/etc/nginx/nginx.conf`、`data/etc/nginx/conf.d/atsflare_routes.conf` 与 `data/etc/nginx/support` 已由 Agent 创建 1. 首次启动后确认 `data/etc/nginx/nginx.conf`、`data/etc/nginx/conf.d/atsflare_routes.conf`、`data/etc/nginx/certs` 与 `data/etc/nginx/lua` 已由 Agent 创建
2. 确认容器实际挂载了主配置、路由目录和证书目录 2. 确认容器实际挂载了主配置、路由目录、证书目录和 Lua 目录
3. 确认宿主机本地可访问 `http://127.0.0.1:18081/atsflare/observability` 与 `http://127.0.0.1:18081/atsflare/stub_status` 3. 确认宿主机本地可访问 `http://127.0.0.1:18081/atsflare/observability` 与 `http://127.0.0.1:18081/atsflare/stub_status`
3. 在管理端发布一次新版本后,确认节点 `current_version` 追平激活版本 3. 在管理端发布一次新版本后,确认节点 `current_version` 追平激活版本
4. 在节点详情查看“当前目标版本”与“最近应用”,确认主配置/路由配置快照和 checksum 已可见 4. 在节点详情查看“当前目标版本”与“最近应用”,确认主配置/路由配置快照和 checksum 已可见
@@ -300,7 +300,7 @@ docker exec atsflare-openresty openresty -t
说明: 说明:
* `docker inspect` 重点确认主配置文件、`conf.d` 目录和证书目录都来自 Agent 受管路径 * `docker inspect` 重点确认主配置文件、`conf.d` 目录、证书目录和 Lua 目录都来自 Agent 受管路径
* 观测端口默认只绑定 `127.0.0.1`;若节点已有冲突,可在 `agent.json` 中调整 `openresty_observability_port` * 观测端口默认只绑定 `127.0.0.1`;若节点已有冲突,可在 `agent.json` 中调整 `openresty_observability_port`
* 若容器名使用默认值,请将上述命令中的名称替换为 `atsflare-openresty` * 若容器名使用默认值,请将上述命令中的名称替换为 `atsflare-openresty`
@@ -315,8 +315,10 @@ docker exec atsflare-openresty openresty -t
"openresty_path": "/usr/local/openresty/nginx/sbin/openresty", "openresty_path": "/usr/local/openresty/nginx/sbin/openresty",
"main_config_path": "/usr/local/openresty/nginx/conf/nginx.conf", "main_config_path": "/usr/local/openresty/nginx/conf/nginx.conf",
"route_config_path": "/usr/local/openresty/nginx/conf/conf.d/atsflare_routes.conf", "route_config_path": "/usr/local/openresty/nginx/conf/conf.d/atsflare_routes.conf",
"support_dir": "/usr/local/openresty/nginx/conf/support", "cert_dir": "/usr/local/openresty/nginx/conf/certs",
"openresty_support_dir": "/usr/local/openresty/nginx/conf/support", "openresty_cert_dir": "/usr/local/openresty/nginx/conf/certs",
"lua_dir": "/usr/local/openresty/nginx/conf/lua",
"openresty_lua_dir": "/usr/local/openresty/nginx/conf/lua",
"openresty_observability_port": 18081 "openresty_observability_port": 18081
} }
``` ```