diff --git a/pmolog/webloger.go b/pmolog/webloger.go index 9da89120..92d7d9bc 100644 --- a/pmolog/webloger.go +++ b/pmolog/webloger.go @@ -2,331 +2,174 @@ package pmolog import ( "container/ring" + "encoding/json" "fmt" "net/http" "sync" + "time" "github.com/sirupsen/logrus" ) -// Structure pour gérer les connexions SSE +const bufferSize = 1000 + +// ---------- SSE Broker ---------- + type SSEBroker struct { clients map[chan string]bool - mutex sync.Mutex + mu sync.RWMutex } -// Buffer circulaire pour stocker les 1000 derniers messages var ( - logBuffer = ring.New(1000) + broker = &SSEBroker{clients: make(map[chan string]bool)} + logBuffer = ring.New(bufferSize) bufferMutex sync.Mutex ) -// Initialiser le broker SSE -var broker = &SSEBroker{ - clients: make(map[chan string]bool), -} +// ---------- Hook for Logrus ---------- -// Hook personnalisé pour Logrus qui envoie les logs aux clients SSE type SSELogHook struct{} -func (hook *SSELogHook) Levels() []logrus.Level { - return logrus.AllLevels -} +func (SSELogHook) Levels() []logrus.Level { return logrus.AllLevels } -func (hook *SSELogHook) Fire(entry *logrus.Entry) error { - // Formater le log avec le niveau et le message - logLine := fmt.Sprintf("[%s] %s", entry.Level.String(), entry.Message) +func (SSELogHook) Fire(entry *logrus.Entry) error { + msg := map[string]string{ + "time": time.Now().Format(time.RFC3339), + "level": entry.Level.String(), + "content": entry.Message, + } + b, _ := json.Marshal(msg) - // Ajouter le message au buffer circulaire + // add to buffer bufferMutex.Lock() - logBuffer.Value = logLine + logBuffer.Value = string(b) logBuffer = logBuffer.Next() bufferMutex.Unlock() - // Envoyer le log à tous les clients connectés - broker.mutex.Lock() - for client := range broker.clients { + // broadcast + broker.mu.RLock() + for ch := range broker.clients { select { - case client <- logLine: - default: - // Client saturé, on skip + case ch <- string(b): + default: // skip if full } } - broker.mutex.Unlock() - + broker.mu.RUnlock() return nil } -// Handler SSE qui envoie les logs en temps réel +// ---------- SSE Handler ---------- + func sseHandler(w http.ResponseWriter, r *http.Request) { - // Configurer les en-têtes SSE w.Header().Set("Content-Type", "text/event-stream") w.Header().Set("Cache-Control", "no-cache") w.Header().Set("Connection", "keep-alive") w.Header().Set("Access-Control-Allow-Origin", "*") - // Créer un canal pour ce client - messageChan := make(chan string, 10) + flusher, ok := w.(http.Flusher) + if !ok { + http.Error(w, "Streaming unsupported", http.StatusInternalServerError) + return + } - // Ajouter le client au broker - broker.mutex.Lock() - broker.clients[messageChan] = true - broker.mutex.Unlock() + ch := make(chan string, 20) + broker.mu.Lock() + broker.clients[ch] = true + broker.mu.Unlock() - // Envoyer d'abord les 1000 derniers messages stockés + // replay buffer bufferMutex.Lock() - logBuffer.Do(func(value interface{}) { - if value != nil { - if msg, ok := value.(string); ok { - // Déterminer le niveau de log pour le style CSS - level := "info" - if len(msg) > 7 { - switch msg[1:6] { - case "ERROR": - level = "error" - case "WARNI": - level = "warning" - case "DEBUG": - level = "debug" - } - } - - // Formater le message en JSON pour inclure le niveau - jsonMsg := fmt.Sprintf("{\"content\": \"%s\", \"level\": \"%s\"}", escapeJSONString(msg), level) - fmt.Fprintf(w, "event: message\ndata: %s\n\n", jsonMsg) - } + logBuffer.Do(func(v interface{}) { + if v != nil { + fmt.Fprintf(w, "event: message\ndata: %s\n\n", v.(string)) } }) - w.(http.Flusher).Flush() + flusher.Flush() bufferMutex.Unlock() - // Envoyer les logs au client au fur et à mesure + // stream new messages for { select { - case msg := <-messageChan: - // Déterminer le niveau de log pour le style CSS - level := "info" - if len(msg) > 7 { - switch msg[1:6] { - case "ERROR": - level = "error" - case "WARNI": - level = "warning" - case "DEBUG": - level = "debug" - } - } - - // Formater le message en JSON pour inclure le niveau - jsonMsg := fmt.Sprintf("{\"content\": \"%s\", \"level\": \"%s\"}", escapeJSONString(msg), level) - fmt.Fprintf(w, "event: message\ndata: %s\n\n", jsonMsg) - w.(http.Flusher).Flush() + case msg := <-ch: + fmt.Fprintf(w, "event: message\ndata: %s\n\n", msg) + flusher.Flush() case <-r.Context().Done(): - // Supprimer le client quand la connexion est fermée - broker.mutex.Lock() - delete(broker.clients, messageChan) - broker.mutex.Unlock() - close(messageChan) + broker.mu.Lock() + delete(broker.clients, ch) + broker.mu.Unlock() + close(ch) return } } } -// Fonction pour échapper les chaînes JSON -func escapeJSONString(s string) string { - // Échapper les guillemets et les antislashes - escaped := "" - for _, c := range s { - switch c { - case '"': - escaped += "\\\"" - case '\\': - escaped += "\\\\" - case '\n': - escaped += "\\n" - case '\r': - escaped += "\\r" - case '\t': - escaped += "\\t" - default: - escaped += string(c) - } - } - return escaped -} +// ---------- HTML ---------- -// Page HTML pour afficher les logs avec support Markdown -var indexHTML = ` - +var indexHTML = ` - Logs en temps réel - - - + + 🚀 Real-Time Logs + + + + -
-

📝 Logs en temps réel (1000 derniers messages)

-
-
- - - - -` +

📝 Logs en temps réel

+
+ + +` + +// ---------- Handlers ---------- -// Handler pour servir la page HTML func indexHandler(w http.ResponseWriter, r *http.Request) { w.Header().Set("Content-Type", "text/html; charset=utf-8") fmt.Fprint(w, indexHTML) } func LoggerWeb(mux *http.ServeMux) { - // Configurer Logrus pour le développement - logrus.SetFormatter(&logrus.TextFormatter{ - ForceColors: true, - FullTimestamp: true, - }) - - // Ajouter le hook SSE à Logrus - logrus.AddHook(&SSELogHook{}) - + logrus.SetFormatter(&logrus.TextFormatter{ForceColors: true, FullTimestamp: true}) + logrus.AddHook(SSELogHook{}) mux.HandleFunc("/log", indexHandler) mux.HandleFunc("/log-sse", sseHandler) - logrus.Info("Web logger connected") - }