vibe-proxy/backend/internal/logging/home_app_log_forwarder.go
2026-08-24 00:10:41 +02:00

296 lines
6.3 KiB
Go

package logging
import (
"context"
"encoding/json"
"errors"
"strings"
"sync"
"sync/atomic"
"time"
"github.com/router-for-me/CLIProxyAPI/v7/internal/home"
log "github.com/sirupsen/logrus"
)
const defaultHomeAppLogQueueSize = 1024
type homeAppLogClient interface {
HeartbeatOK() bool
RPushAppLog(ctx context.Context, payload []byte) error
}
type homeAppLogPayload struct {
Line string `json:"line"`
Level string `json:"level,omitempty"`
Timestamp string `json:"timestamp,omitempty"`
RequestID string `json:"request_id,omitempty"`
client homeAppLogClient
}
// HomeAppLogForwarder forwards application logs to Home after the control connection is healthy.
type HomeAppLogForwarder struct {
formatter log.Formatter
queue chan homeAppLogPayload
stop chan struct{}
stopOnce sync.Once
wg sync.WaitGroup
enabled atomic.Bool
stopped atomic.Bool
ownerMu sync.Mutex
owner homeAppLogClient
}
type homeAppLogMux struct {
mu sync.Mutex
targets map[*HomeAppLogForwarder]struct{}
}
func (h *homeAppLogMux) Levels() []log.Level {
return log.AllLevels
}
func (h *homeAppLogMux) Fire(entry *log.Entry) error {
h.mu.Lock()
targets := make([]*HomeAppLogForwarder, 0, len(h.targets))
for target := range h.targets {
targets = append(targets, target)
}
h.mu.Unlock()
for _, target := range targets {
if errFire := target.Fire(entry); errFire != nil {
return errFire
}
}
return nil
}
func (h *homeAppLogMux) register(target *HomeAppLogForwarder) {
if target == nil {
return
}
h.mu.Lock()
defer h.mu.Unlock()
if h.targets == nil {
h.targets = make(map[*HomeAppLogForwarder]struct{})
}
h.targets[target] = struct{}{}
}
func (h *homeAppLogMux) unregister(target *HomeAppLogForwarder) {
if target == nil {
return
}
h.mu.Lock()
delete(h.targets, target)
h.mu.Unlock()
}
var (
homeAppLogMuxHook = &homeAppLogMux{}
homeAppLogMuxInstallOnce sync.Once
)
func registerHomeAppLogForwarder(forwarder *HomeAppLogForwarder) {
homeAppLogMuxInstallOnce.Do(func() {
log.AddHook(homeAppLogMuxHook)
})
homeAppLogMuxHook.register(forwarder)
}
// StartHomeAppLogForwarder registers a Home log forwarding target with the process-wide logrus hook.
func StartHomeAppLogForwarder(queueSize int) *HomeAppLogForwarder {
if queueSize <= 0 {
queueSize = defaultHomeAppLogQueueSize
}
forwarder := &HomeAppLogForwarder{
formatter: &LogFormatter{},
queue: make(chan homeAppLogPayload, queueSize),
stop: make(chan struct{}),
}
forwarder.enabled.Store(true)
forwarder.wg.Add(1)
go forwarder.run()
registerHomeAppLogForwarder(forwarder)
return forwarder
}
// Stop disables forwarding and waits for the background sender to exit.
func (f *HomeAppLogForwarder) Stop() {
if f == nil {
return
}
f.stopOnce.Do(func() {
f.stopped.Store(true)
f.ownerMu.Lock()
f.owner = nil
f.ownerMu.Unlock()
f.enabled.Store(false)
homeAppLogMuxHook.unregister(f)
close(f.stop)
f.wg.Wait()
})
}
// Bind activates forwarding to client.
func (f *HomeAppLogForwarder) Bind(client *home.Client) {
f.bind(client)
}
func (f *HomeAppLogForwarder) bind(client homeAppLogClient) {
if f == nil || client == nil || f.stopped.Load() {
return
}
f.ownerMu.Lock()
defer f.ownerMu.Unlock()
if f.stopped.Load() {
return
}
f.owner = client
f.enabled.Store(true)
}
// Deactivate stops forwarding only when client owns the forwarder.
func (f *HomeAppLogForwarder) Deactivate(client *home.Client) {
f.deactivate(client)
}
func (f *HomeAppLogForwarder) deactivate(client homeAppLogClient) {
if f == nil || client == nil {
return
}
f.ownerMu.Lock()
if f.owner == client {
f.owner = nil
}
f.ownerMu.Unlock()
}
func (f *HomeAppLogForwarder) client() homeAppLogClient {
f.ownerMu.Lock()
defer f.ownerMu.Unlock()
return f.owner
}
// Levels implements logrus.Hook.
func (f *HomeAppLogForwarder) Levels() []log.Level {
return log.AllLevels
}
// Fire implements logrus.Hook.
func (f *HomeAppLogForwarder) Fire(entry *log.Entry) error {
if f == nil || entry == nil || !f.enabled.Load() {
return nil
}
client := f.client()
if client == nil || !client.HeartbeatOK() {
return nil
}
line, errFormat := f.formatEntry(entry)
if errFormat != nil || strings.TrimSpace(line) == "" {
return nil
}
payload := homeAppLogPayload{
Line: line,
Level: entry.Level.String(),
Timestamp: entry.Time.Format(time.RFC3339Nano),
RequestID: appLogRequestID(entry),
client: client,
}
select {
case f.queue <- payload:
default:
}
return nil
}
func appLogRequestID(entry *log.Entry) string {
if entry == nil {
return ""
}
requestID, _ := entry.Data["request_id"].(string)
requestID = strings.TrimSpace(requestID)
if requestID == "--------" {
return ""
}
return requestID
}
func (f *HomeAppLogForwarder) formatEntry(entry *log.Entry) (string, error) {
formatter := f.formatter
if formatter == nil {
formatter = &LogFormatter{}
}
raw, errFormat := formatter.Format(entry)
if errFormat != nil {
return "", errFormat
}
return string(raw), nil
}
func (f *HomeAppLogForwarder) run() {
defer f.wg.Done()
for {
select {
case <-f.stop:
return
case payload := <-f.queue:
f.forward(payload)
}
}
}
func (f *HomeAppLogForwarder) forward(payload homeAppLogPayload) {
client := payload.client
if client == nil {
client = f.client()
}
if !f.enabled.Load() || client == nil || f.client() != client {
return
}
if !client.HeartbeatOK() {
return
}
raw, errMarshal := json.Marshal(&payload)
if errMarshal != nil {
return
}
if errPush := client.RPushAppLog(context.Background(), raw); errPush != nil && isHomeAppLogUnsupported(errPush) {
f.disableIfCurrentOwner(client)
}
}
func (f *HomeAppLogForwarder) disableIfCurrentOwner(client homeAppLogClient) {
f.ownerMu.Lock()
defer f.ownerMu.Unlock()
if f.owner != client {
return
}
f.enabled.Store(false)
}
func isHomeAppLogUnsupported(err error) bool {
if err == nil {
return false
}
msg := strings.ToLower(strings.TrimSpace(err.Error()))
if msg == "" {
return false
}
for {
switch {
case strings.Contains(msg, "unsupported key"):
return true
case strings.Contains(msg, "unknown command"):
return true
case strings.Contains(msg, "unsupported command"):
return true
}
err = errors.Unwrap(err)
if err == nil {
return false
}
msg = strings.ToLower(strings.TrimSpace(err.Error()))
}
}