Files
goclaw/cmd/gateway_lifecycle.go
Conner MoandConner Mo c21499a7f5 feat(security): let operators un-block CIDRs behind a transparent proxy (#1465)
SSRF protection resolves a hostname and judges the resulting IP. That
model assumes DNS resolution describes where the traffic actually goes,
which stops being true behind a TUN/fake-IP proxy: every query is answered
with a synthetic address out of a reserved range, and the proxy then
routes that address to the real public host. The IP is a handle, not a
destination.

In that environment web_fetch rejects ordinary public sites — observed
with news.sina.cn resolving to 198.18.0.236 — and no configuration can
fix it, because the block list is compiled in. The agent then burns
iterations retrying URLs that can never succeed.

GOCLAW_SSRF_ALLOWED_CIDRS lets an operator name the ranges their proxy
hands out. Empty by default, so nothing changes for deployments that do
not set it, and the accepted and refused entries are both logged at
startup — this widens what LLM- and admin-supplied URLs can reach, so it
should be visible.

Ranges an SSRF actually targets can never be allowlisted: link-local
(including cloud metadata at 169.254.169.254), multicast and unspecified
are refused at parse time, in either direction, so neither an exact entry
nor a wider range that swallows one gets through.

Applied inside isBlocked rather than only in validate() because
NewSafeClient re-checks the pinned IP at dial time through the same
function — relaxing just the pre-flight check would pass validation and
then fail to connect.

internal/tools carries its own private-range list for web_fetch and
web_search, separate from this package and not identical to it. It has to
consult the same setting, or relaxing one gate leaves the other rejecting
the very traffic the operator permitted. Unifying the two lists is left
alone here; it is a wider change than this one.

Co-authored-by: Conner Mo <connermo@ConnerdeMacBook-Pro.local>
2026-07-31 21:30:01 +07:00

421 lines
15 KiB
Go

package cmd
import (
"context"
"fmt"
"log/slog"
"os"
"strings"
"time"
"github.com/nextlevelbuilder/goclaw/internal/bus"
"github.com/nextlevelbuilder/goclaw/internal/cache"
"github.com/nextlevelbuilder/goclaw/internal/channels"
"github.com/nextlevelbuilder/goclaw/internal/channels/bitrix24"
"github.com/nextlevelbuilder/goclaw/internal/config"
"github.com/nextlevelbuilder/goclaw/internal/edition"
"github.com/nextlevelbuilder/goclaw/internal/heartbeat"
"github.com/nextlevelbuilder/goclaw/internal/orchestration"
"github.com/nextlevelbuilder/goclaw/internal/sandbox"
"github.com/nextlevelbuilder/goclaw/internal/scheduler"
"github.com/nextlevelbuilder/goclaw/internal/security"
"github.com/nextlevelbuilder/goclaw/internal/store"
"github.com/nextlevelbuilder/goclaw/internal/tasks"
"github.com/nextlevelbuilder/goclaw/internal/tools"
"github.com/nextlevelbuilder/goclaw/internal/webhooks"
"github.com/nextlevelbuilder/goclaw/pkg/protocol"
)
// lifecycleDeps bundles the extra parameters needed by runLifecycle that are not in gatewayDeps.
type lifecycleDeps struct {
sched *scheduler.Scheduler
heartbeatTicker *heartbeat.Ticker
quotaChecker *channels.QuotaChecker
webFetchTool *tools.WebFetchTool
ttsTool *tools.TtsTool
sandboxMgr sandbox.Manager
postTurn tools.PostTurnProcessor
subagentMgr *tools.SubagentManager
childRunAdmission *orchestration.ChildRunAdmission
consumerTeamStore store.TeamStore
auditCh chan bus.AuditEventPayload
sigCh chan os.Signal
terminateProcess func(int)
}
func drainChildRunsWithRetry(
admission *orchestration.ChildRunAdmission,
firstTimeout time.Duration,
retryTimeout time.Duration,
) error {
if admission == nil {
return nil
}
for attempt, timeout := range []time.Duration{firstTimeout, retryTimeout} {
drainCtx, drainCancel := context.WithTimeout(context.Background(), timeout)
err := admission.Close(drainCtx)
drainCancel()
if err == nil {
return nil
}
slog.Error("gateway: child-run drain attempt failed",
"attempt", attempt+1, "timeout", timeout, "error", err)
}
return fmt.Errorf("%w after retry", orchestration.ErrChildRunDrainTimeout)
}
func drainSubagentManagerWithRetry(
manager *tools.SubagentManager,
firstTimeout time.Duration,
retryTimeout time.Duration,
) error {
if manager == nil {
return nil
}
for attempt, timeout := range []time.Duration{firstTimeout, retryTimeout} {
drainCtx, drainCancel := context.WithTimeout(context.Background(), timeout)
err := manager.CloseContext(drainCtx)
drainCancel()
if err == nil {
return nil
}
slog.Error("gateway: subagent lifecycle drain attempt failed",
"attempt", attempt+1, "timeout", timeout, "error", err)
}
return fmt.Errorf("%w after retry", tools.ErrSubagentLifecycleDrainTimeout)
}
func drainDelegateToolWithRetry(
tool interface{ CloseContext(context.Context) error },
firstTimeout time.Duration,
retryTimeout time.Duration,
) error {
for attempt, timeout := range []time.Duration{firstTimeout, retryTimeout} {
drainCtx, drainCancel := context.WithTimeout(context.Background(), timeout)
err := tool.CloseContext(drainCtx)
drainCancel()
if err == nil {
return nil
}
slog.Error("gateway: delegate completion drain attempt failed",
"attempt", attempt+1, "timeout", timeout, "error", err)
}
return fmt.Errorf("delegate completion drain failed after retry")
}
// runLifecycle wires config-reload subscribers, starts consumers, task recovery,
// the signal handler goroutine, and finally starts the gateway server.
// This is the last phase of runGateway() — called after all setup is complete.
func (d *gatewayDeps) runLifecycle(
ctx context.Context,
cancel context.CancelFunc,
deps lifecycleDeps,
) {
// Reload quota config on config changes via pub/sub.
if deps.quotaChecker != nil {
d.msgBus.Subscribe("quota-config-reload", func(evt bus.Event) {
if evt.Name != bus.TopicConfigChanged {
return
}
updatedCfg, ok := evt.Payload.(*config.Config)
if !ok || updatedCfg.Gateway.Quota == nil {
return
}
config.MergeChannelGroupQuotas(updatedCfg)
deps.quotaChecker.UpdateConfig(*updatedCfg.Gateway.Quota)
slog.Info("quota config reloaded via pub/sub")
})
}
// Reload cron default timezone on config changes via pub/sub.
d.msgBus.Subscribe("cron-config-reload", func(evt bus.Event) {
if evt.Name != bus.TopicConfigChanged {
return
}
updatedCfg, ok := evt.Payload.(*config.Config)
if !ok {
return
}
d.pgStores.Cron.SetDefaultTimezone(updatedCfg.Cron.DefaultTimezone)
})
// Reload web_fetch domain policy on config changes via pub/sub.
d.msgBus.Subscribe("webfetch-config-reload", func(evt bus.Event) {
if evt.Name != bus.TopicConfigChanged {
return
}
updatedCfg, ok := evt.Payload.(*config.Config)
if !ok {
return
}
deps.webFetchTool.UpdatePolicy(updatedCfg.Tools.WebFetch.Policy, updatedCfg.Tools.WebFetch.AllowedDomains, updatedCfg.Tools.WebFetch.BlockedDomains)
})
// Reload global shell deny-group toggles on config changes via pub/sub
// so /config edits apply without a process restart.
subscribeShellDenyGroupsReload(d.msgBus, d.toolsReg)
var providerStore store.ProviderStore
var mcpStore store.MCPServerStore
if d.pgStores != nil {
providerStore = d.pgStores.Providers
mcpStore = d.pgStores.MCP
}
subscribeProviderShellDenyGroupsReload(d.msgBus, d.providerRegistry, providerStore, mcpStore)
// Reload TTS providers on config changes via pub/sub.
d.msgBus.Subscribe("tts-config-reload", func(evt bus.Event) {
if evt.Name != bus.TopicConfigChanged {
return
}
updatedCfg, ok := evt.Payload.(*config.Config)
if !ok {
return
}
if d.pgStores.ConfigSecrets != nil {
// Use master tenant context to load global TTS secrets
masterCtx := store.WithTenantID(context.Background(), store.MasterTenantID)
if secrets, err := d.pgStores.ConfigSecrets.GetAll(masterCtx); err == nil && len(secrets) > 0 {
updatedCfg.ApplyDBSecrets(secrets)
}
}
newMgr := setupTTS(updatedCfg)
if newMgr == nil {
return
}
deps.ttsTool.UpdateManager(newMgr)
if d.ttsHandler != nil {
d.ttsHandler.UpdateManager(newMgr)
}
slog.Info("tts config reloaded", "provider", newMgr.PrimaryProvider(), "auto", string(newMgr.AutoMode()))
})
// Note: vault enrichment provider is resolved per-tenant at runtime,
// no hot-reload handler needed here
// Log orphaned providers on agent deletion. Auto-delete is unsafe because
// providers can be referenced by heartbeats (FK), OAuth tokens, media chains.
d.msgBus.Subscribe("agent-deleted-provider-log", func(evt bus.Event) {
if evt.Name != bus.TopicAgentDeleted {
return
}
payload, ok := evt.Payload.(bus.AgentDeletedPayload)
if !ok || payload.Provider == "" {
return
}
slog.Info("agent deleted, provider may be orphaned — verify via UI",
"agent", payload.AgentKey, "provider", payload.Provider)
})
// Contact collector: auto-collect user info from channels with in-memory dedup cache.
var contactCollector *store.ContactCollector
if d.pgStores.Contacts != nil {
contactCollector = store.NewContactCollector(d.pgStores.Contacts, cache.NewInMemoryCache[bool]())
d.channelMgr.SetContactCollector(contactCollector)
}
go consumeInboundMessages(ctx, d.msgBus, d.agentRouter, d.cfg, deps.sched, d.channelMgr, deps.consumerTeamStore, d.pgStores.AgentLinks, deps.quotaChecker, d.pgStores.Sessions, d.pgStores.Agents, contactCollector, deps.postTurn, deps.subagentMgr, d.usageCapSvc, d.providerRegistry, d.teamWorkEmbedder)
// Webhook callback worker — delivers async webhook_calls rows to receiver callback_url.
// Runs in both editions: Standard (PG, concurrency=4) and Lite (SQLite, concurrency=1).
// sqliteonly: single callback worker — SQLite lacks SKIP LOCKED; BEGIN IMMEDIATE serializes.
var webhookWorkerCancel context.CancelFunc
if d.pgStores != nil &&
d.pgStores.WebhookCalls != nil &&
d.pgStores.Webhooks != nil &&
d.pgStores.Tenants != nil &&
d.agentRouter != nil {
workerConcurrency := 4
if edition.Current().IsLimited() {
// sqliteonly: single callback worker — SQLite lacks SKIP LOCKED; BEGIN IMMEDIATE serializes.
workerConcurrency = 1
}
ww := webhooks.NewWebhookWorker(
d.pgStores.WebhookCalls,
d.pgStores.Webhooks,
d.pgStores.Tenants,
d.agentRouter,
nil, // limiter: created internally with default per-tenant cap (4)
webhooks.WorkerConfig{
WorkerConcurrency: workerConcurrency,
PerTenantConcurrency: 4,
AsyncAgentTimeout: webhooks.ResolveTimeoutSec(d.cfg.Gateway.WebhookAsyncTimeoutSec),
Stream: webhooks.ResolveStream(d.cfg.Gateway.WebhookStream),
},
)
// K6: decrypt raw secret for outbound HMAC signing using the same key as inbound verify.
ww.SetEncKey(os.Getenv("GOCLAW_ENCRYPTION_KEY"))
var workerCtx context.Context
workerCtx, webhookWorkerCancel = context.WithCancel(ctx)
go ww.Run(workerCtx)
}
// Task recovery ticker: re-dispatches stale/pending team tasks on startup and periodically.
var taskTicker *tasks.TaskTicker
if d.pgStores.Teams != nil {
taskTicker = tasks.NewTaskTicker(d.pgStores.Teams, d.pgStores.Agents, d.msgBus, d.cfg.Gateway.TaskRecoveryIntervalSec)
taskTicker.Start()
}
go func() {
sig := <-deps.sigCh
slog.Info("graceful shutdown initiated", "signal", sig)
// Broadcast shutdown event
d.server.BroadcastEvent(*protocol.NewEvent(protocol.EventShutdown, nil))
// Close child-run intake first. A drain timeout must terminate without
// unwinding runGateway defers under a still-live child callback.
if deps.childRunAdmission != nil {
if err := drainChildRunsWithRetry(deps.childRunAdmission, 30*time.Second, 5*time.Second); err != nil {
slog.Error("gateway: terminating after child-run drain failure", "error", err)
terminate := deps.terminateProcess
if terminate == nil {
terminate = os.Exit
}
terminate(1)
return
}
}
// Stop channels, cron, heartbeat, and task ticker
d.channelMgr.StopAll(context.Background())
d.pgStores.Cron.Stop()
deps.heartbeatTicker.Stop()
if taskTicker != nil {
taskTicker.Stop()
}
// Stop webhook callback worker — signals Run() to drain in-flight and exit.
if webhookWorkerCancel != nil {
webhookWorkerCancel()
}
// Drain audit log queue before closing DB
if deps.auditCh != nil {
close(deps.auditCh)
}
if delegate, ok := d.toolsReg.Get("delegate"); ok {
if closer, ok := delegate.(interface {
CloseContext(context.Context) error
}); ok {
if err := drainDelegateToolWithRetry(closer, 65*time.Second, 10*time.Second); err != nil {
slog.Error("gateway: terminating after delegate completion drain failure", "error", err)
terminate := deps.terminateProcess
if terminate == nil {
terminate = os.Exit
}
terminate(1)
return
}
} else if closer, ok := delegate.(interface{ Close() }); ok {
closer.Close()
}
}
if deps.subagentMgr != nil {
if err := drainSubagentManagerWithRetry(deps.subagentMgr, 65*time.Second, 10*time.Second); err != nil {
slog.Error("gateway: terminating after subagent lifecycle drain failure", "error", err)
terminate := deps.terminateProcess
if terminate == nil {
terminate = os.Exit
}
terminate(1)
return
}
}
// Close provider resources (e.g. Claude CLI temp files)
d.providerRegistry.Close()
// Stop permission cache sweep goroutines so they don't leak past shutdown.
if d.permCache != nil {
d.permCache.Close()
}
// Stop sandbox pruning + release containers
if deps.sandboxMgr != nil {
deps.sandboxMgr.Stop()
slog.Info("releasing sandbox containers...")
deps.sandboxMgr.ReleaseAll(context.Background())
}
if deps.sched != nil {
slog.Info("gateway: draining active runs", "timeout", "5s")
deps.sched.Stop() // MarkDraining + StopAll
time.Sleep(5 * time.Second)
}
cancel()
}()
slog.Info("goclaw gateway starting",
"version", Version,
"protocol", protocol.ProtocolVersion,
"agents", d.agentRouter.List(),
"tools", d.toolsReg.Count(),
"channels", d.channelMgr.GetEnabledChannels(),
)
// Tailscale listener: build the mux first, then pass it to initTailscale
// so the same routes are served on both the main listener and Tailscale.
// Compiled via build tags: `go build -tags tsnet` to enable.
mux := d.server.BuildMux()
// Mount channel webhook handlers on the main mux (e.g. Feishu /feishu/events).
// This allows webhook-based channels to share the main server port.
for _, route := range d.channelMgr.WebhookHandlers() {
mux.Handle(route.Path, route.Handler)
slog.Info("webhook route mounted on gateway", "path", route.Path)
}
// Bitrix24: also claim+mount the shared webhook router directly, even if
// no channel_instances row has finished setup yet (bot_code/bot_name
// still empty, or the portal hasn't completed OAuth). /bitrix24/install
// is what completes portal OAuth — gating the route behind a
// fully-configured bot creates a deadlock where an admin can never
// finish installing the first portal on a fresh gateway. ClaimWebhookRoute
// is idempotent (first-claim-wins via CompareAndSwap), so this is a no-op
// if a bitrix24 Channel already claimed the route in the loop above.
if router := bitrix24.WebhookRouter(); router != nil {
if path, handler := router.ClaimWebhookRoute(); path != "" && handler != nil {
mux.Handle(path, handler)
slog.Info("webhook route mounted on gateway", "path", path)
}
}
tsCleanup := initTailscale(ctx, d.cfg, mux)
if tsCleanup != nil {
defer tsCleanup()
}
// Phase 1: suggest localhost binding when Tailscale is active
if d.cfg.Tailscale.Hostname != "" && d.cfg.Gateway.Host == "0.0.0.0" {
slog.Info("Tailscale enabled. Consider setting GOCLAW_HOST=127.0.0.1 for localhost-only + Tailscale access")
}
// Security warnings
if strings.Contains(d.cfg.Database.PostgresDSN, ":goclaw@") {
slog.Warn("security.default_db_password: using default Postgres password — run ./prepare-env.sh to generate a strong one")
}
if len(d.cfg.Gateway.AllowedOrigins) > 0 {
slog.Info("cors: allowed_origins configured", "origins", d.cfg.Gateway.AllowedOrigins)
} else if !edition.Current().IsLimited() {
slog.Warn("security.cors_open: no allowed_origins configured — all WebSocket origins accepted. Set gateway.allowed_origins or GOCLAW_ALLOWED_ORIGINS for production")
}
if allowed, rejected := security.OperatorAllowlistStatus(); len(allowed) > 0 || len(rejected) > 0 {
if len(allowed) > 0 {
slog.Warn("security.ssrf_allowlist: SSRF protection relaxed for operator-configured ranges — tool- and admin-supplied URLs may reach them",
"env", security.SSRFAllowedCIDRsEnv, "allowed", allowed)
}
if len(rejected) > 0 {
slog.Warn("security.ssrf_allowlist_rejected: entries refused; cloud-metadata, multicast and unspecified ranges can never be allowlisted",
"env", security.SSRFAllowedCIDRsEnv, "rejected", rejected)
}
}
if err := d.server.Start(ctx); err != nil {
slog.Error("gateway error", "error", err)
os.Exit(1)
}
}