This commit replaces usages of the standard library `log` package (specifically `log.Printf` and `log.Fatalf`) with the custom `jane/pkg/logger` structured logger across multiple backend files (`web/backend/main.go`, `web/backend/api/gateway.go`, `web/backend/api/oauth.go`, `web/backend/embed.go`, `cmd/picoclaw/internal/gateway/helpers.go`, and `pkg/agent/instance.go`). This fulfills the pending "Structured Logging" task outlined in `docs/design/ETL_TODO.md` to ensure a unified JSON structured logging pattern throughout the codebase. Co-authored-by: hobbyistlabs-coder <267281733+hobbyistlabs-coder@users.noreply.github.com>
265 lines
7.5 KiB
Go
265 lines
7.5 KiB
Go
package gateway
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"os"
|
|
"os/signal"
|
|
"path/filepath"
|
|
"time"
|
|
|
|
"jane/cmd/picoclaw/internal"
|
|
"jane/pkg/agent"
|
|
"jane/pkg/bus"
|
|
"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"
|
|
_ "jane/pkg/channels/matrix"
|
|
_ "jane/pkg/channels/onebot"
|
|
_ "jane/pkg/channels/pico"
|
|
_ "jane/pkg/channels/qq"
|
|
_ "jane/pkg/channels/slack"
|
|
_ "jane/pkg/channels/telegram"
|
|
_ "jane/pkg/channels/whatsapp"
|
|
_ "jane/pkg/channels/whatsapp_native"
|
|
"jane/pkg/config"
|
|
"jane/pkg/cron"
|
|
"jane/pkg/devices"
|
|
"jane/pkg/health"
|
|
"jane/pkg/heartbeat"
|
|
"jane/pkg/logger"
|
|
"jane/pkg/media"
|
|
"jane/pkg/providers"
|
|
"jane/pkg/state"
|
|
"jane/pkg/tools"
|
|
"jane/pkg/voice"
|
|
)
|
|
|
|
func gatewayCmd(debug bool) error {
|
|
if debug {
|
|
logger.SetLevel(logger.DEBUG)
|
|
fmt.Println("🔍 Debug mode enabled")
|
|
}
|
|
|
|
cfg, err := internal.LoadConfig()
|
|
if err != nil {
|
|
return fmt.Errorf("error loading config: %w", err)
|
|
}
|
|
|
|
if cfg.Logger.TimeFormat != "" {
|
|
logger.SetTimeFormat(cfg.Logger.TimeFormat)
|
|
}
|
|
|
|
provider, modelID, err := providers.CreateProvider(cfg)
|
|
if err != nil {
|
|
return fmt.Errorf("error creating provider: %w", err)
|
|
}
|
|
|
|
// Use the resolved model ID from provider creation
|
|
if modelID != "" {
|
|
cfg.Agents.Defaults.ModelName = modelID
|
|
}
|
|
|
|
msgBus := bus.NewMessageBus()
|
|
agentLoop := agent.NewAgentLoop(cfg, msgBus, provider)
|
|
|
|
// Print agent startup info
|
|
fmt.Println("\n📦 Agent Status:")
|
|
startupInfo := agentLoop.GetStartupInfo()
|
|
toolsInfo := startupInfo["tools"].(map[string]any)
|
|
skillsInfo := startupInfo["skills"].(map[string]any)
|
|
fmt.Printf(" • Tools: %d loaded\n", toolsInfo["count"])
|
|
fmt.Printf(" • Skills: %d/%d available\n",
|
|
skillsInfo["available"],
|
|
skillsInfo["total"])
|
|
|
|
// Log to file as well
|
|
logger.InfoCF("agent", "Agent initialized",
|
|
map[string]any{
|
|
"tools_count": toolsInfo["count"],
|
|
"skills_total": skillsInfo["total"],
|
|
"skills_available": skillsInfo["available"],
|
|
})
|
|
|
|
// Setup cron tool and service
|
|
execTimeout := time.Duration(cfg.Tools.Cron.ExecTimeoutMinutes) * time.Minute
|
|
cronService := setupCronTool(
|
|
agentLoop,
|
|
msgBus,
|
|
cfg.WorkspacePath(),
|
|
cfg.Agents.Defaults.RestrictToWorkspace,
|
|
execTimeout,
|
|
cfg,
|
|
)
|
|
|
|
heartbeatService := heartbeat.NewHeartbeatService(
|
|
cfg.WorkspacePath(),
|
|
cfg.Heartbeat.Interval,
|
|
cfg.Heartbeat.Enabled,
|
|
)
|
|
heartbeatService.SetBus(msgBus)
|
|
heartbeatService.SetHandler(func(prompt, channel, chatID string) *tools.ToolResult {
|
|
// Use cli:direct as fallback if no valid channel
|
|
if channel == "" || chatID == "" {
|
|
channel, chatID = "cli", "direct"
|
|
}
|
|
// Use ProcessHeartbeat - no session history, each heartbeat is independent
|
|
var response string
|
|
response, err = agentLoop.ProcessHeartbeat(context.Background(), prompt, channel, chatID)
|
|
if err != nil {
|
|
return tools.ErrorResult(fmt.Sprintf("Heartbeat error: %v", err))
|
|
}
|
|
if response == "HEARTBEAT_OK" {
|
|
return tools.SilentResult("Heartbeat OK")
|
|
}
|
|
// For heartbeat, always return silent - the subagent result will be
|
|
// sent to user via processSystemMessage when the async task completes
|
|
return tools.SilentResult(response)
|
|
})
|
|
|
|
// Create media store for file lifecycle management with TTL cleanup
|
|
mediaStore := media.NewFileMediaStoreWithCleanup(media.MediaCleanerConfig{
|
|
Enabled: cfg.Tools.MediaCleanup.Enabled,
|
|
MaxAge: time.Duration(cfg.Tools.MediaCleanup.MaxAge) * time.Minute,
|
|
Interval: time.Duration(cfg.Tools.MediaCleanup.Interval) * time.Minute,
|
|
})
|
|
mediaStore.Start()
|
|
|
|
channelManager, err := channels.NewManager(cfg, msgBus, mediaStore)
|
|
if err != nil {
|
|
mediaStore.Stop()
|
|
return fmt.Errorf("error creating channel manager: %w", err)
|
|
}
|
|
|
|
// Inject channel manager and media store into agent loop
|
|
agentLoop.SetChannelManager(channelManager)
|
|
agentLoop.SetMediaStore(mediaStore)
|
|
|
|
// Wire up voice transcription if a supported provider is configured.
|
|
if transcriber := voice.DetectTranscriber(cfg); transcriber != nil {
|
|
agentLoop.SetTranscriber(transcriber)
|
|
logger.InfoCF("voice", "Transcription enabled (agent-level)", map[string]any{"provider": transcriber.Name()})
|
|
}
|
|
|
|
enabledChannels := channelManager.GetEnabledChannels()
|
|
if len(enabledChannels) > 0 {
|
|
fmt.Printf("✓ Channels enabled: %s\n", enabledChannels)
|
|
} else {
|
|
fmt.Println("⚠ Warning: No channels enabled")
|
|
}
|
|
|
|
fmt.Printf("✓ Gateway started on %s:%d\n", cfg.Gateway.Host, cfg.Gateway.Port)
|
|
fmt.Println("Press Ctrl+C to stop")
|
|
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
defer cancel()
|
|
|
|
if err := cronService.Start(); err != nil {
|
|
fmt.Printf("Error starting cron service: %v\n", err)
|
|
}
|
|
fmt.Println("✓ Cron service started")
|
|
|
|
if err := heartbeatService.Start(); err != nil {
|
|
fmt.Printf("Error starting heartbeat service: %v\n", err)
|
|
}
|
|
fmt.Println("✓ Heartbeat service started")
|
|
|
|
stateManager := state.NewManager(cfg.WorkspacePath())
|
|
deviceService := devices.NewService(devices.Config{
|
|
Enabled: cfg.Devices.Enabled,
|
|
MonitorUSB: cfg.Devices.MonitorUSB,
|
|
}, stateManager)
|
|
deviceService.SetBus(msgBus)
|
|
if err := deviceService.Start(ctx); err != nil {
|
|
fmt.Printf("Error starting device service: %v\n", err)
|
|
} else if cfg.Devices.Enabled {
|
|
fmt.Println("✓ Device event service started")
|
|
}
|
|
|
|
// Setup shared HTTP server with health endpoints and webhook handlers
|
|
healthServer := health.NewServer(cfg.Gateway.Host, cfg.Gateway.Port)
|
|
addr := fmt.Sprintf("%s:%d", cfg.Gateway.Host, cfg.Gateway.Port)
|
|
channelManager.SetupHTTPServer(addr, healthServer)
|
|
|
|
// Start resource tracker for ETL Visibility
|
|
resourceTracker := health.NewResourceTracker(1 * time.Minute)
|
|
resourceTracker.Start(ctx)
|
|
fmt.Println("✓ Resource tracker started")
|
|
|
|
if err := channelManager.StartAll(ctx); err != nil {
|
|
fmt.Printf("Error starting channels: %v\n", err)
|
|
return err
|
|
}
|
|
|
|
fmt.Printf("✓ Health endpoints available at http://%s:%d/health and /ready\n", cfg.Gateway.Host, cfg.Gateway.Port)
|
|
|
|
go agentLoop.Run(ctx)
|
|
|
|
sigChan := make(chan os.Signal, 1)
|
|
signal.Notify(sigChan, os.Interrupt)
|
|
<-sigChan
|
|
|
|
fmt.Println("\nShutting down...")
|
|
if cp, ok := provider.(providers.StatefulProvider); ok {
|
|
cp.Close()
|
|
}
|
|
cancel()
|
|
msgBus.Close()
|
|
|
|
// Use a fresh context with timeout for graceful shutdown,
|
|
// since the original ctx is already canceled.
|
|
shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), 15*time.Second)
|
|
defer shutdownCancel()
|
|
|
|
channelManager.StopAll(shutdownCtx)
|
|
deviceService.Stop()
|
|
heartbeatService.Stop()
|
|
cronService.Stop()
|
|
resourceTracker.Stop()
|
|
mediaStore.Stop()
|
|
agentLoop.Stop()
|
|
agentLoop.Close()
|
|
fmt.Println("✓ Gateway stopped")
|
|
|
|
return nil
|
|
}
|
|
|
|
func setupCronTool(
|
|
agentLoop *agent.AgentLoop,
|
|
msgBus *bus.MessageBus,
|
|
workspace string,
|
|
restrict bool,
|
|
execTimeout time.Duration,
|
|
cfg *config.Config,
|
|
) *cron.CronService {
|
|
cronStorePath := filepath.Join(workspace, "cron", "jobs.json")
|
|
|
|
// Create cron service
|
|
cronService := cron.NewCronService(cronStorePath, nil)
|
|
|
|
// Create and register CronTool if enabled
|
|
var cronTool *tools.CronTool
|
|
if cfg.Tools.IsToolEnabled("cron") {
|
|
var err error
|
|
cronTool, err = tools.NewCronTool(cronService, agentLoop, msgBus, workspace, restrict, execTimeout, cfg)
|
|
if err != nil {
|
|
logger.FatalCF("gateway", "Critical error during CronTool initialization", map[string]any{"error": err.Error()})
|
|
}
|
|
|
|
agentLoop.RegisterTool(cronTool)
|
|
}
|
|
|
|
// Set onJob handler
|
|
if cronTool != nil {
|
|
cronService.SetOnJob(func(job *cron.CronJob) (string, error) {
|
|
result := cronTool.ExecuteJob(context.Background(), job)
|
|
return result, nil
|
|
})
|
|
}
|
|
|
|
return cronService
|
|
}
|