feat: enhance observability metrics collection and reporting

- Updated BuildSnapshot function to include OpenResty metrics.
- Enhanced BuildTrafficReport to utilize managed OpenResty metrics.
- Introduced new functions for parsing access logs in both JSON and combined formats.
- Added tests for traffic report generation from combined access logs.
- Implemented local observability metrics collection from OpenResty.
- Created Lua scripts for OpenResty to gather observability data.
- Updated configuration documentation to include new observability port and settings.
- Added support for OpenResty observability in the server configuration.
This commit is contained in:
ryan
2026-03-14 17:07:47 +08:00
parent 8cc839669c
commit bdfa80f214
19 changed files with 800 additions and 184 deletions
@@ -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
@@ -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
}
@@ -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)
}
}
+74 -13
View File
@@ -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 {
@@ -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)
}
}