handleEvent 处理事件推送
(ctx context.Context, payload json.RawMessage)
| 306 | |
| 307 | // handleEvent 处理事件推送 |
| 308 | func (h *RemoteAgentHandler) handleEvent(ctx context.Context, payload json.RawMessage) error { |
| 309 | var evt EventPayload |
| 310 | if err := json.Unmarshal(payload, &evt); err != nil { |
| 311 | return fmt.Errorf("failed to parse event payload: %w", err) |
| 312 | } |
| 313 | |
| 314 | logging.Info(ctx, "remote_agent.event.received", map[string]any{ |
| 315 | "agent_id": evt.AgentID, |
| 316 | "event_type": fmt.Sprintf("%T", evt.Envelope.Event), |
| 317 | }) |
| 318 | |
| 319 | h.mu.RLock() |
| 320 | remoteAgent, exists := h.agents[evt.AgentID] |
| 321 | h.mu.RUnlock() |
| 322 | |
| 323 | if !exists { |
| 324 | return fmt.Errorf("agent %s not registered", evt.AgentID) |
| 325 | } |
| 326 | |
| 327 | // 推送事件到 RemoteAgent 的 EventBus |
| 328 | if err := remoteAgent.PushEvent(evt.Envelope); err != nil { |
| 329 | return fmt.Errorf("failed to push event: %w", err) |
| 330 | } |
| 331 | |
| 332 | logging.Info(ctx, "remote_agent.event.pushed", map[string]any{ |
| 333 | "agent_id": evt.AgentID, |
| 334 | }) |
| 335 | |
| 336 | return nil |
| 337 | } |
| 338 | |
| 339 | // handleUnregister 处理 Agent 注销 |
| 340 | func (h *RemoteAgentHandler) handleUnregister(ctx context.Context, payload json.RawMessage) error { |
no test coverage detected