feat(logging): add real-time log streaming via websocket
- Implement WebSocket endpoint for streaming logs to clients - Add log broadcaster to manage client connections and history - Update documentation with new WebSocket API details - Include gorilla/websocket dependency for WebSocket support
This commit is contained in:
@@ -0,0 +1,96 @@
|
||||
package postLog
|
||||
|
||||
import (
|
||||
"sync"
|
||||
)
|
||||
|
||||
type LogMessage struct {
|
||||
Level int `json:"level"`
|
||||
Content string `json:"content"`
|
||||
Timestamp string `json:"timestamp"`
|
||||
}
|
||||
|
||||
type Client struct {
|
||||
ID string
|
||||
Messages chan LogMessage
|
||||
}
|
||||
|
||||
type LogBroadcaster struct {
|
||||
mu sync.RWMutex
|
||||
clients map[string]chan LogMessage
|
||||
history []LogMessage
|
||||
historyM sync.RWMutex
|
||||
}
|
||||
|
||||
var broadcaster *LogBroadcaster
|
||||
|
||||
func InitLogBroadcaster() {
|
||||
broadcaster = &LogBroadcaster{
|
||||
clients: make(map[string]chan LogMessage),
|
||||
history: make([]LogMessage, 0),
|
||||
}
|
||||
}
|
||||
|
||||
func GetLogBroadcaster() *LogBroadcaster {
|
||||
return broadcaster
|
||||
}
|
||||
|
||||
func (lb *LogBroadcaster) AddClient(id string) chan LogMessage {
|
||||
lb.mu.Lock()
|
||||
defer lb.mu.Unlock()
|
||||
|
||||
ch := make(chan LogMessage, 100)
|
||||
lb.clients[id] = ch
|
||||
return ch
|
||||
}
|
||||
|
||||
func (lb *LogBroadcaster) RemoveClient(id string) {
|
||||
lb.mu.Lock()
|
||||
defer lb.mu.Unlock()
|
||||
|
||||
if ch, exists := lb.clients[id]; exists {
|
||||
close(ch)
|
||||
delete(lb.clients, id)
|
||||
}
|
||||
}
|
||||
|
||||
func (lb *LogBroadcaster) Broadcast(msg LogMessage) {
|
||||
lb.mu.RLock()
|
||||
defer lb.mu.RUnlock()
|
||||
|
||||
for _, ch := range lb.clients {
|
||||
select {
|
||||
case ch <- msg:
|
||||
default:
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (lb *LogBroadcaster) GetHistory() []LogMessage {
|
||||
lb.historyM.RLock()
|
||||
defer lb.historyM.RUnlock()
|
||||
|
||||
result := make([]LogMessage, len(lb.history))
|
||||
copy(result, lb.history)
|
||||
return result
|
||||
}
|
||||
|
||||
func (lb *LogBroadcaster) AddToHistory(msg LogMessage) {
|
||||
lb.historyM.Lock()
|
||||
defer lb.historyM.Unlock()
|
||||
|
||||
lb.history = append(lb.history, msg)
|
||||
if len(lb.history) > 100 {
|
||||
lb.history = lb.history[1:]
|
||||
}
|
||||
}
|
||||
|
||||
func (lb *LogBroadcaster) SendHistory(clientCh chan LogMessage) {
|
||||
history := lb.GetHistory()
|
||||
for _, msg := range history {
|
||||
select {
|
||||
case clientCh <- msg:
|
||||
default:
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,66 @@
|
||||
package postLog
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"log"
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"github.com/gorilla/websocket"
|
||||
)
|
||||
|
||||
var upgrader = websocket.Upgrader{
|
||||
ReadBufferSize: 1024,
|
||||
WriteBufferSize: 1024,
|
||||
CheckOrigin: func(r *http.Request) bool {
|
||||
return true
|
||||
},
|
||||
}
|
||||
|
||||
type LogSocketHandler struct {
|
||||
broadcaster *LogBroadcaster
|
||||
}
|
||||
|
||||
func NewLogSocketHandler(b *LogBroadcaster) *LogSocketHandler {
|
||||
return &LogSocketHandler{
|
||||
broadcaster: b,
|
||||
}
|
||||
}
|
||||
|
||||
func (h *LogSocketHandler) Handle(w http.ResponseWriter, r *http.Request) {
|
||||
conn, err := upgrader.Upgrade(w, r, nil)
|
||||
if err != nil {
|
||||
log.Printf("Failed to upgrade connection: %v", err)
|
||||
return
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
clientID := conn.RemoteAddr().String() + "-" + time.Now().Format("20060102150405")
|
||||
clientCh := h.broadcaster.AddClient(clientID)
|
||||
defer h.broadcaster.RemoveClient(clientID)
|
||||
|
||||
h.broadcaster.SendHistory(clientCh)
|
||||
|
||||
done := make(chan struct{})
|
||||
go func() {
|
||||
defer close(done)
|
||||
for {
|
||||
msg := <-clientCh
|
||||
data, err := json.Marshal(msg)
|
||||
if err != nil {
|
||||
log.Printf("Failed to marshal log message: %v", err)
|
||||
return
|
||||
}
|
||||
if err := conn.WriteMessage(websocket.TextMessage, data); err != nil {
|
||||
return
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
for {
|
||||
_, _, err := conn.ReadMessage()
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
+14
-8
@@ -31,7 +31,6 @@ func getSystemTime() string {
|
||||
return now.Format("2006-01-02 15:04:05.000")
|
||||
}
|
||||
|
||||
// SetDebugMode sets the debug mode for the logger
|
||||
func SetDebugMode(debug bool) {
|
||||
loggerMutex.Lock()
|
||||
defer loggerMutex.Unlock()
|
||||
@@ -43,31 +42,38 @@ func PostLog(message string, level int) {
|
||||
|
||||
idx := level
|
||||
if idx < 0 || idx >= len(levelNames) {
|
||||
idx = INFO // default to INFO
|
||||
idx = INFO
|
||||
}
|
||||
|
||||
// Lock to make reading debug flag and all output atomic across threads
|
||||
loggerMutex.Lock()
|
||||
|
||||
// Skip DEBUG when debug is off
|
||||
if idx == DEBUG && !isDebug {
|
||||
loggerMutex.Unlock()
|
||||
return
|
||||
}
|
||||
|
||||
// Colored output
|
||||
levelDisplay := colorOut_256(levelNames[idx], levelColors[idx])
|
||||
|
||||
// Copy logsDB to avoid holding the lock during DB operation
|
||||
db := logsDB
|
||||
|
||||
loggerMutex.Unlock()
|
||||
|
||||
fmt.Printf("[%s - %s] %s\n", timeNow, levelDisplay, message)
|
||||
insertLogToDB(db, level, message, timeNow)
|
||||
|
||||
if broadcaster != nil {
|
||||
broadcaster.AddToHistory(LogMessage{
|
||||
Level: level,
|
||||
Content: message,
|
||||
Timestamp: timeNow,
|
||||
})
|
||||
broadcaster.Broadcast(LogMessage{
|
||||
Level: level,
|
||||
Content: message,
|
||||
Timestamp: timeNow,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// Helper functions for different log levels
|
||||
func Debug(message string) {
|
||||
PostLog(message, DEBUG)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user