Merge remote-tracking branch 'origin/main'

This commit is contained in:
ryan
2026-03-13 16:34:27 +08:00
6 changed files with 67 additions and 39 deletions
+3 -3
View File
@@ -69,13 +69,13 @@ func (r *Runner) Run(ctx context.Context) error {
if heartbeatResult == nil {
heartbeatResult = &protocol.HeartbeatResult{}
}
slog.Info("agent startup heartbeat succeeded", "node_id", nodeID)
slog.Debug("agent startup heartbeat succeeded", "node_id", nodeID)
r.applySettings(heartbeatResult.AgentSettings)
if err = r.SyncService.SyncOnStartup(ctx, heartbeatResult.ActiveConfig); err != nil {
r.recordSyncError(err)
slog.Error("agent startup sync failed", "error", err)
} else {
slog.Info("agent startup sync completed")
slog.Debug("agent startup sync completed")
}
r.tryRestartOpenresty(ctx)
r.tryAutoUpdate(ctx)
@@ -232,7 +232,7 @@ func (r *Runner) tryRegister(ctx context.Context, nodeID *string) error {
r.recordSyncError(err)
slog.Error("agent post-register startup sync failed", "error", err)
} else {
slog.Info("agent post-register startup sync completed")
slog.Debug("agent post-register startup sync completed")
}
r.tryRestartOpenresty(ctx)
r.tryAutoUpdate(ctx)
+5 -5
View File
@@ -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) {
slog.Info("http register node request", "node_id", payload.NodeID, "current_version", payload.CurrentVersion)
slog.Debug("http register node request", "node_id", payload.NodeID, "current_version", payload.CurrentVersion)
resp := protocol.APIResponse[protocol.RegisterNodeResponse]{}
if err := c.postJSON(ctx, "/api/agent/nodes/register", payload, &resp); err != nil {
return nil, err
@@ -38,7 +38,7 @@ func (c *Client) RegisterNode(ctx context.Context, payload protocol.NodePayload)
if !resp.Success {
return nil, errors.New(resp.Message)
}
slog.Info("http register node response", "node_id", resp.Data.NodeID)
slog.Debug("http register node response", "node_id", resp.Data.NodeID)
return &resp.Data, nil
}
@@ -64,18 +64,18 @@ func (c *Client) GetActiveConfig(ctx context.Context) (*protocol.ActiveConfigRes
if !resp.Success {
return nil, errors.New(resp.Message)
}
slog.Info("http get active config response", "version", resp.Data.Version, "checksum", resp.Data.Checksum, "support_files", len(resp.Data.SupportFiles))
slog.Debug("http get active config response", "version", resp.Data.Version, "checksum", resp.Data.Checksum, "support_files", len(resp.Data.SupportFiles))
return &resp.Data, nil
}
func (c *Client) ReportApplyLog(ctx context.Context, payload protocol.ApplyLogPayload) error {
slog.Info("http report apply log request", "node_id", payload.NodeID, "version", payload.Version, "result", payload.Result)
slog.Debug("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)
}
func (c *Client) SetToken(token string) {
c.token = strings.TrimSpace(token)
slog.Info("http client token updated")
slog.Debug("http client token updated")
}
func (c *Client) getJSON(ctx context.Context, path string, target any) error {
+9 -9
View File
@@ -50,22 +50,22 @@ type PathExecutor struct {
}
func (e *PathExecutor) Test(ctx context.Context) error {
slog.Info("running openresty test with binary", "path", e.Path)
slog.Debug("running openresty test with binary", "path", e.Path)
output, err := e.Runner.Run(ctx, e.Path, "-t")
if err != nil {
return fmt.Errorf("openresty -t failed: %w: %s", err, string(output))
}
slog.Info("openresty test succeeded with binary", "path", e.Path)
slog.Debug("openresty test succeeded with binary", "path", e.Path)
return nil
}
func (e *PathExecutor) Reload(ctx context.Context) error {
slog.Info("running openresty reload with binary", "path", e.Path)
slog.Debug("running openresty reload with binary", "path", e.Path)
output, err := e.Runner.Run(ctx, e.Path, "-s", "reload")
if err != nil {
return fmt.Errorf("openresty reload failed: %w: %s", err, string(output))
}
slog.Info("openresty reload succeeded with binary", "path", e.Path)
slog.Debug("openresty reload succeeded with binary", "path", e.Path)
return nil
}
@@ -106,12 +106,12 @@ type DockerExecutor struct {
}
func (e *DockerExecutor) Test(ctx context.Context) error {
slog.Info("running docker openresty test", "container", e.ContainerName, "image", e.Image)
slog.Debug("running docker openresty test", "container", e.ContainerName, "image", e.Image)
output, err := e.runEphemeralRuntimeCommand(ctx, "-t")
if err != nil {
return fmt.Errorf("docker %s -t failed: %w: %s", dockerRuntimeCommand, err, string(output))
}
slog.Info("docker openresty test succeeded", "container", e.ContainerName, "runtime", dockerRuntimeCommand)
slog.Debug("docker openresty test succeeded", "container", e.ContainerName, "runtime", dockerRuntimeCommand)
return nil
}
@@ -130,7 +130,7 @@ func (e *DockerExecutor) EnsureRuntime(ctx context.Context, recreate bool) error
return e.runContainer(ctx)
}
if strings.TrimSpace(string(output)) == "true" {
slog.Info("docker openresty runtime already healthy", "container", e.ContainerName)
slog.Debug("docker openresty runtime already healthy", "container", e.ContainerName)
return nil
}
if err := e.removeContainer(ctx); err != nil {
@@ -294,7 +294,7 @@ func (m *Manager) CurrentChecksum() (string, error) {
return "", err
}
result := bundleChecksum(normalizedMain, normalizedRoute, files)
slog.Info("openresty current checksum calculated", "main_config", m.MainConfigPath, "route_config", m.RouteConfigPath, "checksum", result, "support_files", len(files))
slog.Debug("openresty current checksum calculated", "main_config", m.MainConfigPath, "route_config", m.RouteConfigPath, "checksum", result, "support_files", len(files))
return result, nil
}
@@ -466,7 +466,7 @@ func (m *Manager) backup() (*backupState, error) {
return nil, err
}
state.Files = files
slog.Info("backup captured", "main_exists", state.MainExisted, "route_exists", state.RouteExisted, "support_files", len(state.Files))
slog.Debug("backup captured", "main_exists", state.MainExisted, "route_exists", state.RouteExisted, "support_files", len(state.Files))
return state, nil
}
+10 -10
View File
@@ -73,7 +73,7 @@ func (s *Service) sync(ctx context.Context, startup bool, target *protocol.Activ
slog.Debug("skipping sync because heartbeat returned no active config summary", "mode", mode)
return nil
}
slog.Info("sync startup fallback: active config summary unavailable, fetching active config directly")
slog.Debug("sync startup fallback: active config summary unavailable, fetching active config directly")
config, fetchErr := s.client.GetActiveConfig(ctx)
if fetchErr != nil {
slog.Error("fetch active config failed", "mode", mode, "error", fetchErr)
@@ -87,23 +87,23 @@ func (s *Service) sync(ctx context.Context, startup bool, target *protocol.Activ
}
if currentChecksum == target.Checksum {
slog.Info("local openresty config already up to date", "mode", mode, "version", target.Version)
slog.Debug("local openresty config already up to date", "mode", mode, "version", target.Version)
if startup {
slog.Info("ensuring openresty runtime on startup", "version", target.Version)
slog.Debug("ensuring openresty runtime on startup", "version", target.Version)
if err = s.nginxManager.EnsureRuntime(ctx, true); err != nil {
snapshot.OpenrestyStatus = protocol.OpenrestyStatusUnhealthy
snapshot.OpenrestyMessage = err.Error()
_ = s.stateStore.Save(snapshot)
return err
}
slog.Info("openresty runtime ensured on startup", "version", target.Version)
slog.Debug("openresty runtime ensured on startup", "version", target.Version)
snapshot.OpenrestyStatus = protocol.OpenrestyStatusHealthy
snapshot.OpenrestyMessage = ""
}
snapshot.CurrentVersion = target.Version
snapshot.CurrentChecksum = target.Checksum
snapshot.LastError = ""
slog.Info("sync finished without changes", "mode", mode, "version", target.Version)
slog.Debug("sync finished without changes", "mode", mode, "version", target.Version)
return s.stateStore.Save(snapshot)
}
if snapshot.CurrentVersion == target.Version && snapshot.CurrentChecksum == target.Checksum && !startup {
@@ -121,23 +121,23 @@ 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 {
if currentChecksum == config.Checksum {
slog.Info("local openresty config already up to date", "mode", mode, "version", config.Version)
slog.Debug("local openresty config already up to date", "mode", mode, "version", config.Version)
if startup {
slog.Info("ensuring openresty runtime on startup", "version", config.Version)
slog.Debug("ensuring openresty runtime on startup", "version", config.Version)
if err := s.nginxManager.EnsureRuntime(ctx, true); err != nil {
snapshot.OpenrestyStatus = protocol.OpenrestyStatusUnhealthy
snapshot.OpenrestyMessage = err.Error()
_ = s.stateStore.Save(snapshot)
return err
}
slog.Info("openresty runtime ensured on startup", "version", config.Version)
slog.Debug("openresty runtime ensured on startup", "version", config.Version)
snapshot.OpenrestyStatus = protocol.OpenrestyStatusHealthy
snapshot.OpenrestyMessage = ""
}
snapshot.CurrentVersion = config.Version
snapshot.CurrentChecksum = config.Checksum
snapshot.LastError = ""
slog.Info("sync finished without changes", "mode", mode, "version", config.Version)
slog.Debug("sync finished without changes", "mode", mode, "version", config.Version)
return s.stateStore.Save(snapshot)
}
if target != nil && (target.Version != config.Version || target.Checksum != config.Checksum) {
@@ -199,7 +199,7 @@ func (s *Service) applyIfNeeded(ctx context.Context, mode string, startup bool,
slog.Error("report successful apply log failed", "version", config.Version, "error", err)
return err
}
slog.Info("successful apply log reported", "version", config.Version)
slog.Debug("successful apply log reported", "version", config.Version)
return nil
}
+3 -12
View File
@@ -179,7 +179,7 @@ func GetActiveConfigForAgent() (*AgentConfigResponse, error) {
return nil, err
}
}
slog.Info("agent fetched active config", "version", version.Version, "checksum", version.Checksum)
slog.Debug("agent fetched active config", "version", version.Version, "checksum", version.Checksum)
return &AgentConfigResponse{
Version: version.Version,
Checksum: version.Checksum,
@@ -209,7 +209,7 @@ func ReportApplyLog(payload ApplyLogPayload) (*model.ApplyLog, error) {
if payload.Result != ApplyResultOK && payload.Result != ApplyResultFailed {
return nil, errors.New("result 仅支持 success 或 failed")
}
slog.Info("agent apply log received", "node_id", payload.NodeID, "version", payload.Version, "result", payload.Result)
slog.Debug("agent apply log received", "node_id", payload.NodeID, "version", payload.Version, "result", payload.Result)
log := &model.ApplyLog{
NodeID: payload.NodeID,
@@ -244,7 +244,7 @@ func ReportApplyLog(payload ApplyLogPayload) (*model.ApplyLog, error) {
return nil, err
}
if payload.Result == ApplyResultOK {
slog.Info("agent apply reported success", "node_id", payload.NodeID, "version", payload.Version)
slog.Debug("agent apply reported success", "node_id", payload.NodeID, "version", payload.Version)
} else {
slog.Error("agent apply reported failure", "node_id", payload.NodeID, "version", payload.Version, "message", payload.Message)
}
@@ -267,15 +267,6 @@ func ListNodeViews() ([]*NodeView, error) {
views := make([]*NodeView, 0, len(nodes))
for _, node := range nodes {
computedStatus := computeNodeStatus(node)
if node.Status != computedStatus {
if computedStatus == NodeStatusOffline {
slog.Error("node offline", "node_id", node.NodeID, "name", node.Name, "ip", node.IP, "last_seen_at", node.LastSeenAt.Format(time.RFC3339))
} else if computedStatus == NodeStatusOnline {
slog.Info("node online", "node_id", node.NodeID, "name", node.Name, "ip", node.IP)
}
_ = model.DB.Model(node).Update("status", computedStatus).Error
node.Status = computedStatus
}
view := buildNodeView(node)
view.Status = computedStatus
if log, ok := latestLogs[node.NodeID]; ok {
+37
View File
@@ -288,3 +288,40 @@ func TestCollectNodeHeartbeatChangesOnlyReturnsChangedFields(t *testing.T) {
t.Fatal("did not expect unchanged ip to be included")
}
}
func TestListNodeViewsDoesNotPersistComputedStatus(t *testing.T) {
setupServiceTestDB(t)
node := &model.Node{
NodeID: "node-offline-view",
Name: "edge-offline",
IP: "10.0.0.21",
AgentToken: "token-offline",
AgentVersion: "v0.5.0",
NginxVersion: "1.27.1.2",
Status: NodeStatusOnline,
LastSeenAt: time.Now().Add(-common.NodeOfflineThreshold - time.Minute),
}
if err := node.Insert(); err != nil {
t.Fatalf("failed to insert node: %v", err)
}
views, err := ListNodeViews()
if err != nil {
t.Fatalf("ListNodeViews failed: %v", err)
}
if len(views) != 1 {
t.Fatalf("expected 1 node view, got %d", len(views))
}
if views[0].Status != NodeStatusOffline {
t.Fatalf("expected computed offline status in view, got %s", views[0].Status)
}
storedNode, err := model.GetNodeByID(node.ID)
if err != nil {
t.Fatalf("failed to reload node: %v", err)
}
if storedNode.Status != NodeStatusOnline {
t.Fatalf("expected list query to avoid persisting computed status, got %s", storedNode.Status)
}
}