diff --git a/cmd/picoclaw/internal/gateway/helpers.go b/cmd/picoclaw/internal/gateway/helpers.go index 353883f1d..e8ab00adc 100644 --- a/cmd/picoclaw/internal/gateway/helpers.go +++ b/cmd/picoclaw/internal/gateway/helpers.go @@ -3,7 +3,6 @@ package gateway import ( "context" "fmt" - "log" "os" "os/signal" "path/filepath" @@ -15,6 +14,7 @@ import ( "jane/pkg/channels" _ "jane/pkg/channels/dingtalk" _ "jane/pkg/channels/discord" + _ "jane/pkg/channels/gmessages" _ "jane/pkg/channels/irc" _ "jane/pkg/channels/line" _ "jane/pkg/channels/maixcam" @@ -23,7 +23,6 @@ import ( _ "jane/pkg/channels/pico" _ "jane/pkg/channels/qq" _ "jane/pkg/channels/slack" - _ "jane/pkg/channels/gmessages" _ "jane/pkg/channels/telegram" _ "jane/pkg/channels/whatsapp" _ "jane/pkg/channels/whatsapp_native" @@ -248,7 +247,7 @@ func setupCronTool( var err error cronTool, err = tools.NewCronTool(cronService, agentLoop, msgBus, workspace, restrict, execTimeout, cfg) if err != nil { - log.Fatalf("Critical error during CronTool initialization: %v", err) + logger.FatalCF("gateway", "Critical error during CronTool initialization", map[string]any{"error": err.Error()}) } agentLoop.RegisterTool(cronTool) diff --git a/docs/design/ETL_TODO.md b/docs/design/ETL_TODO.md index 28ed97a9b..413676e90 100644 --- a/docs/design/ETL_TODO.md +++ b/docs/design/ETL_TODO.md @@ -4,7 +4,7 @@ This document tracks the tasks required to implement the "Ultimate Visibility" E ## 1. Extract (Ingestion & Telemetry Collection) -- [ ] **Structured Logging:** Ensure `zerolog` is used consistently across the codebase for structured JSON logging. Add context to logs where missing (session IDs, tool inputs/outputs). +- [x] **Structured Logging:** Ensure `zerolog` is used consistently across the codebase for structured JSON logging. Add context to logs where missing (session IDs, tool inputs/outputs). - [x] **Basic Metrics Implementation:** Introduce a metrics package (e.g., using `expvar` or a Prometheus client) to expose basic application metrics. - [x] **Goroutine Tracking:** Implement a metric to track the number of active Goroutines. - [x] **Memory Tracking:** Implement a metric to track heap allocation and GC pauses. diff --git a/pkg/agent/instance.go b/pkg/agent/instance.go index f706efff2..73ff7decd 100644 --- a/pkg/agent/instance.go +++ b/pkg/agent/instance.go @@ -3,13 +3,13 @@ package agent import ( "context" "fmt" - "log" "os" "path/filepath" "regexp" "strings" "jane/pkg/config" + "jane/pkg/logger" "jane/pkg/memory" "jane/pkg/providers" "jane/pkg/routing" @@ -86,7 +86,7 @@ func NewAgentInstance( if cfg.Tools.IsToolEnabled("exec") { execTool, err := tools.NewExecToolWithConfig(workspace, restrict, cfg) if err != nil { - log.Fatalf("Critical error: unable to initialize exec tool: %v", err) + logger.FatalCF("agent", "Critical error: unable to initialize exec tool", map[string]any{"error": err.Error()}) } toolsRegistry.Register(execTool) } @@ -225,8 +225,7 @@ func NewAgentInstance( }) lightCandidates = resolved } else { - log.Printf("routing: light_model %q not found in model_list — routing disabled for agent %q", - rc.LightModel, agentID) + logger.WarnCF("agent", "routing: light_model not found in model_list — routing disabled for agent", map[string]any{"lightModel": rc.LightModel, "agentID": agentID}) } } @@ -314,7 +313,7 @@ func (a *AgentInstance) Close() error { func initSessionStore(dir string) session.SessionStore { store, err := memory.NewJSONLStore(dir) if err != nil { - log.Printf("memory: init store: %v; using json sessions", err) + logger.WarnCF("agent", "memory: init store failed; using json sessions", map[string]any{"error": err.Error()}) return session.NewSessionManager(dir) } @@ -322,11 +321,11 @@ func initSessionStore(dir string) session.SessionStore { // Migration failure means the store could not write data. // Fall back to SessionManager to avoid a split state where // some sessions are in JSONL and others remain in JSON. - log.Printf("memory: migration failed: %v; falling back to json sessions", merr) + logger.WarnCF("agent", "memory: migration failed; falling back to json sessions", map[string]any{"error": merr.Error()}) store.Close() return session.NewSessionManager(dir) } else if n > 0 { - log.Printf("memory: migrated %d session(s) to jsonl", n) + logger.InfoCF("agent", "memory: migrated session(s) to jsonl", map[string]any{"count": n}) } return session.NewJSONLBackend(store) diff --git a/pkg/channels/gmessages/client.go b/pkg/channels/gmessages/client.go index f86d20a96..0d32d7fbc 100644 --- a/pkg/channels/gmessages/client.go +++ b/pkg/channels/gmessages/client.go @@ -8,9 +8,9 @@ import ( "path/filepath" "time" - "jane/pkg/logger" "github.com/mdp/qrterminal/v3" "github.com/rs/zerolog" + "jane/pkg/logger" "go.mau.fi/mautrix-gmessages/pkg/libgm" "go.mau.fi/mautrix-gmessages/pkg/libgm/events" @@ -221,7 +221,6 @@ func (c *GMessagesChannel) handlePairing(ctx context.Context, client *GMClient) qrterminal.GenerateHalfBlock(qrURL, qrterminal.L, os.Stdout) fmt.Println("Waiting for pairing...") - select { case <-pairingCh: if pairErr != nil { diff --git a/pkg/channels/gmessages/events.go b/pkg/channels/gmessages/events.go index bd577400b..cd23c4148 100644 --- a/pkg/channels/gmessages/events.go +++ b/pkg/channels/gmessages/events.go @@ -53,7 +53,7 @@ func (h *EventHandler) Handle(rawEvt any) { func (h *EventHandler) handleClientReady(evt *events.ClientReady) { logger.InfoCF("channels.gmessages", "Client ready", map[string]any{ - "session_id": evt.SessionID, + "session_id": evt.SessionID, "conversations": len(evt.Conversations), }) @@ -190,7 +190,7 @@ func (h *EventHandler) handleMessage(evt *libgm.WrappedMessage) { // determine if it's a group from DB (we do it in the next steps) peer := bus.Peer{ - ID: chatID, + ID: chatID, } senderInfo := bus.SenderInfo{ @@ -218,7 +218,7 @@ func (h *EventHandler) handleMessage(evt *libgm.WrappedMessage) { logger.DebugCF("channels.gmessages", "Stored message", map[string]any{ "msg_id": dbMsg.MessageID, - "from": senderName, + "from": senderName, "is_old": evt.IsOld, }) } @@ -298,14 +298,14 @@ func ExtractMessageBody(msg *gmproto.Message) string { } type MediaInfo struct { - MediaID string - MimeType string - MediaName string - DecryptionKey []byte - Size int64 - ThumbnailMediaID string - ThumbnailDecryptionKey []byte - InlineData []byte + MediaID string + MimeType string + MediaName string + DecryptionKey []byte + Size int64 + ThumbnailMediaID string + ThumbnailDecryptionKey []byte + InlineData []byte } func ExtractMediaInfo(msg *gmproto.Message) *MediaInfo { @@ -323,13 +323,13 @@ func ExtractMediaInfo(msg *gmproto.Message) *MediaInfo { mi := &MediaInfo{ MediaID: mc.GetMediaID(), - MimeType: mime, - MediaName: mc.GetMediaName(), - DecryptionKey: mc.GetDecryptionKey(), - Size: mc.GetSize(), - ThumbnailMediaID: mc.GetThumbnailMediaID(), + MimeType: mime, + MediaName: mc.GetMediaName(), + DecryptionKey: mc.GetDecryptionKey(), + Size: mc.GetSize(), + ThumbnailMediaID: mc.GetThumbnailMediaID(), ThumbnailDecryptionKey: mc.GetThumbnailDecryptionKey(), - InlineData: mc.GetMediaData(), + InlineData: mc.GetMediaData(), } if mi.MediaID == "" && mi.ThumbnailMediaID != "" { diff --git a/pkg/config/channels.go b/pkg/config/channels.go index fa702d382..6012e2530 100644 --- a/pkg/config/channels.go +++ b/pkg/config/channels.go @@ -1,15 +1,15 @@ package config type ChannelsConfig struct { - WhatsApp WhatsAppConfig `json:"whatsapp"` - Telegram TelegramConfig `json:"telegram"` - Discord DiscordConfig `json:"discord"` - MaixCam MaixCamConfig `json:"maixcam"` - QQ QQConfig `json:"qq"` - DingTalk DingTalkConfig `json:"dingtalk"` - Slack SlackConfig `json:"slack"` - Matrix MatrixConfig `json:"matrix"` - LINE LINEConfig `json:"line"` + WhatsApp WhatsAppConfig `json:"whatsapp"` + Telegram TelegramConfig `json:"telegram"` + Discord DiscordConfig `json:"discord"` + MaixCam MaixCamConfig `json:"maixcam"` + QQ QQConfig `json:"qq"` + DingTalk DingTalkConfig `json:"dingtalk"` + Slack SlackConfig `json:"slack"` + Matrix MatrixConfig `json:"matrix"` + LINE LINEConfig `json:"line"` OneBot OneBotConfig `json:"onebot"` Pico PicoConfig `json:"pico"` IRC IRCConfig `json:"irc"` diff --git a/pkg/config/defaults.go b/pkg/config/defaults.go index 357353dc6..1e510921c 100644 --- a/pkg/config/defaults.go +++ b/pkg/config/defaults.go @@ -43,27 +43,27 @@ func DefaultConfig() *Config { Workspace: filepath.Join(homePath, "Obsidian_Vault", "Patients"), }, { - ID: "coding", - Name: "Coding Persona", - Workspace: filepath.Join(homePath, "workspace", "code"), + ID: "coding", + Name: "Coding Persona", + Workspace: filepath.Join(homePath, "workspace", "code"), MCPServers: []string{"github", "bash"}, }, { - ID: "google", - Name: "Google Persona", - Workspace: filepath.Join(homePath, "workspace", "google"), + ID: "google", + Name: "Google Persona", + Workspace: filepath.Join(homePath, "workspace", "google"), MCPServers: []string{"gmail", "calendar"}, }, { - ID: "communications", - Name: "Communications Persona", - Workspace: filepath.Join(homePath, "workspace", "communications"), + ID: "communications", + Name: "Communications Persona", + Workspace: filepath.Join(homePath, "workspace", "communications"), MCPServers: []string{"slack", "discord"}, }, { - ID: "financial", - Name: "Financial Persona", - Workspace: filepath.Join(homePath, "workspace", "finance"), + ID: "financial", + Name: "Financial Persona", + Workspace: filepath.Join(homePath, "workspace", "finance"), MCPServers: []string{"alpaca"}, }, }, diff --git a/pkg/tools/registry_test.go b/pkg/tools/registry_test.go index bc9b61caf..e3e062f11 100644 --- a/pkg/tools/registry_test.go +++ b/pkg/tools/registry_test.go @@ -359,6 +359,6 @@ func TestToolRegistry_ConcurrentAccess(t *testing.T) { } } -func (t *mockRegistryTool) RequiresApproval() bool { return false } -func (t *mockContextAwareTool) RequiresApproval() bool { return false } +func (t *mockRegistryTool) RequiresApproval() bool { return false } +func (t *mockContextAwareTool) RequiresApproval() bool { return false } func (t *mockAsyncRegistryTool) RequiresApproval() bool { return false } diff --git a/web/backend/api/gateway.go b/web/backend/api/gateway.go index f52faf3c5..2d3890288 100644 --- a/web/backend/api/gateway.go +++ b/web/backend/api/gateway.go @@ -5,7 +5,6 @@ import ( "encoding/json" "fmt" "io" - "log" "net" "net/http" "os" @@ -18,6 +17,7 @@ import ( "time" "jane/pkg/config" + "jane/pkg/logger" "jane/web/backend/utils" ) @@ -57,20 +57,20 @@ func (h *Handler) TryAutoStartGateway() { ready, reason, err := h.gatewayStartReady() if err != nil { - log.Printf("Skip auto-starting gateway: %v", err) + logger.WarnCF("gateway", "Skip auto-starting gateway", map[string]any{"error": err.Error()}) return } if !ready { - log.Printf("Skip auto-starting gateway: %s", reason) + logger.WarnCF("gateway", "Skip auto-starting gateway", map[string]any{"reason": reason}) return } pid, err := h.startGatewayLocked() if err != nil { - log.Printf("Failed to auto-start gateway: %v", err) + logger.ErrorCF("gateway", "Failed to auto-start gateway", map[string]any{"error": err.Error()}) return } - log.Printf("Gateway auto-started (PID: %d)", pid) + logger.InfoCF("gateway", "Gateway auto-started", map[string]any{"pid": pid}) } // gatewayStartReady validates whether current config can start the gateway. @@ -162,7 +162,7 @@ func (h *Handler) startGatewayLocked() (int, error) { // Ensure Pico Channel is configured before starting gateway if _, err := h.ensurePicoChannel(); err != nil { - log.Printf("Warning: failed to ensure pico channel: %v", err) + logger.WarnCF("gateway", "failed to ensure pico channel", map[string]any{"error": err.Error()}) // Non-fatal: gateway can still start without pico channel } @@ -172,7 +172,7 @@ func (h *Handler) startGatewayLocked() (int, error) { gateway.cmd = cmd pid := cmd.Process.Pid - log.Printf("Started picoclaw gateway (PID: %d) from %s", pid, execPath) + logger.InfoCF("gateway", "Started picoclaw gateway", map[string]any{"pid": pid, "execPath": execPath}) // Broadcast starting event gateway.events.Broadcast(GatewayEvent{Status: "starting", PID: pid}) @@ -184,9 +184,9 @@ func (h *Handler) startGatewayLocked() (int, error) { // Wait for exit in background and clean up go func() { if err := cmd.Wait(); err != nil { - log.Printf("Gateway process exited: %v", err) + logger.WarnCF("gateway", "Gateway process exited with error", map[string]any{"error": err.Error()}) } else { - log.Printf("Gateway process exited normally") + logger.InfoCF("gateway", "Gateway process exited normally", nil) } gateway.mu.Lock() @@ -317,7 +317,7 @@ func (h *Handler) handleGatewayStop(w http.ResponseWriter, r *http.Request) { return } - log.Printf("Sent stop signal to gateway (PID: %d)", pid) + logger.InfoCF("gateway", "Sent stop signal to gateway", map[string]any{"pid": pid}) w.Header().Set("Content-Type", "application/json") json.NewEncoder(w).Encode(map[string]any{ diff --git a/web/backend/api/oauth.go b/web/backend/api/oauth.go index c734fad73..319b7fae4 100644 --- a/web/backend/api/oauth.go +++ b/web/backend/api/oauth.go @@ -7,13 +7,13 @@ import ( "fmt" "html" "io" - "log" "net/http" "strings" "time" "jane/pkg/auth" "jane/pkg/config" + "jane/pkg/logger" "jane/pkg/providers" ) @@ -714,7 +714,7 @@ func (h *Handler) persistCredentialAndConfig(provider, authMethod string, cred * if cp.Email == "" { email, err := oauthFetchGoogleUserEmailFunc(cp.AccessToken) if err != nil { - log.Printf("oauth warning: could not fetch google email: %v", err) + logger.WarnCF("oauth", "could not fetch google email", map[string]any{"error": err.Error()}) } else { cp.Email = email } @@ -722,7 +722,7 @@ func (h *Handler) persistCredentialAndConfig(provider, authMethod string, cred * if cp.ProjectID == "" { projectID, err := oauthFetchAntigravityProject(cp.AccessToken) if err != nil { - log.Printf("oauth warning: could not fetch antigravity project id: %v", err) + logger.WarnCF("oauth", "could not fetch antigravity project id", map[string]any{"error": err.Error()}) } else { cp.ProjectID = projectID } diff --git a/web/backend/embed.go b/web/backend/embed.go index 2b28f84b9..dd3c96b86 100644 --- a/web/backend/embed.go +++ b/web/backend/embed.go @@ -3,11 +3,12 @@ package main import ( "embed" "io/fs" - "log" "mime" "net/http" "path" "strings" + + "jane/pkg/logger" ) //go:embed all:dist @@ -19,18 +20,14 @@ func registerEmbedRoutes(mux *http.ServeMux) { // Go's built-in mime.TypeByExtension returns "image/svg" which is incorrect // The correct MIME type per RFC 6838 is "image/svg+xml" if err := mime.AddExtensionType(".svg", "image/svg+xml"); err != nil { - log.Printf("Warning: failed to register SVG MIME type: %v", err) + logger.WarnCF("embed", "failed to register SVG MIME type", map[string]any{"error": err.Error()}) } // Attempt to get the subdirectory 'dist' where Vite usually builds subFS, err := fs.Sub(frontendFS, "dist") if err != nil { // Log a warning if dist doesn't exist yet (e.g., during development before a frontend build) - log.Printf( - "Warning: no 'dist' folder found in embedded frontend. " + - "Ensure you run `pnpm build:backend` in the frontend directory " + - "before building the Go backend.", - ) + logger.WarnCF("embed", "no 'dist' folder found in embedded frontend. Ensure you run `pnpm build:backend` in the frontend directory before building the Go backend.", nil) return } diff --git a/web/backend/main.go b/web/backend/main.go index ee88e7fd0..380e37dc8 100644 --- a/web/backend/main.go +++ b/web/backend/main.go @@ -15,13 +15,13 @@ import ( "errors" "flag" "fmt" - "log" "net/http" "os" "path/filepath" "strconv" "time" + "jane/pkg/logger" "jane/web/backend/api" "jane/web/backend/launcherconfig" "jane/web/backend/middleware" @@ -59,11 +59,11 @@ func main() { absPath, err := filepath.Abs(configPath) if err != nil { - log.Fatalf("Failed to resolve config path: %v", err) + logger.FatalCF("main", "Failed to resolve config path", map[string]any{"error": err.Error()}) } err = utils.EnsureOnboarded(absPath) if err != nil { - log.Printf("Warning: Failed to initialize PicoClaw config automatically: %v", err) + logger.WarnCF("main", "Failed to initialize PicoClaw config automatically", map[string]any{"error": err.Error()}) } var explicitPort bool @@ -80,7 +80,7 @@ func main() { launcherPath := launcherconfig.PathForAppConfig(absPath) launcherCfg, err := launcherconfig.Load(launcherPath, launcherconfig.Default()) if err != nil { - log.Printf("Warning: Failed to load %s: %v", launcherPath, err) + logger.WarnCF("main", "Failed to load launcher config", map[string]any{"path": launcherPath, "error": err.Error()}) launcherCfg = launcherconfig.Default() } @@ -98,7 +98,7 @@ func main() { if err == nil { err = errors.New("must be in range 1-65535") } - log.Fatalf("Invalid port %q: %v", effectivePort, err) + logger.FatalCF("main", "Invalid port", map[string]any{"port": effectivePort, "error": err.Error()}) } // Determine listen address @@ -122,7 +122,7 @@ func main() { accessControlledMux, err := middleware.IPAllowlist(launcherCfg.AllowedCIDRs, mux) if err != nil { - log.Fatalf("Invalid allowed CIDR configuration: %v", err) + logger.FatalCF("main", "Invalid allowed CIDR configuration", map[string]any{"error": err.Error()}) } // Apply middleware stack @@ -151,7 +151,7 @@ func main() { time.Sleep(500 * time.Millisecond) url := "http://localhost:" + effectivePort if err := utils.OpenBrowser(url); err != nil { - log.Printf("Warning: Failed to auto-open browser: %v", err) + logger.WarnCF("main", "Failed to auto-open browser", map[string]any{"error": err.Error()}) } }() } @@ -173,6 +173,6 @@ func main() { } if err := server.ListenAndServe(); err != nil && err != http.ErrServerClosed { - log.Fatalf("Server failed to start: %v", err) + logger.FatalCF("main", "Server failed to start", map[string]any{"error": err.Error()}) } }