[优化] 增强流量可观测性,添加请求长度和字节发送字段;更新支持文件权限设置

This commit is contained in:
ryan
2026-03-15 12:13:07 +08:00
parent eb9a2a8814
commit e1efbf3868
5 changed files with 142 additions and 36 deletions
+4 -1
View File
@@ -320,8 +320,11 @@ func (r *Runner) nodePayload(nodeID string) protocol.NodePayload {
} }
profile := observability.BuildProfile(r.Config, r.StateStore) profile := observability.BuildProfile(r.Config, r.StateStore)
managedOpenRestyMetrics := observability.CollectManagedOpenRestyMetrics(r.Config) managedOpenRestyMetrics := observability.CollectManagedOpenRestyMetrics(r.Config)
trafficReport, accessLogs, fallbackMetrics := observability.BuildTrafficObservability(r.Config, r.StateStore, managedOpenRestyMetrics)
if managedOpenRestyMetrics == nil {
managedOpenRestyMetrics = fallbackMetrics
}
metricSnapshot := observability.BuildSnapshot(r.Config, r.StateStore, managedOpenRestyMetrics) metricSnapshot := observability.BuildSnapshot(r.Config, r.StateStore, managedOpenRestyMetrics)
trafficReport, accessLogs := observability.BuildTrafficObservability(r.Config, r.StateStore, managedOpenRestyMetrics)
healthEvents := observability.BuildHealthEvents(snapshot) healthEvents := observability.BuildHealthEvents(snapshot)
return protocol.NodePayload{ return protocol.NodePayload{
NodeID: nodeID, NodeID: nodeID,
+14 -2
View File
@@ -6,6 +6,7 @@ import (
"encoding/hex" "encoding/hex"
"errors" "errors"
"fmt" "fmt"
"io/fs"
"log/slog" "log/slog"
"os" "os"
"os/exec" "os/exec"
@@ -529,7 +530,7 @@ func (m *Manager) restore(state *backupState) error {
if err := os.MkdirAll(filepath.Dir(targetPath), 0o755); err != nil { if err := os.MkdirAll(filepath.Dir(targetPath), 0o755); err != nil {
return err return err
} }
if err := os.WriteFile(targetPath, []byte(file.Content), 0o600); err != nil { if err := os.WriteFile(targetPath, []byte(file.Content), supportFileMode(file.Path)); err != nil {
return err return err
} }
} }
@@ -554,7 +555,7 @@ func (m *Manager) writeSupportFiles(supportFiles []protocol.SupportFile) error {
if err := os.MkdirAll(filepath.Dir(targetPath), 0o755); err != nil { if err := os.MkdirAll(filepath.Dir(targetPath), 0o755); err != nil {
return err return err
} }
if err := os.WriteFile(targetPath, []byte(file.Content), 0o600); err != nil { if err := os.WriteFile(targetPath, []byte(file.Content), supportFileMode(file.Path)); err != nil {
return err return err
} }
} }
@@ -628,6 +629,17 @@ func (m *Manager) supportFileTargetPath(relativePath string) (string, error) {
return targetPath, nil return targetPath, nil
} }
func supportFileMode(relativePath string) fs.FileMode {
switch strings.ToLower(filepath.Ext(strings.TrimSpace(relativePath))) {
case ".lua", ".crt", ".pem":
return 0o644
case ".key":
return 0o600
default:
return 0o644
}
}
func (m *Manager) renderRouteConfig(content string) string { func (m *Manager) renderRouteConfig(content string) string {
if m.NginxSupportDir == "" { if m.NginxSupportDir == "" {
return content return content
+58
View File
@@ -497,6 +497,64 @@ func TestManagerApplyWritesSupportFilesAndReplacesPlaceholder(t *testing.T) {
if string(certData) != "cert-data" { if string(certData) != "cert-data" {
t.Fatalf("unexpected cert file content: %s", string(certData)) t.Fatalf("unexpected cert file content: %s", string(certData))
} }
luaInfo, err := os.Stat(filepath.Join(manager.SupportDir, "observability", "log.lua"))
if err == nil {
t.Fatalf("expected no lua file in this test, got %v", luaInfo)
}
}
func TestSupportFileMode(t *testing.T) {
testCases := []struct {
path string
want os.FileMode
}{
{path: "observability/log.lua", want: 0o644},
{path: "1.crt", want: 0o644},
{path: "1.pem", want: 0o644},
{path: "1.key", want: 0o600},
{path: "misc.txt", want: 0o644},
}
for _, testCase := range testCases {
if got := supportFileMode(testCase.path); got != testCase.want {
t.Fatalf("unexpected mode for %s: got %o want %o", testCase.path, got, testCase.want)
}
}
}
func TestManagerApplyWritesLuaSupportFilesReadable(t *testing.T) {
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{},
}
err := manager.Apply(context.Background(), "main", "route", []protocol.SupportFile{
{Path: "observability/log.lua", Content: "return"},
{Path: "1.key", Content: "secret"},
})
if err != nil {
t.Fatalf("Apply failed: %v", err)
}
luaInfo, err := os.Stat(filepath.Join(manager.SupportDir, "observability", "log.lua"))
if err != nil {
t.Fatalf("failed to stat lua file: %v", err)
}
if luaInfo.Mode().Perm() != 0o644 {
t.Fatalf("unexpected lua mode: %o", luaInfo.Mode().Perm())
}
keyInfo, err := os.Stat(filepath.Join(manager.SupportDir, "1.key"))
if err != nil {
t.Fatalf("failed to stat key file: %v", err)
}
if keyInfo.Mode().Perm() != 0o600 {
t.Fatalf("unexpected key mode: %o", keyInfo.Mode().Perm())
}
} }
func TestManagerRollbackRestoresSupportFiles(t *testing.T) { func TestManagerRollbackRestoresSupportFiles(t *testing.T) {
+60 -30
View File
@@ -18,37 +18,41 @@ import (
) )
type accessLogRecord struct { type accessLogRecord struct {
Timestamp string `json:"ts"` Timestamp string `json:"ts"`
Host string `json:"host"` Host string `json:"host"`
RemoteAddr string `json:"remote_addr"` RemoteAddr string `json:"remote_addr"`
Path string `json:"path"` Path string `json:"path"`
Status int `json:"status"` Status int `json:"status"`
BytesSent int64 `json:"bytes_sent"`
RequestLength int64 `json:"request_length"`
} }
var combinedAccessLogPattern = regexp.MustCompile(`^(\S+)\s+\S+\s+\S+\s+\[([^\]]+)\]\s+"(?:\S+)\s+(\S+)(?:\s+[^"]*)?"\s+(\d{3})\s+\S+`) var combinedAccessLogPattern = regexp.MustCompile(`^(\S+)\s+\S+\s+\S+\s+\[([^\]]+)\]\s+"(?:\S+)\s+(\S+)(?:\s+[^"]*)?"\s+(\d{3})\s+\S+`)
type trafficAggregate struct { type trafficAggregate struct {
windowStartedAt time.Time windowStartedAt time.Time
windowEndedAt time.Time windowEndedAt time.Time
requestCount int64 requestCount int64
errorCount int64 errorCount int64
statusCodes map[string]int64 openrestyRxBytes int64
topDomains map[string]int64 openrestyTxBytes int64
visitors map[string]struct{} statusCodes map[string]int64
logs []protocol.NodeAccessLog topDomains map[string]int64
visitors map[string]struct{}
logs []protocol.NodeAccessLog
} }
func BuildTrafficReport(cfg *config.Config, stateStore *state.Store, managed *managedOpenRestyMetrics) *protocol.NodeTrafficReport { func BuildTrafficReport(cfg *config.Config, stateStore *state.Store, managed *managedOpenRestyMetrics) *protocol.NodeTrafficReport {
report, _ := BuildTrafficObservability(cfg, stateStore, managed) report, _, _ := BuildTrafficObservability(cfg, stateStore, managed)
return report return report
} }
func BuildTrafficObservability(cfg *config.Config, stateStore *state.Store, managed *managedOpenRestyMetrics) (*protocol.NodeTrafficReport, []protocol.NodeAccessLog) { func BuildTrafficObservability(cfg *config.Config, stateStore *state.Store, managed *managedOpenRestyMetrics) (*protocol.NodeTrafficReport, []protocol.NodeAccessLog, *managedOpenRestyMetrics) {
if cfg == nil || stateStore == nil { if cfg == nil || stateStore == nil {
if managed != nil && managed.TrafficReport != nil { if managed != nil && managed.TrafficReport != nil {
return managed.TrafficReport, nil return managed.TrafficReport, nil, managed
} }
return nil, nil return nil, nil, managed
} }
aggregate := readAccessLogDelta(cfg, stateStore) aggregate := readAccessLogDelta(cfg, stateStore)
@@ -57,12 +61,13 @@ func BuildTrafficObservability(cfg *config.Config, stateStore *state.Store, mana
accessLogs = aggregate.accessLogs() accessLogs = aggregate.accessLogs()
} }
if managed != nil && managed.TrafficReport != nil { if managed != nil && managed.TrafficReport != nil {
return managed.TrafficReport, accessLogs return managed.TrafficReport, accessLogs, managed
} }
if aggregate == nil { if aggregate == nil {
return nil, accessLogs return nil, accessLogs, managed
} }
return aggregate.report(), accessLogs fallbackManaged := aggregate.managedMetrics()
return aggregate.report(), accessLogs, fallbackManaged
} }
func readAccessLogDelta(cfg *config.Config, stateStore *state.Store) *trafficAggregate { func readAccessLogDelta(cfg *config.Config, stateStore *state.Store) *trafficAggregate {
@@ -162,6 +167,12 @@ func (aggregate *trafficAggregate) consume(line []byte) {
if record.Status > 0 { if record.Status > 0 {
aggregate.statusCodes[strconv.Itoa(record.Status)]++ aggregate.statusCodes[strconv.Itoa(record.Status)]++
} }
if record.RequestLength > 0 {
aggregate.openrestyRxBytes += record.RequestLength
}
if record.BytesSent > 0 {
aggregate.openrestyTxBytes += record.BytesSent
}
if host := strings.TrimSpace(record.Host); host != "" { if host := strings.TrimSpace(record.Host); host != "" {
aggregate.topDomains[host]++ aggregate.topDomains[host]++
} }
@@ -178,11 +189,13 @@ func (aggregate *trafficAggregate) consume(line []byte) {
} }
type parsedAccessLogRecord struct { type parsedAccessLogRecord struct {
Timestamp time.Time Timestamp time.Time
Host string Host string
RemoteAddr string RemoteAddr string
Path string Path string
Status int Status int
BytesSent int64
RequestLength int64
} }
func parseAccessLogRecord(raw string) (parsedAccessLogRecord, bool) { func parseAccessLogRecord(raw string) (parsedAccessLogRecord, bool) {
@@ -203,11 +216,13 @@ func parseJSONAccessLogRecord(raw string) (parsedAccessLogRecord, bool) {
return parsedAccessLogRecord{}, false return parsedAccessLogRecord{}, false
} }
return parsedAccessLogRecord{ return parsedAccessLogRecord{
Timestamp: timestamp, Timestamp: timestamp,
Host: strings.TrimSpace(record.Host), Host: strings.TrimSpace(record.Host),
RemoteAddr: strings.TrimSpace(record.RemoteAddr), RemoteAddr: strings.TrimSpace(record.RemoteAddr),
Path: normalizeAccessLogPath(record.Path), Path: normalizeAccessLogPath(record.Path),
Status: record.Status, Status: record.Status,
BytesSent: record.BytesSent,
RequestLength: record.RequestLength,
}, true }, true
} }
@@ -256,6 +271,21 @@ func (aggregate *trafficAggregate) accessLogs() []protocol.NodeAccessLog {
return append([]protocol.NodeAccessLog(nil), aggregate.logs...) return append([]protocol.NodeAccessLog(nil), aggregate.logs...)
} }
func (aggregate *trafficAggregate) managedMetrics() *managedOpenRestyMetrics {
if aggregate == nil {
return nil
}
report := aggregate.report()
if report == nil && aggregate.openrestyRxBytes <= 0 && aggregate.openrestyTxBytes <= 0 {
return nil
}
return &managedOpenRestyMetrics{
TrafficReport: report,
OpenrestyRxBytes: aggregate.openrestyRxBytes,
OpenrestyTxBytes: aggregate.openrestyTxBytes,
}
}
func parseAccessLogTime(value string) (time.Time, error) { func parseAccessLogTime(value string) (time.Time, error) {
trimmed := strings.TrimSpace(value) trimmed := strings.TrimSpace(value)
if trimmed == "" { if trimmed == "" {
@@ -85,21 +85,24 @@ func TestBuildTrafficObservabilityReturnsAccessLogs(t *testing.T) {
} }
logPath := filepath.Join(filepath.Dir(routeConfigPath), "atsflare_access.log") logPath := filepath.Join(filepath.Dir(routeConfigPath), "atsflare_access.log")
content := []byte( content := []byte(
"{\"ts\":\"2026-03-14T08:00:00Z\",\"host\":\"app.example.com\",\"path\":\"/login\",\"remote_addr\":\"10.0.0.1\",\"status\":200}\n" + "{\"ts\":\"2026-03-14T08:00:00Z\",\"host\":\"app.example.com\",\"path\":\"/login\",\"remote_addr\":\"10.0.0.1\",\"status\":200,\"request_length\":128,\"bytes_sent\":512}\n" +
"{\"ts\":\"2026-03-14T08:00:05Z\",\"host\":\"api.example.com\",\"path\":\"/v1/ping\",\"remote_addr\":\"10.0.0.2\",\"status\":502}\n", "{\"ts\":\"2026-03-14T08:00:05Z\",\"host\":\"api.example.com\",\"path\":\"/v1/ping\",\"remote_addr\":\"10.0.0.2\",\"status\":502,\"request_length\":64,\"bytes_sent\":256}\n",
) )
if err := os.WriteFile(logPath, content, 0o644); err != nil { if err := os.WriteFile(logPath, content, 0o644); err != nil {
t.Fatalf("WriteFile failed: %v", err) t.Fatalf("WriteFile failed: %v", err)
} }
stateStore := state.NewStore(filepath.Join(tempDir, "state.json")) stateStore := state.NewStore(filepath.Join(tempDir, "state.json"))
report, accessLogs := BuildTrafficObservability(&config.Config{RouteConfigPath: routeConfigPath}, stateStore, nil) report, accessLogs, fallbackMetrics := BuildTrafficObservability(&config.Config{RouteConfigPath: routeConfigPath}, stateStore, nil)
if report == nil || report.RequestCount != 2 { if report == nil || report.RequestCount != 2 {
t.Fatalf("expected traffic report, got %+v", report) t.Fatalf("expected traffic report, got %+v", report)
} }
if len(accessLogs) != 2 { if len(accessLogs) != 2 {
t.Fatalf("expected access logs, got %+v", accessLogs) t.Fatalf("expected access logs, got %+v", accessLogs)
} }
if fallbackMetrics == nil || fallbackMetrics.OpenrestyRxBytes != 192 || fallbackMetrics.OpenrestyTxBytes != 768 {
t.Fatalf("expected fallback throughput metrics, got %+v", fallbackMetrics)
}
if accessLogs[0].Path != "/login" || accessLogs[1].Path != "/v1/ping" { if accessLogs[0].Path != "/login" || accessLogs[1].Path != "/v1/ping" {
t.Fatalf("unexpected access log paths: %+v", accessLogs) t.Fatalf("unexpected access log paths: %+v", accessLogs)
} }