mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-10 17:26:38 +08:00
[优化] 数据库结构优化
This commit is contained in:
@@ -74,7 +74,7 @@ func (r *Runner) Run(ctx context.Context) error {
|
||||
return err
|
||||
}
|
||||
slog.Info("agent runner started", "node_id", nodeID, "node", r.Config.NodeName, "ip", r.Config.NodeIP)
|
||||
if r.hasAgentToken() {
|
||||
if r.hasAccessToken() {
|
||||
if _, hbErr := r.performHeartbeatCycle(ctx, nodeID, true); hbErr != nil {
|
||||
slog.Error("agent startup heartbeat failed", "error", hbErr)
|
||||
}
|
||||
@@ -119,7 +119,7 @@ func (r *Runner) Run(ctx context.Context) error {
|
||||
delay := wsBackoff.Next()
|
||||
nextWSAttempt = time.Now().Add(delay)
|
||||
slog.Debug("agent ws disconnected; resuming http heartbeat", "retry_after", delay, "error", wsErr)
|
||||
if r.hasAgentToken() {
|
||||
if r.hasAccessToken() {
|
||||
if _, hbErr := r.performHeartbeatCycle(ctx, nodeID, false); hbErr != nil {
|
||||
slog.Error("agent heartbeat after ws disconnect failed", "error", hbErr)
|
||||
}
|
||||
@@ -128,7 +128,7 @@ func (r *Runner) Run(ctx context.Context) error {
|
||||
if wsDone != nil {
|
||||
continue
|
||||
}
|
||||
if !r.hasAgentToken() {
|
||||
if !r.hasAccessToken() {
|
||||
if err = r.tryRegister(ctx, &nodeID); err != nil {
|
||||
slog.Error("agent discovery register failed", "error", err)
|
||||
}
|
||||
@@ -181,7 +181,7 @@ func (r *Runner) performHeartbeatCycle(ctx context.Context, nodeID string, start
|
||||
}
|
||||
|
||||
func (r *Runner) shouldUseWebSocket() bool {
|
||||
enabled := r.WebSocketService != nil && r.websocketUpgradeEnabled && r.hasAgentToken()
|
||||
enabled := r.WebSocketService != nil && r.websocketUpgradeEnabled && r.hasAccessToken()
|
||||
slog.Debug("agent ws upgrade eligibility checked", "enabled", enabled, "server_enabled", r.websocketUpgradeEnabled, "url", r.websocketURL())
|
||||
return enabled
|
||||
}
|
||||
@@ -365,8 +365,8 @@ func (backoff *webSocketBackoff) Reset() {
|
||||
}
|
||||
}
|
||||
|
||||
func (r *Runner) hasAgentToken() bool {
|
||||
return strings.TrimSpace(r.Config.AgentToken) != ""
|
||||
func (r *Runner) hasAccessToken() bool {
|
||||
return strings.TrimSpace(r.Config.AccessToken) != ""
|
||||
}
|
||||
|
||||
func (r *Runner) applySettings(settings *protocol.AgentSettings) bool {
|
||||
@@ -447,7 +447,7 @@ func (r *Runner) tryRegister(ctx context.Context, nodeID *string) error {
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if response == nil || strings.TrimSpace(response.AgentToken) == "" || strings.TrimSpace(response.NodeID) == "" {
|
||||
if response == nil || strings.TrimSpace(response.AccessToken) == "" || strings.TrimSpace(response.NodeID) == "" {
|
||||
return errors.New("discovery register response 缺少 node_id 或 agent_token")
|
||||
}
|
||||
snapshot, err := r.StateStore.Load()
|
||||
@@ -458,14 +458,14 @@ func (r *Runner) tryRegister(ctx context.Context, nodeID *string) error {
|
||||
if err = r.StateStore.Save(snapshot); err != nil {
|
||||
return err
|
||||
}
|
||||
r.Config.AgentToken = response.AgentToken
|
||||
r.Config.AccessToken = response.AccessToken
|
||||
r.Config.DiscoveryToken = ""
|
||||
if err = r.Config.Save(); err != nil {
|
||||
return err
|
||||
}
|
||||
r.HeartbeatService.SetToken(response.AgentToken)
|
||||
r.HeartbeatService.SetToken(response.AccessToken)
|
||||
if r.WebSocketService != nil {
|
||||
r.WebSocketService.SetToken(response.AgentToken)
|
||||
r.WebSocketService.SetToken(response.AccessToken)
|
||||
}
|
||||
*nodeID = response.NodeID
|
||||
slog.Info("agent discovery registration succeeded", "node_id", response.NodeID)
|
||||
@@ -573,23 +573,25 @@ func (r *Runner) nodePayload(nodeID string) protocol.NodePayload {
|
||||
if managedOpenRestyMetrics == nil {
|
||||
managedOpenRestyMetrics = fallbackMetrics
|
||||
}
|
||||
metricSnapshot := observability.BuildSnapshot(r.Config, r.StateStore, managedOpenRestyMetrics)
|
||||
metricSnapshot := observability.BuildSnapshot(r.Config, r.StateStore)
|
||||
openrestyObservation := observability.BuildOpenrestyObservation(managedOpenRestyMetrics)
|
||||
healthEvents := observability.BuildHealthEvents(snapshot)
|
||||
payload := protocol.NodePayload{
|
||||
NodeID: nodeID,
|
||||
Name: r.Config.NodeName,
|
||||
IP: r.Config.NodeIP,
|
||||
AgentVersion: r.Config.AgentVersion,
|
||||
NginxVersion: r.Config.NginxVersion,
|
||||
CurrentVersion: snapshot.CurrentVersion,
|
||||
LastError: snapshot.LastError,
|
||||
OpenrestyStatus: openrestyStatus,
|
||||
OpenrestyMessage: snapshot.OpenrestyMessage,
|
||||
Profile: profile,
|
||||
Snapshot: metricSnapshot,
|
||||
TrafficReport: trafficReport,
|
||||
AccessLogs: accessLogs,
|
||||
HealthEvents: healthEvents,
|
||||
NodeID: nodeID,
|
||||
Name: r.Config.NodeName,
|
||||
IP: r.Config.NodeIP,
|
||||
Version: r.Config.Version,
|
||||
ExtVersion: r.Config.ExtVersion,
|
||||
CurrentVersion: snapshot.CurrentVersion,
|
||||
LastError: snapshot.LastError,
|
||||
OpenrestyStatus: openrestyStatus,
|
||||
OpenrestyMessage: snapshot.OpenrestyMessage,
|
||||
Profile: profile,
|
||||
Snapshot: metricSnapshot,
|
||||
OpenrestyObservation: openrestyObservation,
|
||||
TrafficReport: trafficReport,
|
||||
AccessLogs: accessLogs,
|
||||
HealthEvents: healthEvents,
|
||||
}
|
||||
if r.SyncService != nil {
|
||||
checksums, err := r.SyncService.WAFIPGroupChecksums()
|
||||
@@ -619,17 +621,18 @@ func (r *Runner) prepareHeartbeatPayload(nodeID string) (protocol.NodePayload, [
|
||||
}
|
||||
now := time.Now().UTC()
|
||||
retainAfterUnix := now.Add(-time.Duration(r.Config.ObservabilityReplayMinutes) * time.Minute).Unix()
|
||||
windowStartedAtUnix := state.ObservabilityWindowStartedAt(payload.Snapshot, payload.TrafficReport)
|
||||
windowStartedAtUnix := state.ObservabilityWindowStartedAt(payload.Snapshot, payload.OpenrestyObservation, payload.TrafficReport)
|
||||
if windowStartedAtUnix <= 0 {
|
||||
return payload, nil
|
||||
}
|
||||
|
||||
record := state.ObservabilityBufferRecord{
|
||||
WindowStartedAtUnix: windowStartedAtUnix,
|
||||
Snapshot: payload.Snapshot,
|
||||
TrafficReport: payload.TrafficReport,
|
||||
AccessLogs: payload.AccessLogs,
|
||||
QueuedAtUnix: now.Unix(),
|
||||
WindowStartedAtUnix: windowStartedAtUnix,
|
||||
Snapshot: payload.Snapshot,
|
||||
OpenrestyObservation: payload.OpenrestyObservation,
|
||||
TrafficReport: payload.TrafficReport,
|
||||
AccessLogs: payload.AccessLogs,
|
||||
QueuedAtUnix: now.Unix(),
|
||||
}
|
||||
if err := r.ObservabilityBuffer.Upsert(record, retainAfterUnix); err != nil {
|
||||
slog.Error("upsert observability buffer failed", "error", err)
|
||||
@@ -649,10 +652,11 @@ func (r *Runner) prepareHeartbeatPayload(nodeID string) (protocol.NodePayload, [
|
||||
continue
|
||||
}
|
||||
buffered = append(buffered, protocol.BufferedObservabilityRecord{
|
||||
WindowStartedAtUnix: item.WindowStartedAtUnix,
|
||||
Snapshot: item.Snapshot,
|
||||
TrafficReport: item.TrafficReport,
|
||||
AccessLogs: item.AccessLogs,
|
||||
WindowStartedAtUnix: item.WindowStartedAtUnix,
|
||||
Snapshot: item.Snapshot,
|
||||
OpenrestyObservation: item.OpenrestyObservation,
|
||||
TrafficReport: item.TrafficReport,
|
||||
AccessLogs: item.AccessLogs,
|
||||
})
|
||||
ackWindows = append(ackWindows, item.WindowStartedAtUnix)
|
||||
}
|
||||
|
||||
@@ -194,11 +194,11 @@ func TestRunnerKeepsHeartbeatWhenStartupSyncFails(t *testing.T) {
|
||||
}
|
||||
runner := &Runner{
|
||||
Config: &config.Config{
|
||||
AgentToken: "agent-token",
|
||||
AccessToken: "agent-token",
|
||||
NodeName: "edge-01",
|
||||
NodeIP: "10.0.0.8",
|
||||
AgentVersion: config.AgentVersion,
|
||||
NginxVersion: "1.27.1.2",
|
||||
Version: config.Version,
|
||||
ExtVersion: "1.27.1.2",
|
||||
HeartbeatInterval: config.MillisecondDuration(10 * time.Millisecond),
|
||||
},
|
||||
StateStore: stateStore,
|
||||
@@ -247,11 +247,11 @@ func TestRunnerDoesNotExitOnHeartbeatOrSyncError(t *testing.T) {
|
||||
}
|
||||
runner := &Runner{
|
||||
Config: &config.Config{
|
||||
AgentToken: "agent-token",
|
||||
AccessToken: "agent-token",
|
||||
NodeName: "edge-01",
|
||||
NodeIP: "10.0.0.8",
|
||||
AgentVersion: config.AgentVersion,
|
||||
NginxVersion: "1.27.1.2",
|
||||
Version: config.Version,
|
||||
ExtVersion: "1.27.1.2",
|
||||
HeartbeatInterval: config.MillisecondDuration(10 * time.Millisecond),
|
||||
},
|
||||
StateStore: stateStore,
|
||||
@@ -305,11 +305,11 @@ func TestRunnerReportsOpenrestyHealthAndExecutesRestart(t *testing.T) {
|
||||
}
|
||||
runner := &Runner{
|
||||
Config: &config.Config{
|
||||
AgentToken: "agent-token",
|
||||
AccessToken: "agent-token",
|
||||
NodeName: "edge-01",
|
||||
NodeIP: "10.0.0.8",
|
||||
AgentVersion: config.AgentVersion,
|
||||
NginxVersion: "1.27.1.2",
|
||||
Version: config.Version,
|
||||
ExtVersion: "1.27.1.2",
|
||||
HeartbeatInterval: config.MillisecondDuration(10 * time.Millisecond),
|
||||
},
|
||||
StateStore: stateStore,
|
||||
@@ -361,8 +361,8 @@ func TestRunnerHeartbeatPayloadIncludesObservabilityExtensions(t *testing.T) {
|
||||
Config: &config.Config{
|
||||
NodeName: "edge-observe-1",
|
||||
NodeIP: "10.0.0.51",
|
||||
AgentVersion: config.AgentVersion,
|
||||
NginxVersion: "1.27.1.2",
|
||||
Version: config.Version,
|
||||
ExtVersion: "1.27.1.2",
|
||||
DataDir: tempDir,
|
||||
RouteConfigPath: filepath.Join(tempDir, "conf.d", "openflare_routes.conf"),
|
||||
AccessLogPath: filepath.Join(tempDir, "var", "log", "openflare", "access.log"),
|
||||
@@ -441,11 +441,11 @@ func TestRunnerReplaysBufferedObservabilityAfterHeartbeatRecovery(t *testing.T)
|
||||
}
|
||||
runner := &Runner{
|
||||
Config: &config.Config{
|
||||
AgentToken: "agent-token",
|
||||
AccessToken: "agent-token",
|
||||
NodeName: "edge-buffer-01",
|
||||
NodeIP: "10.0.0.52",
|
||||
AgentVersion: config.AgentVersion,
|
||||
NginxVersion: "1.27.1.2",
|
||||
Version: config.Version,
|
||||
ExtVersion: "1.27.1.2",
|
||||
DataDir: tempDir,
|
||||
RouteConfigPath: filepath.Join(tempDir, "conf.d", "openflare_routes.conf"),
|
||||
HeartbeatInterval: config.MillisecondDuration(10 * time.Millisecond),
|
||||
@@ -498,9 +498,9 @@ func TestRunnerDiscoveryRegisterUpdatesTokenAndNodeID(t *testing.T) {
|
||||
stateStore := state.NewStore(filepath.Join(t.TempDir(), "state.json"))
|
||||
heartbeatService := &fakeHeartbeatService{
|
||||
registerResp: &protocol.RegisterNodeResponse{
|
||||
NodeID: "node-server-assigned",
|
||||
AgentToken: "agent-token-issued",
|
||||
Name: "edge-01",
|
||||
NodeID: "node-server-assigned",
|
||||
AccessToken: "agent-token-issued",
|
||||
Name: "edge-01",
|
||||
},
|
||||
heartbeatResults: []*protocol.HeartbeatResult{{}},
|
||||
onHeartbeat: func(callCount int) {
|
||||
@@ -524,8 +524,8 @@ func TestRunnerDiscoveryRegisterUpdatesTokenAndNodeID(t *testing.T) {
|
||||
DiscoveryToken: cfg.DiscoveryToken,
|
||||
NodeName: cfg.NodeName,
|
||||
NodeIP: cfg.NodeIP,
|
||||
AgentVersion: config.AgentVersion,
|
||||
NginxVersion: "1.27.1.2",
|
||||
Version: config.Version,
|
||||
ExtVersion: "1.27.1.2",
|
||||
HeartbeatInterval: config.MillisecondDuration(10 * time.Millisecond),
|
||||
},
|
||||
StateStore: stateStore,
|
||||
@@ -533,8 +533,8 @@ func TestRunnerDiscoveryRegisterUpdatesTokenAndNodeID(t *testing.T) {
|
||||
SyncService: syncService,
|
||||
}
|
||||
runner.Config = cfg
|
||||
runner.Config.AgentVersion = config.AgentVersion
|
||||
runner.Config.NginxVersion = "1.27.1.2"
|
||||
runner.Config.Version = config.Version
|
||||
runner.Config.ExtVersion = "1.27.1.2"
|
||||
runner.Config.HeartbeatInterval = config.MillisecondDuration(10 * time.Millisecond)
|
||||
|
||||
err = runner.Run(ctx)
|
||||
@@ -554,7 +554,7 @@ func TestRunnerDiscoveryRegisterUpdatesTokenAndNodeID(t *testing.T) {
|
||||
if snapshot.NodeID != "node-server-assigned" {
|
||||
t.Fatalf("expected node id to be replaced, got %q", snapshot.NodeID)
|
||||
}
|
||||
if runner.Config.AgentToken != "agent-token-issued" || runner.Config.DiscoveryToken != "" {
|
||||
if runner.Config.AccessToken != "agent-token-issued" || runner.Config.DiscoveryToken != "" {
|
||||
t.Fatal("expected config token rotation to complete")
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user