Refactor logging to use slog package across the application

- Replaced standard log package with log/slog in httpclient, nginx manager, sync service, updater, and other components for structured logging.
- Introduced environment variable `LOG_LEVEL` to control logging levels (debug, info, warn, error).
- Updated documentation to reflect changes in logging configuration and requirements.
- Added a new logging setup function in the ats_agent internal package to initialize the slog logger.
This commit is contained in:
ryan
2026-03-13 14:53:51 +08:00
parent 06d4831d55
commit aeb7118b30
17 changed files with 321 additions and 209 deletions
+19 -6
View File
@@ -3,7 +3,8 @@ package main
import ( import (
"context" "context"
"flag" "flag"
"log" "log/slog"
"os"
"os/signal" "os/signal"
"syscall" "syscall"
@@ -11,6 +12,7 @@ import (
"atsflare-agent/internal/config" "atsflare-agent/internal/config"
"atsflare-agent/internal/heartbeat" "atsflare-agent/internal/heartbeat"
"atsflare-agent/internal/httpclient" "atsflare-agent/internal/httpclient"
"atsflare-agent/internal/logging"
"atsflare-agent/internal/nginx" "atsflare-agent/internal/nginx"
"atsflare-agent/internal/state" "atsflare-agent/internal/state"
syncservice "atsflare-agent/internal/sync" syncservice "atsflare-agent/internal/sync"
@@ -18,12 +20,15 @@ import (
) )
func main() { func main() {
logging.Setup()
configPath := flag.String("config", "./agent.json", "agent config path") configPath := flag.String("config", "./agent.json", "agent config path")
flag.Parse() flag.Parse()
cfg, err := config.Load(*configPath) cfg, err := config.Load(*configPath)
if err != nil { if err != nil {
log.Fatal(err) slog.Error("load agent config failed", "error", err)
os.Exit(1)
} }
cfg.NginxVersion = nginx.DetectVersion( cfg.NginxVersion = nginx.DetectVersion(
context.Background(), context.Background(),
@@ -38,7 +43,14 @@ func main() {
NginxCertDir: cfg.OpenrestyCertDir, NginxCertDir: cfg.OpenrestyCertDir,
}, },
) )
log.Printf("agent config loaded: server=%s node=%s ip=%s heartbeat_interval=%s route_config=%s cert_dir=%s", cfg.ServerURL, cfg.NodeName, cfg.NodeIP, cfg.HeartbeatInterval, cfg.RouteConfigPath, cfg.CertDir) slog.Info("agent config loaded",
"server", cfg.ServerURL,
"node", cfg.NodeName,
"ip", cfg.NodeIP,
"heartbeat_interval", cfg.HeartbeatInterval,
"route_config", cfg.RouteConfigPath,
"cert_dir", cfg.CertDir,
)
client := httpclient.New(cfg.ServerURL, cfg.InitialAuthToken(), cfg.RequestTimeout.Duration()) client := httpclient.New(cfg.ServerURL, cfg.InitialAuthToken(), cfg.RequestTimeout.Duration())
stateStore := state.NewStore(cfg.StatePath) stateStore := state.NewStore(cfg.StatePath)
@@ -74,10 +86,11 @@ func main() {
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer stop() defer stop()
log.Printf("agent process started") slog.Info("agent process started")
if err = runner.Run(ctx); err != nil && err != context.Canceled { if err = runner.Run(ctx); err != nil && err != context.Canceled {
log.Fatal(err) slog.Error("agent process exited with error", "error", err)
os.Exit(1)
} }
log.Printf("agent process stopped") slog.Info("agent process stopped")
} }
+1 -1
View File
@@ -1,3 +1,3 @@
module atsflare-agent module atsflare-agent
go 1.18 go 1.23.0
+28 -28
View File
@@ -3,7 +3,7 @@ package agent
import ( import (
"context" "context"
"errors" "errors"
"log" "log/slog"
"strings" "strings"
"time" "time"
@@ -59,29 +59,29 @@ func (r *Runner) Run(ctx context.Context) error {
if err != nil { if err != nil {
return err return err
} }
log.Printf("agent runner started: node_id=%s node=%s ip=%s", nodeID, r.Config.NodeName, r.Config.NodeIP) slog.Info("agent runner started", "node_id", nodeID, "node", r.Config.NodeName, "ip", r.Config.NodeIP)
if r.hasAgentToken() { if r.hasAgentToken() {
r.refreshOpenrestyHealth(ctx) r.refreshOpenrestyHealth(ctx)
heartbeatResult, hbErr := r.HeartbeatService.Heartbeat(ctx, r.nodePayload(nodeID)) heartbeatResult, hbErr := r.HeartbeatService.Heartbeat(ctx, r.nodePayload(nodeID))
if hbErr != nil { if hbErr != nil {
log.Printf("agent startup heartbeat failed: %v", hbErr) slog.Error("agent startup heartbeat failed", "error", hbErr)
} else { } else {
if heartbeatResult == nil { if heartbeatResult == nil {
heartbeatResult = &protocol.HeartbeatResult{} heartbeatResult = &protocol.HeartbeatResult{}
} }
log.Printf("agent startup heartbeat succeeded: node_id=%s", nodeID) slog.Info("agent startup heartbeat succeeded", "node_id", nodeID)
r.applySettings(heartbeatResult.AgentSettings) r.applySettings(heartbeatResult.AgentSettings)
if err = r.SyncService.SyncOnStartup(ctx, heartbeatResult.ActiveConfig); err != nil { if err = r.SyncService.SyncOnStartup(ctx, heartbeatResult.ActiveConfig); err != nil {
r.recordSyncError(err) r.recordSyncError(err)
log.Printf("agent startup sync failed: %v", err) slog.Error("agent startup sync failed", "error", err)
} else { } else {
log.Printf("agent startup sync completed") slog.Info("agent startup sync completed")
} }
r.tryRestartOpenresty(ctx) r.tryRestartOpenresty(ctx)
r.tryAutoUpdate(ctx) r.tryAutoUpdate(ctx)
} }
} else if err = r.tryRegister(ctx, &nodeID); err != nil { } else if err = r.tryRegister(ctx, &nodeID); err != nil {
log.Printf("agent initial discovery register failed: %v", err) slog.Error("agent initial discovery register failed", "error", err)
} }
heartbeatTicker := time.NewTicker(r.Config.HeartbeatInterval.Duration()) heartbeatTicker := time.NewTicker(r.Config.HeartbeatInterval.Duration())
@@ -90,19 +90,19 @@ func (r *Runner) Run(ctx context.Context) error {
for { for {
select { select {
case <-ctx.Done(): case <-ctx.Done():
log.Printf("agent runner shutting down: %v", ctx.Err()) slog.Info("agent runner shutting down", "error", ctx.Err())
return ctx.Err() return ctx.Err()
case <-heartbeatTicker.C: case <-heartbeatTicker.C:
if !r.hasAgentToken() { if !r.hasAgentToken() {
if err = r.tryRegister(ctx, &nodeID); err != nil { if err = r.tryRegister(ctx, &nodeID); err != nil {
log.Printf("agent discovery register failed: %v", err) slog.Error("agent discovery register failed", "error", err)
} }
continue continue
} }
r.refreshOpenrestyHealth(ctx) r.refreshOpenrestyHealth(ctx)
heartbeatResult, hbErr := r.HeartbeatService.Heartbeat(ctx, r.nodePayload(nodeID)) heartbeatResult, hbErr := r.HeartbeatService.Heartbeat(ctx, r.nodePayload(nodeID))
if hbErr != nil { if hbErr != nil {
log.Printf("agent heartbeat failed: %v", hbErr) slog.Error("agent heartbeat failed", "error", hbErr)
} else { } else {
if heartbeatResult == nil { if heartbeatResult == nil {
heartbeatResult = &protocol.HeartbeatResult{} heartbeatResult = &protocol.HeartbeatResult{}
@@ -112,7 +112,7 @@ func (r *Runner) Run(ctx context.Context) error {
} }
if err = r.SyncService.SyncOnce(ctx, heartbeatResult.ActiveConfig); err != nil { if err = r.SyncService.SyncOnce(ctx, heartbeatResult.ActiveConfig); err != nil {
r.recordSyncError(err) r.recordSyncError(err)
log.Printf("agent sync failed: %v", err) slog.Error("agent sync failed", "error", err)
} }
r.tryRestartOpenresty(ctx) r.tryRestartOpenresty(ctx)
r.tryAutoUpdate(ctx) r.tryAutoUpdate(ctx)
@@ -133,7 +133,7 @@ func (r *Runner) applySettings(settings *protocol.AgentSettings) bool {
if settings.HeartbeatInterval > 0 { if settings.HeartbeatInterval > 0 {
newInterval := config.MillisecondDuration(time.Duration(settings.HeartbeatInterval) * time.Millisecond) newInterval := config.MillisecondDuration(time.Duration(settings.HeartbeatInterval) * time.Millisecond)
if newInterval != r.Config.HeartbeatInterval { if newInterval != r.Config.HeartbeatInterval {
log.Printf("agent heartbeat interval updated: %s -> %s", r.Config.HeartbeatInterval, newInterval) slog.Info("agent heartbeat interval updated", "from", r.Config.HeartbeatInterval, "to", newInterval)
r.Config.HeartbeatInterval = newInterval r.Config.HeartbeatInterval = newInterval
changed = true changed = true
} }
@@ -155,13 +155,13 @@ func (r *Runner) tryRestartOpenresty(ctx context.Context) {
if r.RuntimeManager == nil { if r.RuntimeManager == nil {
return return
} }
log.Printf("agent openresty restart requested by server") slog.Info("agent openresty restart requested by server")
if err := r.RuntimeManager.Restart(ctx); err != nil { if err := r.RuntimeManager.Restart(ctx); err != nil {
log.Printf("agent openresty restart failed: %v", err) slog.Error("agent openresty restart failed", "error", err)
r.recordOpenrestyUnhealthy(err, false) r.recordOpenrestyUnhealthy(err, false)
return return
} }
log.Printf("agent openresty restart succeeded") slog.Info("agent openresty restart succeeded")
r.recordOpenrestyHealthy() r.recordOpenrestyHealthy()
} }
@@ -182,7 +182,7 @@ func (r *Runner) tryAutoUpdate(ctx context.Context) {
TagName: r.updateTag, TagName: r.updateTag,
Force: force, Force: force,
}); err != nil { }); err != nil {
log.Printf("agent update check failed: %v", err) slog.Error("agent update check failed", "error", err)
} }
if force { if force {
r.updateTag = "" r.updateTag = ""
@@ -194,7 +194,7 @@ func (r *Runner) tryRegister(ctx context.Context, nodeID *string) error {
if strings.TrimSpace(r.Config.DiscoveryToken) == "" { if strings.TrimSpace(r.Config.DiscoveryToken) == "" {
return errors.New("agent_token 为空且未配置 discovery_token") return errors.New("agent_token 为空且未配置 discovery_token")
} }
log.Printf("agent discovery registration started") slog.Info("agent discovery registration started")
response, err := r.HeartbeatService.Register(ctx, r.nodePayload(*nodeID)) response, err := r.HeartbeatService.Register(ctx, r.nodePayload(*nodeID))
if err != nil { if err != nil {
return err return err
@@ -217,11 +217,11 @@ func (r *Runner) tryRegister(ctx context.Context, nodeID *string) error {
} }
r.HeartbeatService.SetToken(response.AgentToken) r.HeartbeatService.SetToken(response.AgentToken)
*nodeID = response.NodeID *nodeID = response.NodeID
log.Printf("agent discovery registration succeeded: node_id=%s", response.NodeID) slog.Info("agent discovery registration succeeded", "node_id", response.NodeID)
r.refreshOpenrestyHealth(ctx) r.refreshOpenrestyHealth(ctx)
heartbeatResult, heartbeatErr := r.HeartbeatService.Heartbeat(ctx, r.nodePayload(*nodeID)) heartbeatResult, heartbeatErr := r.HeartbeatService.Heartbeat(ctx, r.nodePayload(*nodeID))
if heartbeatErr != nil { if heartbeatErr != nil {
log.Printf("agent post-register heartbeat failed: %v", heartbeatErr) slog.Error("agent post-register heartbeat failed", "error", heartbeatErr)
return nil return nil
} }
if heartbeatResult == nil { if heartbeatResult == nil {
@@ -230,9 +230,9 @@ func (r *Runner) tryRegister(ctx context.Context, nodeID *string) error {
r.applySettings(heartbeatResult.AgentSettings) r.applySettings(heartbeatResult.AgentSettings)
if err = r.SyncService.SyncOnStartup(ctx, heartbeatResult.ActiveConfig); err != nil { if err = r.SyncService.SyncOnStartup(ctx, heartbeatResult.ActiveConfig); err != nil {
r.recordSyncError(err) r.recordSyncError(err)
log.Printf("agent post-register startup sync failed: %v", err) slog.Error("agent post-register startup sync failed", "error", err)
} else { } else {
log.Printf("agent post-register startup sync completed") slog.Info("agent post-register startup sync completed")
} }
r.tryRestartOpenresty(ctx) r.tryRestartOpenresty(ctx)
r.tryAutoUpdate(ctx) r.tryAutoUpdate(ctx)
@@ -245,13 +245,13 @@ func (r *Runner) recordSyncError(err error) {
} }
snapshot, loadErr := r.StateStore.Load() snapshot, loadErr := r.StateStore.Load()
if loadErr != nil { if loadErr != nil {
log.Printf("load state before recording sync error failed: %v", loadErr) slog.Error("load state before recording sync error failed", "error", loadErr)
return return
} }
snapshot.LastError = err.Error() snapshot.LastError = err.Error()
log.Printf("recording sync error into state: %s", snapshot.LastError) slog.Warn("recording sync error into state", "error", snapshot.LastError)
if saveErr := r.StateStore.Save(snapshot); saveErr != nil { if saveErr := r.StateStore.Save(snapshot); saveErr != nil {
log.Printf("save state after sync error failed: %v", saveErr) slog.Error("save state after sync error failed", "error", saveErr)
} }
} }
@@ -272,7 +272,7 @@ func (r *Runner) recordOpenrestyHealthy() {
} }
snapshot, err := r.StateStore.Load() snapshot, err := r.StateStore.Load()
if err != nil { if err != nil {
log.Printf("load state before recording openresty health failed: %v", err) slog.Error("load state before recording openresty health failed", "error", err)
return return
} }
if snapshot.OpenrestyStatus == protocol.OpenrestyStatusHealthy && strings.TrimSpace(snapshot.OpenrestyMessage) == "" { if snapshot.OpenrestyStatus == protocol.OpenrestyStatusHealthy && strings.TrimSpace(snapshot.OpenrestyMessage) == "" {
@@ -281,7 +281,7 @@ func (r *Runner) recordOpenrestyHealthy() {
snapshot.OpenrestyStatus = protocol.OpenrestyStatusHealthy snapshot.OpenrestyStatus = protocol.OpenrestyStatusHealthy
snapshot.OpenrestyMessage = "" snapshot.OpenrestyMessage = ""
if err = r.StateStore.Save(snapshot); err != nil { if err = r.StateStore.Save(snapshot); err != nil {
log.Printf("save state after recording openresty health failed: %v", err) slog.Error("save state after recording openresty health failed", "error", err)
} }
} }
@@ -291,7 +291,7 @@ func (r *Runner) recordOpenrestyUnhealthy(err error, fallbackOnly bool) {
} }
snapshot, loadErr := r.StateStore.Load() snapshot, loadErr := r.StateStore.Load()
if loadErr != nil { if loadErr != nil {
log.Printf("load state before recording openresty error failed: %v", loadErr) slog.Error("load state before recording openresty error failed", "error", loadErr)
return return
} }
message := strings.TrimSpace(err.Error()) message := strings.TrimSpace(err.Error())
@@ -300,7 +300,7 @@ func (r *Runner) recordOpenrestyUnhealthy(err error, fallbackOnly bool) {
} }
snapshot.OpenrestyStatus = protocol.OpenrestyStatusUnhealthy snapshot.OpenrestyStatus = protocol.OpenrestyStatusUnhealthy
if saveErr := r.StateStore.Save(snapshot); saveErr != nil { if saveErr := r.StateStore.Save(snapshot); saveErr != nil {
log.Printf("save state after recording openresty error failed: %v", saveErr) slog.Error("save state after recording openresty error failed", "error", saveErr)
} }
} }
+11 -11
View File
@@ -5,7 +5,7 @@ import (
"context" "context"
"encoding/json" "encoding/json"
"errors" "errors"
"log" "log/slog"
"net/http" "net/http"
"strings" "strings"
"time" "time"
@@ -30,7 +30,7 @@ func New(baseURL string, token string, timeout time.Duration) *Client {
} }
func (c *Client) RegisterNode(ctx context.Context, payload protocol.NodePayload) (*protocol.RegisterNodeResponse, error) { func (c *Client) RegisterNode(ctx context.Context, payload protocol.NodePayload) (*protocol.RegisterNodeResponse, error) {
log.Printf("http register node request: node_id=%s current_version=%s", payload.NodeID, payload.CurrentVersion) slog.Info("http register node request", "node_id", payload.NodeID, "current_version", payload.CurrentVersion)
resp := protocol.APIResponse[protocol.RegisterNodeResponse]{} resp := protocol.APIResponse[protocol.RegisterNodeResponse]{}
if err := c.postJSON(ctx, "/api/agent/nodes/register", payload, &resp); err != nil { if err := c.postJSON(ctx, "/api/agent/nodes/register", payload, &resp); err != nil {
return nil, err return nil, err
@@ -38,7 +38,7 @@ func (c *Client) RegisterNode(ctx context.Context, payload protocol.NodePayload)
if !resp.Success { if !resp.Success {
return nil, errors.New(resp.Message) return nil, errors.New(resp.Message)
} }
log.Printf("http register node response: node_id=%s", resp.Data.NodeID) slog.Info("http register node response", "node_id", resp.Data.NodeID)
return &resp.Data, nil return &resp.Data, nil
} }
@@ -64,18 +64,18 @@ func (c *Client) GetActiveConfig(ctx context.Context) (*protocol.ActiveConfigRes
if !resp.Success { if !resp.Success {
return nil, errors.New(resp.Message) return nil, errors.New(resp.Message)
} }
log.Printf("http get active config response: version=%s checksum=%s support_files=%d", resp.Data.Version, resp.Data.Checksum, len(resp.Data.SupportFiles)) slog.Info("http get active config response", "version", resp.Data.Version, "checksum", resp.Data.Checksum, "support_files", len(resp.Data.SupportFiles))
return &resp.Data, nil return &resp.Data, nil
} }
func (c *Client) ReportApplyLog(ctx context.Context, payload protocol.ApplyLogPayload) error { func (c *Client) ReportApplyLog(ctx context.Context, payload protocol.ApplyLogPayload) error {
log.Printf("http report apply log request: node_id=%s version=%s result=%s", payload.NodeID, payload.Version, payload.Result) slog.Info("http report apply log request", "node_id", payload.NodeID, "version", payload.Version, "result", payload.Result)
return c.postJSON(ctx, "/api/agent/apply-logs", payload, nil) return c.postJSON(ctx, "/api/agent/apply-logs", payload, nil)
} }
func (c *Client) SetToken(token string) { func (c *Client) SetToken(token string) {
c.token = strings.TrimSpace(token) c.token = strings.TrimSpace(token)
log.Printf("http client token updated") slog.Info("http client token updated")
} }
func (c *Client) getJSON(ctx context.Context, path string, target any) error { func (c *Client) getJSON(ctx context.Context, path string, target any) error {
@@ -104,28 +104,28 @@ func (c *Client) postJSON(ctx context.Context, path string, body any, target any
func (c *Client) do(req *http.Request, target any) error { func (c *Client) do(req *http.Request, target any) error {
res, err := c.httpClient.Do(req) res, err := c.httpClient.Do(req)
if err != nil { if err != nil {
log.Printf("http request failed: method=%s path=%s error=%v", req.Method, req.URL.Path, err) slog.Error("http request failed", "method", req.Method, "path", req.URL.Path, "error", err)
return err return err
} }
defer res.Body.Close() defer res.Body.Close()
if res.StatusCode != http.StatusOK { if res.StatusCode != http.StatusOK {
log.Printf("http request returned non-200: method=%s path=%s status=%s", req.Method, req.URL.Path, res.Status) slog.Warn("http request returned non-200", "method", req.Method, "path", req.URL.Path, "status", res.Status)
return errors.New(res.Status) return errors.New(res.Status)
} }
if target == nil { if target == nil {
var wrapper protocol.APIResponse[json.RawMessage] var wrapper protocol.APIResponse[json.RawMessage]
if err = json.NewDecoder(res.Body).Decode(&wrapper); err != nil { if err = json.NewDecoder(res.Body).Decode(&wrapper); err != nil {
log.Printf("http response decode failed: method=%s path=%s error=%v", req.Method, req.URL.Path, err) slog.Error("http response decode failed", "method", req.Method, "path", req.URL.Path, "error", err)
return err return err
} }
if !wrapper.Success { if !wrapper.Success {
log.Printf("http api response failed: method=%s path=%s message=%s", req.Method, req.URL.Path, wrapper.Message) slog.Warn("http api response failed", "method", req.Method, "path", req.URL.Path, "message", wrapper.Message)
return errors.New(wrapper.Message) return errors.New(wrapper.Message)
} }
return nil return nil
} }
if err = json.NewDecoder(res.Body).Decode(target); err != nil { if err = json.NewDecoder(res.Body).Decode(target); err != nil {
log.Printf("http response decode failed: method=%s path=%s error=%v", req.Method, req.URL.Path, err) slog.Error("http response decode failed", "method", req.Method, "path", req.URL.Path, "error", err)
return err return err
} }
return nil return nil
+27
View File
@@ -0,0 +1,27 @@
package logging
import (
"log/slog"
"os"
"strings"
)
func Setup() {
handler := slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{
Level: parseLevel(os.Getenv("LOG_LEVEL")),
})
slog.SetDefault(slog.New(handler))
}
func parseLevel(value string) slog.Level {
switch strings.ToLower(strings.TrimSpace(value)) {
case "debug":
return slog.LevelDebug
case "warn", "warning":
return slog.LevelWarn
case "error":
return slog.LevelError
default:
return slog.LevelInfo
}
}
+30 -30
View File
@@ -6,7 +6,7 @@ import (
"encoding/hex" "encoding/hex"
"errors" "errors"
"fmt" "fmt"
"log" "log/slog"
"os" "os"
"os/exec" "os/exec"
"path/filepath" "path/filepath"
@@ -50,22 +50,22 @@ type PathExecutor struct {
} }
func (e *PathExecutor) Test(ctx context.Context) error { func (e *PathExecutor) Test(ctx context.Context) error {
log.Printf("running openresty test with binary: %s", e.Path) slog.Info("running openresty test with binary", "path", e.Path)
output, err := e.Runner.Run(ctx, e.Path, "-t") output, err := e.Runner.Run(ctx, e.Path, "-t")
if err != nil { if err != nil {
return fmt.Errorf("openresty -t failed: %w: %s", err, string(output)) return fmt.Errorf("openresty -t failed: %w: %s", err, string(output))
} }
log.Printf("openresty test succeeded with binary: %s", e.Path) slog.Info("openresty test succeeded with binary", "path", e.Path)
return nil return nil
} }
func (e *PathExecutor) Reload(ctx context.Context) error { func (e *PathExecutor) Reload(ctx context.Context) error {
log.Printf("running openresty reload with binary: %s", e.Path) slog.Info("running openresty reload with binary", "path", e.Path)
output, err := e.Runner.Run(ctx, e.Path, "-s", "reload") output, err := e.Runner.Run(ctx, e.Path, "-s", "reload")
if err != nil { if err != nil {
return fmt.Errorf("openresty reload failed: %w: %s", err, string(output)) return fmt.Errorf("openresty reload failed: %w: %s", err, string(output))
} }
log.Printf("openresty reload succeeded with binary: %s", e.Path) slog.Info("openresty reload succeeded with binary", "path", e.Path)
return nil return nil
} }
@@ -78,7 +78,7 @@ func (e *PathExecutor) CheckHealth(ctx context.Context) error {
} }
func (e *PathExecutor) Restart(ctx context.Context) error { func (e *PathExecutor) Restart(ctx context.Context) error {
log.Printf("restarting openresty with binary: %s", e.Path) slog.Info("restarting openresty with binary", "path", e.Path)
output, err := e.Runner.Run(ctx, e.Path, "-s", "quit") output, err := e.Runner.Run(ctx, e.Path, "-s", "quit")
if err != nil { if err != nil {
text := string(output) text := string(output)
@@ -90,7 +90,7 @@ func (e *PathExecutor) Restart(ctx context.Context) error {
if err != nil { if err != nil {
return fmt.Errorf("openresty start failed: %w: %s", err, string(output)) return fmt.Errorf("openresty start failed: %w: %s", err, string(output))
} }
log.Printf("openresty restart succeeded with binary: %s", e.Path) slog.Info("openresty restart succeeded with binary", "path", e.Path)
return nil return nil
} }
@@ -106,12 +106,12 @@ type DockerExecutor struct {
} }
func (e *DockerExecutor) Test(ctx context.Context) error { func (e *DockerExecutor) Test(ctx context.Context) error {
log.Printf("running docker openresty test: container=%s image=%s", e.ContainerName, e.Image) slog.Info("running docker openresty test", "container", e.ContainerName, "image", e.Image)
output, err := e.runEphemeralRuntimeCommand(ctx, "-t") output, err := e.runEphemeralRuntimeCommand(ctx, "-t")
if err != nil { if err != nil {
return fmt.Errorf("docker %s -t failed: %w: %s", dockerRuntimeCommand, err, string(output)) return fmt.Errorf("docker %s -t failed: %w: %s", dockerRuntimeCommand, err, string(output))
} }
log.Printf("docker openresty test succeeded: container=%s runtime=%s", e.ContainerName, dockerRuntimeCommand) slog.Info("docker openresty test succeeded", "container", e.ContainerName, "runtime", dockerRuntimeCommand)
return nil return nil
} }
@@ -120,7 +120,7 @@ func (e *DockerExecutor) Reload(ctx context.Context) error {
} }
func (e *DockerExecutor) EnsureRuntime(ctx context.Context, recreate bool) error { func (e *DockerExecutor) EnsureRuntime(ctx context.Context, recreate bool) error {
log.Printf("ensuring docker openresty runtime: container=%s recreate=%t", e.ContainerName, recreate) slog.Info("ensuring docker openresty runtime", "container", e.ContainerName, "recreate", recreate)
output, err := e.Runner.Run(ctx, e.DockerBinary, "inspect", "-f", "{{.State.Running}}", e.ContainerName) output, err := e.Runner.Run(ctx, e.DockerBinary, "inspect", "-f", "{{.State.Running}}", e.ContainerName)
if err == nil { if err == nil {
if recreate { if recreate {
@@ -130,7 +130,7 @@ func (e *DockerExecutor) EnsureRuntime(ctx context.Context, recreate bool) error
return e.runContainer(ctx) return e.runContainer(ctx)
} }
if strings.TrimSpace(string(output)) == "true" { if strings.TrimSpace(string(output)) == "true" {
log.Printf("docker openresty runtime already healthy: container=%s", e.ContainerName) slog.Info("docker openresty runtime already healthy", "container", e.ContainerName)
return nil return nil
} }
if err := e.removeContainer(ctx); err != nil { if err := e.removeContainer(ctx); err != nil {
@@ -142,7 +142,7 @@ func (e *DockerExecutor) EnsureRuntime(ctx context.Context, recreate bool) error
} }
func (e *DockerExecutor) CheckHealth(ctx context.Context) error { func (e *DockerExecutor) CheckHealth(ctx context.Context) error {
log.Printf("checking docker openresty runtime health: container=%s", e.ContainerName) slog.Debug("checking docker openresty runtime health", "container", e.ContainerName)
output, err := e.Runner.Run(ctx, e.DockerBinary, "inspect", "-f", "{{.State.Running}}", e.ContainerName) output, err := e.Runner.Run(ctx, e.DockerBinary, "inspect", "-f", "{{.State.Running}}", e.ContainerName)
if err != nil { if err != nil {
return fmt.Errorf("docker inspect openresty failed: %w: %s", err, string(output)) return fmt.Errorf("docker inspect openresty failed: %w: %s", err, string(output))
@@ -158,7 +158,7 @@ func (e *DockerExecutor) Restart(ctx context.Context) error {
} }
func (e *DockerExecutor) removeContainer(ctx context.Context) error { func (e *DockerExecutor) removeContainer(ctx context.Context) error {
log.Printf("removing docker openresty container: container=%s", e.ContainerName) slog.Info("removing docker openresty container", "container", e.ContainerName)
output, err := e.Runner.Run(ctx, e.DockerBinary, "rm", "-f", e.ContainerName) output, err := e.Runner.Run(ctx, e.DockerBinary, "rm", "-f", e.ContainerName)
if err != nil { if err != nil {
text := string(output) text := string(output)
@@ -167,12 +167,12 @@ func (e *DockerExecutor) removeContainer(ctx context.Context) error {
} }
return fmt.Errorf("docker rm openresty failed: %w: %s", err, text) return fmt.Errorf("docker rm openresty failed: %w: %s", err, text)
} }
log.Printf("docker openresty container removed: container=%s", e.ContainerName) slog.Info("docker openresty container removed", "container", e.ContainerName)
return nil return nil
} }
func (e *DockerExecutor) runContainer(ctx context.Context) error { func (e *DockerExecutor) runContainer(ctx context.Context) error {
log.Printf("starting docker openresty container: container=%s image=%s", e.ContainerName, e.Image) slog.Info("starting docker openresty container", "container", e.ContainerName, "image", e.Image)
runArgs := []string{ runArgs := []string{
"run", "-d", "run", "-d",
"--name", e.ContainerName, "--name", e.ContainerName,
@@ -187,7 +187,7 @@ func (e *DockerExecutor) runContainer(ctx context.Context) error {
if runErr != nil { if runErr != nil {
return fmt.Errorf("docker run openresty failed: %w: %s", runErr, string(runOutput)) return fmt.Errorf("docker run openresty failed: %w: %s", runErr, string(runOutput))
} }
log.Printf("docker openresty container started: container=%s", e.ContainerName) slog.Info("docker openresty container started", "container", e.ContainerName)
return nil return nil
} }
@@ -201,39 +201,39 @@ type Manager struct {
} }
func (m *Manager) Apply(ctx context.Context, mainConfig string, routeConfig string, supportFiles []protocol.SupportFile) error { func (m *Manager) Apply(ctx context.Context, mainConfig string, routeConfig string, supportFiles []protocol.SupportFile) error {
log.Printf("openresty apply started: main_config=%s route_config=%s support_files=%d", m.MainConfigPath, m.RouteConfigPath, len(supportFiles)) slog.Info("openresty apply started", "main_config", m.MainConfigPath, "route_config", m.RouteConfigPath, "support_files", len(supportFiles))
backup, err := m.backup() backup, err := m.backup()
if err != nil { if err != nil {
return err return err
} }
if err = m.writeSupportFiles(supportFiles); err != nil { if err = m.writeSupportFiles(supportFiles); err != nil {
log.Printf("writing support files failed, restoring backup: error=%v", err) slog.Error("writing support files failed, restoring backup", "error", err)
_ = m.restore(backup) _ = m.restore(backup)
return err return err
} }
renderedMainConfig := m.renderMainConfig(mainConfig) renderedMainConfig := m.renderMainConfig(mainConfig)
if err = os.WriteFile(m.MainConfigPath, []byte(renderedMainConfig), 0o644); err != nil { if err = os.WriteFile(m.MainConfigPath, []byte(renderedMainConfig), 0o644); err != nil {
log.Printf("writing openresty main config failed, restoring backup: error=%v", err) slog.Error("writing openresty main config failed, restoring backup", "error", err)
_ = m.restore(backup) _ = m.restore(backup)
return err return err
} }
renderedRouteConfig := m.renderRouteConfig(routeConfig) renderedRouteConfig := m.renderRouteConfig(routeConfig)
if err = os.WriteFile(m.RouteConfigPath, []byte(renderedRouteConfig), 0o644); err != nil { if err = os.WriteFile(m.RouteConfigPath, []byte(renderedRouteConfig), 0o644); err != nil {
log.Printf("writing openresty route config failed, restoring backup: error=%v", err) slog.Error("writing openresty route config failed, restoring backup", "error", err)
_ = m.restore(backup) _ = m.restore(backup)
return err return err
} }
if err = m.Executor.Test(ctx); err != nil { if err = m.Executor.Test(ctx); err != nil {
log.Printf("openresty test failed after config write, restoring backup: error=%v", err) slog.Error("openresty test failed after config write, restoring backup", "error", err)
_ = m.restore(backup) _ = m.restore(backup)
return err return err
} }
if err = m.Executor.Reload(ctx); err != nil { if err = m.Executor.Reload(ctx); err != nil {
log.Printf("openresty reload failed after config write, restoring backup: error=%v", err) slog.Error("openresty reload failed after config write, restoring backup", "error", err)
_ = m.restore(backup) _ = m.restore(backup)
return err return err
} }
log.Printf("openresty apply completed successfully: main_config=%s route_config=%s", m.MainConfigPath, m.RouteConfigPath) slog.Info("openresty apply completed successfully", "main_config", m.MainConfigPath, "route_config", m.RouteConfigPath)
return nil return nil
} }
@@ -241,7 +241,7 @@ func (m *Manager) EnsureRuntime(ctx context.Context, recreate bool) error {
if m.Executor == nil { if m.Executor == nil {
return errors.New("executor 未配置") return errors.New("executor 未配置")
} }
log.Printf("openresty ensure runtime requested: recreate=%t", recreate) slog.Info("openresty ensure runtime requested", "recreate", recreate)
return m.Executor.EnsureRuntime(ctx, recreate) return m.Executor.EnsureRuntime(ctx, recreate)
} }
@@ -256,7 +256,7 @@ func (m *Manager) Restart(ctx context.Context) error {
if m.Executor == nil { if m.Executor == nil {
return errors.New("executor 未配置") return errors.New("executor 未配置")
} }
log.Printf("openresty restart requested") slog.Info("openresty restart requested")
return m.Executor.Restart(ctx) return m.Executor.Restart(ctx)
} }
@@ -294,7 +294,7 @@ func (m *Manager) CurrentChecksum() (string, error) {
return "", err return "", err
} }
result := bundleChecksum(normalizedMain, normalizedRoute, files) result := bundleChecksum(normalizedMain, normalizedRoute, files)
log.Printf("openresty current checksum calculated: main_config=%s route_config=%s checksum=%s support_files=%d", m.MainConfigPath, m.RouteConfigPath, result, len(files)) slog.Info("openresty current checksum calculated", "main_config", m.MainConfigPath, "route_config", m.RouteConfigPath, "checksum", result, "support_files", len(files))
return result, nil return result, nil
} }
@@ -344,10 +344,10 @@ func NewExecutor(options ExecutorOptions) Executor {
func DetectVersion(ctx context.Context, options ExecutorOptions) string { func DetectVersion(ctx context.Context, options ExecutorOptions) string {
version, err := detectVersion(ctx, options, &OSCommandRunner{}) version, err := detectVersion(ctx, options, &OSCommandRunner{})
if err != nil { if err != nil {
log.Printf("detect openresty version failed: %v", err) slog.Error("detect openresty version failed", "error", err)
return "" return ""
} }
log.Printf("detected openresty version: %s", version) slog.Info("detected openresty version", "version", version)
return version return version
} }
@@ -466,7 +466,7 @@ func (m *Manager) backup() (*backupState, error) {
return nil, err return nil, err
} }
state.Files = files state.Files = files
log.Printf("backup captured: main_exists=%t route_exists=%t support_files=%d", state.MainExisted, state.RouteExisted, len(state.Files)) slog.Info("backup captured", "main_exists", state.MainExisted, "route_exists", state.RouteExisted, "support_files", len(state.Files))
return state, nil return state, nil
} }
@@ -474,7 +474,7 @@ func (m *Manager) restore(state *backupState) error {
if state == nil { if state == nil {
return nil return nil
} }
log.Printf("restoring nginx backup: main_existed=%t route_existed=%t support_files=%d", state.MainExisted, state.RouteExisted, len(state.Files)) slog.Warn("restoring nginx backup", "main_existed", state.MainExisted, "route_existed", state.RouteExisted, "support_files", len(state.Files))
if state.MainExisted { if state.MainExisted {
if err := os.WriteFile(m.MainConfigPath, state.MainData, 0o644); err != nil { if err := os.WriteFile(m.MainConfigPath, state.MainData, 0o644); err != nil {
return err return err
+23 -23
View File
@@ -4,7 +4,7 @@ import (
"context" "context"
"crypto/sha256" "crypto/sha256"
"encoding/hex" "encoding/hex"
"log" "log/slog"
"strings" "strings"
"atsflare-agent/internal/protocol" "atsflare-agent/internal/protocol"
@@ -70,13 +70,13 @@ func (s *Service) sync(ctx context.Context, startup bool, target *protocol.Activ
if target == nil || target.Version == "" || target.Checksum == "" { if target == nil || target.Version == "" || target.Checksum == "" {
if !startup { if !startup {
log.Printf("skipping sync because heartbeat returned no active config summary: mode=%s", mode) slog.Debug("skipping sync because heartbeat returned no active config summary", "mode", mode)
return nil return nil
} }
log.Printf("sync startup fallback: active config summary unavailable, fetching active config directly") slog.Info("sync startup fallback: active config summary unavailable, fetching active config directly")
config, fetchErr := s.client.GetActiveConfig(ctx) config, fetchErr := s.client.GetActiveConfig(ctx)
if fetchErr != nil { if fetchErr != nil {
log.Printf("fetch active config failed: mode=%s error=%v", mode, fetchErr) slog.Error("fetch active config failed", "mode", mode, "error", fetchErr)
return fetchErr return fetchErr
} }
target = &protocol.ActiveConfigMeta{ target = &protocol.ActiveConfigMeta{
@@ -87,33 +87,33 @@ func (s *Service) sync(ctx context.Context, startup bool, target *protocol.Activ
} }
if currentChecksum == target.Checksum { if currentChecksum == target.Checksum {
log.Printf("local openresty config already up to date: mode=%s version=%s", mode, target.Version) slog.Info("local openresty config already up to date", "mode", mode, "version", target.Version)
if startup { if startup {
log.Printf("ensuring openresty runtime on startup: version=%s", target.Version) slog.Info("ensuring openresty runtime on startup", "version", target.Version)
if err = s.nginxManager.EnsureRuntime(ctx, true); err != nil { if err = s.nginxManager.EnsureRuntime(ctx, true); err != nil {
snapshot.OpenrestyStatus = protocol.OpenrestyStatusUnhealthy snapshot.OpenrestyStatus = protocol.OpenrestyStatusUnhealthy
snapshot.OpenrestyMessage = err.Error() snapshot.OpenrestyMessage = err.Error()
_ = s.stateStore.Save(snapshot) _ = s.stateStore.Save(snapshot)
return err return err
} }
log.Printf("openresty runtime ensured on startup: version=%s", target.Version) slog.Info("openresty runtime ensured on startup", "version", target.Version)
snapshot.OpenrestyStatus = protocol.OpenrestyStatusHealthy snapshot.OpenrestyStatus = protocol.OpenrestyStatusHealthy
snapshot.OpenrestyMessage = "" snapshot.OpenrestyMessage = ""
} }
snapshot.CurrentVersion = target.Version snapshot.CurrentVersion = target.Version
snapshot.CurrentChecksum = target.Checksum snapshot.CurrentChecksum = target.Checksum
snapshot.LastError = "" snapshot.LastError = ""
log.Printf("sync finished without changes: mode=%s version=%s", mode, target.Version) slog.Info("sync finished without changes", "mode", mode, "version", target.Version)
return s.stateStore.Save(snapshot) return s.stateStore.Save(snapshot)
} }
if snapshot.CurrentVersion == target.Version && snapshot.CurrentChecksum == target.Checksum && !startup { if snapshot.CurrentVersion == target.Version && snapshot.CurrentChecksum == target.Checksum && !startup {
log.Printf("skipping config fetch because state already records target version/checksum: version=%s checksum=%s", target.Version, target.Checksum) slog.Debug("skipping config fetch because state already records target version/checksum", "version", target.Version, "checksum", target.Checksum)
return nil return nil
} }
config, err := s.client.GetActiveConfig(ctx) config, err := s.client.GetActiveConfig(ctx)
if err != nil { if err != nil {
log.Printf("fetch active config failed: mode=%s error=%v", mode, err) slog.Error("fetch active config failed", "mode", mode, "error", err)
return err return err
} }
return s.applyIfNeeded(ctx, mode, startup, snapshot, currentChecksum, target, config) return s.applyIfNeeded(ctx, mode, startup, snapshot, currentChecksum, target, config)
@@ -121,30 +121,30 @@ func (s *Service) sync(ctx context.Context, startup bool, target *protocol.Activ
func (s *Service) applyIfNeeded(ctx context.Context, mode string, startup bool, snapshot *state.Snapshot, currentChecksum string, target *protocol.ActiveConfigMeta, config *protocol.ActiveConfigResponse) error { func (s *Service) applyIfNeeded(ctx context.Context, mode string, startup bool, snapshot *state.Snapshot, currentChecksum string, target *protocol.ActiveConfigMeta, config *protocol.ActiveConfigResponse) error {
if currentChecksum == config.Checksum { if currentChecksum == config.Checksum {
log.Printf("local openresty config already up to date: mode=%s version=%s", mode, config.Version) slog.Info("local openresty config already up to date", "mode", mode, "version", config.Version)
if startup { if startup {
log.Printf("ensuring openresty runtime on startup: version=%s", config.Version) slog.Info("ensuring openresty runtime on startup", "version", config.Version)
if err := s.nginxManager.EnsureRuntime(ctx, true); err != nil { if err := s.nginxManager.EnsureRuntime(ctx, true); err != nil {
snapshot.OpenrestyStatus = protocol.OpenrestyStatusUnhealthy snapshot.OpenrestyStatus = protocol.OpenrestyStatusUnhealthy
snapshot.OpenrestyMessage = err.Error() snapshot.OpenrestyMessage = err.Error()
_ = s.stateStore.Save(snapshot) _ = s.stateStore.Save(snapshot)
return err return err
} }
log.Printf("openresty runtime ensured on startup: version=%s", config.Version) slog.Info("openresty runtime ensured on startup", "version", config.Version)
snapshot.OpenrestyStatus = protocol.OpenrestyStatusHealthy snapshot.OpenrestyStatus = protocol.OpenrestyStatusHealthy
snapshot.OpenrestyMessage = "" snapshot.OpenrestyMessage = ""
} }
snapshot.CurrentVersion = config.Version snapshot.CurrentVersion = config.Version
snapshot.CurrentChecksum = config.Checksum snapshot.CurrentChecksum = config.Checksum
snapshot.LastError = "" snapshot.LastError = ""
log.Printf("sync finished without changes: mode=%s version=%s", mode, config.Version) slog.Info("sync finished without changes", "mode", mode, "version", config.Version)
return s.stateStore.Save(snapshot) return s.stateStore.Save(snapshot)
} }
if target != nil && (target.Version != config.Version || target.Checksum != config.Checksum) { if target != nil && (target.Version != config.Version || target.Checksum != config.Checksum) {
log.Printf("active config changed between heartbeat and fetch: heartbeat_version=%s heartbeat_checksum=%s fetched_version=%s fetched_checksum=%s", target.Version, target.Checksum, config.Version, config.Checksum) slog.Warn("active config changed between heartbeat and fetch", "heartbeat_version", target.Version, "heartbeat_checksum", target.Checksum, "fetched_version", config.Version, "fetched_checksum", config.Checksum)
} }
if snapshot.CurrentVersion == config.Version && snapshot.CurrentChecksum == config.Checksum && !startup { if snapshot.CurrentVersion == config.Version && snapshot.CurrentChecksum == config.Checksum && !startup {
log.Printf("skipping apply because state already records target version/checksum: version=%s checksum=%s", config.Version, config.Checksum) slog.Debug("skipping apply because state already records target version/checksum", "version", config.Version, "checksum", config.Checksum)
return nil return nil
} }
routeConfig := config.RouteConfig routeConfig := config.RouteConfig
@@ -153,9 +153,9 @@ func (s *Service) applyIfNeeded(ctx context.Context, mode string, startup bool,
} }
mainConfigChecksum := checksumString(config.MainConfig) mainConfigChecksum := checksumString(config.MainConfig)
routeConfigChecksum := checksumString(routeConfig) routeConfigChecksum := checksumString(routeConfig)
log.Printf("applying new openresty config: mode=%s from_version=%s to_version=%s old_checksum=%s new_checksum=%s", mode, snapshot.CurrentVersion, config.Version, currentChecksum, config.Checksum) slog.Info("applying new openresty config", "mode", mode, "from_version", snapshot.CurrentVersion, "to_version", config.Version, "old_checksum", currentChecksum, "new_checksum", config.Checksum)
if err := s.nginxManager.Apply(ctx, config.MainConfig, routeConfig, config.SupportFiles); err != nil { if err := s.nginxManager.Apply(ctx, config.MainConfig, routeConfig, config.SupportFiles); err != nil {
log.Printf("apply openresty config failed: mode=%s version=%s error=%v", mode, config.Version, err) slog.Error("apply openresty config failed", "mode", mode, "version", config.Version, "error", err)
snapshot.LastError = err.Error() snapshot.LastError = err.Error()
snapshot.OpenrestyStatus = protocol.OpenrestyStatusUnhealthy snapshot.OpenrestyStatus = protocol.OpenrestyStatusUnhealthy
snapshot.OpenrestyMessage = err.Error() snapshot.OpenrestyMessage = err.Error()
@@ -171,13 +171,13 @@ func (s *Service) applyIfNeeded(ctx context.Context, mode string, startup bool,
SupportFileCount: len(config.SupportFiles), SupportFileCount: len(config.SupportFiles),
}) })
if reportErr != nil { if reportErr != nil {
log.Printf("report failed apply log failed: version=%s error=%v", config.Version, reportErr) slog.Error("report failed apply log failed", "version", config.Version, "error", reportErr)
return reportErr return reportErr
} }
log.Printf("failed apply log reported: version=%s", config.Version) slog.Warn("failed apply log reported", "version", config.Version)
return err return err
} }
log.Printf("openresty config applied successfully: mode=%s version=%s", mode, config.Version) slog.Info("openresty config applied successfully", "mode", mode, "version", config.Version)
snapshot.CurrentVersion = config.Version snapshot.CurrentVersion = config.Version
snapshot.CurrentChecksum = config.Checksum snapshot.CurrentChecksum = config.Checksum
snapshot.LastError = "" snapshot.LastError = ""
@@ -196,10 +196,10 @@ func (s *Service) applyIfNeeded(ctx context.Context, mode string, startup bool,
RouteConfigChecksum: routeConfigChecksum, RouteConfigChecksum: routeConfigChecksum,
SupportFileCount: len(config.SupportFiles), SupportFileCount: len(config.SupportFiles),
}); err != nil { }); err != nil {
log.Printf("report successful apply log failed: version=%s error=%v", config.Version, err) slog.Error("report successful apply log failed", "version", config.Version, "error", err)
return err return err
} }
log.Printf("successful apply log reported: version=%s", config.Version) slog.Info("successful apply log reported", "version", config.Version)
return nil return nil
} }
+3 -3
View File
@@ -5,7 +5,7 @@ import (
"encoding/json" "encoding/json"
"fmt" "fmt"
"io" "io"
"log" "log/slog"
"net/http" "net/http"
"os" "os"
"runtime" "runtime"
@@ -64,7 +64,7 @@ func (s *Service) CheckAndUpdate(ctx context.Context, repo string, options agent
return nil return nil
} }
log.Printf("agent update available: %s -> %s", localVersion, remoteVersion) slog.Info("agent update available", "from", localVersion, "to", remoteVersion)
assetName := assetNameForGOOSGOARCH(runtime.GOOS, runtime.GOARCH) assetName := assetNameForGOOSGOARCH(runtime.GOOS, runtime.GOARCH)
var downloadURL string var downloadURL string
@@ -219,7 +219,7 @@ func (s *Service) downloadAndRestart(ctx context.Context, url string, targetPath
} }
tmpFile.Close() tmpFile.Close()
log.Printf("agent binary updated, restarting...") slog.Info("agent binary updated, restarting")
return replaceAndRestart(targetPath, tmpPath) return replaceAndRestart(targetPath, tmpPath)
} }
+2 -3
View File
@@ -3,7 +3,6 @@ package common
import ( import (
"flag" "flag"
"fmt" "fmt"
"log"
"os" "os"
"path/filepath" "path/filepath"
"strings" "strings"
@@ -59,12 +58,12 @@ func init() {
var err error var err error
*LogDir, err = filepath.Abs(*LogDir) *LogDir, err = filepath.Abs(*LogDir)
if err != nil { if err != nil {
log.Fatal(err) FatalLog(err)
} }
if _, err := os.Stat(*LogDir); os.IsNotExist(err) { if _, err := os.Stat(*LogDir); os.IsNotExist(err) {
err = os.Mkdir(*LogDir, 0777) err = os.Mkdir(*LogDir, 0777)
if err != nil { if err != nil {
log.Fatal(err) FatalLog(err)
} }
} }
} }
+117 -60
View File
@@ -1,15 +1,14 @@
package common package common
import ( import (
"fmt" "context"
"github.com/gin-gonic/gin" "github.com/gin-gonic/gin"
"io" "io"
"log" "log/slog"
"os" "os"
"path/filepath" "path/filepath"
"strings" "strings"
"time" )
)
type logLevel int type logLevel int
@@ -20,19 +19,76 @@ const (
logLevelError logLevelError
) )
var currentLogLevel = logLevelInfo var currentLogLevel = logLevelInfo
var currentLogLevelName = "info" var currentLogLevelName = "info"
var commonLogWriter io.Writer = os.Stdout var commonLogWriter io.Writer = os.Stdout
var errorLogWriter io.Writer = os.Stderr var errorLogWriter io.Writer = os.Stderr
var defaultLogger *slog.Logger
type levelRouterHandler struct {
commonHandler slog.Handler
errorHandler slog.Handler
}
func (h *levelRouterHandler) Enabled(ctx context.Context, level slog.Level) bool {
return h.commonHandler.Enabled(ctx, level) || h.errorHandler.Enabled(ctx, level)
}
func (h *levelRouterHandler) Handle(ctx context.Context, record slog.Record) error {
if record.Level >= slog.LevelError {
return h.errorHandler.Handle(ctx, record)
}
return h.commonHandler.Handle(ctx, record)
}
func (h *levelRouterHandler) WithAttrs(attrs []slog.Attr) slog.Handler {
return &levelRouterHandler{
commonHandler: h.commonHandler.WithAttrs(attrs),
errorHandler: h.errorHandler.WithAttrs(attrs),
}
}
func (h *levelRouterHandler) WithGroup(name string) slog.Handler {
return &levelRouterHandler{
commonHandler: h.commonHandler.WithGroup(name),
errorHandler: h.errorHandler.WithGroup(name),
}
}
func configureGinWriters() { func configureGinWriters() {
if shouldLog(logLevelDebug) { if shouldLog(logLevelDebug) {
gin.DefaultWriter = commonLogWriter gin.DefaultWriter = commonLogWriter
} else { } else {
gin.DefaultWriter = io.Discard gin.DefaultWriter = io.Discard
} }
gin.DefaultErrorWriter = errorLogWriter gin.DefaultErrorWriter = errorLogWriter
} }
func slogLevel() slog.Level {
switch currentLogLevel {
case logLevelDebug:
return slog.LevelDebug
case logLevelWarn:
return slog.LevelWarn
case logLevelError:
return slog.LevelError
default:
return slog.LevelInfo
}
}
func ensureLogger() *slog.Logger {
if defaultLogger != nil {
return defaultLogger
}
handlerOptions := &slog.HandlerOptions{Level: slogLevel()}
defaultLogger = slog.New(&levelRouterHandler{
commonHandler: slog.NewTextHandler(commonLogWriter, handlerOptions),
errorHandler: slog.NewTextHandler(errorLogWriter, handlerOptions),
})
slog.SetDefault(defaultLogger)
return defaultLogger
}
func SetLogLevel(level string) { func SetLogLevel(level string) {
normalized := strings.TrimSpace(strings.ToLower(level)) normalized := strings.TrimSpace(strings.ToLower(level))
@@ -61,42 +117,43 @@ func shouldLog(level logLevel) bool {
return level >= currentLogLevel return level >= currentLogLevel
} }
func SetupGinLog() { func SetupGinLog() {
if *LogDir != "" { if *LogDir != "" {
commonLogPath := filepath.Join(*LogDir, "common.log") commonLogPath := filepath.Join(*LogDir, "common.log")
errorLogPath := filepath.Join(*LogDir, "error.log") errorLogPath := filepath.Join(*LogDir, "error.log")
commonFd, err := os.OpenFile(commonLogPath, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0644) commonFd, err := os.OpenFile(commonLogPath, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0644)
if err != nil { if err != nil {
log.Fatal("failed to open log file") _, _ = io.WriteString(os.Stderr, "failed to open common log file\n")
} os.Exit(1)
errorFd, err := os.OpenFile(errorLogPath, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0644) }
if err != nil { errorFd, err := os.OpenFile(errorLogPath, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0644)
log.Fatal("failed to open log file") if err != nil {
} _, _ = io.WriteString(os.Stderr, "failed to open error log file\n")
commonLogWriter = io.MultiWriter(os.Stdout, commonFd) os.Exit(1)
errorLogWriter = io.MultiWriter(os.Stderr, errorFd) }
} commonLogWriter = io.MultiWriter(os.Stdout, commonFd)
configureGinWriters() errorLogWriter = io.MultiWriter(os.Stderr, errorFd)
} }
configureGinWriters()
func SysLog(s string) { defaultLogger = nil
if !shouldLog(logLevelInfo) { ensureLogger()
return }
}
t := time.Now() func SysLog(s string) {
_, _ = fmt.Fprintf(commonLogWriter, "[SYS] %v | %s \n", t.Format("2006/01/02 - 15:04:05"), s) if !shouldLog(logLevelInfo) {
} return
}
func SysError(s string) { ensureLogger().Info(s)
if !shouldLog(logLevelError) { }
return
} func SysError(s string) {
t := time.Now() if !shouldLog(logLevelError) {
_, _ = fmt.Fprintf(errorLogWriter, "[SYS] %v | %s \n", t.Format("2006/01/02 - 15:04:05"), s) return
} }
ensureLogger().Error(s)
func FatalLog(v ...any) { }
t := time.Now()
_, _ = fmt.Fprintf(errorLogWriter, "[FATAL] %v | %v \n", t.Format("2006/01/02 - 15:04:05"), v) func FatalLog(v ...any) {
os.Exit(1) ensureLogger().Error("fatal error", "details", v)
} os.Exit(1)
}
+2 -3
View File
@@ -4,7 +4,6 @@ import (
"fmt" "fmt"
"github.com/google/uuid" "github.com/google/uuid"
"html/template" "html/template"
"log"
"net" "net"
"os/exec" "os/exec"
"runtime" "runtime"
@@ -24,14 +23,14 @@ func OpenBrowser(url string) {
err = exec.Command("open", url).Start() err = exec.Command("open", url).Start()
} }
if err != nil { if err != nil {
log.Println(err) SysError(err.Error())
} }
} }
func GetIp() (ip string) { func GetIp() (ip string) {
ips, err := net.InterfaceAddrs() ips, err := net.InterfaceAddrs()
if err != nil { if err != nil {
log.Println(err) SysError(err.Error())
return ip return ip
} }
+1 -2
View File
@@ -12,7 +12,6 @@ import (
"github.com/gin-contrib/sessions/cookie" "github.com/gin-contrib/sessions/cookie"
"github.com/gin-contrib/sessions/redis" "github.com/gin-contrib/sessions/redis"
"github.com/gin-gonic/gin" "github.com/gin-gonic/gin"
"log"
"os" "os"
"strconv" "strconv"
) )
@@ -87,7 +86,7 @@ func main() {
common.SysLog(fmt.Sprintf("server listening on :%s", port)) common.SysLog(fmt.Sprintf("server listening on :%s", port))
err = server.Run(":" + port) err = server.Run(":" + port)
if err != nil { if err != nil {
log.Println(err) common.SysError(err.Error())
} }
} }
+3 -4
View File
@@ -3,7 +3,6 @@ package middleware
import ( import (
"atsflare/common" "atsflare/common"
"context" "context"
"fmt"
"github.com/gin-gonic/gin" "github.com/gin-gonic/gin"
"net/http" "net/http"
"time" "time"
@@ -19,7 +18,7 @@ func redisRateLimiter(c *gin.Context, maxRequestNum int, duration int64, mark st
key := "rateLimit:" + mark + c.ClientIP() key := "rateLimit:" + mark + c.ClientIP()
listLength, err := rdb.LLen(ctx, key).Result() listLength, err := rdb.LLen(ctx, key).Result()
if err != nil { if err != nil {
fmt.Println(err.Error()) common.SysError(err.Error())
c.Status(http.StatusInternalServerError) c.Status(http.StatusInternalServerError)
c.Abort() c.Abort()
return return
@@ -31,7 +30,7 @@ func redisRateLimiter(c *gin.Context, maxRequestNum int, duration int64, mark st
oldTimeStr, _ := rdb.LIndex(ctx, key, -1).Result() oldTimeStr, _ := rdb.LIndex(ctx, key, -1).Result()
oldTime, err := time.Parse(timeFormat, oldTimeStr) oldTime, err := time.Parse(timeFormat, oldTimeStr)
if err != nil { if err != nil {
fmt.Println(err) common.SysError(err.Error())
c.Status(http.StatusInternalServerError) c.Status(http.StatusInternalServerError)
c.Abort() c.Abort()
return return
@@ -39,7 +38,7 @@ func redisRateLimiter(c *gin.Context, maxRequestNum int, duration int64, mark st
nowTimeStr := time.Now().Format(timeFormat) nowTimeStr := time.Now().Format(timeFormat)
nowTime, err := time.Parse(timeFormat, nowTimeStr) nowTime, err := time.Parse(timeFormat, nowTimeStr)
if err != nil { if err != nil {
fmt.Println(err) common.SysError(err.Error())
c.Status(http.StatusInternalServerError) c.Status(http.StatusInternalServerError)
c.Abort() c.Abort()
return return
+2 -3
View File
@@ -8,7 +8,6 @@ import (
"encoding/json" "encoding/json"
"fmt" "fmt"
"io" "io"
"log"
"net/http" "net/http"
"os" "os"
"os/exec" "os/exec"
@@ -137,7 +136,7 @@ func ScheduleServerUpgrade(channel string) (*LatestServerRelease, error) {
go func(task *preparedServerUpgrade) { go func(task *preparedServerUpgrade) {
time.Sleep(serverUpgradeDispatchDelay) time.Sleep(serverUpgradeDispatchDelay)
if err := executeServerUpgrade(task); err != nil { if err := executeServerUpgrade(task); err != nil {
log.Printf("server self-update failed: %v", err) common.SysError(fmt.Sprintf("server self-update failed: %v", err))
serverUpgradeState.Lock() serverUpgradeState.Lock()
serverUpgradeState.inProgress = false serverUpgradeState.inProgress = false
serverUpgradeState.Unlock() serverUpgradeState.Unlock()
@@ -251,7 +250,7 @@ func ConfirmManualServerUpgrade(uploadToken string) (*UploadedServerBinary, erro
go func(task *manualServerBinaryCandidate) { go func(task *manualServerBinaryCandidate) {
time.Sleep(serverUpgradeDispatchDelay) time.Sleep(serverUpgradeDispatchDelay)
if err := executeManualServerUpgrade(task); err != nil { if err := executeManualServerUpgrade(task); err != nil {
log.Printf("server manual upgrade failed: %v", err) common.SysError(fmt.Sprintf("server manual upgrade failed: %v", err))
serverUpgradeState.Lock() serverUpgradeState.Lock()
serverUpgradeState.inProgress = false serverUpgradeState.inProgress = false
serverUpgradeState.Unlock() serverUpgradeState.Unlock()
+16 -3
View File
@@ -188,9 +188,22 @@ Agent 当前支持两类启动配置:
1. 命令行参数 1. 命令行参数
2. `agent.json` 配置文件 2. `agent.json` 配置文件
当前 Agent **没有额外环境变量作为正式配置入口**,启动行为主要由 `-config` 参数和配置文件字段决定。 当前 Agent 仅支持少量环境变量用于日志输出控制;核心启动行为仍由 `-config` 参数和配置文件字段决定。
### 2.1 Agent 命令行参数 ### 2.1.1 Agent 环境变量
支持的环境变量:
| 环境变量 | 作用 | 默认值 | 示例 |
| --- | --- | --- | --- |
| `LOG_LEVEL` | 指定 Agent 的 `slog` 日志等级;支持 `debug`/`info`/`warn`/`error`(`warning` 兼容为 `warn`) | `info` | `LOG_LEVEL=debug` |
说明:
* `LOG_LEVEL` 只影响 Agent 本地日志输出,不改变心跳、同步或配置行为
* 未设置或设置为不支持的值时,将回退为 `info`
### 2.1 Agent 命令行参数
启动示例: 启动示例:
+33 -29
View File
@@ -6,18 +6,18 @@
## 1. 前置条件 ## 1. 前置条件
### 1.1 Server ### 1.1 Server
* Go 1.23+
* Node.js 18+
* 可写 SQLite 文件目录
* Go 1.18+ ### 1.2 Agent
* Node.js 18+
* 可写 SQLite 文件目录 * Go 1.23+
* 对 Agent 数据目录有写权限
### 1.2 Agent * 若使用独立 OpenResty 模式:可执行 `openresty -t` 与 `openresty -s reload`
* 若使用 Docker 模式:具备 Docker 执行权限
* Go 1.18+
* 对 Agent 数据目录有写权限
* 若使用独立 OpenResty 模式:可执行 `openresty -t` 与 `openresty -s reload`
* 若使用 Docker 模式:具备 Docker 执行权限
* 第五版主配置接管模式下,Agent 对 OpenResty 主配置目标路径必须具备写权限 * 第五版主配置接管模式下,Agent 对 OpenResty 主配置目标路径必须具备写权限
--- ---
@@ -41,18 +41,20 @@ pnpm build
### 2.2 启动服务 ### 2.2 启动服务
```bash ```bash
cd atsf_server cd atsf_server
export SESSION_SECRET='replace-with-random-string' export SESSION_SECRET='replace-with-random-string'
export SQLITE_PATH='./atsflare.db' export SQLITE_PATH='./atsflare.db'
go run . export LOG_LEVEL='info'
``` go run .
```
说明: 说明:
* 默认不依赖全局 `AGENT_TOKEN` * 默认不依赖全局 `AGENT_TOKEN`
* 节点接入凭证由数据库维护:节点专属 `agent_token` + 全局 `discovery_token` * 节点接入凭证由数据库维护:节点专属 `agent_token` + 全局 `discovery_token`
* 默认监听端口为 `3000` * 默认监听端口为 `3000`
* `LOG_LEVEL` 支持 `debug` / `info` / `warn` / `error`
### 2.3 使用 docker-compose 启动 Server ### 2.3 使用 docker-compose 启动 Server
@@ -196,18 +198,20 @@ swag init -g main.go -o docs
### 4.1 直接运行 ### 4.1 直接运行
```bash ```bash
cd atsf_agent cd atsf_agent
go run ./cmd/agent -config /path/to/agent.json export LOG_LEVEL='info'
``` go run ./cmd/agent -config /path/to/agent.json
```
### 4.2 编译后二进制运行 ### 4.2 编译后二进制运行
```bash ```bash
cd atsf_agent cd atsf_agent
go build -o atsflare-agent ./cmd/agent go build -o atsflare-agent ./cmd/agent
./atsflare-agent -config /path/to/agent.json export LOG_LEVEL='info'
``` ./atsflare-agent -config /path/to/agent.json
```
--- ---
+3
View File
@@ -21,6 +21,7 @@
`atsf_server` 继续作为单体控制面: `atsf_server` 继续作为单体控制面:
* Go 1.23+
* Gin * Gin
* GORM * GORM
* SQLite * SQLite
@@ -36,6 +37,7 @@
`atsf_agent` 继续作为 Go 单体程序: `atsf_agent` 继续作为 Go 单体程序:
* Go 1.23+
* 单二进制 * 单二进制
* 节点本地执行 * 节点本地执行
* `openresty_path` 优先 * `openresty_path` 优先
@@ -275,6 +277,7 @@ Agent 必须满足:
要求: 要求:
* Server 与 Agent 统一使用 `slog` 输出结构化日志
* 日志足够定位问题 * 日志足够定位问题
* 不打印敏感凭证完整值 * 不打印敏感凭证完整值