feat: 添加可观察性监听地址支持,优化配置和测试用例

This commit is contained in:
ryan
2026-03-14 17:42:06 +08:00
parent 4be44733e8
commit 7628397785
8 changed files with 111 additions and 58 deletions
+7 -6
View File
@@ -60,12 +60,13 @@ func main() {
runtimeRouteConfigPath = nginx.DockerRouteConfigPath runtimeRouteConfigPath = nginx.DockerRouteConfigPath
} }
runtimeManager := &nginx.Manager{ runtimeManager := &nginx.Manager{
MainConfigPath: cfg.MainConfigPath, MainConfigPath: cfg.MainConfigPath,
RouteConfigPath: cfg.RouteConfigPath, RouteConfigPath: cfg.RouteConfigPath,
RuntimeRouteConfigPath: runtimeRouteConfigPath, RuntimeRouteConfigPath: runtimeRouteConfigPath,
SupportDir: cfg.SupportDir, SupportDir: cfg.SupportDir,
NginxSupportDir: cfg.OpenrestySupportDir, NginxSupportDir: cfg.OpenrestySupportDir,
OpenrestyObservabilityPort: cfg.OpenrestyObservabilityPort, OpenrestyObservabilityListen: nginx.ObservabilityListenAddress(cfg.OpenrestyPath, cfg.OpenrestyObservabilityPort),
OpenrestyObservabilityPort: cfg.OpenrestyObservabilityPort,
Executor: nginx.NewExecutor(nginx.ExecutorOptions{ Executor: nginx.NewExecutor(nginx.ExecutorOptions{
NginxPath: cfg.OpenrestyPath, NginxPath: cfg.OpenrestyPath,
DockerBinary: cfg.DockerBinary, DockerBinary: cfg.DockerBinary,
+25 -7
View File
@@ -21,6 +21,7 @@ const SupportDirPlaceholder = "__ATSF_SUPPORT_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__"
const ObservabilityListenPlaceholder = "__ATSF_OBSERVABILITY_LISTEN__"
const ObservabilityPortPlaceholder = "__ATSF_OBSERVABILITY_PORT__" const ObservabilityPortPlaceholder = "__ATSF_OBSERVABILITY_PORT__"
const DockerMainConfigPath = "/usr/local/openresty/nginx/conf/nginx.conf" const DockerMainConfigPath = "/usr/local/openresty/nginx/conf/nginx.conf"
const DockerRouteConfigPath = "/etc/nginx/conf.d/atsflare_routes.conf" const DockerRouteConfigPath = "/etc/nginx/conf.d/atsflare_routes.conf"
@@ -198,13 +199,14 @@ func (e *DockerExecutor) runContainer(ctx context.Context) error {
} }
type Manager struct { type Manager struct {
MainConfigPath string MainConfigPath string
RouteConfigPath string RouteConfigPath string
RuntimeRouteConfigPath string RuntimeRouteConfigPath string
SupportDir string SupportDir string
NginxSupportDir string NginxSupportDir string
OpenrestyObservabilityPort int OpenrestyObservabilityListen string
Executor Executor OpenrestyObservabilityPort int
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 {
@@ -298,6 +300,9 @@ func (m *Manager) CurrentChecksum() (string, error) {
if luaDir := m.luaRuntimePath(); luaDir != "" { if luaDir := m.luaRuntimePath(); luaDir != "" {
normalizedMain = strings.ReplaceAll(normalizedMain, luaDir, LuaDirPlaceholder) normalizedMain = strings.ReplaceAll(normalizedMain, luaDir, LuaDirPlaceholder)
} }
if listen := strings.TrimSpace(m.OpenrestyObservabilityListen); listen != "" {
normalizedMain = strings.ReplaceAll(normalizedMain, listen, ObservabilityListenPlaceholder)
}
if m.OpenrestyObservabilityPort > 0 { if m.OpenrestyObservabilityPort > 0 {
normalizedMain = strings.ReplaceAll(normalizedMain, fmt.Sprintf("%d", m.OpenrestyObservabilityPort), ObservabilityPortPlaceholder) normalizedMain = strings.ReplaceAll(normalizedMain, fmt.Sprintf("%d", m.OpenrestyObservabilityPort), ObservabilityPortPlaceholder)
} }
@@ -641,12 +646,25 @@ func (m *Manager) renderMainConfig(content string) string {
if luaDir := m.luaRuntimePath(); luaDir != "" { if luaDir := m.luaRuntimePath(); luaDir != "" {
rendered = strings.ReplaceAll(rendered, LuaDirPlaceholder, luaDir) rendered = strings.ReplaceAll(rendered, LuaDirPlaceholder, luaDir)
} }
if listen := strings.TrimSpace(m.OpenrestyObservabilityListen); listen != "" {
rendered = strings.ReplaceAll(rendered, ObservabilityListenPlaceholder, listen)
}
if m.OpenrestyObservabilityPort > 0 { if m.OpenrestyObservabilityPort > 0 {
rendered = strings.ReplaceAll(rendered, ObservabilityPortPlaceholder, fmt.Sprintf("%d", m.OpenrestyObservabilityPort)) rendered = strings.ReplaceAll(rendered, ObservabilityPortPlaceholder, fmt.Sprintf("%d", m.OpenrestyObservabilityPort))
} }
return rendered return rendered
} }
func ObservabilityListenAddress(openrestyPath string, port int) string {
if port <= 0 {
return ""
}
if strings.TrimSpace(openrestyPath) != "" {
return fmt.Sprintf("127.0.0.1:%d", port)
}
return fmt.Sprintf("%d", port)
}
func (m *Manager) routeConfigIncludePath() string { func (m *Manager) routeConfigIncludePath() string {
if strings.TrimSpace(m.RuntimeRouteConfigPath) != "" { if strings.TrimSpace(m.RuntimeRouteConfigPath) != "" {
return strings.TrimSpace(m.RuntimeRouteConfigPath) return strings.TrimSpace(m.RuntimeRouteConfigPath)
+23 -6
View File
@@ -460,14 +460,15 @@ func TestParseNginxVersionIgnoresDockerEntrypointPaths(t *testing.T) {
func TestManagerApplyWritesSupportFilesAndReplacesPlaceholder(t *testing.T) { func TestManagerApplyWritesSupportFilesAndReplacesPlaceholder(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"), SupportDir: filepath.Join(tempDir, "support"),
NginxSupportDir: "/etc/nginx/atsflare-support", NginxSupportDir: "/etc/nginx/atsflare-support",
Executor: &fakeExecutor{}, OpenrestyObservabilityListen: "18081",
Executor: &fakeExecutor{},
} }
err := manager.Apply(context.Background(), "include __ATSF_ROUTE_CONFIG__;", "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_SUPPORT_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"},
}) })
@@ -482,6 +483,13 @@ func TestManagerApplyWritesSupportFilesAndReplacesPlaceholder(t *testing.T) {
if !strings.Contains(string(routeData), "/etc/nginx/atsflare-support/1.crt") { if !strings.Contains(string(routeData), "/etc/nginx/atsflare-support/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)
if err != nil {
t.Fatalf("failed to read main config: %v", err)
}
if !strings.Contains(string(mainData), "listen 18081;") {
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.SupportDir, "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)
@@ -608,3 +616,12 @@ func TestManagerApplyRejectsSupportFilePathTraversal(t *testing.T) {
t.Fatalf("expected escaped file to not exist, stat err = %v", statErr) t.Fatalf("expected escaped file to not exist, stat err = %v", statErr)
} }
} }
func TestObservabilityListenAddress(t *testing.T) {
if got := ObservabilityListenAddress("", 18081); got != "18081" {
t.Fatalf("unexpected docker observability listen address: %s", got)
}
if got := ObservabilityListenAddress("/usr/local/openresty/nginx/sbin/openresty", 18081); got != "127.0.0.1:18081" {
t.Fatalf("unexpected path observability listen address: %s", got)
}
}
+1 -1
View File
@@ -37,7 +37,7 @@ type trafficAggregate struct {
} }
func BuildTrafficReport(cfg *config.Config, stateStore *state.Store, managed *managedOpenRestyMetrics) *protocol.NodeTrafficReport { func BuildTrafficReport(cfg *config.Config, stateStore *state.Store, managed *managedOpenRestyMetrics) *protocol.NodeTrafficReport {
if managed != nil && managed.TrafficReport != nil && managed.TrafficReport.RequestCount > 0 { if managed != nil && managed.TrafficReport != nil {
return managed.TrafficReport return managed.TrafficReport
} }
if cfg == nil || stateStore == nil { if cfg == nil || stateStore == nil {
@@ -6,6 +6,7 @@ import (
"testing" "testing"
"atsflare-agent/internal/config" "atsflare-agent/internal/config"
"atsflare-agent/internal/protocol"
"atsflare-agent/internal/state" "atsflare-agent/internal/state"
) )
@@ -107,3 +108,24 @@ func TestBuildTrafficReportParsesCombinedAccessLog(t *testing.T) {
t.Fatalf("expected combined access log to omit top domains when host is unavailable, got %+v", report.TopDomains) t.Fatalf("expected combined access log to omit top domains when host is unavailable, got %+v", report.TopDomains)
} }
} }
func TestBuildTrafficReportReturnsManagedWindowEvenWhenRequestCountZero(t *testing.T) {
report := BuildTrafficReport(nil, nil, &managedOpenRestyMetrics{
TrafficReport: &protocol.NodeTrafficReport{
WindowStartedAtUnix: 1710403200,
WindowEndedAtUnix: 1710403260,
RequestCount: 0,
ErrorCount: 0,
UniqueVisitorCount: 0,
StatusCodes: map[string]int64{},
TopDomains: map[string]int64{},
SourceCountries: map[string]int64{},
},
})
if report == nil {
t.Fatal("expected managed traffic report to be returned even when request count is zero")
}
if report.RequestCount != 0 || report.WindowStartedAtUnix != 1710403200 || report.WindowEndedAtUnix != 1710403260 {
t.Fatalf("unexpected managed traffic report: %+v", report)
}
}
+6 -5
View File
@@ -115,11 +115,12 @@ type configBundle struct {
} }
const ( const (
nginxSupportDirPlaceholder = "__ATSF_SUPPORT_DIR__" nginxSupportDirPlaceholder = "__ATSF_SUPPORT_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__"
nginxObservabilityPortPlaceholder = "__ATSF_OBSERVABILITY_PORT__" nginxObservabilityListenPlaceholder = "__ATSF_OBSERVABILITY_LISTEN__"
nginxObservabilityPortPlaceholder = "__ATSF_OBSERVABILITY_PORT__"
) )
var requiredMainConfigTemplatePlaceholders = []string{ var requiredMainConfigTemplatePlaceholders = []string{
+11 -2
View File
@@ -60,8 +60,11 @@ func TestCreateTLSCertificateAndRenderHTTPSConfig(t *testing.T) {
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__/observability/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 127.0.0.1:__ATSF_OBSERVABILITY_PORT__;") { if !strings.Contains(result.Version.MainConfig, "listen __ATSF_OBSERVABILITY_LISTEN__;") {
t.Fatal("expected main config to include managed openresty observability port placeholder") t.Fatal("expected main config to include managed openresty observability listen placeholder")
}
if strings.Contains(result.Version.MainConfig, "allow 127.0.0.1;") {
t.Fatal("expected main config to avoid hard-coded allow rules on observability server")
} }
if !strings.Contains(result.Version.RenderedConfig, "listen 443 ssl;") { if !strings.Contains(result.Version.RenderedConfig, "listen 443 ssl;") {
t.Fatal("expected rendered config to include https server block") t.Fatal("expected rendered config to include https server block")
@@ -78,6 +81,12 @@ func TestCreateTLSCertificateAndRenderHTTPSConfig(t *testing.T) {
if !strings.Contains(result.Version.SupportFilesJSON, "observability/log.lua") || !strings.Contains(result.Version.SupportFilesJSON, "observability/read.lua") { 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") 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) {
@@ -8,6 +8,7 @@ const (
openRestyObservabilityLogLuaPath = openRestyObservabilitySupportDir + "/log.lua" openRestyObservabilityLogLuaPath = openRestyObservabilitySupportDir + "/log.lua"
openRestyObservabilityReadLuaPath = openRestyObservabilitySupportDir + "/read.lua" openRestyObservabilityReadLuaPath = openRestyObservabilitySupportDir + "/read.lua"
openRestyObservabilityWindowTTL = 7200 openRestyObservabilityWindowTTL = 7200
openRestyObservabilityWindowSize = 60
) )
const openRestyObservabilityInitLua = `local dict = ngx.shared.atsflare_observability const openRestyObservabilityInitLua = `local dict = ngx.shared.atsflare_observability
@@ -15,12 +16,7 @@ if not dict then
return return
end end
local now = ngx.time() return
local current_window = dict:get("current_window")
if not current_window then
dict:set("current_window", now)
dict:set("window_started_at:" .. now, now)
end
` `
const openRestyObservabilityLogLua = `local dict = ngx.shared.atsflare_observability const openRestyObservabilityLogLua = `local dict = ngx.shared.atsflare_observability
@@ -28,15 +24,16 @@ if not dict then
return return
end end
local ttl = ` + "7200" + ` local request_uri = tostring(ngx.var.uri or "")
local current_window = dict:get("current_window") if request_uri == "/atsflare/observability" or request_uri == "/atsflare/stub_status" then
local now = ngx.time() return
if not current_window then
current_window = now
dict:set("current_window", current_window)
dict:set("window_started_at:" .. current_window, now)
end end
local ttl = ` + "7200" + `
local now = ngx.time()
local window_size = ` + "60" + `
local window_start = now - (now % window_size)
local function ensure_counter(key) local function ensure_counter(key)
dict:add(key, 0, ttl) dict:add(key, 0, ttl)
end end
@@ -64,7 +61,7 @@ local function remember_value(list_key, marker_key, value)
dict:set(list_key, existing .. "\n" .. value, ttl) dict:set(list_key, existing .. "\n" .. value, ttl)
end end
local window_prefix = tostring(current_window) local window_prefix = tostring(window_start)
incr("request_count:" .. window_prefix, 1) incr("request_count:" .. window_prefix, 1)
local status = tostring(ngx.status or 0) local status = tostring(ngx.status or 0)
@@ -116,12 +113,9 @@ if not dict then
end end
local now = ngx.time() local now = ngx.time()
local current_window = dict:get("current_window") local window_size = ` + "60" + `
if not current_window then local window_start = now - (now % window_size)
current_window = now local current_window = tostring(window_start)
dict:set("current_window", current_window)
dict:set("window_started_at:" .. current_window, now)
end
local function read_counter(key) local function read_counter(key)
return tonumber(dict:get(key) or 0) or 0 return tonumber(dict:get(key) or 0) or 0
@@ -140,7 +134,7 @@ local function read_map(window_id, prefix, list_key)
end end
local payload = { local payload = {
window_started_at_unix = read_counter("window_started_at:" .. current_window), window_started_at_unix = window_start,
window_ended_at_unix = now, window_ended_at_unix = now,
request_count = read_counter("request_count:" .. current_window), request_count = read_counter("request_count:" .. current_window),
error_count = read_counter("error_count:" .. current_window), error_count = read_counter("error_count:" .. current_window),
@@ -152,13 +146,6 @@ local payload = {
openresty_tx_bytes = read_counter("openresty_tx_bytes:" .. current_window) openresty_tx_bytes = read_counter("openresty_tx_bytes:" .. current_window)
} }
local next_window = now
if next_window <= current_window then
next_window = current_window + 1
end
dict:set("current_window", next_window)
dict:set("window_started_at:" .. next_window, now)
ngx.header.content_type = "application/json" ngx.header.content_type = "application/json"
ngx.say(cjson.encode(payload)) ngx.say(cjson.encode(payload))
` `
@@ -178,11 +165,9 @@ func renderOpenRestyObservabilityTemplateBlock() string {
fmt.Sprintf(" log_by_lua_file %s/%s;", nginxLuaDirPlaceholder, openRestyObservabilityLogLuaPath), fmt.Sprintf(" log_by_lua_file %s/%s;", nginxLuaDirPlaceholder, openRestyObservabilityLogLuaPath),
"", "",
fmt.Sprintf(" server {"), fmt.Sprintf(" server {"),
fmt.Sprintf(" listen 127.0.0.1:%s;", nginxObservabilityPortPlaceholder), fmt.Sprintf(" listen %s;", nginxObservabilityListenPlaceholder),
" server_name atsflare-observability;", " server_name atsflare-observability;",
" access_log off;", " access_log off;",
" allow 127.0.0.1;",
" deny all;",
"", "",
" location = /atsflare/observability {", " location = /atsflare/observability {",
" default_type application/json;", " default_type application/json;",