diff --git a/atsf_agent/cmd/agent/main.go b/atsf_agent/cmd/agent/main.go index 0dd1bec6..34fd808f 100644 --- a/atsf_agent/cmd/agent/main.go +++ b/atsf_agent/cmd/agent/main.go @@ -60,12 +60,13 @@ func main() { runtimeRouteConfigPath = nginx.DockerRouteConfigPath } runtimeManager := &nginx.Manager{ - MainConfigPath: cfg.MainConfigPath, - RouteConfigPath: cfg.RouteConfigPath, - RuntimeRouteConfigPath: runtimeRouteConfigPath, - SupportDir: cfg.SupportDir, - NginxSupportDir: cfg.OpenrestySupportDir, - OpenrestyObservabilityPort: cfg.OpenrestyObservabilityPort, + MainConfigPath: cfg.MainConfigPath, + RouteConfigPath: cfg.RouteConfigPath, + RuntimeRouteConfigPath: runtimeRouteConfigPath, + SupportDir: cfg.SupportDir, + NginxSupportDir: cfg.OpenrestySupportDir, + OpenrestyObservabilityListen: nginx.ObservabilityListenAddress(cfg.OpenrestyPath, cfg.OpenrestyObservabilityPort), + OpenrestyObservabilityPort: cfg.OpenrestyObservabilityPort, Executor: nginx.NewExecutor(nginx.ExecutorOptions{ NginxPath: cfg.OpenrestyPath, DockerBinary: cfg.DockerBinary, diff --git a/atsf_agent/internal/nginx/manager.go b/atsf_agent/internal/nginx/manager.go index 9ee49a74..1f5ffe02 100644 --- a/atsf_agent/internal/nginx/manager.go +++ b/atsf_agent/internal/nginx/manager.go @@ -21,6 +21,7 @@ const SupportDirPlaceholder = "__ATSF_SUPPORT_DIR__" const RouteConfigPlaceholder = "__ATSF_ROUTE_CONFIG__" const AccessLogPlaceholder = "__ATSF_ACCESS_LOG__" const LuaDirPlaceholder = "__ATSF_LUA_DIR__" +const ObservabilityListenPlaceholder = "__ATSF_OBSERVABILITY_LISTEN__" const ObservabilityPortPlaceholder = "__ATSF_OBSERVABILITY_PORT__" const DockerMainConfigPath = "/usr/local/openresty/nginx/conf/nginx.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 { - MainConfigPath string - RouteConfigPath string - RuntimeRouteConfigPath string - SupportDir string - NginxSupportDir string - OpenrestyObservabilityPort int - Executor Executor + MainConfigPath string + RouteConfigPath string + RuntimeRouteConfigPath string + SupportDir string + NginxSupportDir string + OpenrestyObservabilityListen string + OpenrestyObservabilityPort int + Executor Executor } 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 != "" { normalizedMain = strings.ReplaceAll(normalizedMain, luaDir, LuaDirPlaceholder) } + if listen := strings.TrimSpace(m.OpenrestyObservabilityListen); listen != "" { + normalizedMain = strings.ReplaceAll(normalizedMain, listen, ObservabilityListenPlaceholder) + } if m.OpenrestyObservabilityPort > 0 { 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 != "" { rendered = strings.ReplaceAll(rendered, LuaDirPlaceholder, luaDir) } + if listen := strings.TrimSpace(m.OpenrestyObservabilityListen); listen != "" { + rendered = strings.ReplaceAll(rendered, ObservabilityListenPlaceholder, listen) + } if m.OpenrestyObservabilityPort > 0 { rendered = strings.ReplaceAll(rendered, ObservabilityPortPlaceholder, fmt.Sprintf("%d", m.OpenrestyObservabilityPort)) } 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 { if strings.TrimSpace(m.RuntimeRouteConfigPath) != "" { return strings.TrimSpace(m.RuntimeRouteConfigPath) diff --git a/atsf_agent/internal/nginx/manager_test.go b/atsf_agent/internal/nginx/manager_test.go index f414c3d2..4051945c 100644 --- a/atsf_agent/internal/nginx/manager_test.go +++ b/atsf_agent/internal/nginx/manager_test.go @@ -460,14 +460,15 @@ func TestParseNginxVersionIgnoresDockerEntrypointPaths(t *testing.T) { func TestManagerApplyWritesSupportFilesAndReplacesPlaceholder(t *testing.T) { tempDir := t.TempDir() manager := &Manager{ - MainConfigPath: filepath.Join(tempDir, "nginx.conf"), - RouteConfigPath: filepath.Join(tempDir, "routes.conf"), - SupportDir: filepath.Join(tempDir, "support"), - NginxSupportDir: "/etc/nginx/atsflare-support", - Executor: &fakeExecutor{}, + MainConfigPath: filepath.Join(tempDir, "nginx.conf"), + RouteConfigPath: filepath.Join(tempDir, "routes.conf"), + SupportDir: filepath.Join(tempDir, "support"), + NginxSupportDir: "/etc/nginx/atsflare-support", + 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.key", Content: "key-data"}, }) @@ -482,6 +483,13 @@ func TestManagerApplyWritesSupportFilesAndReplacesPlaceholder(t *testing.T) { if !strings.Contains(string(routeData), "/etc/nginx/atsflare-support/1.crt") { 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")) if err != nil { 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) } } + +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) + } +} diff --git a/atsf_agent/internal/observability/traffic.go b/atsf_agent/internal/observability/traffic.go index c69e7fdb..32922d7f 100644 --- a/atsf_agent/internal/observability/traffic.go +++ b/atsf_agent/internal/observability/traffic.go @@ -37,7 +37,7 @@ type trafficAggregate struct { } 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 } if cfg == nil || stateStore == nil { diff --git a/atsf_agent/internal/observability/traffic_test.go b/atsf_agent/internal/observability/traffic_test.go index 447e0eff..6cbb1f37 100644 --- a/atsf_agent/internal/observability/traffic_test.go +++ b/atsf_agent/internal/observability/traffic_test.go @@ -6,6 +6,7 @@ import ( "testing" "atsflare-agent/internal/config" + "atsflare-agent/internal/protocol" "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) } } + +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) + } +} diff --git a/atsf_server/service/config_version.go b/atsf_server/service/config_version.go index a1a14bef..6bc1a842 100644 --- a/atsf_server/service/config_version.go +++ b/atsf_server/service/config_version.go @@ -115,11 +115,12 @@ type configBundle struct { } const ( - nginxSupportDirPlaceholder = "__ATSF_SUPPORT_DIR__" - nginxRouteConfigPlaceholder = "__ATSF_ROUTE_CONFIG__" - nginxAccessLogPlaceholder = "__ATSF_ACCESS_LOG__" - nginxLuaDirPlaceholder = "__ATSF_LUA_DIR__" - nginxObservabilityPortPlaceholder = "__ATSF_OBSERVABILITY_PORT__" + nginxSupportDirPlaceholder = "__ATSF_SUPPORT_DIR__" + nginxRouteConfigPlaceholder = "__ATSF_ROUTE_CONFIG__" + nginxAccessLogPlaceholder = "__ATSF_ACCESS_LOG__" + nginxLuaDirPlaceholder = "__ATSF_LUA_DIR__" + nginxObservabilityListenPlaceholder = "__ATSF_OBSERVABILITY_LISTEN__" + nginxObservabilityPortPlaceholder = "__ATSF_OBSERVABILITY_PORT__" ) var requiredMainConfigTemplatePlaceholders = []string{ diff --git a/atsf_server/service/https_phase1_test.go b/atsf_server/service/https_phase1_test.go index 8c377350..6c3f7895 100644 --- a/atsf_server/service/https_phase1_test.go +++ b/atsf_server/service/https_phase1_test.go @@ -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;") { 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__;") { - t.Fatal("expected main config to include managed openresty observability port placeholder") + if !strings.Contains(result.Version.MainConfig, "listen __ATSF_OBSERVABILITY_LISTEN__;") { + 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;") { 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") { 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) { diff --git a/atsf_server/service/openresty_observability_assets.go b/atsf_server/service/openresty_observability_assets.go index cc99a0c0..e962c3ce 100644 --- a/atsf_server/service/openresty_observability_assets.go +++ b/atsf_server/service/openresty_observability_assets.go @@ -8,6 +8,7 @@ const ( openRestyObservabilityLogLuaPath = openRestyObservabilitySupportDir + "/log.lua" openRestyObservabilityReadLuaPath = openRestyObservabilitySupportDir + "/read.lua" openRestyObservabilityWindowTTL = 7200 + openRestyObservabilityWindowSize = 60 ) const openRestyObservabilityInitLua = `local dict = ngx.shared.atsflare_observability @@ -15,12 +16,7 @@ if not dict then return end -local now = ngx.time() -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 +return ` const openRestyObservabilityLogLua = `local dict = ngx.shared.atsflare_observability @@ -28,15 +24,16 @@ if not dict then return end -local ttl = ` + "7200" + ` -local current_window = dict:get("current_window") -local now = ngx.time() -if not current_window then - current_window = now - dict:set("current_window", current_window) - dict:set("window_started_at:" .. current_window, now) +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 @@ -64,7 +61,7 @@ local function remember_value(list_key, marker_key, value) dict:set(list_key, existing .. "\n" .. value, ttl) end -local window_prefix = tostring(current_window) +local window_prefix = tostring(window_start) incr("request_count:" .. window_prefix, 1) local status = tostring(ngx.status or 0) @@ -116,12 +113,9 @@ if not dict then end local now = ngx.time() -local current_window = dict:get("current_window") -if not current_window then - current_window = now - dict:set("current_window", current_window) - dict:set("window_started_at:" .. current_window, now) -end +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 @@ -140,7 +134,7 @@ local function read_map(window_id, prefix, list_key) end local payload = { - window_started_at_unix = read_counter("window_started_at:" .. current_window), + 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), @@ -152,13 +146,6 @@ local payload = { 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.say(cjson.encode(payload)) ` @@ -178,11 +165,9 @@ func renderOpenRestyObservabilityTemplateBlock() string { fmt.Sprintf(" log_by_lua_file %s/%s;", nginxLuaDirPlaceholder, openRestyObservabilityLogLuaPath), "", fmt.Sprintf(" server {"), - fmt.Sprintf(" listen 127.0.0.1:%s;", nginxObservabilityPortPlaceholder), + fmt.Sprintf(" listen %s;", nginxObservabilityListenPlaceholder), " server_name atsflare-observability;", " access_log off;", - " allow 127.0.0.1;", - " deny all;", "", " location = /atsflare/observability {", " default_type application/json;",