fix gateway reload will cause pico stop working issue

This commit is contained in:
Cytown 2026-03-27 16:21:09 +08:00
parent 9cbb4ab7ad
commit e3d4b2a74a
7 changed files with 165 additions and 43 deletions

View file

@ -2343,7 +2343,7 @@ func TestProcessMessage_PublishesReasoningContentToReasoningChannel(t *testing.T
if outbound.Content != "thinking trace" {
t.Fatalf("reasoning content = %q, want %q", outbound.Content, "thinking trace")
}
case <-time.After(2 * time.Second):
case <-time.After(3 * time.Second):
t.Fatal("expected reasoning content to be published to reasoning channel")
}
}

View file

@ -0,0 +1,74 @@
package channels
import (
"net/http"
"strings"
"sync"
)
// dynamicServeMux is an http.Handler that supports dynamic registration
// and unregistration of handlers without recreating the server.
type dynamicServeMux struct {
mu sync.RWMutex
handlers map[string]http.Handler
}
func newDynamicServeMux() *dynamicServeMux {
return &dynamicServeMux{
handlers: make(map[string]http.Handler),
}
}
// Handle registers the handler for the given pattern.
func (dm *dynamicServeMux) Handle(pattern string, handler http.Handler) {
dm.mu.Lock()
defer dm.mu.Unlock()
dm.handlers[pattern] = handler
}
// HandleFunc registers the handler function for the given pattern.
func (dm *dynamicServeMux) HandleFunc(pattern string, handler func(http.ResponseWriter, *http.Request)) {
dm.Handle(pattern, http.HandlerFunc(handler))
}
// Unhandle removes the handler for the given pattern.
func (dm *dynamicServeMux) Unhandle(pattern string) {
dm.mu.Lock()
defer dm.mu.Unlock()
delete(dm.handlers, pattern)
}
// ServeHTTP dispatches the request to the handler whose pattern best matches
// the request URL path. It supports both exact path matches and subtree
// (trailing-slash) prefix matches, choosing the longest prefix on collision.
func (dm *dynamicServeMux) ServeHTTP(w http.ResponseWriter, r *http.Request) {
dm.mu.RLock()
defer dm.mu.RUnlock()
path := r.URL.Path
// Exact match first.
if h, ok := dm.handlers[path]; ok {
h.ServeHTTP(w, r)
return
}
// Longest subtree prefix match (patterns ending with "/").
var bestLen int
var bestHandler http.Handler
for pattern, handler := range dm.handlers {
if strings.HasSuffix(pattern, "/") && strings.HasPrefix(path, pattern) {
if len(pattern) > bestLen {
bestLen = len(pattern)
bestHandler = handler
}
}
}
if bestHandler != nil {
bestHandler.ServeHTTP(w, r)
return
}
http.NotFound(w, r)
}

View file

@ -83,7 +83,7 @@ type Manager struct {
config *config.Config
mediaStore media.MediaStore
dispatchTask *asyncTask
mux *http.ServeMux
mux *dynamicServeMux
httpServer *http.Server
mu sync.RWMutex
placeholders sync.Map // "channel:chatID" → placeholderID (string)
@ -436,7 +436,7 @@ func (m *Manager) initChannels(channels *config.ChannelsConfig) error {
// It registers health endpoints from the health server and discovers channels
// that implement WebhookHandler and/or HealthChecker to register their handlers.
func (m *Manager) SetupHTTPServer(addr string, healthServer *health.Server) {
m.mux = http.NewServeMux()
m.mux = newDynamicServeMux()
// Register health endpoints
if healthServer != nil {
@ -444,7 +444,28 @@ func (m *Manager) SetupHTTPServer(addr string, healthServer *health.Server) {
}
// Discover and register webhook handlers and health checkers
m.registerHTTPHandlersLocked()
m.httpServer = &http.Server{
Addr: addr,
Handler: m.mux,
ReadTimeout: 30 * time.Second,
WriteTimeout: 30 * time.Second,
}
}
// registerHTTPHandlersLocked registers webhook and health-check handlers for
// all channels currently in m.channels. Caller must hold m.mu (or ensure
// exclusive access).
func (m *Manager) registerHTTPHandlersLocked() {
for name, ch := range m.channels {
m.registerChannelHTTPHandler(name, ch)
}
}
// registerChannelHTTPHandler registers the webhook/health handlers for a
// single channel onto m.mux.
func (m *Manager) registerChannelHTTPHandler(name string, ch Channel) {
if wh, ok := ch.(WebhookHandler); ok {
m.mux.Handle(wh.WebhookPath(), wh)
logger.InfoCF("channels", "Webhook handler registered", map[string]any{
@ -461,11 +482,22 @@ func (m *Manager) SetupHTTPServer(addr string, healthServer *health.Server) {
}
}
m.httpServer = &http.Server{
Addr: addr,
Handler: m.mux,
ReadTimeout: 30 * time.Second,
WriteTimeout: 30 * time.Second,
// unregisterChannelHTTPHandler removes the webhook/health handlers for a
// single channel from m.mux.
func (m *Manager) unregisterChannelHTTPHandler(name string, ch Channel) {
if wh, ok := ch.(WebhookHandler); ok {
m.mux.Unhandle(wh.WebhookPath())
logger.InfoCF("channels", "Webhook handler unregistered", map[string]any{
"channel": name,
"path": wh.WebhookPath(),
})
}
if hc, ok := ch.(HealthChecker); ok {
m.mux.Unhandle(hc.HealthPath())
logger.InfoCF("channels", "Health endpoint unregistered", map[string]any{
"channel": name,
"path": hc.HealthPath(),
})
}
}
@ -984,8 +1016,13 @@ func (m *Manager) GetEnabledChannels() []string {
func (m *Manager) Reload(ctx context.Context, cfg *config.Config) error {
m.mu.Lock()
defer m.mu.Unlock()
m.config = cfg
list := toChannelHashes(cfg)
added, removed := compareChannels(m.channelHashes, list)
deferFuncs := make([]func(), 0, len(removed)+len(added))
for _, name := range removed {
// Stop all channels
channel := m.channels[name]
@ -998,9 +1035,9 @@ func (m *Manager) Reload(ctx context.Context, cfg *config.Config) error {
"error": err.Error(),
})
}
go func() {
deferFuncs = append(deferFuncs, func() {
m.UnregisterChannel(name)
}()
})
}
dispatchCtx, cancel := context.WithCancel(ctx)
m.dispatchTask = &asyncTask{cancel: cancel}
@ -1031,13 +1068,17 @@ func (m *Manager) Reload(ctx context.Context, cfg *config.Config) error {
m.workers[name] = w
go m.runWorker(dispatchCtx, name, w)
go m.runMediaWorker(dispatchCtx, name, w)
go func() {
deferFuncs = append(deferFuncs, func() {
m.RegisterChannel(name, channel)
}()
})
}
m.config = cfg
m.channelHashes = toChannelHashes(cfg)
m.channelHashes = list
go func() {
for _, f := range deferFuncs {
f()
}
}()
return nil
}
@ -1045,11 +1086,17 @@ func (m *Manager) RegisterChannel(name string, channel Channel) {
m.mu.Lock()
defer m.mu.Unlock()
m.channels[name] = channel
if m.mux != nil {
m.registerChannelHTTPHandler(name, channel)
}
}
func (m *Manager) UnregisterChannel(name string) {
m.mu.Lock()
defer m.mu.Unlock()
if ch, ok := m.channels[name]; ok && m.mux != nil {
m.unregisterChannelHTTPHandler(name, ch)
}
if w, ok := m.workers[name]; ok && w != nil {
close(w.queue)
<-w.done

View file

@ -490,12 +490,13 @@ func restartServices(
}
al.SetMediaStore(runningServices.MediaStore)
runningServices.ChannelManager, err = channels.NewManager(cfg, msgBus, runningServices.MediaStore)
if err != nil {
return fmt.Errorf("error recreating channel manager: %w", err)
}
al.SetChannelManager(runningServices.ChannelManager)
if err = runningServices.ChannelManager.Reload(context.Background(), cfg); err != nil {
return fmt.Errorf("error reload channels: %w", err)
}
fmt.Println(" ✓ Channels restarted.")
enabledChannels := runningServices.ChannelManager.GetEnabledChannels()
if len(enabledChannels) > 0 {
fmt.Printf(" ✓ Channels enabled: %s\n", enabledChannels)
@ -503,18 +504,6 @@ func restartServices(
fmt.Println(" ⚠ Warning: No channels enabled")
}
addr := fmt.Sprintf("%s:%d", cfg.Gateway.Host, cfg.Gateway.Port)
// Reuse existing HealthServer to preserve reloadFunc
if runningServices.HealthServer == nil {
runningServices.HealthServer = health.NewServer(cfg.Gateway.Host, cfg.Gateway.Port)
}
runningServices.ChannelManager.SetupHTTPServer(addr, runningServices.HealthServer)
if err = runningServices.ChannelManager.Reload(context.Background(), cfg); err != nil {
return fmt.Errorf("error reload channels: %w", err)
}
fmt.Println(" ✓ Channels restarted.")
stateManager := state.NewManager(cfg.WorkspacePath())
runningServices.DeviceService = devices.NewService(devices.Config{
Enabled: cfg.Devices.Enabled,

View file

@ -198,9 +198,17 @@ func (s *Server) readyHandler(w http.ResponseWriter, r *http.Request) {
})
}
// HandlerMux is the interface for registering HTTP handlers, used by
// RegisterOnMux so that callers can pass any mux implementation
// (e.g. *http.ServeMux or a custom dynamic mux).
type HandlerMux interface {
Handle(pattern string, handler http.Handler)
HandleFunc(pattern string, handler func(http.ResponseWriter, *http.Request))
}
// RegisterOnMux registers /health, /ready and /reload handlers onto the given mux.
// This allows the health endpoints to be served by a shared HTTP server.
func (s *Server) RegisterOnMux(mux *http.ServeMux) {
func (s *Server) RegisterOnMux(mux HandlerMux) {
mux.HandleFunc("/health", s.healthHandler)
mux.HandleFunc("/ready", s.readyHandler)
mux.HandleFunc("/reload", s.reloadHandler)

View file

@ -353,6 +353,10 @@ func WarnCF(component string, message string, fields map[string]any) {
logMessage(WARN, component, message, fields)
}
func Warnf(message string, ss ...any) {
logMessage(WARN, "", fmt.Sprintf(message, ss...), nil)
}
func Error(message string) {
logMessage(ERROR, "", message, nil)
}

View file

@ -730,8 +730,8 @@ func (h *Handler) gatewayStatusData() map[string]any {
gateway.mu.Unlock()
logger.ErrorC("gateway", fmt.Sprintf("Gateway health check failed: %v", err))
} else {
logger.InfoC("gateway", fmt.Sprintf("Gateway health status: %d", statusCode))
if statusCode != http.StatusOK {
logger.WarnC("gateway", fmt.Sprintf("Gateway health status: %d", statusCode))
gateway.mu.Lock()
setGatewayRuntimeStatusLocked("error")
gateway.mu.Unlock()