docs(gateway): add Chinese comments to gateway.go and add web-architecture.md

This commit is contained in:
ddshuai 2026-04-03 17:19:37 +08:00
parent 84a52428b5
commit e1e3598a62
2 changed files with 182 additions and 16 deletions

View file

@ -43,28 +43,31 @@ import (
) )
const ( const (
serviceShutdownTimeout = 30 * time.Second serviceShutdownTimeout = 30 * time.Second // 服务关闭超时时间
providerReloadTimeout = 30 * time.Second providerReloadTimeout = 30 * time.Second // Provider 重载超时时间
gracefulShutdownTimeout = 15 * time.Second gracefulShutdownTimeout = 15 * time.Second // 优雅关停超时时间
logPath = "logs" logPath = "logs"
panicFile = "gateway_panic.log" panicFile = "gateway_panic.log"
logFile = "gateway.log" logFile = "gateway.log"
) )
// services 网关运行时管理的所有服务集合
type services struct { type services struct {
CronService *cron.CronService CronService *cron.CronService // 定时任务服务
HeartbeatService *heartbeat.HeartbeatService HeartbeatService *heartbeat.HeartbeatService // 心跳服务
MediaStore media.MediaStore MediaStore media.MediaStore // 媒体文件存储
ChannelManager *channels.Manager ChannelManager *channels.Manager // 渠道管理器
DeviceService *devices.Service DeviceService *devices.Service // 设备事件服务
HealthServer *health.Server HealthServer *health.Server // 健康检查 HTTP 服务
manualReloadChan chan struct{} manualReloadChan chan struct{} // 手动重载信号通道
reloading atomic.Bool reloading atomic.Bool // 重载进行中标记(原子操作,防止并发重载)
} }
// startupBlockedProvider 启动受限模式的占位 Provider
// 当没有配置默认模型时使用,所有请求直接返回错误
type startupBlockedProvider struct { type startupBlockedProvider struct {
reason string reason string // 受限原因
} }
func (p *startupBlockedProvider) Chat( func (p *startupBlockedProvider) Chat(
@ -81,8 +84,9 @@ func (p *startupBlockedProvider) GetDefaultModel() string {
return "" return ""
} }
// Run starts the gateway runtime using the configuration loaded from configPath. // Run 启动网关运行时,从 configPath 加载配置并初始化所有服务
func Run(debug bool, homePath, configPath string, allowEmptyStartup bool) error { func Run(debug bool, homePath, configPath string, allowEmptyStartup bool) error {
// 初始化 panic 日志,捕获未恢复的异常
panicPath := filepath.Join(homePath, logPath, panicFile) panicPath := filepath.Join(homePath, logPath, panicFile)
panicFunc, err := logger.InitPanic(panicPath) panicFunc, err := logger.InitPanic(panicPath)
if err != nil { if err != nil {
@ -90,16 +94,19 @@ func Run(debug bool, homePath, configPath string, allowEmptyStartup bool) error
} }
defer panicFunc() defer panicFunc()
// 启用文件日志
if err = logger.EnableFileLogging(filepath.Join(homePath, logPath, logFile)); err != nil { if err = logger.EnableFileLogging(filepath.Join(homePath, logPath, logFile)); err != nil {
panic(fmt.Sprintf("error enabling file logging: %v", err)) panic(fmt.Sprintf("error enabling file logging: %v", err))
} }
defer logger.DisableFileLogging() defer logger.DisableFileLogging()
// 加载配置文件
cfg, err := config.LoadConfig(configPath) cfg, err := config.LoadConfig(configPath)
if err != nil { if err != nil {
return fmt.Errorf("error loading config: %w", err) return fmt.Errorf("error loading config: %w", err)
} }
// 设置日志级别
logger.SetLevelFromString(cfg.Gateway.LogLevel) logger.SetLevelFromString(cfg.Gateway.LogLevel)
if debug { if debug {
@ -107,6 +114,7 @@ func Run(debug bool, homePath, configPath string, allowEmptyStartup bool) error
fmt.Println("🔍 Debug mode enabled") fmt.Println("🔍 Debug mode enabled")
} }
// 创建 LLM ProviderAI 模型提供者)
provider, modelID, err := createStartupProvider(cfg, allowEmptyStartup) provider, modelID, err := createStartupProvider(cfg, allowEmptyStartup)
if err != nil { if err != nil {
return fmt.Errorf("error creating provider: %w", err) return fmt.Errorf("error creating provider: %w", err)
@ -116,9 +124,11 @@ func Run(debug bool, homePath, configPath string, allowEmptyStartup bool) error
cfg.Agents.Defaults.ModelName = modelID cfg.Agents.Defaults.ModelName = modelID
} }
// 创建消息总线和 Agent 循环
msgBus := bus.NewMessageBus() msgBus := bus.NewMessageBus()
agentLoop := agent.NewAgentLoop(cfg, msgBus, provider) agentLoop := agent.NewAgentLoop(cfg, msgBus, provider)
// 打印 Agent 启动信息
fmt.Println("\n📦 Agent Status:") fmt.Println("\n📦 Agent Status:")
startupInfo := agentLoop.GetStartupInfo() startupInfo := agentLoop.GetStartupInfo()
toolsInfo := startupInfo["tools"].(map[string]any) toolsInfo := startupInfo["tools"].(map[string]any)
@ -133,14 +143,16 @@ func Run(debug bool, homePath, configPath string, allowEmptyStartup bool) error
"skills_available": skillsInfo["available"], "skills_available": skillsInfo["available"],
}) })
// 初始化并启动所有子服务(定时任务、心跳、渠道、设备等)
runningServices, err := setupAndStartServices(cfg, agentLoop, msgBus) runningServices, err := setupAndStartServices(cfg, agentLoop, msgBus)
if err != nil { if err != nil {
return err return err
} }
// Setup manual reload channel for /reload endpoint // 设置手动重载通道,用于 /reload HTTP 端点
manualReloadChan := make(chan struct{}, 1) manualReloadChan := make(chan struct{}, 1)
runningServices.manualReloadChan = manualReloadChan runningServices.manualReloadChan = manualReloadChan
// reloadTrigger: 重载触发函数,确保同一时间只有一个重载任务执行
reloadTrigger := func() error { reloadTrigger := func() error {
if !runningServices.reloading.CompareAndSwap(false, true) { if !runningServices.reloading.CompareAndSwap(false, true) {
return fmt.Errorf("reload already in progress") return fmt.Errorf("reload already in progress")
@ -149,7 +161,7 @@ func Run(debug bool, homePath, configPath string, allowEmptyStartup bool) error
case manualReloadChan <- struct{}{}: case manualReloadChan <- struct{}{}:
return nil return nil
default: default:
// Should not happen, but reset flag if channel is full // 通道已满(不应发生),重置标记
runningServices.reloading.Store(false) runningServices.reloading.Store(false)
return fmt.Errorf("reload already queued") return fmt.Errorf("reload already queued")
} }
@ -163,8 +175,10 @@ func Run(debug bool, homePath, configPath string, allowEmptyStartup bool) error
ctx, cancel := context.WithCancel(context.Background()) ctx, cancel := context.WithCancel(context.Background())
defer cancel() defer cancel()
// 启动 Agent 循环处理协程
go agentLoop.Run(ctx) go agentLoop.Run(ctx)
// 配置文件热重载监控
var configReloadChan <-chan *config.Config var configReloadChan <-chan *config.Config
stopWatch := func() {} stopWatch := func() {}
if cfg.Gateway.HotReload { if cfg.Gateway.HotReload {
@ -173,16 +187,20 @@ func Run(debug bool, homePath, configPath string, allowEmptyStartup bool) error
} }
defer stopWatch() defer stopWatch()
// 监听系统信号
sigChan := make(chan os.Signal, 1) sigChan := make(chan os.Signal, 1)
signal.Notify(sigChan, os.Interrupt, syscall.SIGTERM) signal.Notify(sigChan, os.Interrupt, syscall.SIGTERM)
// 主事件循环:监听系统信号、配置变更、手动重载
for { for {
select { select {
case <-sigChan: case <-sigChan:
// 收到中断信号,执行优雅关停
logger.Info("Shutting down...") logger.Info("Shutting down...")
shutdownGateway(runningServices, agentLoop, provider, true) shutdownGateway(runningServices, agentLoop, provider, true)
return nil return nil
case newCfg := <-configReloadChan: case newCfg := <-configReloadChan:
// 配置文件变更,触发热重载
if !runningServices.reloading.CompareAndSwap(false, true) { if !runningServices.reloading.CompareAndSwap(false, true) {
logger.Warn("Config reload skipped: another reload is in progress") logger.Warn("Config reload skipped: another reload is in progress")
continue continue
@ -192,6 +210,7 @@ func Run(debug bool, homePath, configPath string, allowEmptyStartup bool) error
logger.Errorf("Config reload failed: %v", err) logger.Errorf("Config reload failed: %v", err)
} }
case <-manualReloadChan: case <-manualReloadChan:
// 手动重载(通过 /reload HTTP 端点触发)
logger.Info("Manual reload triggered via /reload endpoint") logger.Info("Manual reload triggered via /reload endpoint")
newCfg, err := config.LoadConfig(configPath) newCfg, err := config.LoadConfig(configPath)
if err != nil { if err != nil {
@ -214,6 +233,7 @@ func Run(debug bool, homePath, configPath string, allowEmptyStartup bool) error
} }
} }
// executeReload 执行重载操作,确保重载完成后重置标记
func executeReload( func executeReload(
ctx context.Context, ctx context.Context,
agentLoop *agent.AgentLoop, agentLoop *agent.AgentLoop,
@ -227,6 +247,8 @@ func executeReload(
return handleConfigReload(ctx, agentLoop, newCfg, provider, runningServices, msgBus, allowEmptyStartup) return handleConfigReload(ctx, agentLoop, newCfg, provider, runningServices, msgBus, allowEmptyStartup)
} }
// createStartupProvider 根据配置创建启动时的 LLM Provider
// 当 allowEmptyStartup 为 true 且未配置模型时,返回受限模式的占位 Provider
func createStartupProvider( func createStartupProvider(
cfg *config.Config, cfg *config.Config,
allowEmptyStartup bool, allowEmptyStartup bool,
@ -244,6 +266,8 @@ func createStartupProvider(
return providers.CreateProvider(cfg) return providers.CreateProvider(cfg)
} }
// setupAndStartServices 初始化并启动所有子服务
// 包括:定时任务、心跳、媒体存储、渠道管理、健康检查、设备事件等
func setupAndStartServices( func setupAndStartServices(
cfg *config.Config, cfg *config.Config,
agentLoop *agent.AgentLoop, agentLoop *agent.AgentLoop,
@ -251,6 +275,7 @@ func setupAndStartServices(
) (*services, error) { ) (*services, error) {
runningServices := &services{} runningServices := &services{}
// 初始化定时任务服务
execTimeout := time.Duration(cfg.Tools.Cron.ExecTimeoutMinutes) * time.Minute execTimeout := time.Duration(cfg.Tools.Cron.ExecTimeoutMinutes) * time.Minute
var err error var err error
runningServices.CronService, err = setupCronTool( runningServices.CronService, err = setupCronTool(
@ -269,6 +294,7 @@ func setupAndStartServices(
} }
fmt.Println("✓ Cron service started") fmt.Println("✓ Cron service started")
// 初始化心跳服务
runningServices.HeartbeatService = heartbeat.NewHeartbeatService( runningServices.HeartbeatService = heartbeat.NewHeartbeatService(
cfg.WorkspacePath(), cfg.WorkspacePath(),
cfg.Heartbeat.Interval, cfg.Heartbeat.Interval,
@ -281,6 +307,7 @@ func setupAndStartServices(
} }
fmt.Println("✓ Heartbeat service started") fmt.Println("✓ Heartbeat service started")
// 初始化媒体文件存储(带自动清理)
runningServices.MediaStore = media.NewFileMediaStoreWithCleanup(media.MediaCleanerConfig{ runningServices.MediaStore = media.NewFileMediaStoreWithCleanup(media.MediaCleanerConfig{
Enabled: cfg.Tools.MediaCleanup.Enabled, Enabled: cfg.Tools.MediaCleanup.Enabled,
MaxAge: time.Duration(cfg.Tools.MediaCleanup.MaxAge) * time.Minute, MaxAge: time.Duration(cfg.Tools.MediaCleanup.MaxAge) * time.Minute,
@ -290,6 +317,7 @@ func setupAndStartServices(
fms.Start() fms.Start()
} }
// 初始化渠道管理器
runningServices.ChannelManager, err = channels.NewManager(cfg, msgBus, runningServices.MediaStore) runningServices.ChannelManager, err = channels.NewManager(cfg, msgBus, runningServices.MediaStore)
if err != nil { if err != nil {
if fms, ok := runningServices.MediaStore.(*media.FileMediaStore); ok { if fms, ok := runningServices.MediaStore.(*media.FileMediaStore); ok {
@ -301,6 +329,7 @@ func setupAndStartServices(
agentLoop.SetChannelManager(runningServices.ChannelManager) agentLoop.SetChannelManager(runningServices.ChannelManager)
agentLoop.SetMediaStore(runningServices.MediaStore) agentLoop.SetMediaStore(runningServices.MediaStore)
// 检测并设置语音转文字能力
if transcriber := voice.DetectTranscriber(cfg); transcriber != nil { if transcriber := voice.DetectTranscriber(cfg); transcriber != nil {
agentLoop.SetTranscriber(transcriber) agentLoop.SetTranscriber(transcriber)
logger.InfoCF("voice", "Transcription enabled (agent-level)", map[string]any{"provider": transcriber.Name()}) logger.InfoCF("voice", "Transcription enabled (agent-level)", map[string]any{"provider": transcriber.Name()})
@ -313,10 +342,12 @@ func setupAndStartServices(
fmt.Println("⚠ Warning: No channels enabled") fmt.Println("⚠ Warning: No channels enabled")
} }
// 启动 HTTP 服务器和健康检查端点
addr := fmt.Sprintf("%s:%d", cfg.Gateway.Host, cfg.Gateway.Port) addr := fmt.Sprintf("%s:%d", cfg.Gateway.Host, cfg.Gateway.Port)
runningServices.HealthServer = health.NewServer(cfg.Gateway.Host, cfg.Gateway.Port) runningServices.HealthServer = health.NewServer(cfg.Gateway.Host, cfg.Gateway.Port)
runningServices.ChannelManager.SetupHTTPServer(addr, runningServices.HealthServer) runningServices.ChannelManager.SetupHTTPServer(addr, runningServices.HealthServer)
// 启动所有已启用的渠道
if err = runningServices.ChannelManager.StartAll(context.Background()); err != nil { if err = runningServices.ChannelManager.StartAll(context.Background()); err != nil {
return nil, fmt.Errorf("error starting channels: %w", err) return nil, fmt.Errorf("error starting channels: %w", err)
} }
@ -327,6 +358,7 @@ func setupAndStartServices(
cfg.Gateway.Port, cfg.Gateway.Port,
) )
// 初始化设备事件服务
stateManager := state.NewManager(cfg.WorkspacePath()) stateManager := state.NewManager(cfg.WorkspacePath())
runningServices.DeviceService = devices.NewService(devices.Config{ runningServices.DeviceService = devices.NewService(devices.Config{
Enabled: cfg.Devices.Enabled, Enabled: cfg.Devices.Enabled,
@ -342,11 +374,13 @@ func setupAndStartServices(
return runningServices, nil return runningServices, nil
} }
// stopAndCleanupServices 按顺序停止并清理所有服务
// 重载时不会停止渠道管理器isReload=true 时跳过)
func stopAndCleanupServices(runningServices *services, shutdownTimeout time.Duration, isReload bool) { func stopAndCleanupServices(runningServices *services, shutdownTimeout time.Duration, isReload bool) {
shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), shutdownTimeout) shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), shutdownTimeout)
defer shutdownCancel() defer shutdownCancel()
// reload should not stop channel manager // 重载时不停渠道管理器
if !isReload && runningServices.ChannelManager != nil { if !isReload && runningServices.ChannelManager != nil {
runningServices.ChannelManager.StopAll(shutdownCtx) runningServices.ChannelManager.StopAll(shutdownCtx)
} }
@ -366,12 +400,14 @@ func stopAndCleanupServices(runningServices *services, shutdownTimeout time.Dura
} }
} }
// shutdownGateway 完整关闭网关:关闭 Provider → 停止所有服务 → 停止 Agent 循环
func shutdownGateway( func shutdownGateway(
runningServices *services, runningServices *services,
agentLoop *agent.AgentLoop, agentLoop *agent.AgentLoop,
provider providers.LLMProvider, provider providers.LLMProvider,
fullShutdown bool, fullShutdown bool,
) { ) {
// 如果是完整关闭且 Provider 有状态,先关闭 Provider
if cp, ok := provider.(providers.StatefulProvider); ok && fullShutdown { if cp, ok := provider.(providers.StatefulProvider); ok && fullShutdown {
cp.Close() cp.Close()
} }
@ -384,6 +420,9 @@ func shutdownGateway(
logger.Info("✓ Gateway stopped") logger.Info("✓ Gateway stopped")
} }
// handleConfigReload 处理配置文件热重载
// 流程:停止服务 → 创建新 Provider → 重载 Agent → 重启服务
// 如果任何步骤失败,会尝试用旧配置重启服务
func handleConfigReload( func handleConfigReload(
ctx context.Context, ctx context.Context,
al *agent.AgentLoop, al *agent.AgentLoop,
@ -399,12 +438,15 @@ func handleConfigReload(
logger.Infof(" New model is '%s', recreating provider...", newModel) logger.Infof(" New model is '%s', recreating provider...", newModel)
// 第一步:停止所有服务
logger.Info(" Stopping all services...") logger.Info(" Stopping all services...")
stopAndCleanupServices(runningServices, serviceShutdownTimeout, true) stopAndCleanupServices(runningServices, serviceShutdownTimeout, true)
// 第二步:用新配置创建 Provider
newProvider, newModelID, err := createStartupProvider(newCfg, allowEmptyStartup) newProvider, newModelID, err := createStartupProvider(newCfg, allowEmptyStartup)
if err != nil { if err != nil {
logger.Errorf(" ⚠ Error creating new provider: %v", err) logger.Errorf(" ⚠ Error creating new provider: %v", err)
// 创建失败,尝试用旧配置恢复服务
logger.Warn(" Attempting to restart services with old provider and config...") logger.Warn(" Attempting to restart services with old provider and config...")
if restartErr := restartServices(al, runningServices, msgBus); restartErr != nil { if restartErr := restartServices(al, runningServices, msgBus); restartErr != nil {
logger.Errorf(" ⚠ Failed to restart services: %v", restartErr) logger.Errorf(" ⚠ Failed to restart services: %v", restartErr)
@ -419,11 +461,13 @@ func handleConfigReload(
reloadCtx, reloadCancel := context.WithTimeout(context.Background(), providerReloadTimeout) reloadCtx, reloadCancel := context.WithTimeout(context.Background(), providerReloadTimeout)
defer reloadCancel() defer reloadCancel()
// 第三步:重载 Agent 循环(更新 Provider 和配置)
if err := al.ReloadProviderAndConfig(reloadCtx, newProvider, newCfg); err != nil { if err := al.ReloadProviderAndConfig(reloadCtx, newProvider, newCfg); err != nil {
logger.Errorf(" ⚠ Error reloading agent loop: %v", err) logger.Errorf(" ⚠ Error reloading agent loop: %v", err)
if cp, ok := newProvider.(providers.StatefulProvider); ok { if cp, ok := newProvider.(providers.StatefulProvider); ok {
cp.Close() cp.Close()
} }
// 重载失败,尝试用旧配置恢复服务
logger.Warn(" Attempting to restart services with old provider and config...") logger.Warn(" Attempting to restart services with old provider and config...")
if restartErr := restartServices(al, runningServices, msgBus); restartErr != nil { if restartErr := restartServices(al, runningServices, msgBus); restartErr != nil {
logger.Errorf(" ⚠ Failed to restart services: %v", restartErr) logger.Errorf(" ⚠ Failed to restart services: %v", restartErr)
@ -433,6 +477,7 @@ func handleConfigReload(
*providerRef = newProvider *providerRef = newProvider
// 第四步:用新配置重启所有服务
logger.Info(" Restarting all services with new configuration...") logger.Info(" Restarting all services with new configuration...")
if err := restartServices(al, runningServices, msgBus); err != nil { if err := restartServices(al, runningServices, msgBus); err != nil {
logger.Errorf(" ⚠ Error restarting services: %v", err) logger.Errorf(" ⚠ Error restarting services: %v", err)
@ -443,6 +488,8 @@ func handleConfigReload(
return nil return nil
} }
// restartServices 用当前配置重新创建并启动所有服务
// 用于配置热重载后的服务恢复
func restartServices( func restartServices(
al *agent.AgentLoop, al *agent.AgentLoop,
runningServices *services, runningServices *services,
@ -450,6 +497,7 @@ func restartServices(
) error { ) error {
cfg := al.GetConfig() cfg := al.GetConfig()
// 重建定时任务服务
execTimeout := time.Duration(cfg.Tools.Cron.ExecTimeoutMinutes) * time.Minute execTimeout := time.Duration(cfg.Tools.Cron.ExecTimeoutMinutes) * time.Minute
var err error var err error
runningServices.CronService, err = setupCronTool( runningServices.CronService, err = setupCronTool(
@ -468,6 +516,7 @@ func restartServices(
} }
fmt.Println(" ✓ Cron service restarted") fmt.Println(" ✓ Cron service restarted")
// 重建心跳服务
runningServices.HeartbeatService = heartbeat.NewHeartbeatService( runningServices.HeartbeatService = heartbeat.NewHeartbeatService(
cfg.WorkspacePath(), cfg.WorkspacePath(),
cfg.Heartbeat.Interval, cfg.Heartbeat.Interval,
@ -480,6 +529,7 @@ func restartServices(
} }
fmt.Println(" ✓ Heartbeat service restarted") fmt.Println(" ✓ Heartbeat service restarted")
// 重建媒体文件存储
runningServices.MediaStore = media.NewFileMediaStoreWithCleanup(media.MediaCleanerConfig{ runningServices.MediaStore = media.NewFileMediaStoreWithCleanup(media.MediaCleanerConfig{
Enabled: cfg.Tools.MediaCleanup.Enabled, Enabled: cfg.Tools.MediaCleanup.Enabled,
MaxAge: time.Duration(cfg.Tools.MediaCleanup.MaxAge) * time.Minute, MaxAge: time.Duration(cfg.Tools.MediaCleanup.MaxAge) * time.Minute,
@ -492,6 +542,7 @@ func restartServices(
al.SetChannelManager(runningServices.ChannelManager) al.SetChannelManager(runningServices.ChannelManager)
// 重载渠道管理器
if err = runningServices.ChannelManager.Reload(context.Background(), cfg); err != nil { if err = runningServices.ChannelManager.Reload(context.Background(), cfg); err != nil {
return fmt.Errorf("error reload channels: %w", err) return fmt.Errorf("error reload channels: %w", err)
} }
@ -504,6 +555,7 @@ func restartServices(
fmt.Println(" ⚠ Warning: No channels enabled") fmt.Println(" ⚠ Warning: No channels enabled")
} }
// 重建设备事件服务
stateManager := state.NewManager(cfg.WorkspacePath()) stateManager := state.NewManager(cfg.WorkspacePath())
runningServices.DeviceService = devices.NewService(devices.Config{ runningServices.DeviceService = devices.NewService(devices.Config{
Enabled: cfg.Devices.Enabled, Enabled: cfg.Devices.Enabled,
@ -516,6 +568,7 @@ func restartServices(
fmt.Println(" ✓ Device event service restarted") fmt.Println(" ✓ Device event service restarted")
} }
// 重新检测语音转文字能力
transcriber := voice.DetectTranscriber(cfg) transcriber := voice.DetectTranscriber(cfg)
al.SetTranscriber(transcriber) al.SetTranscriber(transcriber)
if transcriber != nil { if transcriber != nil {
@ -527,6 +580,8 @@ func restartServices(
return nil return nil
} }
// setupConfigWatcherPolling 创建配置文件轮询监控器
// 每 2 秒检查配置文件的修改时间和大小,变更时触发重载
func setupConfigWatcherPolling(configPath string, debug bool) (chan *config.Config, func()) { func setupConfigWatcherPolling(configPath string, debug bool) (chan *config.Config, func()) {
configChan := make(chan *config.Config, 1) configChan := make(chan *config.Config, 1)
stop := make(chan struct{}) stop := make(chan struct{})
@ -536,6 +591,7 @@ func setupConfigWatcherPolling(configPath string, debug bool) (chan *config.Conf
go func() { go func() {
defer wg.Done() defer wg.Done()
// 记录文件最后修改时间和大小
lastModTime := getFileModTime(configPath) lastModTime := getFileModTime(configPath)
lastSize := getFileSize(configPath) lastSize := getFileSize(configPath)
@ -548,16 +604,19 @@ func setupConfigWatcherPolling(configPath string, debug bool) (chan *config.Conf
currentModTime := getFileModTime(configPath) currentModTime := getFileModTime(configPath)
currentSize := getFileSize(configPath) currentSize := getFileSize(configPath)
// 检测到文件变更
if currentModTime.After(lastModTime) || currentSize != lastSize { if currentModTime.After(lastModTime) || currentSize != lastSize {
if debug { if debug {
logger.Debugf("🔍 Config file change detected") logger.Debugf("🔍 Config file change detected")
} }
// 等待 500ms 避免文件写入中途读取
time.Sleep(500 * time.Millisecond) time.Sleep(500 * time.Millisecond)
lastModTime = currentModTime lastModTime = currentModTime
lastSize = currentSize lastSize = currentSize
// 加载并验证新配置
newCfg, err := config.LoadConfig(configPath) newCfg, err := config.LoadConfig(configPath)
if err != nil { if err != nil {
logger.Errorf("⚠ Error loading new config: %v", err) logger.Errorf("⚠ Error loading new config: %v", err)
@ -573,6 +632,7 @@ func setupConfigWatcherPolling(configPath string, debug bool) (chan *config.Conf
logger.Info("✓ Config file validated and loaded") logger.Info("✓ Config file validated and loaded")
// 非阻塞发送,如果上一次重载还未处理则跳过
select { select {
case configChan <- newCfg: case configChan <- newCfg:
default: default:
@ -593,6 +653,7 @@ func setupConfigWatcherPolling(configPath string, debug bool) (chan *config.Conf
return configChan, stopFunc return configChan, stopFunc
} }
// getFileModTime 获取文件的最后修改时间
func getFileModTime(path string) time.Time { func getFileModTime(path string) time.Time {
info, err := os.Stat(path) info, err := os.Stat(path)
if err != nil { if err != nil {
@ -601,6 +662,7 @@ func getFileModTime(path string) time.Time {
return info.ModTime() return info.ModTime()
} }
// getFileSize 获取文件大小
func getFileSize(path string) int64 { func getFileSize(path string) int64 {
info, err := os.Stat(path) info, err := os.Stat(path)
if err != nil { if err != nil {
@ -609,6 +671,8 @@ func getFileSize(path string) int64 {
return info.Size() return info.Size()
} }
// setupCronTool 初始化定时任务服务和工具
// 创建 CronService 实例,如果 cron 工具启用则注册到 Agent 循环
func setupCronTool( func setupCronTool(
agentLoop *agent.AgentLoop, agentLoop *agent.AgentLoop,
msgBus *bus.MessageBus, msgBus *bus.MessageBus,
@ -632,6 +696,7 @@ func setupCronTool(
agentLoop.RegisterTool(cronTool) agentLoop.RegisterTool(cronTool)
} }
// 设置定时任务执行回调
if cronTool != nil { if cronTool != nil {
cronService.SetOnJob(func(job *cron.CronJob) (string, error) { cronService.SetOnJob(func(job *cron.CronJob) (string, error) {
result := cronTool.ExecuteJob(context.Background(), job) result := cronTool.ExecuteJob(context.Background(), job)
@ -642,6 +707,8 @@ func setupCronTool(
return cronService, nil return cronService, nil
} }
// createHeartbeatHandler 创建心跳处理函数
// 当渠道和聊天 ID 为空时使用默认值cli/direct
func createHeartbeatHandler(agentLoop *agent.AgentLoop) func(prompt, channel, chatID string) *tools.ToolResult { func createHeartbeatHandler(agentLoop *agent.AgentLoop) func(prompt, channel, chatID string) *tools.ToolResult {
return func(prompt, channel, chatID string) *tools.ToolResult { return func(prompt, channel, chatID string) *tools.ToolResult {
if channel == "" || chatID == "" { if channel == "" || chatID == "" {

99
web-architecture.md Normal file
View file

@ -0,0 +1,99 @@
# Web 模块架构与打包流程
## 项目结构
```
web/
├── frontend/ → React SPA (Vite 构建)
├── backend/ → Go Web 服务器 (嵌入前端产物)
├── Makefile → 统一构建入口
├── build/ → 最终输出目录
├── picoclaw-launcher.desktop
└── picoclaw-launcher.png
```
## 打包流程
`make build` 分两步执行:
### 1. 构建前端 (`pnpm build:backend`)
```bash
tsc -b && vite build --outDir ../backend/dist --emptyOutDir
```
- TypeScript 编译检查
- Vite 打包,产物直接输出到 `backend/dist/`
前端技术栈React 19 + TanStack Router + Tailwind CSS v4 + Radix UI (shadcn)
### 2. 编译 Go 后端
```bash
CGO_ENABLED=0 go build -v -tags stdjson \
-ldflags "-X ...Version=... -X ...GitCommit=... -s -w" \
-o build/picoclaw-launcher ./backend/
```
- `backend/embed.go` 通过 `//go:embed all:dist` 将前端产物嵌入二进制
- `-s -w` 去掉调试信息,减小体积
- `-ldflags -X` 注入版本号、Git commit、构建时间
- macOS 上 `CGO_ENABLED=1`systray 依赖Linux/Windows 上 `CGO_ENABLED=0`
### 最终产物
单个 `build/picoclaw-launcher` 可执行文件,内嵌完整前端 SPA。
## WebSocket 链路
前端通过 WebSocket 与 picoclaw gateway 通信web 后端作为反向代理中转:
```
前端 WebSocket
→ web 后端反向代理 (GET /pico/ws)
→ picoclaw gateway HTTP server
→ PicoChannel.ServeHTTP (pkg/channels/pico/pico.go:225)
→ handleWebSocket (pico.go:318) — 升级连接
→ authenticate (pico.go:383) — 验证 token
→ readLoop (pico.go:425) — 消息读取循环
→ pingLoop (pico.go:486) — 心跳保活
```
### 关键代码位置
| 组件 | 文件 |
|------|------|
| WebSocket 反向代理 | `web/backend/api/pico.go` |
| Gateway 注册 pico channel | `pkg/gateway/gateway.go` (import) |
| Pico Channel 实现 | `pkg/channels/pico/pico.go` |
| 前端代理配置 | `web/frontend/vite.config.ts` |
### API 端点
| 路由 | 用途 |
|------|------|
| `GET /pico/ws` | WebSocket 代理(转发到 gateway |
| `POST /api/pico/token` | 生成/刷新 WebSocket 认证 token |
| `POST /api/pico/setup` | 初始化 Pico Channel |
## 开发模式
`make dev` 分别启动前后端Vite dev server 通过代理转发请求:
```typescript
// vite.config.ts
proxy: {
"/api": { target: "http://localhost:18800" },
"/ws": { target: "ws://localhost:18800", ws: true },
}
```
## 其他命令
```bash
make dev # 启动前后端开发服务器
make build # 构建生产包
make test # 运行测试
make lint # 代码检查
make clean # 清理构建产物
```