From e1efbf3868aad34a1f8e3a1ab3580aeea01ed518 Mon Sep 17 00:00:00 2001 From: ryan Date: Sun, 15 Mar 2026 12:13:07 +0800 Subject: [PATCH] =?UTF-8?q?[=E4=BC=98=E5=8C=96]=20=E5=A2=9E=E5=BC=BA?= =?UTF-8?q?=E6=B5=81=E9=87=8F=E5=8F=AF=E8=A7=82=E6=B5=8B=E6=80=A7=EF=BC=8C?= =?UTF-8?q?=E6=B7=BB=E5=8A=A0=E8=AF=B7=E6=B1=82=E9=95=BF=E5=BA=A6=E5=92=8C?= =?UTF-8?q?=E5=AD=97=E8=8A=82=E5=8F=91=E9=80=81=E5=AD=97=E6=AE=B5=EF=BC=9B?= =?UTF-8?q?=E6=9B=B4=E6=96=B0=E6=94=AF=E6=8C=81=E6=96=87=E4=BB=B6=E6=9D=83?= =?UTF-8?q?=E9=99=90=E8=AE=BE=E7=BD=AE?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- atsf_agent/internal/agent/runner.go | 5 +- atsf_agent/internal/nginx/manager.go | 16 +++- atsf_agent/internal/nginx/manager_test.go | 58 ++++++++++++ atsf_agent/internal/observability/traffic.go | 90 ++++++++++++------- .../internal/observability/traffic_test.go | 9 +- 5 files changed, 142 insertions(+), 36 deletions(-) diff --git a/atsf_agent/internal/agent/runner.go b/atsf_agent/internal/agent/runner.go index 95276e27..00661984 100644 --- a/atsf_agent/internal/agent/runner.go +++ b/atsf_agent/internal/agent/runner.go @@ -320,8 +320,11 @@ func (r *Runner) nodePayload(nodeID string) protocol.NodePayload { } profile := observability.BuildProfile(r.Config, r.StateStore) 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) - trafficReport, accessLogs := observability.BuildTrafficObservability(r.Config, r.StateStore, managedOpenRestyMetrics) healthEvents := observability.BuildHealthEvents(snapshot) return protocol.NodePayload{ NodeID: nodeID, diff --git a/atsf_agent/internal/nginx/manager.go b/atsf_agent/internal/nginx/manager.go index 1f5ffe02..f7bebe43 100644 --- a/atsf_agent/internal/nginx/manager.go +++ b/atsf_agent/internal/nginx/manager.go @@ -6,6 +6,7 @@ import ( "encoding/hex" "errors" "fmt" + "io/fs" "log/slog" "os" "os/exec" @@ -529,7 +530,7 @@ func (m *Manager) restore(state *backupState) error { if err := os.MkdirAll(filepath.Dir(targetPath), 0o755); err != nil { 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 } } @@ -554,7 +555,7 @@ func (m *Manager) writeSupportFiles(supportFiles []protocol.SupportFile) error { if err := os.MkdirAll(filepath.Dir(targetPath), 0o755); err != nil { 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 } } @@ -628,6 +629,17 @@ func (m *Manager) supportFileTargetPath(relativePath string) (string, error) { 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 { if m.NginxSupportDir == "" { return content diff --git a/atsf_agent/internal/nginx/manager_test.go b/atsf_agent/internal/nginx/manager_test.go index 4051945c..b7b21e23 100644 --- a/atsf_agent/internal/nginx/manager_test.go +++ b/atsf_agent/internal/nginx/manager_test.go @@ -497,6 +497,64 @@ func TestManagerApplyWritesSupportFilesAndReplacesPlaceholder(t *testing.T) { if string(certData) != "cert-data" { 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) { diff --git a/atsf_agent/internal/observability/traffic.go b/atsf_agent/internal/observability/traffic.go index 975c3a38..8b7999b8 100644 --- a/atsf_agent/internal/observability/traffic.go +++ b/atsf_agent/internal/observability/traffic.go @@ -18,37 +18,41 @@ import ( ) type accessLogRecord struct { - Timestamp string `json:"ts"` - Host string `json:"host"` - RemoteAddr string `json:"remote_addr"` - Path string `json:"path"` - Status int `json:"status"` + Timestamp string `json:"ts"` + Host string `json:"host"` + RemoteAddr string `json:"remote_addr"` + Path string `json:"path"` + 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+`) type trafficAggregate struct { - windowStartedAt time.Time - windowEndedAt time.Time - requestCount int64 - errorCount int64 - statusCodes map[string]int64 - topDomains map[string]int64 - visitors map[string]struct{} - logs []protocol.NodeAccessLog + windowStartedAt time.Time + windowEndedAt time.Time + requestCount int64 + errorCount int64 + openrestyRxBytes int64 + openrestyTxBytes int64 + statusCodes map[string]int64 + topDomains map[string]int64 + visitors map[string]struct{} + logs []protocol.NodeAccessLog } func BuildTrafficReport(cfg *config.Config, stateStore *state.Store, managed *managedOpenRestyMetrics) *protocol.NodeTrafficReport { - report, _ := BuildTrafficObservability(cfg, stateStore, managed) + report, _, _ := BuildTrafficObservability(cfg, stateStore, managed) 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 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) @@ -57,12 +61,13 @@ func BuildTrafficObservability(cfg *config.Config, stateStore *state.Store, mana accessLogs = aggregate.accessLogs() } if managed != nil && managed.TrafficReport != nil { - return managed.TrafficReport, accessLogs + return managed.TrafficReport, accessLogs, managed } 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 { @@ -162,6 +167,12 @@ func (aggregate *trafficAggregate) consume(line []byte) { if record.Status > 0 { 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 != "" { aggregate.topDomains[host]++ } @@ -178,11 +189,13 @@ func (aggregate *trafficAggregate) consume(line []byte) { } type parsedAccessLogRecord struct { - Timestamp time.Time - Host string - RemoteAddr string - Path string - Status int + Timestamp time.Time + Host string + RemoteAddr string + Path string + Status int + BytesSent int64 + RequestLength int64 } func parseAccessLogRecord(raw string) (parsedAccessLogRecord, bool) { @@ -203,11 +216,13 @@ func parseJSONAccessLogRecord(raw string) (parsedAccessLogRecord, bool) { return parsedAccessLogRecord{}, false } return parsedAccessLogRecord{ - Timestamp: timestamp, - Host: strings.TrimSpace(record.Host), - RemoteAddr: strings.TrimSpace(record.RemoteAddr), - Path: normalizeAccessLogPath(record.Path), - Status: record.Status, + Timestamp: timestamp, + Host: strings.TrimSpace(record.Host), + RemoteAddr: strings.TrimSpace(record.RemoteAddr), + Path: normalizeAccessLogPath(record.Path), + Status: record.Status, + BytesSent: record.BytesSent, + RequestLength: record.RequestLength, }, true } @@ -256,6 +271,21 @@ func (aggregate *trafficAggregate) accessLogs() []protocol.NodeAccessLog { 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) { trimmed := strings.TrimSpace(value) if trimmed == "" { diff --git a/atsf_agent/internal/observability/traffic_test.go b/atsf_agent/internal/observability/traffic_test.go index fca91ebc..cb34535f 100644 --- a/atsf_agent/internal/observability/traffic_test.go +++ b/atsf_agent/internal/observability/traffic_test.go @@ -85,21 +85,24 @@ func TestBuildTrafficObservabilityReturnsAccessLogs(t *testing.T) { } logPath := filepath.Join(filepath.Dir(routeConfigPath), "atsflare_access.log") 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:05Z\",\"host\":\"api.example.com\",\"path\":\"/v1/ping\",\"remote_addr\":\"10.0.0.2\",\"status\":502}\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,\"request_length\":64,\"bytes_sent\":256}\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, accessLogs := BuildTrafficObservability(&config.Config{RouteConfigPath: routeConfigPath}, stateStore, nil) + report, accessLogs, fallbackMetrics := BuildTrafficObservability(&config.Config{RouteConfigPath: routeConfigPath}, stateStore, nil) if report == nil || report.RequestCount != 2 { t.Fatalf("expected traffic report, got %+v", report) } if len(accessLogs) != 2 { 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" { t.Fatalf("unexpected access log paths: %+v", accessLogs) }