diff --git a/README.md b/README.md index 2f9104f3..1473f91f 100644 --- a/README.md +++ b/README.md @@ -46,6 +46,7 @@ ATSFlare 当前定位为内部自用的反向代理控制面,不面向外部 * 配置版本化:支持预览、发布、激活、历史回滚,版本不可变 * 节点接入:支持全局 `discovery_token` 首次接入,也支持节点专属 `agent_token` * Agent 自动应用:周期性同步、落盘、`openresty -t`、`openresty -s reload`、失败自动回滚 +* 节点观测:Agent 会向受管 OpenResty 注入 Lua 观测脚本,按 heartbeat 上报最近窗口请求、错误、UV 与连接指标 * TLS 与域名管理:支持证书托管、域名资产维护、精确匹配与通配符匹配 * 运维能力:配置变更摘要、Agent 运行参数下发、Agent 正式版自动更新与 preview 手动升级、Server 正式版 GitHub 自升级、Server preview 手动检查升级、Server 手动上传二进制确认升级 * 管理端 UI:基于 Next.js App Router + React 19 + Tailwind CSS 4 的新版前端 @@ -213,6 +214,7 @@ Docker 镜像工作流仅构建 `atsf_server`,并产出 `linux/amd64` 与 `lin | `agent_token` | 节点专属认证 Token | | `discovery_token` | 首次自动注册使用的全局 Token | | `data_dir` | Agent 托管数据目录 | +| `openresty_observability_port` | Agent 读取 OpenResty Lua 本地观测指标的 loopback 端口 | | `nginx_path` | 本机 Nginx 路径,设置后走本机模式 | | `nginx_container_name` | Docker 模式下的 Nginx 容器名 | diff --git a/atsf_agent/cmd/agent/main.go b/atsf_agent/cmd/agent/main.go index 5c39bf89..a3efaa7a 100644 --- a/atsf_agent/cmd/agent/main.go +++ b/atsf_agent/cmd/agent/main.go @@ -33,14 +33,15 @@ func main() { cfg.NginxVersion = nginx.DetectVersion( context.Background(), nginx.ExecutorOptions{ - NginxPath: cfg.OpenrestyPath, - DockerBinary: cfg.DockerBinary, - ContainerName: cfg.OpenrestyContainerName, - Image: cfg.OpenrestyDockerImage, - MainConfigPath: cfg.MainConfigPath, - RouteConfigPath: cfg.RouteConfigPath, - CertDir: cfg.CertDir, - NginxCertDir: cfg.OpenrestyCertDir, + NginxPath: cfg.OpenrestyPath, + DockerBinary: cfg.DockerBinary, + ContainerName: cfg.OpenrestyContainerName, + Image: cfg.OpenrestyDockerImage, + MainConfigPath: cfg.MainConfigPath, + RouteConfigPath: cfg.RouteConfigPath, + CertDir: cfg.CertDir, + NginxCertDir: cfg.OpenrestyCertDir, + OpenrestyObservabilityPort: cfg.OpenrestyObservabilityPort, }, ) slog.Info("agent config loaded", @@ -59,20 +60,22 @@ func main() { runtimeRouteConfigPath = nginx.DockerRouteConfigPath } runtimeManager := &nginx.Manager{ - MainConfigPath: cfg.MainConfigPath, - RouteConfigPath: cfg.RouteConfigPath, - RuntimeRouteConfigPath: runtimeRouteConfigPath, - CertDir: cfg.CertDir, - NginxCertDir: cfg.OpenrestyCertDir, + MainConfigPath: cfg.MainConfigPath, + RouteConfigPath: cfg.RouteConfigPath, + RuntimeRouteConfigPath: runtimeRouteConfigPath, + CertDir: cfg.CertDir, + NginxCertDir: cfg.OpenrestyCertDir, + OpenrestyObservabilityPort: cfg.OpenrestyObservabilityPort, Executor: nginx.NewExecutor(nginx.ExecutorOptions{ - NginxPath: cfg.OpenrestyPath, - DockerBinary: cfg.DockerBinary, - ContainerName: cfg.OpenrestyContainerName, - Image: cfg.OpenrestyDockerImage, - MainConfigPath: cfg.MainConfigPath, - RouteConfigPath: cfg.RouteConfigPath, - CertDir: cfg.CertDir, - NginxCertDir: cfg.OpenrestyCertDir, + NginxPath: cfg.OpenrestyPath, + DockerBinary: cfg.DockerBinary, + ContainerName: cfg.OpenrestyContainerName, + Image: cfg.OpenrestyDockerImage, + MainConfigPath: cfg.MainConfigPath, + RouteConfigPath: cfg.RouteConfigPath, + CertDir: cfg.CertDir, + NginxCertDir: cfg.OpenrestyCertDir, + OpenrestyObservabilityPort: cfg.OpenrestyObservabilityPort, }), } runner := &agent.Runner{ diff --git a/atsf_agent/internal/agent/runner.go b/atsf_agent/internal/agent/runner.go index 6ef78e7c..755dcec4 100644 --- a/atsf_agent/internal/agent/runner.go +++ b/atsf_agent/internal/agent/runner.go @@ -312,8 +312,9 @@ func (r *Runner) nodePayload(nodeID string) protocol.NodePayload { openrestyStatus = protocol.OpenrestyStatusUnknown } profile := observability.BuildProfile(r.Config, r.StateStore) - metricSnapshot := observability.BuildSnapshot(r.Config, r.StateStore) - trafficReport := observability.BuildTrafficReport(r.Config, r.StateStore) + managedOpenRestyMetrics := observability.CollectManagedOpenRestyMetrics(r.Config) + metricSnapshot := observability.BuildSnapshot(r.Config, r.StateStore, managedOpenRestyMetrics) + trafficReport := observability.BuildTrafficReport(r.Config, r.StateStore, managedOpenRestyMetrics) healthEvents := observability.BuildHealthEvents(snapshot) return protocol.NodePayload{ NodeID: nodeID, diff --git a/atsf_agent/internal/config/config.go b/atsf_agent/internal/config/config.go index 78763bae..669fe0fb 100644 --- a/atsf_agent/internal/config/config.go +++ b/atsf_agent/internal/config/config.go @@ -17,49 +17,52 @@ const ( defaultCertDirRelativePath = "etc/nginx/certs" defaultDockerStateRelativePath = "var/lib/atsflare/agent-state.json" defaultDockerOpenRestyCertDir = "/etc/nginx/atsflare-certs" + defaultOpenRestyObservabilityPort = 18081 ) type Config struct { - ServerURL string `json:"server_url"` - AgentToken string `json:"agent_token"` - DiscoveryToken string `json:"discovery_token"` - NodeName string `json:"node_name"` - NodeIP string `json:"node_ip"` - AgentVersion string `json:"-"` - NginxVersion string `json:"-"` - OpenrestyPath string `json:"openresty_path"` - OpenrestyContainerName string `json:"openresty_container_name"` - OpenrestyDockerImage string `json:"openresty_docker_image"` - DockerBinary string `json:"docker_binary"` - DataDir string `json:"data_dir"` - MainConfigPath string `json:"main_config_path"` - RouteConfigPath string `json:"route_config_path"` - CertDir string `json:"cert_dir"` - OpenrestyCertDir string `json:"openresty_cert_dir"` - StatePath string `json:"state_path"` - HeartbeatInterval MillisecondDuration `json:"heartbeat_interval"` - RequestTimeout MillisecondDuration `json:"request_timeout"` - configPath string + ServerURL string `json:"server_url"` + AgentToken string `json:"agent_token"` + DiscoveryToken string `json:"discovery_token"` + NodeName string `json:"node_name"` + NodeIP string `json:"node_ip"` + AgentVersion string `json:"-"` + NginxVersion string `json:"-"` + OpenrestyPath string `json:"openresty_path"` + OpenrestyContainerName string `json:"openresty_container_name"` + OpenrestyDockerImage string `json:"openresty_docker_image"` + DockerBinary string `json:"docker_binary"` + DataDir string `json:"data_dir"` + MainConfigPath string `json:"main_config_path"` + RouteConfigPath string `json:"route_config_path"` + CertDir string `json:"cert_dir"` + OpenrestyCertDir string `json:"openresty_cert_dir"` + OpenrestyObservabilityPort int `json:"openresty_observability_port"` + StatePath string `json:"state_path"` + HeartbeatInterval MillisecondDuration `json:"heartbeat_interval"` + RequestTimeout MillisecondDuration `json:"request_timeout"` + configPath string } type configFile struct { - ServerURL string `json:"server_url"` - AgentToken string `json:"agent_token"` - DiscoveryToken string `json:"discovery_token"` - NodeName string `json:"node_name"` - NodeIP string `json:"node_ip"` - OpenrestyPath string `json:"openresty_path"` - OpenrestyContainerName string `json:"openresty_container_name"` - OpenrestyDockerImage string `json:"openresty_docker_image"` - DockerBinary string `json:"docker_binary"` - DataDir string `json:"data_dir"` - MainConfigPath string `json:"main_config_path"` - RouteConfigPath string `json:"route_config_path"` - CertDir string `json:"cert_dir"` - OpenrestyCertDir string `json:"openresty_cert_dir"` - StatePath string `json:"state_path"` - HeartbeatInterval MillisecondDuration `json:"heartbeat_interval"` - RequestTimeout MillisecondDuration `json:"request_timeout"` + ServerURL string `json:"server_url"` + AgentToken string `json:"agent_token"` + DiscoveryToken string `json:"discovery_token"` + NodeName string `json:"node_name"` + NodeIP string `json:"node_ip"` + OpenrestyPath string `json:"openresty_path"` + OpenrestyContainerName string `json:"openresty_container_name"` + OpenrestyDockerImage string `json:"openresty_docker_image"` + DockerBinary string `json:"docker_binary"` + DataDir string `json:"data_dir"` + MainConfigPath string `json:"main_config_path"` + RouteConfigPath string `json:"route_config_path"` + CertDir string `json:"cert_dir"` + OpenrestyCertDir string `json:"openresty_cert_dir"` + OpenrestyObservabilityPort int `json:"openresty_observability_port"` + StatePath string `json:"state_path"` + HeartbeatInterval MillisecondDuration `json:"heartbeat_interval"` + RequestTimeout MillisecondDuration `json:"request_timeout"` } func Load(path string) (*Config, error) { @@ -72,23 +75,24 @@ func Load(path string) (*Config, error) { return nil, err } cfg := &Config{ - ServerURL: file.ServerURL, - AgentToken: file.AgentToken, - DiscoveryToken: file.DiscoveryToken, - NodeName: file.NodeName, - NodeIP: file.NodeIP, - OpenrestyPath: file.OpenrestyPath, - OpenrestyContainerName: file.OpenrestyContainerName, - OpenrestyDockerImage: file.OpenrestyDockerImage, - DockerBinary: file.DockerBinary, - DataDir: file.DataDir, - MainConfigPath: file.MainConfigPath, - RouteConfigPath: file.RouteConfigPath, - CertDir: file.CertDir, - OpenrestyCertDir: file.OpenrestyCertDir, - StatePath: file.StatePath, - HeartbeatInterval: file.HeartbeatInterval, - RequestTimeout: file.RequestTimeout, + ServerURL: file.ServerURL, + AgentToken: file.AgentToken, + DiscoveryToken: file.DiscoveryToken, + NodeName: file.NodeName, + NodeIP: file.NodeIP, + OpenrestyPath: file.OpenrestyPath, + OpenrestyContainerName: file.OpenrestyContainerName, + OpenrestyDockerImage: file.OpenrestyDockerImage, + DockerBinary: file.DockerBinary, + DataDir: file.DataDir, + MainConfigPath: file.MainConfigPath, + RouteConfigPath: file.RouteConfigPath, + CertDir: file.CertDir, + OpenrestyCertDir: file.OpenrestyCertDir, + OpenrestyObservabilityPort: file.OpenrestyObservabilityPort, + StatePath: file.StatePath, + HeartbeatInterval: file.HeartbeatInterval, + RequestTimeout: file.RequestTimeout, } cfg.configPath = path applyDefaults(cfg, filepath.Dir(path)) @@ -144,6 +148,9 @@ func applyDefaults(cfg *Config, baseDir string) { cfg.OpenrestyCertDir = defaultDockerOpenRestyCertDir } } + if cfg.OpenrestyObservabilityPort <= 0 { + cfg.OpenrestyObservabilityPort = defaultOpenRestyObservabilityPort + } if cfg.HeartbeatInterval <= 0 { cfg.HeartbeatInterval = MillisecondDuration(10 * time.Second) } @@ -198,6 +205,9 @@ func validate(cfg *Config) error { if cfg.NodeIP == "" { return errors.New("node_ip 不能为空") } + if cfg.OpenrestyObservabilityPort <= 0 || cfg.OpenrestyObservabilityPort > 65535 { + return errors.New("openresty_observability_port 必须在 1-65535 之间") + } return nil } diff --git a/atsf_agent/internal/config/config_test.go b/atsf_agent/internal/config/config_test.go index 9dd37704..1fe54864 100644 --- a/atsf_agent/internal/config/config_test.go +++ b/atsf_agent/internal/config/config_test.go @@ -54,6 +54,9 @@ func TestLoadDockerModeUsesManagedPaths(t *testing.T) { if cfg.StatePath != filepath.Join(dir, "data", defaultDockerStateRelativePath) { t.Fatalf("unexpected state path: %s", cfg.StatePath) } + if cfg.OpenrestyObservabilityPort != defaultOpenRestyObservabilityPort { + t.Fatalf("unexpected openresty observability port: %d", cfg.OpenrestyObservabilityPort) + } } func TestLoadPathModeKeepsExplicitPaths(t *testing.T) { @@ -94,6 +97,9 @@ func TestLoadPathModeKeepsExplicitPaths(t *testing.T) { if cfg.OpenrestyCertDir != cfg.CertDir { t.Fatalf("expected path mode openresty cert dir to equal cert dir, got %s / %s", cfg.OpenrestyCertDir, cfg.CertDir) } + if cfg.OpenrestyObservabilityPort != defaultOpenRestyObservabilityPort { + t.Fatalf("unexpected path mode openresty observability port: %d", cfg.OpenrestyObservabilityPort) + } } func TestLoadUsesCustomDataDirForGeneratedFiles(t *testing.T) { @@ -204,6 +210,9 @@ func TestSavePersistsMillisecondsAndOmitsRuntimeVersions(t *testing.T) { if decoded["request_timeout"] != float64(7000) { t.Fatalf("unexpected request timeout: %#v", decoded["request_timeout"]) } + if decoded["openresty_observability_port"] != float64(defaultOpenRestyObservabilityPort) { + t.Fatalf("unexpected observability port: %#v", decoded["openresty_observability_port"]) + } if _, ok := decoded["nginx_path"]; ok { t.Fatal("legacy nginx_path should not be persisted") } diff --git a/atsf_agent/internal/nginx/manager.go b/atsf_agent/internal/nginx/manager.go index d2bdd41b..74b51280 100644 --- a/atsf_agent/internal/nginx/manager.go +++ b/atsf_agent/internal/nginx/manager.go @@ -20,6 +20,8 @@ import ( const CertDirPlaceholder = "__ATSF_CERT_DIR__" const RouteConfigPlaceholder = "__ATSF_ROUTE_CONFIG__" const AccessLogPlaceholder = "__ATSF_ACCESS_LOG__" +const LuaDirPlaceholder = "__ATSF_LUA_DIR__" +const ObservabilityPortPlaceholder = "__ATSF_OBSERVABILITY_PORT__" const DockerMainConfigPath = "/usr/local/openresty/nginx/conf/nginx.conf" const DockerRouteConfigPath = "/etc/nginx/conf.d/atsflare_routes.conf" const DockerAccessLogPath = "/etc/nginx/conf.d/atsflare_access.log" @@ -97,14 +99,15 @@ func (e *PathExecutor) Restart(ctx context.Context) error { } type DockerExecutor struct { - DockerBinary string - ContainerName string - Image string - MainConfigPath string - RouteConfigDir string - CertDir string - NginxCertDir string - Runner CommandRunner + DockerBinary string + ContainerName string + Image string + MainConfigPath string + RouteConfigDir string + CertDir string + NginxCertDir string + OpenrestyObservabilityPort int + Runner CommandRunner } func (e *DockerExecutor) Test(ctx context.Context) error { @@ -180,6 +183,7 @@ func (e *DockerExecutor) runContainer(ctx context.Context) error { "--name", e.ContainerName, "-p", "80:80", "-p", "443:443", + "-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:/etc/nginx/conf.d", e.RouteConfigDir), "-v", fmt.Sprintf("%s:%s", e.CertDir, e.NginxCertDir), @@ -194,12 +198,13 @@ func (e *DockerExecutor) runContainer(ctx context.Context) error { } type Manager struct { - MainConfigPath string - RouteConfigPath string - RuntimeRouteConfigPath string - CertDir string - NginxCertDir string - Executor Executor + MainConfigPath string + RouteConfigPath string + RuntimeRouteConfigPath string + CertDir string + NginxCertDir string + OpenrestyObservabilityPort int + Executor Executor } func (m *Manager) Apply(ctx context.Context, mainConfig string, routeConfig string, supportFiles []protocol.SupportFile) error { @@ -290,6 +295,12 @@ func (m *Manager) CurrentChecksum() (string, error) { if accessLogPath := m.accessLogRuntimePath(); accessLogPath != "" { normalizedMain = strings.ReplaceAll(normalizedMain, accessLogPath, AccessLogPlaceholder) } + if luaDir := m.luaRuntimePath(); luaDir != "" { + normalizedMain = strings.ReplaceAll(normalizedMain, luaDir, LuaDirPlaceholder) + } + if m.OpenrestyObservabilityPort > 0 { + normalizedMain = strings.ReplaceAll(normalizedMain, fmt.Sprintf("%d", m.OpenrestyObservabilityPort), ObservabilityPortPlaceholder) + } normalizedRoute := string(data) if m.NginxCertDir != "" { normalizedRoute = strings.ReplaceAll(normalizedRoute, m.NginxCertDir, CertDirPlaceholder) @@ -304,14 +315,15 @@ func (m *Manager) CurrentChecksum() (string, error) { } type ExecutorOptions struct { - NginxPath string - DockerBinary string - ContainerName string - Image string - MainConfigPath string - RouteConfigPath string - CertDir string - NginxCertDir string + NginxPath string + DockerBinary string + ContainerName string + Image string + MainConfigPath string + RouteConfigPath string + CertDir string + NginxCertDir string + OpenrestyObservabilityPort int } func NewExecutor(options ExecutorOptions) Executor { @@ -335,14 +347,15 @@ func NewExecutor(options ExecutorOptions) Executor { certDir = absDir } return &DockerExecutor{ - DockerBinary: options.DockerBinary, - ContainerName: options.ContainerName, - Image: options.Image, - MainConfigPath: mainConfigPath, - RouteConfigDir: routeConfigDir, - CertDir: certDir, - NginxCertDir: options.NginxCertDir, - Runner: runner, + DockerBinary: options.DockerBinary, + ContainerName: options.ContainerName, + Image: options.Image, + MainConfigPath: mainConfigPath, + RouteConfigDir: routeConfigDir, + CertDir: certDir, + NginxCertDir: options.NginxCertDir, + OpenrestyObservabilityPort: options.OpenrestyObservabilityPort, + Runner: runner, } } @@ -588,7 +601,11 @@ func (m *Manager) supportFileTargetPath(relativePath string) (string, error) { if strings.TrimSpace(m.CertDir) == "" { return "", errors.New("cert dir 不能为空") } - normalizedPath := filepath.Clean(filepath.FromSlash(strings.TrimSpace(relativePath))) + candidate := strings.TrimSpace(relativePath) + if strings.Contains(candidate, `\`) { + candidate = strings.ReplaceAll(candidate, `\`, "/") + } + normalizedPath := filepath.Clean(filepath.FromSlash(candidate)) if normalizedPath == "." || normalizedPath == "" { return "", errors.New("support file path 不能为空") } @@ -621,6 +638,12 @@ func (m *Manager) renderMainConfig(content string) string { if accessLogPath := m.accessLogRuntimePath(); accessLogPath != "" { rendered = strings.ReplaceAll(rendered, AccessLogPlaceholder, accessLogPath) } + if luaDir := m.luaRuntimePath(); luaDir != "" { + rendered = strings.ReplaceAll(rendered, LuaDirPlaceholder, luaDir) + } + if m.OpenrestyObservabilityPort > 0 { + rendered = strings.ReplaceAll(rendered, ObservabilityPortPlaceholder, fmt.Sprintf("%d", m.OpenrestyObservabilityPort)) + } return rendered } @@ -639,6 +662,13 @@ func (m *Manager) accessLogRuntimePath() string { return filepath.ToSlash(filepath.Join(filepath.Dir(includePath), "atsflare_access.log")) } +func (m *Manager) luaRuntimePath() string { + if strings.TrimSpace(m.NginxCertDir) == "" { + return "" + } + return filepath.ToSlash(m.NginxCertDir) +} + func checksum(content string) string { sum := sha256.Sum256([]byte(content)) return hex.EncodeToString(sum[:]) diff --git a/atsf_agent/internal/nginx/manager_test.go b/atsf_agent/internal/nginx/manager_test.go index 0dc3fa8a..ff902e65 100644 --- a/atsf_agent/internal/nginx/manager_test.go +++ b/atsf_agent/internal/nginx/manager_test.go @@ -210,14 +210,15 @@ func TestDockerExecutorRunContainerMountsManagedFiles(t *testing.T) { certDir := filepath.Clean("/tmp/managed/certs") runner := &fakeRunner{} executor := &DockerExecutor{ - DockerBinary: "docker", - ContainerName: "atsflare-openresty", - Image: "openresty/openresty:alpine", - MainConfigPath: mainConfigPath, - RouteConfigDir: routeConfigDir, - CertDir: certDir, - NginxCertDir: "/etc/nginx/atsflare-certs", - Runner: runner, + DockerBinary: "docker", + ContainerName: "atsflare-openresty", + Image: "openresty/openresty:alpine", + MainConfigPath: mainConfigPath, + RouteConfigDir: routeConfigDir, + CertDir: certDir, + NginxCertDir: "/etc/nginx/atsflare-certs", + OpenrestyObservabilityPort: 18081, + Runner: runner, } if err := executor.runContainer(context.Background()); err != nil { @@ -233,6 +234,7 @@ func TestDockerExecutorRunContainerMountsManagedFiles(t *testing.T) { "--name", "atsflare-openresty", "-p", "80:80", "-p", "443:443", + "-p", "127.0.0.1:18081:18081", "-v", mainConfigPath + ":" + DockerMainConfigPath, "-v", routeConfigDir + ":/etc/nginx/conf.d", "-v", certDir + ":/etc/nginx/atsflare-certs", @@ -253,14 +255,15 @@ func TestDockerExecutorRecreatesContainerOnStartup(t *testing.T) { }, } executor := &DockerExecutor{ - DockerBinary: "docker", - ContainerName: "atsflare-openresty", - Image: "openresty/openresty:alpine", - MainConfigPath: filepath.Clean("/tmp/nginx.conf"), - RouteConfigDir: filepath.Clean("/tmp/routes"), - CertDir: filepath.Clean("/tmp/certs"), - NginxCertDir: "/etc/nginx/atsflare-certs", - Runner: runner, + DockerBinary: "docker", + ContainerName: "atsflare-openresty", + Image: "openresty/openresty:alpine", + MainConfigPath: filepath.Clean("/tmp/nginx.conf"), + RouteConfigDir: filepath.Clean("/tmp/routes"), + CertDir: filepath.Clean("/tmp/certs"), + NginxCertDir: "/etc/nginx/atsflare-certs", + OpenrestyObservabilityPort: 18081, + Runner: runner, } if err := executor.EnsureRuntime(context.Background(), true); err != nil { @@ -279,13 +282,14 @@ func TestDockerExecutorRecreatesContainerOnStartup(t *testing.T) { func TestNewExecutorUsesAbsoluteDockerMountPath(t *testing.T) { executor := NewExecutor(ExecutorOptions{ - DockerBinary: "docker", - ContainerName: "atsflare-openresty", - Image: "openresty/openresty:alpine", - MainConfigPath: "./data/etc/nginx/nginx.conf", - RouteConfigPath: "./data/etc/nginx/conf.d/atsflare_routes.conf", - CertDir: "./data/etc/nginx/certs", - NginxCertDir: "/etc/nginx/atsflare-certs", + DockerBinary: "docker", + ContainerName: "atsflare-openresty", + Image: "openresty/openresty:alpine", + MainConfigPath: "./data/etc/nginx/nginx.conf", + RouteConfigPath: "./data/etc/nginx/conf.d/atsflare_routes.conf", + CertDir: "./data/etc/nginx/certs", + NginxCertDir: "/etc/nginx/atsflare-certs", + OpenrestyObservabilityPort: 18081, }) dockerExecutor, ok := executor.(*DockerExecutor) diff --git a/atsf_agent/internal/observability/collector.go b/atsf_agent/internal/observability/collector.go index 58273a39..6ea376ec 100644 --- a/atsf_agent/internal/observability/collector.go +++ b/atsf_agent/internal/observability/collector.go @@ -40,7 +40,7 @@ func BuildProfile(cfg *config.Config, stateStore *state.Store) *protocol.NodeSys return profile } -func BuildSnapshot(cfg *config.Config, stateStore *state.Store) *protocol.NodeMetricSnapshot { +func BuildSnapshot(cfg *config.Config, stateStore *state.Store, managed *managedOpenRestyMetrics) *protocol.NodeMetricSnapshot { now := time.Now().UTC() metric := &protocol.NodeMetricSnapshot{ CapturedAtUnix: now.Unix(), @@ -56,6 +56,11 @@ func BuildSnapshot(cfg *config.Config, stateStore *state.Store) *protocol.NodeMe metric.NetworkRxBytes, metric.NetworkTxBytes = readLinuxNetworkTotals() metric.DiskReadBytes, metric.DiskWriteBytes = readLinuxDiskTotals() + if managed != nil { + metric.OpenrestyRxBytes = managed.OpenrestyRxBytes + metric.OpenrestyTxBytes = managed.OpenrestyTxBytes + metric.OpenrestyConnections = managed.OpenrestyConnections + } if stateStore == nil { return metric diff --git a/atsf_agent/internal/observability/openresty_local.go b/atsf_agent/internal/observability/openresty_local.go new file mode 100644 index 00000000..25a0953c --- /dev/null +++ b/atsf_agent/internal/observability/openresty_local.go @@ -0,0 +1,129 @@ +package observability + +import ( + "atsflare-agent/internal/config" + "atsflare-agent/internal/protocol" + "encoding/json" + "fmt" + "io" + "net/http" + "regexp" + "strconv" + "strings" + "time" +) + +const openRestyObservabilityPath = "/atsflare/observability" +const openRestyStubStatusPath = "/atsflare/stub_status" + +var stubStatusActivePattern = regexp.MustCompile(`Active connections:\s+(\d+)`) + +type managedOpenRestyMetrics struct { + TrafficReport *protocol.NodeTrafficReport + OpenrestyRxBytes int64 + OpenrestyTxBytes int64 + OpenrestyConnections int64 +} + +type openRestyObservabilityResponse struct { + WindowStartedAtUnix int64 `json:"window_started_at_unix"` + WindowEndedAtUnix int64 `json:"window_ended_at_unix"` + RequestCount int64 `json:"request_count"` + ErrorCount int64 `json:"error_count"` + UniqueVisitorCount int64 `json:"unique_visitor_count"` + StatusCodes map[string]int64 `json:"status_codes"` + TopDomains map[string]int64 `json:"top_domains"` + SourceCountries map[string]int64 `json:"source_countries"` + OpenrestyRxBytes int64 `json:"openresty_rx_bytes"` + OpenrestyTxBytes int64 `json:"openresty_tx_bytes"` +} + +func CollectManagedOpenRestyMetrics(cfg *config.Config) *managedOpenRestyMetrics { + if cfg == nil || cfg.OpenrestyObservabilityPort <= 0 { + return nil + } + + baseURL := fmt.Sprintf("http://127.0.0.1:%d", cfg.OpenrestyObservabilityPort) + client := &http.Client{Timeout: 1500 * time.Millisecond} + + observabilityResp := openRestyObservabilityResponse{} + if err := fetchLocalJSON(client, baseURL+openRestyObservabilityPath, &observabilityResp); err != nil { + return nil + } + + result := &managedOpenRestyMetrics{ + TrafficReport: &protocol.NodeTrafficReport{ + WindowStartedAtUnix: observabilityResp.WindowStartedAtUnix, + WindowEndedAtUnix: observabilityResp.WindowEndedAtUnix, + RequestCount: observabilityResp.RequestCount, + ErrorCount: observabilityResp.ErrorCount, + UniqueVisitorCount: observabilityResp.UniqueVisitorCount, + StatusCodes: normalizeCountMap(observabilityResp.StatusCodes), + TopDomains: normalizeCountMap(observabilityResp.TopDomains), + SourceCountries: normalizeCountMap(observabilityResp.SourceCountries), + }, + OpenrestyRxBytes: observabilityResp.OpenrestyRxBytes, + OpenrestyTxBytes: observabilityResp.OpenrestyTxBytes, + } + + if text, err := fetchLocalText(client, baseURL+openRestyStubStatusPath); err == nil { + result.OpenrestyConnections = parseStubStatusActiveConnections(text) + } + + return result +} + +func fetchLocalJSON(client *http.Client, url string, target any) error { + resp, err := client.Get(url) + if err != nil { + return err + } + defer resp.Body.Close() + if resp.StatusCode != http.StatusOK { + return fmt.Errorf("unexpected local observability status: %s", resp.Status) + } + return json.NewDecoder(resp.Body).Decode(target) +} + +func fetchLocalText(client *http.Client, url string) (string, error) { + resp, err := client.Get(url) + if err != nil { + return "", err + } + defer resp.Body.Close() + if resp.StatusCode != http.StatusOK { + return "", fmt.Errorf("unexpected local stub status: %s", resp.Status) + } + data, err := io.ReadAll(resp.Body) + if err != nil { + return "", err + } + return string(data), nil +} + +func parseStubStatusActiveConnections(raw string) int64 { + matches := stubStatusActivePattern.FindStringSubmatch(raw) + if len(matches) != 2 { + return 0 + } + value, err := strconv.ParseInt(matches[1], 10, 64) + if err != nil { + return 0 + } + return value +} + +func normalizeCountMap(values map[string]int64) map[string]int64 { + if len(values) == 0 { + return map[string]int64{} + } + result := make(map[string]int64, len(values)) + for key, value := range values { + key = strings.TrimSpace(key) + if key == "" || value <= 0 { + continue + } + result[key] = value + } + return result +} diff --git a/atsf_agent/internal/observability/openresty_local_test.go b/atsf_agent/internal/observability/openresty_local_test.go new file mode 100644 index 00000000..85109681 --- /dev/null +++ b/atsf_agent/internal/observability/openresty_local_test.go @@ -0,0 +1,82 @@ +package observability + +import ( + "net" + "net/http" + "net/http/httptest" + "strings" + "testing" + + "atsflare-agent/internal/config" +) + +func TestCollectManagedOpenRestyMetrics(t *testing.T) { + listener, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("Listen failed: %v", err) + } + port := listener.Addr().(*net.TCPAddr).Port + + mux := http.NewServeMux() + mux.HandleFunc(openRestyObservabilityPath, func(writer http.ResponseWriter, request *http.Request) { + writer.Header().Set("Content-Type", "application/json") + _, _ = writer.Write([]byte(`{"window_started_at_unix":1710403200,"window_ended_at_unix":1710403210,"request_count":12,"error_count":2,"unique_visitor_count":5,"status_codes":{"200":10,"502":2},"top_domains":{"app.example.com":9,"api.example.com":3},"source_countries":{},"openresty_rx_bytes":4096,"openresty_tx_bytes":8192}`)) + }) + mux.HandleFunc(openRestyStubStatusPath, func(writer http.ResponseWriter, request *http.Request) { + _, _ = writer.Write([]byte("Active connections: 7 \nserver accepts handled requests\n 10 10 12 \nReading: 1 Writing: 2 Waiting: 4 \n")) + }) + + server := httptest.NewUnstartedServer(mux) + server.Listener = listener + server.Start() + defer server.Close() + + metrics := CollectManagedOpenRestyMetrics(&config.Config{ + OpenrestyObservabilityPort: port, + }) + if metrics == nil || metrics.TrafficReport == nil { + t.Fatalf("expected managed openresty metrics, got %+v", metrics) + } + if metrics.TrafficReport.RequestCount != 12 || metrics.TrafficReport.ErrorCount != 2 { + t.Fatalf("unexpected traffic report: %+v", metrics.TrafficReport) + } + if metrics.OpenrestyRxBytes != 4096 || metrics.OpenrestyTxBytes != 8192 { + t.Fatalf("unexpected openresty byte counters: %+v", metrics) + } + if metrics.OpenrestyConnections != 7 { + t.Fatalf("unexpected openresty connections: %+v", metrics) + } +} + +func TestParseStubStatusActiveConnections(t *testing.T) { + if value := parseStubStatusActiveConnections("Active connections: 19\n"); value != 19 { + t.Fatalf("unexpected active connections: %d", value) + } +} + +func TestNormalizeCountMapDropsEmptyKeys(t *testing.T) { + normalized := normalizeCountMap(map[string]int64{ + "": 4, + " 200 ": 3, + "app.example.com": 0, + }) + if len(normalized) != 1 || normalized["200"] != 3 { + t.Fatalf("unexpected normalized map: %+v", normalized) + } +} + +func TestCollectManagedOpenRestyMetricsHandlesUnavailableEndpoint(t *testing.T) { + cfg := &config.Config{OpenrestyObservabilityPort: 1} + if metrics := CollectManagedOpenRestyMetrics(cfg); metrics != nil { + t.Fatalf("expected nil metrics for unavailable endpoint, got %+v", metrics) + } +} + +func TestOpenRestyObservabilityPathsAreStable(t *testing.T) { + if !strings.HasPrefix(openRestyObservabilityPath, "/atsflare/") { + t.Fatalf("unexpected observability path: %s", openRestyObservabilityPath) + } + if !strings.HasPrefix(openRestyStubStatusPath, "/atsflare/") { + t.Fatalf("unexpected stub status path: %s", openRestyStubStatusPath) + } +} diff --git a/atsf_agent/internal/observability/traffic.go b/atsf_agent/internal/observability/traffic.go index 49dfd0eb..c69e7fdb 100644 --- a/atsf_agent/internal/observability/traffic.go +++ b/atsf_agent/internal/observability/traffic.go @@ -10,6 +10,7 @@ import ( "io" "os" "path/filepath" + "regexp" "sort" "strconv" "strings" @@ -23,6 +24,8 @@ type accessLogRecord struct { Status int `json:"status"` } +var combinedAccessLogPattern = regexp.MustCompile(`^(\S+)\s+\S+\s+\S+\s+\[([^\]]+)\]\s+"[^"]*"\s+(\d{3})\s+\S+`) + type trafficAggregate struct { windowStartedAt time.Time windowEndedAt time.Time @@ -33,7 +36,10 @@ type trafficAggregate struct { visitors map[string]struct{} } -func BuildTrafficReport(cfg *config.Config, stateStore *state.Store) *protocol.NodeTrafficReport { +func BuildTrafficReport(cfg *config.Config, stateStore *state.Store, managed *managedOpenRestyMetrics) *protocol.NodeTrafficReport { + if managed != nil && managed.TrafficReport != nil && managed.TrafficReport.RequestCount > 0 { + return managed.TrafficReport + } if cfg == nil || stateStore == nil { return nil } @@ -115,21 +121,16 @@ func (aggregate *trafficAggregate) consume(line []byte) { return } - var record accessLogRecord - if err := json.Unmarshal([]byte(trimmed), &record); err != nil { + record, ok := parseAccessLogRecord(trimmed) + if !ok { return } - timestamp, err := parseAccessLogTime(record.Timestamp) - if err != nil { - return + if aggregate.windowStartedAt.IsZero() || record.Timestamp.Before(aggregate.windowStartedAt) { + aggregate.windowStartedAt = record.Timestamp } - - if aggregate.windowStartedAt.IsZero() || timestamp.Before(aggregate.windowStartedAt) { - aggregate.windowStartedAt = timestamp - } - if aggregate.windowEndedAt.IsZero() || timestamp.After(aggregate.windowEndedAt) { - aggregate.windowEndedAt = timestamp + if aggregate.windowEndedAt.IsZero() || record.Timestamp.After(aggregate.windowEndedAt) { + aggregate.windowEndedAt = record.Timestamp } aggregate.requestCount++ @@ -147,6 +148,58 @@ func (aggregate *trafficAggregate) consume(line []byte) { } } +type parsedAccessLogRecord struct { + Timestamp time.Time + Host string + RemoteAddr string + Status int +} + +func parseAccessLogRecord(raw string) (parsedAccessLogRecord, bool) { + record, ok := parseJSONAccessLogRecord(raw) + if ok { + return record, true + } + return parseCombinedAccessLogRecord(raw) +} + +func parseJSONAccessLogRecord(raw string) (parsedAccessLogRecord, bool) { + var record accessLogRecord + if err := json.Unmarshal([]byte(raw), &record); err != nil { + return parsedAccessLogRecord{}, false + } + timestamp, err := parseAccessLogTime(record.Timestamp) + if err != nil { + return parsedAccessLogRecord{}, false + } + return parsedAccessLogRecord{ + Timestamp: timestamp, + Host: strings.TrimSpace(record.Host), + RemoteAddr: strings.TrimSpace(record.RemoteAddr), + Status: record.Status, + }, true +} + +func parseCombinedAccessLogRecord(raw string) (parsedAccessLogRecord, bool) { + matches := combinedAccessLogPattern.FindStringSubmatch(raw) + if len(matches) != 4 { + return parsedAccessLogRecord{}, false + } + timestamp, err := parseAccessLogTime(matches[2]) + if err != nil { + return parsedAccessLogRecord{}, false + } + status, err := strconv.Atoi(matches[3]) + if err != nil { + return parsedAccessLogRecord{}, false + } + return parsedAccessLogRecord{ + Timestamp: timestamp, + RemoteAddr: strings.TrimSpace(matches[1]), + Status: status, + }, true +} + func (aggregate *trafficAggregate) report() *protocol.NodeTrafficReport { if aggregate.requestCount == 0 || aggregate.windowStartedAt.IsZero() || aggregate.windowEndedAt.IsZero() { return nil @@ -165,7 +218,15 @@ func (aggregate *trafficAggregate) report() *protocol.NodeTrafficReport { } func parseAccessLogTime(value string) (time.Time, error) { - return time.Parse(time.RFC3339, strings.TrimSpace(value)) + trimmed := strings.TrimSpace(value) + if trimmed == "" { + return time.Time{}, errors.New("empty access log time") + } + timestamp, err := time.Parse(time.RFC3339, trimmed) + if err == nil { + return timestamp, nil + } + return time.Parse("02/Jan/2006:15:04:05 -0700", trimmed) } func cloneTrafficCounts(values map[string]int64, limit int) map[string]int64 { diff --git a/atsf_agent/internal/observability/traffic_test.go b/atsf_agent/internal/observability/traffic_test.go index 8f7b4987..447e0eff 100644 --- a/atsf_agent/internal/observability/traffic_test.go +++ b/atsf_agent/internal/observability/traffic_test.go @@ -26,7 +26,7 @@ func TestBuildTrafficReportAggregatesManagedAccessLog(t *testing.T) { } stateStore := state.NewStore(filepath.Join(tempDir, "state.json")) - report := BuildTrafficReport(&config.Config{RouteConfigPath: routeConfigPath}, stateStore) + report := BuildTrafficReport(&config.Config{RouteConfigPath: routeConfigPath}, stateStore, nil) if report == nil { t.Fatal("expected traffic report") } @@ -48,7 +48,7 @@ func TestBuildTrafficReportAggregatesManagedAccessLog(t *testing.T) { t.Fatalf("unexpected access log offset: %d", snapshot.AccessLogOffset) } - secondReport := BuildTrafficReport(&config.Config{RouteConfigPath: routeConfigPath}, stateStore) + secondReport := BuildTrafficReport(&config.Config{RouteConfigPath: routeConfigPath}, stateStore, nil) if secondReport != nil { t.Fatalf("expected no report without appended lines, got %+v", secondReport) } @@ -70,8 +70,40 @@ func TestBuildTrafficReportResetsOffsetAfterTruncate(t *testing.T) { t.Fatalf("Save failed: %v", err) } - report := BuildTrafficReport(&config.Config{RouteConfigPath: routeConfigPath}, stateStore) + report := BuildTrafficReport(&config.Config{RouteConfigPath: routeConfigPath}, stateStore, nil) if report == nil || report.RequestCount != 1 { t.Fatalf("expected one request after truncate reset, got %+v", report) } } + +func TestBuildTrafficReportParsesCombinedAccessLog(t *testing.T) { + tempDir := t.TempDir() + routeConfigPath := filepath.Join(tempDir, "conf.d", "atsflare_routes.conf") + if err := os.MkdirAll(filepath.Dir(routeConfigPath), 0o755); err != nil { + t.Fatalf("MkdirAll failed: %v", err) + } + logPath := filepath.Join(filepath.Dir(routeConfigPath), "atsflare_access.log") + content := []byte( + "10.0.0.1 - - [14/Mar/2026:08:00:00 +0000] \"GET / HTTP/1.1\" 200 123 \"-\" \"curl/8.0\"\n" + + "10.0.0.2 - - [14/Mar/2026:08:00:05 +0000] \"GET /healthz HTTP/1.1\" 502 64 \"-\" \"curl/8.0\"\n" + + "10.0.0.1 - - [14/Mar/2026:08:00:10 +0000] \"GET /api HTTP/1.1\" 200 256 \"-\" \"curl/8.0\"\n", + ) + if err := os.WriteFile(logPath, content, 0o644); err != nil { + t.Fatalf("WriteFile failed: %v", err) + } + + stateStore := state.NewStore(filepath.Join(tempDir, "state.json")) + report := BuildTrafficReport(&config.Config{RouteConfigPath: routeConfigPath}, stateStore, nil) + if report == nil { + t.Fatal("expected traffic report from combined access log") + } + if report.RequestCount != 3 || report.ErrorCount != 1 || report.UniqueVisitorCount != 2 { + t.Fatalf("unexpected combined log counters: %+v", report) + } + if report.StatusCodes["200"] != 2 || report.StatusCodes["502"] != 1 { + t.Fatalf("unexpected combined log status codes: %+v", report.StatusCodes) + } + if len(report.TopDomains) != 0 { + t.Fatalf("expected combined access log to omit top domains when host is unavailable, got %+v", report.TopDomains) + } +} diff --git a/atsf_server/service/config_version.go b/atsf_server/service/config_version.go index 66aa066e..5da47bfc 100644 --- a/atsf_server/service/config_version.go +++ b/atsf_server/service/config_version.go @@ -115,9 +115,11 @@ type configBundle struct { } const ( - nginxCertDirPlaceholder = "__ATSF_CERT_DIR__" - nginxRouteConfigPlaceholder = "__ATSF_ROUTE_CONFIG__" - nginxAccessLogPlaceholder = "__ATSF_ACCESS_LOG__" + nginxCertDirPlaceholder = "__ATSF_CERT_DIR__" + nginxRouteConfigPlaceholder = "__ATSF_ROUTE_CONFIG__" + nginxAccessLogPlaceholder = "__ATSF_ACCESS_LOG__" + nginxLuaDirPlaceholder = "__ATSF_LUA_DIR__" + nginxObservabilityPortPlaceholder = "__ATSF_OBSERVABILITY_PORT__" ) var requiredMainConfigTemplatePlaceholders = []string{ @@ -351,6 +353,7 @@ func buildCurrentConfigBundle(requireRoutes bool) (*configBundle, error) { if err != nil { return nil, err } + supportFiles = append(supportFiles, buildOpenRestyObservabilitySupportFiles()...) mainConfig := renderMainConfig(openRestyConfig) return &configBundle{ Routes: routes, @@ -675,17 +678,21 @@ func renderTemplateDirective(enabled bool, statement string) string { } func renderOpenRestyCacheTemplateBlock(cfg openRestyConfigSnapshot) string { + lines := make([]string, 0, 8) if !cfg.CacheEnabled { - return "" + lines = append(lines, renderOpenRestyObservabilityTemplateBlock()) + return strings.Join(lines, "") } - return strings.Join([]string{ + lines = append(lines, strings.Join([]string{ fmt.Sprintf(" proxy_cache_path %s levels=%s keys_zone=atsflare_cache:10m inactive=%s max_size=%s;", cfg.CachePath, cfg.CacheLevels, cfg.CacheInactive, cfg.CacheMaxSize), fmt.Sprintf(" proxy_cache_key \"%s\";", cfg.CacheKeyTemplate), fmt.Sprintf(" proxy_cache_lock %s;", onOff(cfg.CacheLockEnabled)), fmt.Sprintf(" proxy_cache_lock_timeout %s;", cfg.CacheLockTimeout), fmt.Sprintf(" proxy_cache_use_stale %s;", cfg.CacheUseStale), "", - }, "\n") + }, "\n")) + lines = append(lines, renderOpenRestyObservabilityTemplateBlock()) + return strings.Join(lines, "") } func onOff(value bool) string { diff --git a/atsf_server/service/https_phase1_test.go b/atsf_server/service/https_phase1_test.go index 1e0c788d..0cbebacc 100644 --- a/atsf_server/service/https_phase1_test.go +++ b/atsf_server/service/https_phase1_test.go @@ -57,6 +57,12 @@ func TestCreateTLSCertificateAndRenderHTTPSConfig(t *testing.T) { if !strings.Contains(result.Version.MainConfig, "access_log __ATSF_ACCESS_LOG__ atsflare_json;") { 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;") { + 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.RenderedConfig, "listen 443 ssl;") { t.Fatal("expected rendered config to include https server block") } @@ -69,6 +75,9 @@ func TestCreateTLSCertificateAndRenderHTTPSConfig(t *testing.T) { if !strings.Contains(result.Version.SupportFilesJSON, ".crt") || !strings.Contains(result.Version.SupportFilesJSON, ".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") + } } func TestCreateProxyRouteRejectsHTTPSWithoutCertificate(t *testing.T) { @@ -187,6 +196,9 @@ func TestPreviewAndDiffConfigVersion(t *testing.T) { if !strings.Contains(preview.MainConfig, "include __ATSF_ROUTE_CONFIG__;") { 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;") { + t.Fatal("expected preview main config to include managed openresty lua log hook") + } if !strings.Contains(preview.RenderedConfig, `proxy_set_header X-Release "candidate";`) { t.Fatal("expected preview config to include modified custom header") } diff --git a/atsf_server/service/openresty_observability_assets.go b/atsf_server/service/openresty_observability_assets.go new file mode 100644 index 00000000..cc99a0c0 --- /dev/null +++ b/atsf_server/service/openresty_observability_assets.go @@ -0,0 +1,212 @@ +package service + +import "fmt" + +const ( + openRestyObservabilitySupportDir = "observability" + openRestyObservabilityInitLuaPath = openRestyObservabilitySupportDir + "/init.lua" + openRestyObservabilityLogLuaPath = openRestyObservabilitySupportDir + "/log.lua" + openRestyObservabilityReadLuaPath = openRestyObservabilitySupportDir + "/read.lua" + openRestyObservabilityWindowTTL = 7200 +) + +const openRestyObservabilityInitLua = `local dict = ngx.shared.atsflare_observability +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 +` + +const openRestyObservabilityLogLua = `local dict = ngx.shared.atsflare_observability +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) +end + +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(current_window) +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 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 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 = read_counter("window_started_at:" .. current_window), + 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) +} + +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)) +` + +func buildOpenRestyObservabilitySupportFiles() []SupportFile { + return []SupportFile{ + {Path: openRestyObservabilityInitLuaPath, Content: openRestyObservabilityInitLua}, + {Path: openRestyObservabilityLogLuaPath, Content: openRestyObservabilityLogLua}, + {Path: openRestyObservabilityReadLuaPath, Content: openRestyObservabilityReadLua}, + } +} + +func renderOpenRestyObservabilityTemplateBlock() string { + return stringsJoinLines( + " lua_shared_dict atsflare_observability 10m;", + fmt.Sprintf(" init_worker_by_lua_file %s/%s;", nginxLuaDirPlaceholder, openRestyObservabilityInitLuaPath), + fmt.Sprintf(" log_by_lua_file %s/%s;", nginxLuaDirPlaceholder, openRestyObservabilityLogLuaPath), + "", + fmt.Sprintf(" server {"), + fmt.Sprintf(" listen 127.0.0.1:%s;", nginxObservabilityPortPlaceholder), + " server_name atsflare-observability;", + " access_log off;", + " allow 127.0.0.1;", + " deny all;", + "", + " location = /atsflare/observability {", + " default_type application/json;", + fmt.Sprintf(" content_by_lua_file %s/%s;", nginxLuaDirPlaceholder, openRestyObservabilityReadLuaPath), + " }", + "", + " location = /atsflare/stub_status {", + " stub_status;", + " }", + " }", + "", + ) +} + +func stringsJoinLines(lines ...string) string { + if len(lines) == 0 { + return "" + } + result := "" + for index, line := range lines { + if index > 0 { + result += "\n" + } + result += line + } + return result + "\n" +} diff --git a/atsf_server/web/features/nodes/components/node-detail-page.tsx b/atsf_server/web/features/nodes/components/node-detail-page.tsx index 334b9976..e6e9b9d5 100644 --- a/atsf_server/web/features/nodes/components/node-detail-page.tsx +++ b/atsf_server/web/features/nodes/components/node-detail-page.tsx @@ -1350,7 +1350,6 @@ export function NodeDetailPage({ nodeId }: { nodeId: string }) {

{node.latest_apply_checksum ? (
-

目标 Checksum:{node.latest_apply_checksum}

支持文件:{node.latest_support_file_count}

) : null} diff --git a/docs/app-config.md b/docs/app-config.md index f99763ee..e54c0525 100644 --- a/docs/app-config.md +++ b/docs/app-config.md @@ -235,12 +235,13 @@ go run ./cmd/agent -config ./agent.json { "server_url": "http://127.0.0.1:3000", "discovery_token": "replace-with-global-discovery-token", - "data_dir": "./data", - "openresty_container_name": "atsflare-openresty", - "openresty_docker_image": "openresty/openresty:alpine", - "heartbeat_interval": 10000, - "request_timeout": 10000 -} + "data_dir": "./data", + "openresty_container_name": "atsflare-openresty", + "openresty_docker_image": "openresty/openresty:alpine", + "openresty_observability_port": 18081, + "heartbeat_interval": 10000, + "request_timeout": 10000 +} ``` 使用节点专属 Token 的示例: @@ -255,11 +256,12 @@ go run ./cmd/agent -config ./agent.json "openresty_path": "/usr/local/openresty/nginx/sbin/openresty", "main_config_path": "/usr/local/openresty/nginx/conf/nginx.conf", "route_config_path": "/usr/local/openresty/nginx/conf/conf.d/atsflare_routes.conf", - "cert_dir": "/usr/local/openresty/nginx/conf/certs", - "openresty_cert_dir": "/usr/local/openresty/nginx/conf/certs", - "state_path": "./data/agent-state.json", - "heartbeat_interval": 10000, - "request_timeout": 10000 + "cert_dir": "/usr/local/openresty/nginx/conf/certs", + "openresty_cert_dir": "/usr/local/openresty/nginx/conf/certs", + "openresty_observability_port": 18081, + "state_path": "./data/agent-state.json", + "heartbeat_interval": 10000, + "request_timeout": 10000 } ``` @@ -274,8 +276,9 @@ go run ./cmd/agent -config ./agent.json | `node_ip` | 节点 IP | 否 | 自动探测第一个可用 IPv4 | `192.168.1.20` | | `openresty_path` | 本机 OpenResty 可执行文件路径;设置后按本机 OpenResty 模式运行 | 否 | 空;未设置时按 Docker OpenResty 模式处理 | `/usr/local/openresty/nginx/sbin/openresty` | | `openresty_container_name` | Docker 模式下的 OpenResty 容器名 | 否 | `atsflare-openresty` | `atsflare-openresty` | -| `openresty_docker_image` | Docker 模式下用于初始化/管理的 OpenResty 镜像 | 否 | `openresty/openresty:alpine` | `openresty/openresty:alpine` | -| `docker_binary` | Docker 可执行文件名或路径 | 否 | `docker` | `/usr/bin/docker` | +| `openresty_docker_image` | Docker 模式下用于初始化/管理的 OpenResty 镜像 | 否 | `openresty/openresty:alpine` | `openresty/openresty:alpine` | +| `openresty_observability_port` | Agent 注入的 OpenResty 本地观测端口;用于 heartbeat 前读取 Lua 窗口指标和 `stub_status`,默认仅监听 `127.0.0.1` | 否 | `18081` | `18081` | +| `docker_binary` | Docker 可执行文件名或路径 | 否 | `docker` | `/usr/bin/docker` | | `data_dir` | Agent 数据目录,用于存储托管配置、证书和状态文件 | 否 | 配置文件所在目录下的 `data` 子目录 | `./data` | | `main_config_path` | 第五版主配置接管时 OpenResty 主配置文件写入路径 | 第五版本机模式建议必填 | Docker 模式可使用受管默认路径;本机模式建议显式设置 | `/usr/local/openresty/nginx/conf/nginx.conf` | | `route_config_path` | 路由配置文件写入路径 | 否 | 默认为 `data_dir` 下托管路径 | `/etc/nginx/conf.d/atsflare_routes.conf` | @@ -291,9 +294,10 @@ go run ./cmd/agent -config ./agent.json * `heartbeat_interval`、`request_timeout` 支持两种写法: * 毫秒整数,例如 `10000` * Go duration 字符串,例如 `"30s"` -* `node_name` 与 `node_ip` 未填写时会自动探测;若自动探测失败,配置校验会报错 -* 未配置 `openresty_path` 时,默认为 Docker OpenResty 模式 -* 配置保存时,`agent_version`、`nginx_version` 由程序运行时维护,不需要写入 JSON +* `node_name` 与 `node_ip` 未填写时会自动探测;若自动探测失败,配置校验会报错 +* 未配置 `openresty_path` 时,默认为 Docker OpenResty 模式 +* `openresty_observability_port` 默认仅绑定本地回环地址;若节点本机已有端口冲突,可改为其他未占用端口 +* 配置保存时,`agent_version`、`nginx_version` 由程序运行时维护,不需要写入 JSON * 第五版主配置接管完成后,本机模式下应优先通过 `main_config_path` 由 Agent 写入受管主配置,而不是依赖节点手工维护 include 规则 ### 2.4 Agent 托管路径默认值 @@ -311,7 +315,12 @@ Docker OpenResty 模式下: | 字段 | 默认值 | | --- | --- | -| `openresty_cert_dir` | `/etc/nginx/atsflare-certs` | +| `openresty_cert_dir` | `/etc/nginx/atsflare-certs` | + +补充说明: + +* Agent 当前会随受管配置一并向 OpenResty 注入 Lua 观测脚本,并在每次 heartbeat 前通过 `http://127.0.0.1:/atsflare/observability` 读取最近窗口请求指标 +* 同一端口还会暴露仅本机可访问的 `stub_status`,用于采集 OpenResty 活动连接数 ### 2.5 Agent 启动示例 diff --git a/docs/deployment.md b/docs/deployment.md index f84bf9ab..58bbd739 100644 --- a/docs/deployment.md +++ b/docs/deployment.md @@ -154,6 +154,7 @@ swag init -g main.go -o docs "data_dir": "./data", "openresty_container_name": "atsflare-openresty", "openresty_docker_image": "openresty/openresty:alpine", + "openresty_observability_port": 18081, "heartbeat_interval": 10000, "request_timeout": 10000 } @@ -168,6 +169,7 @@ swag init -g main.go -o docs "data_dir": "./data", "openresty_container_name": "atsflare-openresty", "openresty_docker_image": "openresty/openresty:alpine", + "openresty_observability_port": 18081, "heartbeat_interval": 10000, "request_timeout": 10000 } @@ -178,19 +180,21 @@ swag init -g main.go -o docs * `agent_version` 由 Agent 代码内常量提供,升级时同步修改代码 * 为兼容现有 Agent / Server API,运行时版本仍通过 `openresty_version` 字段上报,但其值现在表示 OpenResty 版本 * 时间字段使用毫秒整数 -* `agent_token` 与 `discovery_token` 至少填写一个 -* 若 `agent_token` 为空且 `discovery_token` 存在,Agent 会自动注册并写回新的专属 `agent_token` -* `node_name` 与 `node_ip` 可省略,未填写时自动探测 +* `agent_token` 与 `discovery_token` 至少填写一个 +* 若 `agent_token` 为空且 `discovery_token` 存在,Agent 会自动注册并写回新的专属 `agent_token` +* `node_name` 与 `node_ip` 可省略,未填写时自动探测 * 未配置 `openresty_path` 时,默认使用 Docker OpenResty 容器 +* Agent 会在受管 OpenResty 中注入 Lua 观测脚本,并通过 `openresty_observability_port` 对本机暴露最近窗口指标与 `stub_status` ### 3.3 第五版新增部署约束 第五版开发完成后,OpenResty 主配置将进入 Agent 受管范围。部署与联调时应满足: -* 本机 OpenResty 模式需要为 Agent 显式提供主配置文件写入路径 -* Docker OpenResty 模式需要保证主配置、路由配置和证书目录位于同一套受管挂载路径中 -* 节点现存手工维护的主配置如继续保留,必须先迁移为 Server 渲染模板的等价配置,再切换到受管模式 -* 主配置切换前必须预留回滚副本,并通过一次 `openresty -t` 失败演练验证回滚 +* 本机 OpenResty 模式需要为 Agent 显式提供主配置文件写入路径 +* Docker OpenResty 模式需要保证主配置、路由配置和证书目录位于同一套受管挂载路径中 +* Docker OpenResty 模式会额外挂载一个仅本机可访问的 `127.0.0.1:` 观测端口,用于 Agent 在 heartbeat 前抓取 Lua 窗口指标 +* 节点现存手工维护的主配置如继续保留,必须先迁移为 Server 渲染模板的等价配置,再切换到受管模式 +* 主配置切换前必须预留回滚副本,并通过一次 `openresty -t` 失败演练验证回滚 --- @@ -252,7 +256,8 @@ export LOG_LEVEL='info' "discovery_token": "replace-with-global-discovery-token", "data_dir": "./data", "openresty_container_name": "atsflare-openresty", - "openresty_docker_image": "openresty/openresty:alpine" + "openresty_docker_image": "openresty/openresty:alpine", + "openresty_observability_port": 18081 } ``` @@ -260,6 +265,7 @@ export LOG_LEVEL='info' 1. 首次启动后确认 `data/etc/nginx/nginx.conf`、`data/etc/nginx/conf.d/atsflare_routes.conf` 与 `data/etc/nginx/certs` 已由 Agent 创建 2. 确认容器实际挂载了主配置、路由目录和证书目录 +3. 确认宿主机本地可访问 `http://127.0.0.1:18081/atsflare/observability` 与 `http://127.0.0.1:18081/atsflare/stub_status` 3. 在管理端发布一次新版本后,确认节点 `current_version` 追平激活版本 4. 在节点详情查看“当前目标版本”与“最近应用”,确认主配置/路由配置快照和 checksum 已可见 @@ -273,6 +279,7 @@ docker exec atsflare-openresty openresty -t 说明: * `docker inspect` 重点确认主配置文件、`conf.d` 目录和证书目录都来自 Agent 受管路径 +* 观测端口默认只绑定 `127.0.0.1`;若节点已有冲突,可在 `agent.json` 中调整 `openresty_observability_port` * 若容器名使用默认值,请将上述命令中的名称替换为 `atsflare-openresty` ### 5.3.2 本机 OpenResty 模式最小验证 @@ -287,7 +294,8 @@ docker exec atsflare-openresty openresty -t "main_config_path": "/usr/local/openresty/nginx/conf/nginx.conf", "route_config_path": "/usr/local/openresty/nginx/conf/conf.d/atsflare_routes.conf", "cert_dir": "/usr/local/openresty/nginx/conf/certs", - "openresty_cert_dir": "/usr/local/openresty/nginx/conf/certs" + "openresty_cert_dir": "/usr/local/openresty/nginx/conf/certs", + "openresty_observability_port": 18081 } ``` @@ -295,8 +303,9 @@ docker exec atsflare-openresty openresty -t 1. 发布前先备份 `main_config_path` 与 `route_config_path` 2. 首次发布后执行 `openresty -t`,确认主配置已由 Server 模板接管且 include 指向 Agent 写入的路由文件 -3. 再次发布修改后的规则或 OpenResty 参数,确认 `openresty -s reload` 成功且节点版本更新 -4. 在节点详情与应用记录页确认主配置 checksum、路由配置 checksum 和支持文件数已上报 +3. 确认本机可访问 `http://127.0.0.1:18081/atsflare/observability` 与 `http://127.0.0.1:18081/atsflare/stub_status` +4. 再次发布修改后的规则或 OpenResty 参数,确认 `openresty -s reload` 成功且节点版本更新 +5. 在节点详情与应用记录页确认主配置 checksum、路由配置 checksum 和支持文件数已上报 ### 5.4 验证管理端状态 diff --git a/docs/development-plan.md b/docs/development-plan.md index 45b86700..cd8d87c1 100644 --- a/docs/development-plan.md +++ b/docs/development-plan.md @@ -19,7 +19,7 @@ * 已完成 heartbeat 扩展协议与观测数据分层落地,节点已支持上报 `profile`、`snapshot`、`traffic_report`、`health_events` * 已完成 Server 侧观测模型、入库链路与查询接口,节点观测已拆分为系统画像、资源快照、窗口流量聚合、健康事件 -* 已完成 Agent 侧真实请求窗口聚合采集,当前通过受管 OpenResty 访问日志做增量读取与窗口聚合,不再只是预留服务端通道 +* 已完成 Agent 侧真实请求窗口聚合采集,当前通过受管 OpenResty Lua 观测脚本与本地指标端口输出窗口聚合结果,不再依赖访问日志增量读取 * 已完成节点详情观测接口与页面第一轮改造,当前已支持系统画像、实时资源、运行状态、24 小时趋势、状态码分布、Top Domain 与健康事件时间线 * 已完成首页总览专用聚合接口与首页第一轮改造,当前已支持系统运行总览、风险态势、峰值摘要、24 小时趋势、节点健康列表与活动异常 * 已完成首页风险态势到节点页的轻量筛选联动,支持从总览跳转到节点页查看离线节点、OpenResty 异常节点与配置落后节点