[功能] 添加访问日志折叠、IP汇总和趋势查询功能,并实现日志清理功能

This commit is contained in:
ryan
2026-03-18 10:33:35 +08:00
parent 2e875f583b
commit 29a6fedbe9
10 changed files with 1954 additions and 273 deletions
+136 -3
View File
@@ -1,6 +1,7 @@
package controller
import (
"net/http"
"openflare/service"
"strconv"
@@ -13,17 +14,149 @@ import (
// @Produce json
// @Security BearerAuth
// @Param node_id query string false "Node ID"
// @Param remote_addr query string false "Remote address"
// @Param host query string false "Host"
// @Param path query string false "Path"
// @Param p query int false "Page index"
// @Param page_size query int false "Page size"
// @Param sort_by query string false "Sort by"
// @Param sort_order query string false "Sort order"
// @Success 200 {object} map[string]interface{}
// @Router /api/access-logs/ [get]
func GetAccessLogs(c *gin.Context) {
page, _ := strconv.Atoi(c.DefaultQuery("p", "0"))
pageSize, _ := strconv.Atoi(c.DefaultQuery("page_size", "0"))
logs, err := service.ListAccessLogs(c.Query("node_id"), page, pageSize)
logs, err := service.ListAccessLogs(readAccessLogQuery(c))
if err != nil {
respondFailure(c, err.Error())
return
}
respondSuccess(c, logs)
}
// GetFoldedAccessLogs godoc
// @Summary List folded access logs
// @Tags AccessLogs
// @Produce json
// @Security BearerAuth
// @Param node_id query string false "Node ID"
// @Param remote_addr query string false "Remote address"
// @Param host query string false "Host"
// @Param path query string false "Path"
// @Param p query int false "Page index"
// @Param page_size query int false "Page size"
// @Param sort_by query string false "Sort by"
// @Param sort_order query string false "Sort order"
// @Param fold_minutes query int false "Fold minutes"
// @Success 200 {object} map[string]interface{}
// @Router /api/access-logs/folds [get]
func GetFoldedAccessLogs(c *gin.Context) {
query := readAccessLogQuery(c)
query.FoldMinutes = readQueryInt(c, "fold_minutes")
logs, err := service.ListFoldedAccessLogs(query)
if err != nil {
respondFailure(c, err.Error())
return
}
respondSuccess(c, logs)
}
// GetAccessLogIPSummaries godoc
// @Summary List access log IP summaries
// @Tags AccessLogs
// @Produce json
// @Security BearerAuth
// @Param node_id query string false "Node ID"
// @Param remote_addr query string false "Remote address"
// @Param host query string false "Host"
// @Param p query int false "Page index"
// @Param page_size query int false "Page size"
// @Param sort_by query string false "Sort by"
// @Param sort_order query string false "Sort order"
// @Success 200 {object} map[string]interface{}
// @Router /api/access-logs/ip-summary [get]
func GetAccessLogIPSummaries(c *gin.Context) {
result, err := service.ListAccessLogIPSummaries(service.AccessLogIPSummaryQuery{
NodeID: c.Query("node_id"),
RemoteAddr: c.Query("remote_addr"),
Host: c.Query("host"),
Page: readQueryInt(c, "p"),
PageSize: readQueryInt(c, "page_size"),
SortBy: c.Query("sort_by"),
SortOrder: c.Query("sort_order"),
})
if err != nil {
respondFailure(c, err.Error())
return
}
respondSuccess(c, result)
}
// GetAccessLogIPTrend godoc
// @Summary Get access log IP trend
// @Tags AccessLogs
// @Produce json
// @Security BearerAuth
// @Param node_id query string false "Node ID"
// @Param remote_addr query string true "Remote address"
// @Param host query string false "Host"
// @Param hours query int false "Hours"
// @Param bucket_minutes query int false "Bucket minutes"
// @Success 200 {object} map[string]interface{}
// @Router /api/access-logs/ip-summary/trend [get]
func GetAccessLogIPTrend(c *gin.Context) {
result, err := service.GetAccessLogIPTrend(service.AccessLogIPTrendQuery{
NodeID: c.Query("node_id"),
RemoteAddr: c.Query("remote_addr"),
Host: c.Query("host"),
Hours: readQueryInt(c, "hours"),
BucketMinutes: readQueryInt(c, "bucket_minutes"),
})
if err != nil {
respondFailure(c, err.Error())
return
}
respondSuccess(c, result)
}
// CleanupAccessLogs godoc
// @Summary Cleanup access logs by retention days
// @Tags AccessLogs
// @Accept json
// @Produce json
// @Security BearerAuth
// @Success 200 {object} map[string]interface{}
// @Router /api/access-logs/cleanup [post]
func CleanupAccessLogs(c *gin.Context) {
var input service.AccessLogCleanupInput
if err := c.ShouldBindJSON(&input); err != nil {
c.JSON(http.StatusBadRequest, gin.H{
"success": false,
"message": "参数错误",
"error": err.Error(),
})
return
}
result, err := service.CleanupAccessLogs(input)
if err != nil {
respondFailure(c, err.Error())
return
}
respondSuccess(c, result)
}
func readAccessLogQuery(c *gin.Context) service.AccessLogQuery {
return service.AccessLogQuery{
NodeID: c.Query("node_id"),
RemoteAddr: c.Query("remote_addr"),
Host: c.Query("host"),
Path: c.Query("path"),
Page: readQueryInt(c, "p"),
PageSize: readQueryInt(c, "page_size"),
SortBy: c.Query("sort_by"),
SortOrder: c.Query("sort_order"),
}
}
func readQueryInt(c *gin.Context, key string) int {
value, _ := strconv.Atoi(c.DefaultQuery(key, "0"))
return value
}
+291 -33
View File
@@ -1,16 +1,22 @@
package model
import "time"
import (
"fmt"
"strings"
"time"
"gorm.io/gorm"
)
type NodeAccessLog struct {
ID uint `json:"id" gorm:"primaryKey"`
NodeID string `json:"node_id" gorm:"index;size:64;not null"`
LoggedAt time.Time `json:"logged_at" gorm:"index"`
RemoteAddr string `json:"remote_addr" gorm:"size:128"`
NodeID string `json:"node_id" gorm:"index:idx_node_access_logs_node_logged_at,priority:1;size:64;not null"`
LoggedAt time.Time `json:"logged_at" gorm:"index:idx_node_access_logs_logged_at;index:idx_node_access_logs_node_logged_at,priority:2"`
RemoteAddr string `json:"remote_addr" gorm:"index:idx_node_access_logs_remote_addr;size:128"`
Region string `json:"region" gorm:"size:128"`
Host string `json:"host" gorm:"size:255"`
Host string `json:"host" gorm:"index:idx_node_access_logs_host;size:255"`
Path string `json:"path" gorm:"size:2048"`
StatusCode int `json:"status_code"`
StatusCode int `json:"status_code" gorm:"index:idx_node_access_logs_status_code"`
RawJSON string `json:"raw_json" gorm:"type:text"`
CreatedAt time.Time `json:"created_at"`
}
@@ -20,39 +26,91 @@ type NodeAccessLogRegionCount struct {
Count int64 `json:"count"`
}
func ListNodeAccessLogs(nodeID string, since time.Time, offset int, limit int) (logs []*NodeAccessLog, err error) {
query := DB.Order("logged_at desc, id desc")
if nodeID != "" {
query = query.Where("node_id = ?", nodeID)
}
if !since.IsZero() {
query = query.Where("logged_at >= ?", since)
}
if offset > 0 {
query = query.Offset(offset)
}
if limit > 0 {
query = query.Limit(limit)
}
err = query.Find(&logs).Error
type NodeAccessLogQuery struct {
NodeID string
RemoteAddr string
Host string
Path string
Since time.Time
Page int
PageSize int
SortBy string
SortOrder string
}
type NodeAccessLogBucketQuery struct {
NodeID string
RemoteAddr string
Host string
Path string
Since time.Time
Page int
PageSize int
SortBy string
SortOrder string
FoldMinutes int
}
type NodeAccessLogBucketRow struct {
BucketEpoch int64 `json:"bucket_epoch"`
RequestCount int64 `json:"request_count"`
UniqueIPCount int64 `json:"unique_ip_count"`
UniqueHostCount int64 `json:"unique_host_count"`
SuccessCount int64 `json:"success_count"`
ClientErrorCount int64 `json:"client_error_count"`
ServerErrorCount int64 `json:"server_error_count"`
}
type NodeAccessLogIPSummaryQuery struct {
NodeID string
RemoteAddr string
Host string
Since time.Time
Page int
PageSize int
SortBy string
SortOrder string
}
type NodeAccessLogIPSummaryRow struct {
RemoteAddr string `json:"remote_addr"`
TotalRequests int64 `json:"total_requests"`
RecentRequests int64 `json:"recent_requests"`
LastSeenEpoch int64 `json:"last_seen_epoch"`
}
type NodeAccessLogIPTrendQuery struct {
NodeID string
RemoteAddr string
Host string
Since time.Time
BucketMinutes int
}
type NodeAccessLogTrendPointRow struct {
BucketEpoch int64 `json:"bucket_epoch"`
RequestCount int64 `json:"request_count"`
}
func ListNodeAccessLogs(query NodeAccessLogQuery) (logs []*NodeAccessLog, err error) {
offset := query.Page * query.PageSize
db := buildNodeAccessLogQuery(DB, query).
Order(buildNodeAccessLogSortClause(query.SortBy, query.SortOrder)).
Limit(query.PageSize).
Offset(offset)
err = db.Find(&logs).Error
return logs, err
}
func CountNodeAccessLogs(nodeID string, since time.Time) (totalRecords int64, totalIPs int64, err error) {
query := DB.Model(&NodeAccessLog{})
if nodeID != "" {
query = query.Where("node_id = ?", nodeID)
}
if !since.IsZero() {
query = query.Where("logged_at >= ?", since)
}
if err = query.Count(&totalRecords).Error; err != nil {
func CountNodeAccessLogs(query NodeAccessLogQuery) (totalRecords int64, totalIPs int64, err error) {
base := buildNodeAccessLogQuery(DB.Model(&NodeAccessLog{}), query)
if err = base.Count(&totalRecords).Error; err != nil {
return 0, 0, err
}
if err = query.
distinctQuery := buildNodeAccessLogQuery(DB.Model(&NodeAccessLog{}), query).
Where("remote_addr <> ''").
Distinct("remote_addr").
Count(&totalIPs).Error; err != nil {
Distinct("remote_addr")
if err = distinctQuery.Count(&totalIPs).Error; err != nil {
return 0, 0, err
}
return totalRecords, totalIPs, nil
@@ -75,3 +133,203 @@ func ListNodeAccessLogRegionCounts(nodeID string, since time.Time, limit int) (i
err = query.Scan(&items).Error
return items, err
}
func ListNodeAccessLogBuckets(query NodeAccessLogBucketQuery) (items []*NodeAccessLogBucketRow, err error) {
offset := query.Page * query.PageSize
bucketExpr := accessLogBucketEpochExpr(query.FoldMinutes)
base := buildNodeAccessLogQuery(DB.Model(&NodeAccessLog{}), NodeAccessLogQuery{
NodeID: query.NodeID,
RemoteAddr: query.RemoteAddr,
Host: query.Host,
Path: query.Path,
Since: query.Since,
})
err = base.Select(fmt.Sprintf(
"%s as bucket_epoch, count(*) as request_count, count(distinct remote_addr) as unique_ip_count, count(distinct host) as unique_host_count, sum(case when status_code < 400 then 1 else 0 end) as success_count, sum(case when status_code >= 400 and status_code < 500 then 1 else 0 end) as client_error_count, sum(case when status_code >= 500 then 1 else 0 end) as server_error_count",
bucketExpr,
)).
Group(bucketExpr).
Order(buildNodeAccessLogBucketSortClause(query.SortBy, query.SortOrder)).
Limit(query.PageSize).
Offset(offset).
Scan(&items).Error
return items, err
}
func CountNodeAccessLogBuckets(query NodeAccessLogBucketQuery) (total int64, err error) {
bucketExpr := accessLogBucketEpochExpr(query.FoldMinutes)
base := buildNodeAccessLogQuery(DB.Model(&NodeAccessLog{}), NodeAccessLogQuery{
NodeID: query.NodeID,
RemoteAddr: query.RemoteAddr,
Host: query.Host,
Path: query.Path,
Since: query.Since,
})
rows := []struct {
BucketEpoch int64 `gorm:"column:bucket_epoch"`
}{}
err = base.Select(fmt.Sprintf("%s as bucket_epoch", bucketExpr)).
Group(bucketExpr).
Scan(&rows).Error
if err != nil {
return 0, err
}
return int64(len(rows)), nil
}
func ListNodeAccessLogIPSummaries(query NodeAccessLogIPSummaryQuery, recentSince time.Time) (items []*NodeAccessLogIPSummaryRow, err error) {
offset := query.Page * query.PageSize
base := buildNodeAccessLogQuery(DB.Model(&NodeAccessLog{}), NodeAccessLogQuery{
NodeID: query.NodeID,
RemoteAddr: query.RemoteAddr,
Host: query.Host,
Since: query.Since,
}).Where("remote_addr <> ''")
lastSeenExpr := accessLogEpochExpr("max(logged_at)")
err = base.Select(
"remote_addr as remote_addr, count(*) as total_requests, sum(case when logged_at >= ? then 1 else 0 end) as recent_requests, "+lastSeenExpr+" as last_seen_epoch",
recentSince,
).
Group("remote_addr").
Order(buildNodeAccessLogIPSummarySortClause(query.SortBy, query.SortOrder)).
Limit(query.PageSize).
Offset(offset).
Scan(&items).Error
return items, err
}
func CountNodeAccessLogIPSummaries(query NodeAccessLogIPSummaryQuery) (total int64, err error) {
base := buildNodeAccessLogQuery(DB.Model(&NodeAccessLog{}), NodeAccessLogQuery{
NodeID: query.NodeID,
RemoteAddr: query.RemoteAddr,
Host: query.Host,
Since: query.Since,
}).Where("remote_addr <> ''")
rows := []struct {
RemoteAddr string `gorm:"column:remote_addr"`
}{}
err = base.Select("remote_addr").
Group("remote_addr").
Scan(&rows).Error
if err != nil {
return 0, err
}
return int64(len(rows)), nil
}
func ListNodeAccessLogIPTrend(query NodeAccessLogIPTrendQuery) (items []*NodeAccessLogTrendPointRow, err error) {
bucketExpr := accessLogBucketEpochExpr(query.BucketMinutes)
base := buildNodeAccessLogQuery(DB.Model(&NodeAccessLog{}), NodeAccessLogQuery{
NodeID: query.NodeID,
RemoteAddr: query.RemoteAddr,
Host: query.Host,
Since: query.Since,
}).Where("remote_addr = ?", strings.TrimSpace(query.RemoteAddr))
err = base.Select(fmt.Sprintf("%s as bucket_epoch, count(*) as request_count", bucketExpr)).
Group(bucketExpr).
Order("bucket_epoch asc").
Scan(&items).Error
return items, err
}
func DeleteNodeAccessLogsBefore(before time.Time) (deleted int64, err error) {
result := DB.Where("logged_at < ?", before).Delete(&NodeAccessLog{})
return result.RowsAffected, result.Error
}
func buildNodeAccessLogQuery(db *gorm.DB, query NodeAccessLogQuery) *gorm.DB {
if db == nil {
db = DB.Model(&NodeAccessLog{})
}
if db.Statement == nil || db.Statement.Model == nil {
db = db.Model(&NodeAccessLog{})
}
if trimmed := strings.TrimSpace(query.NodeID); trimmed != "" {
db = db.Where("node_id LIKE ?", "%"+trimmed+"%")
}
if trimmed := strings.TrimSpace(query.RemoteAddr); trimmed != "" {
db = db.Where("remote_addr LIKE ?", "%"+trimmed+"%")
}
if trimmed := strings.TrimSpace(query.Host); trimmed != "" {
db = db.Where("host LIKE ?", "%"+trimmed+"%")
}
if trimmed := strings.TrimSpace(query.Path); trimmed != "" {
db = db.Where("path LIKE ?", "%"+trimmed+"%")
}
if !query.Since.IsZero() {
db = db.Where("logged_at >= ?", query.Since)
}
return db
}
func buildNodeAccessLogSortClause(sortBy string, sortOrder string) string {
column := "logged_at"
switch strings.TrimSpace(sortBy) {
case "status_code":
column = "status_code"
case "remote_addr":
column = "remote_addr"
case "host":
column = "host"
case "path":
column = "path"
}
order := normalizeSortOrder(sortOrder)
if column == "logged_at" {
return fmt.Sprintf("%s %s, id %s", column, order, order)
}
return fmt.Sprintf("%s %s, logged_at desc, id desc", column, order)
}
func buildNodeAccessLogBucketSortClause(sortBy string, sortOrder string) string {
order := normalizeSortOrder(sortOrder)
switch strings.TrimSpace(sortBy) {
case "request_count":
return fmt.Sprintf("request_count %s, bucket_epoch desc", order)
default:
return fmt.Sprintf("bucket_epoch %s", order)
}
}
func buildNodeAccessLogIPSummarySortClause(sortBy string, sortOrder string) string {
order := normalizeSortOrder(sortOrder)
switch strings.TrimSpace(sortBy) {
case "recent_requests":
return fmt.Sprintf("recent_requests %s, last_seen_epoch desc, remote_addr asc", order)
case "last_seen_at":
return fmt.Sprintf("last_seen_epoch %s, total_requests desc, remote_addr asc", order)
case "remote_addr":
return fmt.Sprintf("remote_addr %s", order)
default:
return fmt.Sprintf("total_requests %s, last_seen_epoch desc, remote_addr asc", order)
}
}
func accessLogBucketEpochExpr(bucketMinutes int) string {
bucketSeconds := bucketMinutes * 60
if bucketSeconds <= 0 {
bucketSeconds = 180
}
switch DB.Dialector.Name() {
case "postgres":
return fmt.Sprintf("CAST(floor(extract(epoch from logged_at) / %d) * %d AS BIGINT)", bucketSeconds, bucketSeconds)
default:
return fmt.Sprintf("CAST((strftime('%%s', logged_at) / %d) * %d AS INTEGER)", bucketSeconds, bucketSeconds)
}
}
func accessLogEpochExpr(expression string) string {
switch DB.Dialector.Name() {
case "postgres":
return fmt.Sprintf("CAST(extract(epoch from %s) AS BIGINT)", expression)
default:
return fmt.Sprintf("CAST(strftime('%%s', %s) AS INTEGER)", expression)
}
}
func normalizeSortOrder(sortOrder string) string {
if strings.EqualFold(strings.TrimSpace(sortOrder), "asc") {
return "asc"
}
return "desc"
}
+4
View File
@@ -139,6 +139,10 @@ func SetApiRouter(router *gin.Engine) {
accessLogRoute.Use(middleware.AdminAuth())
{
accessLogRoute.GET("/", controller.GetAccessLogs)
accessLogRoute.GET("/folds", controller.GetFoldedAccessLogs)
accessLogRoute.GET("/ip-summary", controller.GetAccessLogIPSummaries)
accessLogRoute.GET("/ip-summary/trend", controller.GetAccessLogIPTrend)
accessLogRoute.POST("/cleanup", controller.CleanupAccessLogs)
}
agentRoute := apiRouter.Group("/agent")
{
+411 -42
View File
@@ -1,16 +1,36 @@
package service
import (
"errors"
"openflare/model"
"strings"
"time"
)
const (
defaultAccessLogPageSize = 20
maxAccessLogPageSize = 200
defaultAccessLogPageSize = 20
maxAccessLogPageSize = 200
defaultAccessLogSortBy = "logged_at"
defaultAccessLogSortOrder = "desc"
defaultAccessLogFoldMinute = 3
defaultIPTrendHours = 24
defaultIPTrendBucketMinute = 30
maxIPTrendHours = 168
nodeAccessLogRetentionDays = 90
)
type AccessLogQuery struct {
NodeID string `json:"node_id"`
RemoteAddr string `json:"remote_addr"`
Host string `json:"host"`
Path string `json:"path"`
Page int `json:"page"`
PageSize int `json:"page_size"`
SortBy string `json:"sort_by"`
SortOrder string `json:"sort_order"`
FoldMinutes int `json:"fold_minutes"`
}
type AccessLogView struct {
ID uint `json:"id"`
NodeID string `json:"node_id"`
@@ -32,52 +52,99 @@ type AccessLogList struct {
TotalIP int64 `json:"total_ip"`
}
func ListAccessLogs(nodeID string, page int, pageSize int) (*AccessLogList, error) {
normalizedPage := normalizeAccessLogPage(page)
normalizedPageSize := normalizeAccessLogPageSize(pageSize)
offset := normalizedPage * normalizedPageSize
trimmedNodeID := strings.TrimSpace(nodeID)
since := time.Now().Add(-nodeAccessLogRetentionWindow)
logs, err := model.ListNodeAccessLogs(
trimmedNodeID,
since,
offset,
normalizedPageSize+1,
)
type FoldedAccessLogView struct {
BucketStartedAt time.Time `json:"bucket_started_at"`
RequestCount int64 `json:"request_count"`
UniqueIPCount int64 `json:"unique_ip_count"`
UniqueHostCount int64 `json:"unique_host_count"`
SuccessCount int64 `json:"success_count"`
ClientErrorCount int64 `json:"client_error_count"`
ServerErrorCount int64 `json:"server_error_count"`
}
type FoldedAccessLogList struct {
Items []FoldedAccessLogView `json:"items"`
Page int `json:"page"`
PageSize int `json:"page_size"`
HasMore bool `json:"has_more"`
TotalBucket int64 `json:"total_bucket"`
TotalRecord int64 `json:"total_record"`
TotalIP int64 `json:"total_ip"`
FoldMinutes int `json:"fold_minutes"`
}
type AccessLogIPSummaryQuery struct {
NodeID string `json:"node_id"`
RemoteAddr string `json:"remote_addr"`
Host string `json:"host"`
Page int `json:"page"`
PageSize int `json:"page_size"`
SortBy string `json:"sort_by"`
SortOrder string `json:"sort_order"`
}
type AccessLogIPSummaryView struct {
RemoteAddr string `json:"remote_addr"`
TotalRequests int64 `json:"total_requests"`
RecentRequests int64 `json:"recent_requests"`
LastSeenAt time.Time `json:"last_seen_at"`
}
type AccessLogIPSummaryList struct {
Items []AccessLogIPSummaryView `json:"items"`
Page int `json:"page"`
PageSize int `json:"page_size"`
HasMore bool `json:"has_more"`
TotalIP int64 `json:"total_ip"`
SortBy string `json:"sort_by"`
SortOrder string `json:"sort_order"`
}
type AccessLogIPTrendQuery struct {
NodeID string `json:"node_id"`
RemoteAddr string `json:"remote_addr"`
Host string `json:"host"`
Hours int `json:"hours"`
BucketMinutes int `json:"bucket_minutes"`
}
type AccessLogIPTrendPoint struct {
BucketStartedAt time.Time `json:"bucket_started_at"`
RequestCount int64 `json:"request_count"`
}
type AccessLogIPTrendView struct {
RemoteAddr string `json:"remote_addr"`
Hours int `json:"hours"`
BucketMinutes int `json:"bucket_minutes"`
Points []AccessLogIPTrendPoint `json:"points"`
}
type AccessLogCleanupInput struct {
RetentionDays int `json:"retention_days"`
}
type AccessLogCleanupResult struct {
RetentionDays int `json:"retention_days"`
DeletedCount int64 `json:"deleted_count"`
Cutoff time.Time `json:"cutoff"`
}
func ListAccessLogs(input AccessLogQuery) (*AccessLogList, error) {
normalized := normalizeAccessLogQuery(input)
modelQuery := buildModelAccessLogQuery(normalized)
logs, err := model.ListNodeAccessLogs(modelQuery)
if err != nil {
return nil, err
}
totalRecords, totalIPs, err := model.CountNodeAccessLogs(trimmedNodeID, since)
totalRecords, totalIPs, err := model.CountNodeAccessLogs(modelQuery)
if err != nil {
return nil, err
}
nodeIDs := make([]string, 0, len(logs))
seenNodeIDs := make(map[string]struct{}, len(logs))
for _, item := range logs {
if item == nil {
continue
}
if _, exists := seenNodeIDs[item.NodeID]; exists {
continue
}
seenNodeIDs[item.NodeID] = struct{}{}
nodeIDs = append(nodeIDs, item.NodeID)
}
nodes, err := model.ListNodesByNodeIDs(nodeIDs)
nodeNames, err := listNodeNameMap(logs)
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
}
hasMore := len(logs) > normalizedPageSize
if hasMore {
logs = logs[:normalizedPageSize]
}
views := make([]AccessLogView, 0, len(logs))
for _, item := range logs {
if item == nil {
@@ -97,14 +164,270 @@ func ListAccessLogs(nodeID string, page int, pageSize int) (*AccessLogList, erro
}
return &AccessLogList{
Items: views,
Page: normalizedPage,
PageSize: normalizedPageSize,
HasMore: hasMore,
Page: normalized.Page,
PageSize: normalized.PageSize,
HasMore: int64((normalized.Page+1)*normalized.PageSize) < totalRecords,
TotalRecord: totalRecords,
TotalIP: totalIPs,
}, nil
}
func ListFoldedAccessLogs(input AccessLogQuery) (*FoldedAccessLogList, error) {
normalized := normalizeAccessLogQuery(input)
foldMinutes, err := normalizeFoldMinutes(normalized.FoldMinutes)
if err != nil {
return nil, err
}
modelQuery := buildModelAccessLogQuery(normalized)
bucketQuery := model.NodeAccessLogBucketQuery{
NodeID: modelQuery.NodeID,
RemoteAddr: modelQuery.RemoteAddr,
Host: modelQuery.Host,
Path: modelQuery.Path,
Since: modelQuery.Since,
Page: normalized.Page,
PageSize: normalized.PageSize,
SortBy: normalizeFoldSortBy(normalized.SortBy),
SortOrder: normalized.SortOrder,
FoldMinutes: foldMinutes,
}
items, err := model.ListNodeAccessLogBuckets(bucketQuery)
if err != nil {
return nil, err
}
totalBuckets, err := model.CountNodeAccessLogBuckets(bucketQuery)
if err != nil {
return nil, err
}
totalRecords, totalIPs, err := model.CountNodeAccessLogs(modelQuery)
if err != nil {
return nil, err
}
views := make([]FoldedAccessLogView, 0, len(items))
for _, item := range items {
if item == nil {
continue
}
views = append(views, FoldedAccessLogView{
BucketStartedAt: time.Unix(item.BucketEpoch, 0).UTC(),
RequestCount: item.RequestCount,
UniqueIPCount: item.UniqueIPCount,
UniqueHostCount: item.UniqueHostCount,
SuccessCount: item.SuccessCount,
ClientErrorCount: item.ClientErrorCount,
ServerErrorCount: item.ServerErrorCount,
})
}
return &FoldedAccessLogList{
Items: views,
Page: normalized.Page,
PageSize: normalized.PageSize,
HasMore: int64((normalized.Page+1)*normalized.PageSize) < totalBuckets,
TotalBucket: totalBuckets,
TotalRecord: totalRecords,
TotalIP: totalIPs,
FoldMinutes: foldMinutes,
}, nil
}
func ListAccessLogIPSummaries(input AccessLogIPSummaryQuery) (*AccessLogIPSummaryList, error) {
normalized := normalizeAccessLogIPSummaryQuery(input)
since := time.Now().UTC().Add(-nodeAccessLogRetentionWindow)
recentSince := time.Now().UTC().Add(-3 * time.Hour)
query := model.NodeAccessLogIPSummaryQuery{
NodeID: strings.TrimSpace(normalized.NodeID),
RemoteAddr: strings.TrimSpace(normalized.RemoteAddr),
Host: strings.TrimSpace(normalized.Host),
Since: since,
Page: normalized.Page,
PageSize: normalized.PageSize,
SortBy: normalized.SortBy,
SortOrder: normalized.SortOrder,
}
items, err := model.ListNodeAccessLogIPSummaries(query, recentSince)
if err != nil {
return nil, err
}
totalIP, err := model.CountNodeAccessLogIPSummaries(query)
if err != nil {
return nil, err
}
views := make([]AccessLogIPSummaryView, 0, len(items))
for _, item := range items {
if item == nil {
continue
}
views = append(views, AccessLogIPSummaryView{
RemoteAddr: item.RemoteAddr,
TotalRequests: item.TotalRequests,
RecentRequests: item.RecentRequests,
LastSeenAt: time.Unix(item.LastSeenEpoch, 0).UTC(),
})
}
return &AccessLogIPSummaryList{
Items: views,
Page: normalized.Page,
PageSize: normalized.PageSize,
HasMore: int64((normalized.Page+1)*normalized.PageSize) < totalIP,
TotalIP: totalIP,
SortBy: normalized.SortBy,
SortOrder: normalized.SortOrder,
}, nil
}
func GetAccessLogIPTrend(input AccessLogIPTrendQuery) (*AccessLogIPTrendView, error) {
normalized, err := normalizeAccessLogIPTrendQuery(input)
if err != nil {
return nil, err
}
points, err := model.ListNodeAccessLogIPTrend(model.NodeAccessLogIPTrendQuery{
NodeID: strings.TrimSpace(normalized.NodeID),
RemoteAddr: strings.TrimSpace(normalized.RemoteAddr),
Host: strings.TrimSpace(normalized.Host),
Since: time.Now().UTC().Add(-time.Duration(normalized.Hours) * time.Hour),
BucketMinutes: normalized.BucketMinutes,
})
if err != nil {
return nil, err
}
pointMap := make(map[int64]int64, len(points))
for _, item := range points {
if item == nil {
continue
}
pointMap[item.BucketEpoch] = item.RequestCount
}
bucketDuration := time.Duration(normalized.BucketMinutes) * time.Minute
start := time.Now().UTC().Add(-time.Duration(normalized.Hours) * time.Hour).Truncate(bucketDuration)
end := time.Now().UTC().Truncate(bucketDuration)
views := make([]AccessLogIPTrendPoint, 0, int(end.Sub(start)/bucketDuration)+1)
for cursor := start; !cursor.After(end); cursor = cursor.Add(bucketDuration) {
views = append(views, AccessLogIPTrendPoint{
BucketStartedAt: cursor,
RequestCount: pointMap[cursor.Unix()],
})
}
return &AccessLogIPTrendView{
RemoteAddr: normalized.RemoteAddr,
Hours: normalized.Hours,
BucketMinutes: normalized.BucketMinutes,
Points: views,
}, nil
}
func CleanupAccessLogs(input AccessLogCleanupInput) (*AccessLogCleanupResult, error) {
if input.RetentionDays <= 0 || input.RetentionDays > nodeAccessLogRetentionDays {
return nil, errors.New("retention_days 必须在 1 到 90 之间")
}
cutoff := time.Now().UTC().Add(-time.Duration(input.RetentionDays) * 24 * time.Hour)
deleted, err := model.DeleteNodeAccessLogsBefore(cutoff)
if err != nil {
return nil, err
}
return &AccessLogCleanupResult{
RetentionDays: input.RetentionDays,
DeletedCount: deleted,
Cutoff: cutoff,
}, nil
}
func buildModelAccessLogQuery(input AccessLogQuery) model.NodeAccessLogQuery {
return model.NodeAccessLogQuery{
NodeID: strings.TrimSpace(input.NodeID),
RemoteAddr: strings.TrimSpace(input.RemoteAddr),
Host: strings.TrimSpace(input.Host),
Path: strings.TrimSpace(input.Path),
Since: time.Now().UTC().Add(-nodeAccessLogRetentionWindow),
Page: input.Page,
PageSize: input.PageSize,
SortBy: input.SortBy,
SortOrder: input.SortOrder,
}
}
func listNodeNameMap(logs []*model.NodeAccessLog) (map[string]string, error) {
nodeIDs := make([]string, 0, len(logs))
seen := make(map[string]struct{}, len(logs))
for _, item := range logs {
if item == nil || item.NodeID == "" {
continue
}
if _, exists := seen[item.NodeID]; exists {
continue
}
seen[item.NodeID] = struct{}{}
nodeIDs = append(nodeIDs, item.NodeID)
}
nodes, err := model.ListNodesByNodeIDs(nodeIDs)
if err != nil {
return nil, err
}
result := make(map[string]string, len(nodes))
for _, node := range nodes {
if node == nil {
continue
}
result[node.NodeID] = node.Name
}
return result, nil
}
func normalizeAccessLogQuery(input AccessLogQuery) AccessLogQuery {
return AccessLogQuery{
NodeID: strings.TrimSpace(input.NodeID),
RemoteAddr: strings.TrimSpace(input.RemoteAddr),
Host: strings.TrimSpace(input.Host),
Path: strings.TrimSpace(input.Path),
Page: normalizeAccessLogPage(input.Page),
PageSize: normalizeAccessLogPageSize(input.PageSize),
SortBy: normalizeAccessLogSortBy(input.SortBy),
SortOrder: normalizeAccessLogSortOrder(input.SortOrder),
FoldMinutes: input.FoldMinutes,
}
}
func normalizeAccessLogIPSummaryQuery(input AccessLogIPSummaryQuery) AccessLogIPSummaryQuery {
return AccessLogIPSummaryQuery{
NodeID: strings.TrimSpace(input.NodeID),
RemoteAddr: strings.TrimSpace(input.RemoteAddr),
Host: strings.TrimSpace(input.Host),
Page: normalizeAccessLogPage(input.Page),
PageSize: normalizeAccessLogPageSize(input.PageSize),
SortBy: normalizeIPSummarySortBy(input.SortBy),
SortOrder: normalizeAccessLogSortOrder(input.SortOrder),
}
}
func normalizeAccessLogIPTrendQuery(input AccessLogIPTrendQuery) (AccessLogIPTrendQuery, error) {
remoteAddr := strings.TrimSpace(input.RemoteAddr)
if remoteAddr == "" {
return AccessLogIPTrendQuery{}, errors.New("remote_addr 不能为空")
}
hours := input.Hours
if hours <= 0 {
hours = defaultIPTrendHours
}
if hours > maxIPTrendHours {
hours = maxIPTrendHours
}
bucketMinutes := input.BucketMinutes
if bucketMinutes <= 0 {
bucketMinutes = defaultIPTrendBucketMinute
}
switch bucketMinutes {
case 5, 10, 15, 30, 60:
default:
return AccessLogIPTrendQuery{}, errors.New("bucket_minutes 仅支持 5、10、15、30、60")
}
return AccessLogIPTrendQuery{
NodeID: strings.TrimSpace(input.NodeID),
RemoteAddr: remoteAddr,
Host: strings.TrimSpace(input.Host),
Hours: hours,
BucketMinutes: bucketMinutes,
}, nil
}
func normalizeAccessLogPage(page int) int {
if page < 0 {
return 0
@@ -121,3 +444,49 @@ func normalizeAccessLogPageSize(pageSize int) int {
}
return pageSize
}
func normalizeAccessLogSortBy(sortBy string) string {
switch strings.TrimSpace(sortBy) {
case "status_code", "remote_addr", "host", "path":
return strings.TrimSpace(sortBy)
default:
return defaultAccessLogSortBy
}
}
func normalizeAccessLogSortOrder(sortOrder string) string {
if strings.EqualFold(strings.TrimSpace(sortOrder), "asc") {
return "asc"
}
return defaultAccessLogSortOrder
}
func normalizeFoldSortBy(sortBy string) string {
switch strings.TrimSpace(sortBy) {
case "request_count":
return "request_count"
default:
return "bucket_started_at"
}
}
func normalizeIPSummarySortBy(sortBy string) string {
switch strings.TrimSpace(sortBy) {
case "recent_requests", "last_seen_at", "remote_addr":
return strings.TrimSpace(sortBy)
default:
return "total_requests"
}
}
func normalizeFoldMinutes(value int) (int, error) {
if value <= 0 {
return defaultAccessLogFoldMinute, nil
}
switch value {
case 3, 5:
return value, nil
default:
return 0, errors.New("fold_minutes 仅支持 3 或 5")
}
}
+122 -3
View File
@@ -64,7 +64,7 @@ func TestListAccessLogsIncludesSummaryTotals(t *testing.T) {
t.Fatalf("failed to seed access logs: %v", err)
}
result, err := ListAccessLogs("", 0, 2)
result, err := ListAccessLogs(AccessLogQuery{Page: 0, PageSize: 2})
if err != nil {
t.Fatalf("ListAccessLogs failed: %v", err)
}
@@ -84,7 +84,7 @@ func TestListAccessLogsIncludesSummaryTotals(t *testing.T) {
t.Fatal("expected has_more to be true")
}
filtered, err := ListAccessLogs("node-a", 0, 50)
filtered, err := ListAccessLogs(AccessLogQuery{NodeID: "node-a", Page: 0, PageSize: 50})
if err != nil {
t.Fatalf("ListAccessLogs filtered failed: %v", err)
}
@@ -125,7 +125,7 @@ func TestListAccessLogsUsesDefaultPageSize(t *testing.T) {
t.Fatalf("failed to seed access logs: %v", err)
}
result, err := ListAccessLogs("", 0, 0)
result, err := ListAccessLogs(AccessLogQuery{})
if err != nil {
t.Fatalf("ListAccessLogs failed: %v", err)
}
@@ -139,3 +139,122 @@ func TestListAccessLogsUsesDefaultPageSize(t *testing.T) {
t.Fatal("expected has_more to be true")
}
}
func TestListFoldedAccessLogsAndIPSummaries(t *testing.T) {
setupServiceTestDB(t)
now := time.Now().UTC()
if err := model.DB.Create(&model.Node{
NodeID: "node-folded",
Name: "edge-folded",
}).Error; err != nil {
t.Fatalf("failed to seed node: %v", err)
}
logs := []*model.NodeAccessLog{
{
NodeID: "node-folded",
LoggedAt: now.Add(-4 * time.Minute),
RemoteAddr: "203.0.113.1",
Host: "alpha.example.com",
Path: "/first",
StatusCode: 200,
},
{
NodeID: "node-folded",
LoggedAt: now.Add(-3 * time.Minute),
RemoteAddr: "203.0.113.1",
Host: "alpha.example.com",
Path: "/second",
StatusCode: 502,
},
{
NodeID: "node-folded",
LoggedAt: now.Add(-2 * time.Minute),
RemoteAddr: "203.0.113.2",
Host: "beta.example.com",
Path: "/third",
StatusCode: 404,
},
}
if err := model.DB.Create(&logs).Error; err != nil {
t.Fatalf("failed to seed access logs: %v", err)
}
folded, err := ListFoldedAccessLogs(AccessLogQuery{
NodeID: "node-folded",
Page: 0,
PageSize: 10,
FoldMinutes: 5,
})
if err != nil {
t.Fatalf("ListFoldedAccessLogs failed: %v", err)
}
if len(folded.Items) != 2 {
t.Fatalf("expected two folded buckets, got %+v", folded.Items)
}
if folded.TotalRecord != 3 || folded.TotalBucket != 2 {
t.Fatalf("unexpected folded totals: %+v", folded)
}
if folded.Items[0].RequestCount+folded.Items[1].RequestCount != 3 {
t.Fatalf("unexpected folded request count sum: %+v", folded.Items)
}
ipSummaries, err := ListAccessLogIPSummaries(AccessLogIPSummaryQuery{
NodeID: "node-folded",
Page: 0,
PageSize: 10,
SortBy: "total_requests",
SortOrder: "desc",
})
if err != nil {
t.Fatalf("ListAccessLogIPSummaries failed: %v", err)
}
if len(ipSummaries.Items) != 2 {
t.Fatalf("expected two ip summary rows, got %+v", ipSummaries.Items)
}
if ipSummaries.Items[0].RemoteAddr != "203.0.113.1" || ipSummaries.Items[0].TotalRequests != 2 {
t.Fatalf("unexpected top ip summary row: %+v", ipSummaries.Items[0])
}
}
func TestCleanupAccessLogsDeletesExpiredData(t *testing.T) {
setupServiceTestDB(t)
now := time.Now().UTC()
if err := model.DB.Create([]*model.NodeAccessLog{
{
NodeID: "node-cleanup",
LoggedAt: now.Add(-10 * 24 * time.Hour),
RemoteAddr: "203.0.113.9",
Host: "cleanup.example.com",
Path: "/old",
StatusCode: 200,
},
{
NodeID: "node-cleanup",
LoggedAt: now.Add(-2 * 24 * time.Hour),
RemoteAddr: "203.0.113.10",
Host: "cleanup.example.com",
Path: "/recent",
StatusCode: 200,
},
}).Error; err != nil {
t.Fatalf("failed to seed cleanup logs: %v", err)
}
result, err := CleanupAccessLogs(AccessLogCleanupInput{RetentionDays: 7})
if err != nil {
t.Fatalf("CleanupAccessLogs failed: %v", err)
}
if result.DeletedCount != 1 {
t.Fatalf("expected 1 deleted record, got %+v", result)
}
remaining, err := ListAccessLogs(AccessLogQuery{Page: 0, PageSize: 10, NodeID: "node-cleanup"})
if err != nil {
t.Fatalf("ListAccessLogs failed after cleanup: %v", err)
}
if len(remaining.Items) != 1 || remaining.Items[0].Path != "/recent" {
t.Fatalf("unexpected remaining logs after cleanup: %+v", remaining.Items)
}
}
+22 -4
View File
@@ -708,7 +708,12 @@ func TestHeartbeatNodePersistsObservabilityPayload(t *testing.T) {
t.Fatalf("unexpected request reports: %+v", reports)
}
accessLogs, err := model.ListNodeAccessLogs(node.NodeID, time.Time{}, 0, 10)
accessLogs, err := model.ListNodeAccessLogs(model.NodeAccessLogQuery{
NodeID: node.NodeID,
Since: time.Time{},
Page: 0,
PageSize: 10,
})
if err != nil {
t.Fatalf("expected node access logs query to succeed: %v", err)
}
@@ -822,7 +827,12 @@ func TestHeartbeatNodePersistsBufferedObservabilityPayload(t *testing.T) {
t.Fatalf("expected current and buffered reports, got %+v", reports)
}
accessLogs, err := model.ListNodeAccessLogs(node.NodeID, time.Time{}, 0, 10)
accessLogs, err := model.ListNodeAccessLogs(model.NodeAccessLogQuery{
NodeID: node.NodeID,
Since: time.Time{},
Page: 0,
PageSize: 10,
})
if err != nil {
t.Fatalf("expected node access logs query to succeed: %v", err)
}
@@ -930,7 +940,11 @@ func TestListAccessLogsUsesPagination(t *testing.T) {
t.Fatalf("failed to seed access logs: %v", err)
}
pageOne, err := ListAccessLogs(node.NodeID, 0, 2)
pageOne, err := ListAccessLogs(AccessLogQuery{
NodeID: node.NodeID,
Page: 0,
PageSize: 2,
})
if err != nil {
t.Fatalf("ListAccessLogs page 1 failed: %v", err)
}
@@ -944,7 +958,11 @@ func TestListAccessLogsUsesPagination(t *testing.T) {
t.Fatalf("expected paged access log region to be returned, got %+v", pageOne.Items[0])
}
pageTwo, err := ListAccessLogs(node.NodeID, 1, 2)
pageTwo, err := ListAccessLogs(AccessLogQuery{
NodeID: node.NodeID,
Page: 1,
PageSize: 2,
})
if err != nil {
t.Fatalf("ListAccessLogs page 2 failed: %v", err)
}
+1 -1
View File
@@ -17,7 +17,7 @@ const (
NodeHealthSeverityInfo = "info"
NodeHealthSeverityWarning = "warning"
NodeHealthSeverityCritical = "critical"
nodeAccessLogRetentionWindow = 24 * time.Hour
nodeAccessLogRetentionWindow = nodeAccessLogRetentionDays * 24 * time.Hour
)
type AgentNodeSystemProfile struct {
@@ -1,15 +1,56 @@
import { apiRequest } from '@/lib/api/client';
import type { AccessLogList } from '@/features/access-logs/types';
import type {
AccessLogCleanupPayload,
AccessLogCleanupResult,
AccessLogFilters,
AccessLogIPSummaryFilters,
AccessLogIPSummaryList,
AccessLogIPTrend,
AccessLogIPTrendFilters,
AccessLogList,
FoldedAccessLogFilters,
FoldedAccessLogList,
} from '@/features/access-logs/types';
export function getAccessLogs(page: number, nodeId?: string, pageSize = 20) {
const normalizedNodeId = nodeId?.trim();
const searchParams = new URLSearchParams({
p: String(Math.max(page, 0)),
page_size: String(pageSize),
function buildSearchParams(filters: object) {
const searchParams = new URLSearchParams();
Object.entries(filters as Record<string, string | number | undefined>).forEach(([key, value]) => {
if (value === undefined || value === null || value === '') {
return;
}
searchParams.set(key, String(value));
});
return searchParams.toString();
}
export function getAccessLogs(filters: AccessLogFilters) {
const query = buildSearchParams(filters);
return apiRequest<AccessLogList>(`/access-logs/${query ? `?${query}` : ''}`);
}
export function getFoldedAccessLogs(filters: FoldedAccessLogFilters) {
const query = buildSearchParams(filters);
return apiRequest<FoldedAccessLogList>(`/access-logs/folds${query ? `?${query}` : ''}`);
}
export function getAccessLogIPSummaries(filters: AccessLogIPSummaryFilters) {
const query = buildSearchParams(filters);
return apiRequest<AccessLogIPSummaryList>(
`/access-logs/ip-summary${query ? `?${query}` : ''}`,
);
}
export function getAccessLogIPTrend(filters: AccessLogIPTrendFilters) {
const query = buildSearchParams(filters);
return apiRequest<AccessLogIPTrend>(
`/access-logs/ip-summary/trend${query ? `?${query}` : ''}`,
);
}
export function cleanupAccessLogs(payload: AccessLogCleanupPayload) {
return apiRequest<AccessLogCleanupResult>('/access-logs/cleanup', {
method: 'POST',
body: JSON.stringify(payload),
});
if (normalizedNodeId) {
searchParams.set('node_id', normalizedNodeId);
}
return apiRequest<AccessLogList>(`/access-logs/?${searchParams.toString()}`);
}
File diff suppressed because it is too large Load Diff
@@ -1,3 +1,14 @@
export interface AccessLogFilters {
node_id?: string;
remote_addr?: string;
host?: string;
path?: string;
p?: number;
page_size?: number;
sort_by?: string;
sort_order?: 'asc' | 'desc';
}
export interface AccessLogItem {
id: number;
node_id: string;
@@ -18,3 +29,85 @@ export interface AccessLogList {
total_record: number;
total_ip: number;
}
export interface FoldedAccessLogFilters extends AccessLogFilters {
fold_minutes: 3 | 5;
}
export interface FoldedAccessLogItem {
bucket_started_at: string;
request_count: number;
unique_ip_count: number;
unique_host_count: number;
success_count: number;
client_error_count: number;
server_error_count: number;
}
export interface FoldedAccessLogList {
items: FoldedAccessLogItem[];
page: number;
page_size: number;
has_more: boolean;
total_bucket: number;
total_record: number;
total_ip: number;
fold_minutes: number;
}
export interface AccessLogIPSummaryFilters {
node_id?: string;
remote_addr?: string;
host?: string;
p?: number;
page_size?: number;
sort_by?: string;
sort_order?: 'asc' | 'desc';
}
export interface AccessLogIPSummaryItem {
remote_addr: string;
total_requests: number;
recent_requests: number;
last_seen_at: string;
}
export interface AccessLogIPSummaryList {
items: AccessLogIPSummaryItem[];
page: number;
page_size: number;
has_more: boolean;
total_ip: number;
sort_by: string;
sort_order: 'asc' | 'desc';
}
export interface AccessLogIPTrendFilters {
node_id?: string;
remote_addr: string;
host?: string;
hours?: number;
bucket_minutes?: number;
}
export interface AccessLogIPTrendPoint {
bucket_started_at: string;
request_count: number;
}
export interface AccessLogIPTrend {
remote_addr: string;
hours: number;
bucket_minutes: number;
points: AccessLogIPTrendPoint[];
}
export interface AccessLogCleanupPayload {
retention_days: number;
}
export interface AccessLogCleanupResult {
retention_days: number;
deleted_count: number;
cutoff: string;
}