fix: 初验针对修改

This commit is contained in:
zxr
2026-09-09 16:45:01 +08:00
parent 38a86e1fd8
commit b0ffde3100
25 changed files with 1531 additions and 158 deletions

View File

@@ -1,6 +1,7 @@
package ingest
import (
"context"
"encoding/json"
"errors"
"fmt"
@@ -63,17 +64,36 @@ func enqueuePayloadWithDB(db *gorm.DB, logEventID uint, payloadJSON string) (uin
return row.ID, nil
}
func StartAlertDispatcher() {
func StartAlertDispatcher(ctx context.Context) {
owner := dispatcherOwner()
go func() {
startWorker("alert_dispatcher", func() error {
ticker := time.NewTicker(2 * time.Second)
defer ticker.Stop()
for range ticker.C {
if _, err := ProcessAlertOutboxBatch(impl.DBService, 20, owner); err != nil {
log.Printf("logs: alert outbox dispatch: %v", err)
for {
select {
case <-ctx.Done():
return nil
case <-ticker.C:
if _, err := ProcessAlertOutboxBatch(impl.DBService, 20, owner); err != nil {
markWorkerError("alert_dispatcher", err)
log.Printf("logs: alert outbox dispatch: %v", err)
continue
}
if depth, err := alertOutboxQueueDepth(ctx); err == nil {
markWorkerQueueDepth("alert_dispatcher", depth)
}
markWorkerSucceeded("alert_dispatcher")
}
}
}()
})
}
func alertOutboxQueueDepth(ctx context.Context) (int64, error) {
var depth int64
err := impl.DBService.WithContext(ctx).Model(&models.AlertOutbox{}).
Where("status IN ?", []string{outboxStatusPending, outboxStatusRetrying, outboxStatusProcessing}).
Count(&depth).Error
return depth, err
}
func dispatcherOwner() string {

View File

@@ -1,6 +1,7 @@
package ingest
import (
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
@@ -138,19 +139,29 @@ func (e *Engine) Refresh() error {
return nil
}
func StartRefresher() {
func StartRefresher(ctx context.Context) error {
interval := config.Spec.Ingest.RuleRefreshSecs
if interval <= 0 {
interval = 30
if err := Global.Refresh(); err != nil {
return err
}
_ = Global.Refresh()
go func() {
startWorker("rule_refresher", func() error {
markWorkerSucceeded("rule_refresher")
t := time.NewTicker(time.Duration(interval) * time.Second)
defer t.Stop()
for range t.C {
_ = Global.Refresh()
for {
select {
case <-ctx.Done():
return nil
case <-t.C:
if err := Global.Refresh(); err != nil {
markWorkerError("rule_refresher", err)
continue
}
markWorkerSucceeded("rule_refresher")
}
}
}()
})
return nil
}
func normOID(s string) string {

View File

@@ -0,0 +1,106 @@
package ingest
import (
"sync"
"time"
)
// WorkerStatus 描述日志后台 worker 的当前运行状态。
type WorkerStatus struct {
Name string `json:"name"`
Running bool `json:"running"`
LastStartedAt time.Time `json:"last_started_at,omitempty"`
LastSucceededAt time.Time `json:"last_succeeded_at,omitempty"`
LastErrorAt time.Time `json:"last_error_at,omitempty"`
LastError string `json:"last_error,omitempty"`
QueueDepth *int64 `json:"queue_depth,omitempty"`
}
func markWorkerQueueDepth(name string, depth int64) {
lifecycle.Lock()
status := lifecycle.workers[name]
status.Name = name
status.QueueDepth = new(int64)
*status.QueueDepth = depth
lifecycle.workers[name] = status
lifecycle.Unlock()
}
var lifecycle = struct {
sync.RWMutex
wg sync.WaitGroup
workers map[string]WorkerStatus
}{workers: make(map[string]WorkerStatus)}
func startWorker(name string, run func() error) {
lifecycle.Lock()
lifecycle.workers[name] = WorkerStatus{Name: name, Running: true, LastStartedAt: time.Now().UTC()}
lifecycle.wg.Add(1)
lifecycle.Unlock()
go func() {
defer lifecycle.wg.Done()
err := run()
lifecycle.Lock()
status := lifecycle.workers[name]
status.Running = false
if err != nil {
status.LastErrorAt = time.Now().UTC()
status.LastError = err.Error()
}
lifecycle.workers[name] = status
lifecycle.Unlock()
}()
}
func markWorkerSucceeded(name string) {
lifecycle.Lock()
status := lifecycle.workers[name]
status.Name = name
status.LastSucceededAt = time.Now().UTC()
status.LastError = ""
lifecycle.workers[name] = status
lifecycle.Unlock()
}
func markWorkerError(name string, err error) {
if err == nil {
return
}
lifecycle.Lock()
status := lifecycle.workers[name]
status.Name = name
status.LastErrorAt = time.Now().UTC()
status.LastError = err.Error()
lifecycle.workers[name] = status
lifecycle.Unlock()
}
// WorkerStatuses 返回状态快照。
func WorkerStatuses() []WorkerStatus {
lifecycle.RLock()
defer lifecycle.RUnlock()
result := make([]WorkerStatus, 0, len(lifecycle.workers))
for _, status := range lifecycle.workers {
if status.QueueDepth != nil {
depth := *status.QueueDepth
status.QueueDepth = &depth
}
result = append(result, status)
}
return result
}
// Wait 等待后台 worker 退出,返回是否在超时内完成。
func Wait(timeout time.Duration) bool {
done := make(chan struct{})
go func() {
lifecycle.wg.Wait()
close(done)
}()
select {
case <-done:
return true
case <-time.After(timeout):
return false
}
}

View File

@@ -1,41 +1,52 @@
package ingest
import (
"log"
"net"
"strings"
"git.apinb.com/ops/logs/internal/config"
)
func StartSyslogUDP() {
addr := strings.TrimSpace(config.Spec.Ingest.SyslogListenAddr)
if addr == "" {
return
}
go func() {
pc, err := net.ListenPacket("udp", addr)
if err != nil {
log.Printf("logs: syslog UDP listen %s: %v", addr, err)
return
}
defer pc.Close()
log.Printf("logs: syslog listening UDP %s", addr)
buf := make([]byte, 65536)
for {
n, remote, err := pc.ReadFrom(buf)
if err != nil {
log.Printf("logs: syslog read: %v", err)
continue
}
udpAddr, _ := remote.(*net.UDPAddr)
if udpAddr == nil {
continue
}
p := make([]byte, n)
copy(p, buf[:n])
a := *udpAddr
Global.HandleSyslog(&a, p)
}
}()
}
package ingest
import (
"context"
"fmt"
"log"
"net"
"strings"
"git.apinb.com/ops/logs/internal/config"
)
func StartSyslogUDP(ctx context.Context) error {
addr := strings.TrimSpace(config.Spec.Ingest.SyslogListenAddr)
if addr == "" {
return nil
}
pc, err := net.ListenPacket("udp", addr)
if err != nil {
return fmt.Errorf("syslog UDP 监听 %s 失败: %w", addr, err)
}
startWorker("syslog_udp", func() error {
defer pc.Close()
go func() {
<-ctx.Done()
_ = pc.Close()
}()
log.Printf("logs: syslog listening UDP %s", addr)
buf := make([]byte, 65536)
for {
n, remote, err := pc.ReadFrom(buf)
if err != nil {
if ctx.Err() != nil {
return nil
}
markWorkerError("syslog_udp", err)
log.Printf("logs: syslog read: %v", err)
continue
}
udpAddr, _ := remote.(*net.UDPAddr)
if udpAddr == nil {
continue
}
p := make([]byte, n)
copy(p, buf[:n])
a := *udpAddr
Global.HandleSyslog(&a, p)
markWorkerSucceeded("syslog_udp")
}
})
return nil
}

View File

@@ -1,32 +1,58 @@
package ingest
import (
"log"
"net"
"strings"
"git.apinb.com/ops/logs/internal/config"
"github.com/gosnmp/gosnmp"
)
func StartTrapUDP() {
addr := strings.TrimSpace(config.Spec.Ingest.TrapListenAddr)
if addr == "" {
return
}
go func() {
tl := gosnmp.NewTrapListener()
tl.OnNewTrap = func(pkt *gosnmp.SnmpPacket, u *net.UDPAddr) {
if u == nil || pkt == nil {
return
}
ua := *u
Global.HandleTrap(&ua, pkt)
}
tl.Params = gosnmp.Default
tl.Params.Logger = gosnmp.NewLogger(log.Default())
if err := tl.Listen(addr); err != nil {
log.Printf("logs: trap listener %s: %v", addr, err)
}
}()
}
package ingest
import (
"context"
"fmt"
"log"
"net"
"strings"
"time"
"git.apinb.com/ops/logs/internal/config"
"github.com/gosnmp/gosnmp"
)
func StartTrapUDP(ctx context.Context) error {
addr := strings.TrimSpace(config.Spec.Ingest.TrapListenAddr)
if addr == "" {
return nil
}
tl := gosnmp.NewTrapListener()
tl.OnNewTrap = func(pkt *gosnmp.SnmpPacket, u *net.UDPAddr) {
if u == nil || pkt == nil {
return
}
ua := *u
Global.HandleTrap(&ua, pkt)
markWorkerSucceeded("trap_udp")
}
tl.Params = gosnmp.Default
tl.Params.Logger = gosnmp.NewLogger(log.Default())
started := make(chan error, 1)
startWorker("trap_udp", func() error {
go func() {
<-ctx.Done()
tl.Close()
}()
err := tl.Listen(addr)
if ctx.Err() != nil {
return nil
}
return err
})
go func() {
select {
case <-tl.Listening():
started <- nil
case <-ctx.Done():
started <- ctx.Err()
case <-time.After(5 * time.Second):
tl.Close()
started <- fmt.Errorf("Trap UDP 监听 %s 启动超时", addr)
}
}()
if err := <-started; err != nil {
return err
}
return nil
}