package main import ( "encoding/json" "fmt" "net/http" "strconv" "sync" "time" ) // TraceEvent represents a single step in the message processing pipeline. type TraceEvent struct { ID string `json:"id"` Timestamp time.Time `json:"timestamp"` SessionID string `json:"session_id,omitempty"` UserID string `json:"user_id,omitempty"` Hop string `json:"hop"` // message_received, intent, subsession, llm_call, tool_call, vision, synthesis, review, response, think Label string `json:"label"` Status string `json:"status"` // success, error, running Detail string `json:"detail,omitempty"` DurationMs int64 `json:"duration_ms,omitempty"` Data map[string]interface{} `json:"data,omitempty"` } var ( traceMu sync.Mutex traceEvents []TraceEvent traceMax = 200 ) // AddTraceEvent appends a trace event to the in-memory ring buffer. func AddTraceEvent(hop, sessionID, userID, label, status, detail string, durationMs int64, data map[string]interface{}) { traceMu.Lock() defer traceMu.Unlock() ev := TraceEvent{ ID: fmt.Sprintf("%s-%d", hop, time.Now().UnixNano()), Timestamp: time.Now(), SessionID: sessionID, UserID: userID, Hop: hop, Label: label, Status: status, Detail: detail, DurationMs: durationMs, Data: data, } traceEvents = append(traceEvents, ev) if len(traceEvents) > traceMax { traceEvents = traceEvents[len(traceEvents)-traceMax:] } } // GetTraceEvents returns recent trace events, optionally filtered by session. func GetTraceEvents(sessionID string, limit int) []TraceEvent { traceMu.Lock() defer traceMu.Unlock() if limit <= 0 || limit > len(traceEvents) { limit = len(traceEvents) } var result []TraceEvent // Return newest first. for i := len(traceEvents) - 1; i >= 0 && len(result) < limit; i-- { ev := traceEvents[i] if sessionID == "" || ev.SessionID == sessionID { result = append(result, ev) } } return result } // registerTraceEndpoint adds the trace events API. func registerTraceEndpoint(mux *http.ServeMux) { mux.HandleFunc("/api/v1/trace/events", func(w http.ResponseWriter, r *http.Request) { sessionID := r.URL.Query().Get("session_id") limit := 100 if l, err := strconv.Atoi(r.URL.Query().Get("limit")); err == nil && l > 0 && l <= 500 { limit = l } events := GetTraceEvents(sessionID, limit) w.Header().Set("Content-Type", "application/json") json.NewEncoder(w).Encode(map[string]interface{}{ "events": events, "total": len(events), }) }) }