package api import ( "encoding/json" "strings" fleetai "crypto-miner-server/internal/ai" "crypto-miner-server/internal/spreadrouter" ) // PathTraceSurgicalAdapter resolves persisted Path Tracer sessions for surgical replay. type PathTraceSurgicalAdapter struct { Hub *WSHub PathTrace *PathTracerHandler } func (a *PathTraceSurgicalAdapter) TraceForAgent(agentID string) (fleetai.SurgicalTraceContext, bool) { if a == nil || agentID == "" { return fleetai.SurgicalTraceContext{}, false } sessions := traceSessionsSnapshot(a.PathTrace) if len(sessions) == 0 && a.Hub != nil && a.Hub.db != nil { rows, err := a.Hub.db.ListPathTraceSessions() if err != nil { return fleetai.SurgicalTraceContext{}, false } for _, row := range rows { var rec pathTraceSessionPersist if json.Unmarshal(row.Payload, &rec) != nil { continue } sess := &TraceSession{ ID: rec.ID, AgentIDs: rec.AgentIDs, Hops: rec.Hops, Error: rec.Error, DiscoverError: rec.DiscoverError, } sessions = append(sessions, sess) } } for _, sess := range sessions { if ctx, ok := surgicalTraceFromSession(sess, agentID); ok { return ctx, true } } return fleetai.SurgicalTraceContext{}, false } func surgicalTraceFromSession(sess *TraceSession, agentID string) (fleetai.SurgicalTraceContext, bool) { if sess == nil || agentID == "" { return fleetai.SurgicalTraceContext{}, false } hopIndex := -1 for i, hop := range sess.Hops { if hop != nil && hop.AgentID == agentID { hopIndex = i break } } if hopIndex < 0 { for _, id := range sess.AgentIDs { if id == agentID { hopIndex = 0 break } } } if hopIndex < 0 { return fleetai.SurgicalTraceContext{}, false } ctx := fleetai.SurgicalTraceContext{ SessionID: sess.ID, HopIndex: hopIndex, HopCount: len(sess.Hops), SessionError: strings.TrimSpace(sess.Error), DiscoverError: strings.TrimSpace(sess.DiscoverError), } if len(sess.Hops) > 0 { last := sess.Hops[len(sess.Hops)-1] if last != nil { ctx.EgressAgentID = last.AgentID } } for _, host := range serviceGraphList(sess.ServiceGraph) { sub := spreadrouter.NormalizeSubnet(host.Subnet) if sub == "" { sub = spreadrouter.SubnetFromIP(host.Host) } if sub != "" { ctx.TargetSubnets = append(ctx.TargetSubnets, sub) } } return ctx, true } // DatabaseStrainMemoryAdapter persists surgical replay outcomes. type DatabaseStrainMemoryAdapter struct { DB interface { InsertStrainMemory(agentID, sessionID, failedTier, strain, fixType, fixArgs, outcome string) error } } func (a *DatabaseStrainMemoryAdapter) InsertStrainMemory(agentID, sessionID, failedTier, strain, fixType, fixArgs, outcome string) error { if a == nil || a.DB == nil { return nil } return a.DB.InsertStrainMemory(agentID, sessionID, failedTier, strain, fixType, fixArgs, outcome) } // HubSeerEmitter broadcasts and persists Seer feed events. type HubSeerEmitter struct { Hub *WSHub DB interface { InsertSeerEvent(eventType, agentID string, payload []byte) (int64, error) } } func (e *HubSeerEmitter) EmitSeerEvent(eventType, agentID string, payload map[string]interface{}) error { if e == nil { return nil } raw, _ := json.Marshal(payload) var id int64 if e.DB != nil { var err error id, err = e.DB.InsertSeerEvent(eventType, agentID, raw) if err != nil { return err } } if e.Hub != nil { e.Hub.BroadcastSeerEvent(map[string]interface{}{ "id": id, "event_type": eventType, "agent_id": agentID, "payload": json.RawMessage(raw), }) } return nil } // StrainFromAgent returns spread_strain for strain memory rows. func StrainFromAgent(hub *WSHub, agentID string) string { if hub == nil || hub.db == nil || agentID == "" { return "" } ag, err := hub.db.GetAgent(agentID) if err != nil || ag == nil { return "" } return strings.TrimSpace(ag.SpreadStrain) }