mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-08 08:36:37 +08:00
feat: add access log functionality and related components
- Introduced NodeAccessLog model and corresponding database migrations. - Implemented access log retrieval in the service layer. - Created API endpoint for accessing logs with appropriate security measures. - Developed frontend components for displaying access logs, including filtering and summary statistics. - Updated observability buffer to include access logs and ensure proper merging and retention. - Enhanced traffic report and observability tests to validate new access log features. - Updated documentation to reflect changes in access log handling and data retention policies.
This commit is contained in:
@@ -0,0 +1,55 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"atsflare/model"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
const accessLogListLimit = 500
|
||||
|
||||
type AccessLogView struct {
|
||||
ID uint `json:"id"`
|
||||
NodeID string `json:"node_id"`
|
||||
NodeName string `json:"node_name"`
|
||||
LoggedAt time.Time `json:"logged_at"`
|
||||
RemoteAddr string `json:"remote_addr"`
|
||||
Host string `json:"host"`
|
||||
Path string `json:"path"`
|
||||
StatusCode int `json:"status_code"`
|
||||
}
|
||||
|
||||
func ListAccessLogs(nodeID string) ([]AccessLogView, error) {
|
||||
logs, err := model.ListNodeAccessLogs(strings.TrimSpace(nodeID), time.Now().Add(-nodeAccessLogRetentionWindow), accessLogListLimit)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
nodes, err := model.ListNodes()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
nodeNames := make(map[string]string, len(nodes))
|
||||
for _, node := range nodes {
|
||||
if node == nil {
|
||||
continue
|
||||
}
|
||||
nodeNames[node.NodeID] = node.Name
|
||||
}
|
||||
views := make([]AccessLogView, 0, len(logs))
|
||||
for _, item := range logs {
|
||||
if item == nil {
|
||||
continue
|
||||
}
|
||||
views = append(views, AccessLogView{
|
||||
ID: item.ID,
|
||||
NodeID: item.NodeID,
|
||||
NodeName: nodeNames[item.NodeID],
|
||||
LoggedAt: item.LoggedAt,
|
||||
RemoteAddr: item.RemoteAddr,
|
||||
Host: item.Host,
|
||||
Path: item.Path,
|
||||
StatusCode: item.StatusCode,
|
||||
})
|
||||
}
|
||||
return views, nil
|
||||
}
|
||||
@@ -36,6 +36,7 @@ type AgentNodePayload struct {
|
||||
Profile *AgentNodeSystemProfile `json:"profile,omitempty"`
|
||||
Snapshot *AgentNodeMetricSnapshot `json:"snapshot,omitempty"`
|
||||
TrafficReport *AgentNodeTrafficReport `json:"traffic_report,omitempty"`
|
||||
AccessLogs []AgentNodeAccessLog `json:"access_logs,omitempty"`
|
||||
BufferedObservability []AgentBufferedObservabilityRecord `json:"buffered_observability,omitempty"`
|
||||
HealthEvents []AgentNodeHealthEvent `json:"health_events"`
|
||||
}
|
||||
|
||||
@@ -620,6 +620,22 @@ func TestHeartbeatNodePersistsObservabilityPayload(t *testing.T) {
|
||||
TopDomains: map[string]int64{"example.com": 900},
|
||||
SourceCountries: map[string]int64{"CN": 700, "US": 200},
|
||||
},
|
||||
AccessLogs: []AgentNodeAccessLog{
|
||||
{
|
||||
LoggedAtUnix: time.Now().Add(-45 * time.Second).Unix(),
|
||||
RemoteAddr: "203.0.113.10",
|
||||
Host: "example.com",
|
||||
Path: "/login",
|
||||
StatusCode: 200,
|
||||
},
|
||||
{
|
||||
LoggedAtUnix: time.Now().Add(-40 * time.Second).Unix(),
|
||||
RemoteAddr: "198.51.100.20",
|
||||
Host: "api.example.com",
|
||||
Path: "/v1/ping",
|
||||
StatusCode: 502,
|
||||
},
|
||||
},
|
||||
HealthEvents: []AgentNodeHealthEvent{
|
||||
{
|
||||
EventType: "openresty_unhealthy",
|
||||
@@ -657,6 +673,14 @@ func TestHeartbeatNodePersistsObservabilityPayload(t *testing.T) {
|
||||
t.Fatalf("unexpected request reports: %+v", reports)
|
||||
}
|
||||
|
||||
accessLogs, err := model.ListNodeAccessLogs(node.NodeID, time.Time{}, 10)
|
||||
if err != nil {
|
||||
t.Fatalf("expected node access logs query to succeed: %v", err)
|
||||
}
|
||||
if len(accessLogs) != 2 || accessLogs[0].Path == "" {
|
||||
t.Fatalf("unexpected access logs: %+v", accessLogs)
|
||||
}
|
||||
|
||||
events, err := model.ListNodeHealthEvents(node.NodeID, true, 10)
|
||||
if err != nil {
|
||||
t.Fatalf("expected node health events query to succeed: %v", err)
|
||||
@@ -724,6 +748,15 @@ func TestHeartbeatNodePersistsBufferedObservabilityPayload(t *testing.T) {
|
||||
TopDomains: map[string]int64{"edge.example.com": 40},
|
||||
SourceCountries: map[string]int64{"CN": 20},
|
||||
},
|
||||
AccessLogs: []AgentNodeAccessLog{
|
||||
{
|
||||
LoggedAtUnix: now.Add(-110 * time.Second).Unix(),
|
||||
RemoteAddr: "203.0.113.21",
|
||||
Host: "edge.example.com",
|
||||
Path: "/buffered",
|
||||
StatusCode: 200,
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
})
|
||||
@@ -747,6 +780,14 @@ func TestHeartbeatNodePersistsBufferedObservabilityPayload(t *testing.T) {
|
||||
t.Fatalf("expected current and buffered reports, got %+v", reports)
|
||||
}
|
||||
|
||||
accessLogs, err := model.ListNodeAccessLogs(node.NodeID, time.Time{}, 10)
|
||||
if err != nil {
|
||||
t.Fatalf("expected node access logs query to succeed: %v", err)
|
||||
}
|
||||
if len(accessLogs) != 1 || accessLogs[0].Path != "/buffered" {
|
||||
t.Fatalf("expected buffered access logs to persist, got %+v", accessLogs)
|
||||
}
|
||||
|
||||
_, err = HeartbeatNode(node, AgentNodePayload{
|
||||
NodeID: node.NodeID,
|
||||
Name: node.Name,
|
||||
|
||||
@@ -17,6 +17,7 @@ const (
|
||||
NodeHealthSeverityInfo = "info"
|
||||
NodeHealthSeverityWarning = "warning"
|
||||
NodeHealthSeverityCritical = "critical"
|
||||
nodeAccessLogRetentionWindow = 24 * time.Hour
|
||||
)
|
||||
|
||||
type AgentNodeSystemProfile struct {
|
||||
@@ -60,10 +61,19 @@ type AgentNodeTrafficReport struct {
|
||||
SourceCountries map[string]int64 `json:"source_countries"`
|
||||
}
|
||||
|
||||
type AgentNodeAccessLog struct {
|
||||
LoggedAtUnix int64 `json:"logged_at_unix"`
|
||||
RemoteAddr string `json:"remote_addr"`
|
||||
Host string `json:"host"`
|
||||
Path string `json:"path"`
|
||||
StatusCode int `json:"status_code"`
|
||||
}
|
||||
|
||||
type AgentBufferedObservabilityRecord struct {
|
||||
WindowStartedAtUnix int64 `json:"window_started_at_unix"`
|
||||
Snapshot *AgentNodeMetricSnapshot `json:"snapshot,omitempty"`
|
||||
TrafficReport *AgentNodeTrafficReport `json:"traffic_report,omitempty"`
|
||||
AccessLogs []AgentNodeAccessLog `json:"access_logs,omitempty"`
|
||||
}
|
||||
|
||||
type AgentNodeHealthEvent struct {
|
||||
@@ -78,7 +88,7 @@ func persistHeartbeatObservability(nodeID string, payload AgentNodePayload, repo
|
||||
if strings.TrimSpace(nodeID) == "" {
|
||||
return
|
||||
}
|
||||
if payload.Profile == nil && payload.Snapshot == nil && payload.TrafficReport == nil && len(payload.BufferedObservability) == 0 && payload.HealthEvents == nil {
|
||||
if payload.Profile == nil && payload.Snapshot == nil && payload.TrafficReport == nil && len(payload.AccessLogs) == 0 && len(payload.BufferedObservability) == 0 && payload.HealthEvents == nil {
|
||||
return
|
||||
}
|
||||
|
||||
@@ -95,6 +105,9 @@ func persistHeartbeatObservability(nodeID string, payload AgentNodePayload, repo
|
||||
if err := persistNodeTrafficReport(tx, nodeID, payload.TrafficReport, reportedAt); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := persistNodeAccessLogs(tx, nodeID, payload.AccessLogs, reportedAt); err != nil {
|
||||
return err
|
||||
}
|
||||
if payload.HealthEvents != nil {
|
||||
if err := reconcileNodeHealthEvents(tx, nodeID, payload.HealthEvents, reportedAt); err != nil {
|
||||
return err
|
||||
@@ -114,6 +127,9 @@ func persistBufferedObservability(tx *gorm.DB, nodeID string, records []AgentBuf
|
||||
if err := persistNodeTrafficReport(tx, nodeID, record.TrafficReport, reportedAt); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := persistNodeAccessLogs(tx, nodeID, record.AccessLogs, reportedAt); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -186,6 +202,35 @@ func persistNodeTrafficReport(tx *gorm.DB, nodeID string, report *AgentNodeTraff
|
||||
return tx.Where("node_id = ? AND window_started_at = ? AND window_ended_at = ?", nodeID, record.WindowStartedAt, record.WindowEndedAt).Assign(record).FirstOrCreate(record).Error
|
||||
}
|
||||
|
||||
func persistNodeAccessLogs(tx *gorm.DB, nodeID string, logs []AgentNodeAccessLog, reportedAt time.Time) error {
|
||||
if len(logs) == 0 {
|
||||
return nil
|
||||
}
|
||||
for _, item := range logs {
|
||||
record := &model.NodeAccessLog{
|
||||
NodeID: nodeID,
|
||||
LoggedAt: timeFromUnix(item.LoggedAtUnix, reportedAt),
|
||||
RemoteAddr: strings.TrimSpace(item.RemoteAddr),
|
||||
Host: strings.TrimSpace(item.Host),
|
||||
Path: strings.TrimSpace(item.Path),
|
||||
StatusCode: item.StatusCode,
|
||||
RawJSON: marshalJSON(item),
|
||||
}
|
||||
if err := tx.Where(
|
||||
"node_id = ? AND logged_at = ? AND remote_addr = ? AND host = ? AND path = ? AND status_code = ?",
|
||||
nodeID,
|
||||
record.LoggedAt,
|
||||
record.RemoteAddr,
|
||||
record.Host,
|
||||
record.Path,
|
||||
record.StatusCode,
|
||||
).Assign(record).FirstOrCreate(record).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return tx.Where("node_id = ? AND logged_at < ?", nodeID, reportedAt.Add(-nodeAccessLogRetentionWindow)).Delete(&model.NodeAccessLog{}).Error
|
||||
}
|
||||
|
||||
func reconcileNodeHealthEvents(tx *gorm.DB, nodeID string, events []AgentNodeHealthEvent, reportedAt time.Time) error {
|
||||
activeTypes := make(map[string]AgentNodeHealthEvent, len(events))
|
||||
for _, event := range events {
|
||||
|
||||
Reference in New Issue
Block a user