Files
goclaw/cmd/gateway.go
T
bd5adc61c8 feat(bitrix24): imbot.v2 migration, 2-way media, openline sender-tag echo, and hardening (#1236)
* refactor(bitrix24): rename "Path B" framing to maintainer-specified naming [B24:2794]

Per maintainer hard rule #10 (no generic "Path A/B" framing) from PR #1061
review. The Bitrix24 MCP auto-onboard flow is Bitrix-specific glue
("Bitrix24 OAuth -> existing mcp_user_credentials bridge"), NOT a generic
MCP architecture pattern.

Naming convention applied consistently:
- First mention per file: full "Bitrix24 OAuth -> existing
  mcp_user_credentials bridge" (matches maintainer comment verbatim).
- Subsequent mentions in same file: shortened "mcp_user_credentials bridge".
- Test/log context referencing literal endpoint /api/auto-onboard: keep
  "auto-onboard" reference (it's the actual API endpoint name).

Changes are documentation-only:
- Rename in code comments + test descriptions + plan docs.
- Clarify framing in mcp_client.go + provisioner.go doc comments to
  emphasize Bitrix-specific glue (not generic MCP infra).
- Reuse existing mcp_user_credentials table + MCPServerStore methods
  (no schema / store / abstraction change).

Files:
- cmd/gateway.go (factory registration doc)
- internal/channels/bitrix24/{channel,factory,mcp_client,provisioner}.go
- internal/channels/bitrix24/{mcp_client,provisioner}_test.go
- plan/goclaw-mcp-integration.md (21 occurrences)

Verified: go build + MCP-related tests pass (TestProvision*,
TestInitMCPProvisioner*, TestMCPClient*).

Phase 1 of Path C execution per
plans/reports/decision-log-260519-1555-bitrix24-pr-fork-decision.md.

* fix: confine outbound media paths to agent workspace [B24:2794]

Tool MEDIA:<path> output reached channel file-upload sinks (Bitrix
imbot.v2.File.upload, Telegram sendDocument, etc.) verbatim via
parseMediaResult, with no workspace-boundary check. A malicious or buggy
tool emitting MEDIA:/etc/passwd could exfiltrate arbitrary files to chat.

Extract the EvalSymlinks+Rel containment from extractMediaFromContent into
a shared confineToWorkspace helper and apply it at the parseMediaResult
sink in processToolResult. Fixing at the source/egress boundary protects
every channel at once rather than per-channel. Paths that escape the
workspace are dropped and logged (security.media_path_rejected).

Add TestConfineToWorkspace (boundary unit) and
TestParseMediaResultConfinedToWorkspace (sink regression for H2).

* feat(bitrix24): support inbound + outbound media via imbot.v2 File API [B24:2794]

Bitrix24 channel was text-only; attachments were parsed but dropped.
- Inbound: download chat files via imbot.v2.File.download (one-time URL),
  forward to the agent with MIME preserved (internal/channels/bitrix24/download.go).
- Outbound: upload agent media to the chat via imbot.v2.File.upload
  (internal/channels/bitrix24/send_media.go).
- Add BaseChannel.HandleMessageMedia to preserve MIME/filename through the bus.
- Per-channel media_max_mb cap (default 20) applies to both directions.

Tests: 92 pass (internal/channels/bitrix24 + internal/channels), go vet clean (PG + sqliteonly).

* refactor(bitrix24): migrate messaging/bot-list/unregister to imbot v2 API [B24:2794]

Move outbound REST calls to the imbot v2 family (keeps register on v1):
- imbot.message.add -> imbot.v2.Chat.Message.send (fields.message shape, live-verified)
- imbot.bot.list (+ legacy imbot.list fallback) -> imbot.v2.Bot.list; add botListRows
  to normalize the v2 {bots:[...]} envelope, legacy array, and id-keyed map forms
- imbot.unregister -> imbot.v2.Bot.unregister

Bot registration stays on v1 imbot.register: v2 imbot.v2.Bot.register changes the
event-delivery model (per-event handler URLs -> eventMode), which would require
rewriting the inbound event parser. No user-facing behavior change.

Tests: bitrix24 package green; go vet ./... clean.

* feat(bitrix24): route whisper via v1 SKIP_CONNECTOR + add v2 replyId [B24:2794]

Bot was leaking HiddenMessage (whisper) replies to the external Zalo
connector because every outbound call went through imbot.v2.Chat.Message.send,
which has no equivalent of the v1 SKIP_CONNECTOR flag. Branch the outbound
path on inbound visibility:

  whisper → imbot.message.add + SKIP_CONNECTOR=Y  (v1, send_v1.go)
  public  → imbot.v2.Chat.Message.send + fields.replyId  (v2, send_v2.go)

Pipeline:
  events.go        parse data[PARAMS][PARAMS][COMPONENT_ID]=HiddenMessage
                   into EventParams.IsHiddenMessage (form + JSON variants)
  handle.go        set bitrix_visibility on InboundMessage.Metadata
  consumer         forward visibility + message_id into OutboundMessage
  send.go          resolveSendOptions + sendChunk dispatcher +
                   shared callWithRateLimitRetry helper
  metadata_keys.go single source of truth for the keys + values

Defaults preserve pre-refactor behaviour: callers that don't populate
bitrix_visibility still go through v2 public, and replyId is omitted
unless a numeric bitrix_message_id arrives in metadata.

Tests:
  TestParseEvent_FormURLEncoded_IsHiddenMessage  (3 cases)
  TestParseEvent_JSON_IsHiddenMessage             (3 cases)
  TestResolveSendOptions                          (8 cases)
  TestSend_BranchesOnVisibility                   (4 cases)

* feat(bitrix24): openline sender-tag echo on replies [B24:2794]

Openline sender-tag echo (this change):
- Capture the connector sender tag ("[name #id]:" or "[name] #id:") from
  inbound openline group messages, strip it from the body the agent sees,
  and re-prepend the canonical "[name] #id:" form to the reply so the Open
  Channel connector routes the answer back to the right external user.
- New sender_prefix.go helper (+ test) accepts both inbound layouts and
  emits one canonical form; scoped to messages carrying the tag, so plain
  chats are unaffected.
- metadata_keys.go: MetaKeySenderPrefix; handle.go capture/strip/stash;
  gateway_consumer_normal.go forwards the key; send.go prepends it on the
  first chunk before chunking.

Bundled bitrix24 channel-core work already on this branch:
- handle.go: @mention is the sole trigger for both staff and connector
  customers; unmentioned traffic is dropped (was: drop all connector msgs).
- isGroupMessageType: treat SONET_GROUP "B" as a group.
- handle_test.go, mcp_client_test.go: cover the above.

* feat(bitrix24): accept colon-less openline sender tag, echo [name] #id [B24:2794]

The Open Channel connector dropped the trailing colon from its sender tag:
inbound now arrives as "[Name] #id <msg>" (was "[Name] #id: <msg>"). The
id-bearing patterns required the colon, so the tag fell through to the
name-only branch and the reply echoed "[Name]" — dropping the #id the
connector needs to route the answer back.

- sender_prefix.go: make the trailing ":" optional on both id layouts
  ([name #id] / [name] #id, with or without colon) and echo the canonical
  "[name] #id" (no colon) to match the connector's current format. Bare
  "[name]" (no id) still echoes "[name]" for Open Channel only.
- handle.go: gate the bare name-only layout to Open Channel (isOpenChannel)
  so ordinary group chats starting with "[x] ..." are left untouched.
- sender_prefix_test.go: cover colon/no-colon x id-inside/id-outside, the
  name-only openline case, and the non-openline no-op.

* fix: security and robustness fixes from the bitrix24 channel review [B24:2794]

- download.go: block redirect-based SSRF on inbound media. CheckRedirect
  re-validates each hop (http(s) only, reject private/loopback/link-local
  hosts, cap hops); the initial portal-domain pin is no longer bypassable
  via a 3xx to an internal service. Public-host redirects still allowed.
- handle.go: extract/echo the openline sender tag only for Open Channel
  sessions (was: any group chat), removing bogus prefixes in CRM group
  chats and narrowing the forged-tag misroute surface.
- loop_tools.go + loop_media.go: confine result.Media to the agent / team /
  tenant-allowed roots (new confineToAnyRoot) before a channel uploads it,
  so a prompt-injected out-of-workspace path (e.g. /etc/passwd) cannot
  exfiltrate, while legitimate cross-workspace media (team files, delegatee
  output) still flows.
- send_media.go: bounded outbound read via io.LimitReader replaces the
  os.Stat + os.ReadFile pair, closing the TOCTOU size-cap bypass; cap a
  single message's outbound attachments at 10 (mirrors inbound).
- register.go: paginate imbot.v2.Bot.list (limit/offset + hasNextPage,
  capped at 40 pages) so verify/lookup see bots past the first 50.
- mcp_client.go: redact access_token / refresh_token / client_secret from an
  echoed MCP error body before it is logged or returned (+ test).

* fix(security): validate resolved dial IP on Bitrix media redirects [B24:2794]

The inbound media download redirect guard only string-checked the redirect
hostname (isPrivateOrLoopback on req.URL.Hostname()), so a redirect to a public
hostname that resolves to 127.0.0.1 / 169.254.169.254 / an RFC1918 address — or a
DNS-rebinding swap between check and dial — still passed the guard and the client
would connect. Reported in PR review.

Add security.NewRedirectFollowingSafeClient: it follows redirects but validates
the RESOLVED destination IP of every hop at dial time via net.Dialer.Control,
reusing the existing blocked-CIDR list. The IP it checks is the IP actually
dialed, so both redirect-to-internal and DNS rebinding are refused, while
legitimate public CDN redirects still succeed. download.go now uses it instead of
the hostname-string guard.

Tests: deterministic dial-control table (loopback / link-local / private /
multicast / unspecified / public, v4 + v6), malformed/non-IP addr, test bypass,
loopback-dial-blocked client wiring, and redirect cap + scheme checks.

* feat(bitrix24): per-participant Zalo openline identity from 3-token sender tag [B24:2794]

Parse the connector's "[Name] #uid #msgId" sender tag so each external
customer in a shared Open Channel group gets its own contact + USER.md
instead of collapsing onto the connector proxy id. Identity minting is
gated on IS_CONNECTOR=Y to reject operator forged tags. Echo back the
msgId only ("#msgId") on replies; keep the legacy single-number and
name-only layouts unchanged. Zero DB migration.

- sender_prefix.go: parseOpenlineSenderTag() classifies 3-token / legacy / name-only
- handle.go: synthetic senderID "openlines:{instance}:{chat}:{uid}" + participant_user_id metadata, gated on FromIsConnector
- gateway_consumer_normal.go: deriveGroupUserID() routes participant -> per-person scope, group fallback otherwise
- send.go: buildAddressMention numeric-id guard so synthetic ids don't emit invalid [USER=...] BBCode
- MetaKeyMessageID kept as Bitrix MESSAGE_ID (drives v2 fields.replyId); connector msgId surfaced only via echo prefix

---------

Co-authored-by: DangTinh311 <dangtinh31193@gmail.com>
Co-authored-by: Chinh Dang <chinhdang@192.168.68.104>
2026-06-22 14:23:34 +07:00

817 lines
34 KiB
Go

package cmd
import (
"context"
"fmt"
"io"
"log/slog"
"os"
"os/signal"
"path/filepath"
"strings"
"syscall"
"time"
"github.com/google/uuid"
"github.com/nextlevelbuilder/goclaw/internal/agent"
"github.com/nextlevelbuilder/goclaw/internal/bgalert"
"github.com/nextlevelbuilder/goclaw/internal/bootstrap"
"github.com/nextlevelbuilder/goclaw/internal/bus"
"github.com/nextlevelbuilder/goclaw/internal/cache"
"github.com/nextlevelbuilder/goclaw/internal/channelmemory"
"github.com/nextlevelbuilder/goclaw/internal/channels"
"github.com/nextlevelbuilder/goclaw/internal/channels/bitrix24"
"github.com/nextlevelbuilder/goclaw/internal/channels/discord"
"github.com/nextlevelbuilder/goclaw/internal/channels/facebook"
"github.com/nextlevelbuilder/goclaw/internal/channels/feishu"
"github.com/nextlevelbuilder/goclaw/internal/channels/pancake"
slackchannel "github.com/nextlevelbuilder/goclaw/internal/channels/slack"
"github.com/nextlevelbuilder/goclaw/internal/channels/telegram"
"github.com/nextlevelbuilder/goclaw/internal/channels/whatsapp"
"github.com/nextlevelbuilder/goclaw/internal/channels/zalo"
zalopersonal "github.com/nextlevelbuilder/goclaw/internal/channels/zalo/personal"
"github.com/nextlevelbuilder/goclaw/internal/config"
"github.com/nextlevelbuilder/goclaw/internal/consolidation"
"github.com/nextlevelbuilder/goclaw/internal/edition"
"github.com/nextlevelbuilder/goclaw/internal/eventbus"
"github.com/nextlevelbuilder/goclaw/internal/gateway"
"github.com/nextlevelbuilder/goclaw/internal/gateway/methods"
"github.com/nextlevelbuilder/goclaw/internal/hooks"
httpapi "github.com/nextlevelbuilder/goclaw/internal/http"
kg "github.com/nextlevelbuilder/goclaw/internal/knowledgegraph"
mcpbridge "github.com/nextlevelbuilder/goclaw/internal/mcp"
mcpoauth "github.com/nextlevelbuilder/goclaw/internal/mcp/oauth"
"github.com/nextlevelbuilder/goclaw/internal/media"
"github.com/nextlevelbuilder/goclaw/internal/providers"
"github.com/nextlevelbuilder/goclaw/internal/scheduler"
"github.com/nextlevelbuilder/goclaw/internal/security"
"github.com/nextlevelbuilder/goclaw/internal/skills"
"github.com/nextlevelbuilder/goclaw/internal/store"
"github.com/nextlevelbuilder/goclaw/internal/tools"
usagecaps "github.com/nextlevelbuilder/goclaw/internal/usage/caps"
"github.com/nextlevelbuilder/goclaw/internal/vault"
"github.com/nextlevelbuilder/goclaw/pkg/protocol"
// Register workstation backend factories via init().
_ "github.com/nextlevelbuilder/goclaw/internal/workstation/backends"
)
func gatewayLogOutput() io.Writer {
logFile := strings.TrimSpace(os.Getenv("GOCLAW_LOG_FILE"))
if logFile == "" {
if st, err := os.Stat("/var/log/goclaw"); err == nil && st.IsDir() {
logFile = "/var/log/goclaw/goclaw.log"
}
}
if logFile == "" {
return os.Stdout
}
f, err := os.OpenFile(logFile, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0o644)
if err != nil {
fmt.Fprintf(os.Stderr, "warning: cannot open GOCLAW_LOG_FILE=%q: %v\n", logFile, err)
return os.Stdout
}
fmt.Fprintf(os.Stderr, "logging to %s\n", logFile)
return io.MultiWriter(os.Stdout, f)
}
func runGateway() {
// Setup structured logging
logLevel := slog.LevelInfo
if verbose {
logLevel = slog.LevelDebug
}
// Env override (docker/K8s friendly, default: info): GOCLAW_LOG_LEVEL=debug|info|warn|error
if lvl := os.Getenv("GOCLAW_LOG_LEVEL"); lvl != "" {
switch strings.ToLower(lvl) {
case "debug":
logLevel = slog.LevelDebug
case "info":
logLevel = slog.LevelInfo
case "warn":
logLevel = slog.LevelWarn
case "error":
logLevel = slog.LevelError
default:
fmt.Fprintf(os.Stderr, "warning: unknown GOCLAW_LOG_LEVEL=%q, using info\n", lvl)
}
}
logOutput := gatewayLogOutput()
textHandler := slog.NewTextHandler(logOutput, &slog.HandlerOptions{Level: logLevel})
logTee := gateway.NewLogTee(textHandler)
slog.SetDefault(slog.New(logTee))
// Load config
cfgPath := resolveConfigPath()
cfg, err := config.Load(cfgPath)
if err != nil {
slog.Error("failed to load config", "error", err)
os.Exit(1)
}
if err := config.ValidateGatewayAuth(cfg.Gateway); err != nil {
slog.Error("unsafe gateway auth configuration", "error", err)
os.Exit(1)
}
// Edition override: explicit GOCLAW_EDITION takes precedence over auto-detection.
// Auto-detection happens later in setupStoresAndTracing (sqlite → lite).
if edName := os.Getenv("GOCLAW_EDITION"); edName != "" {
switch edName {
case "lite":
edition.SetCurrent(edition.Lite)
slog.Info("edition: lite (explicit)")
case "standard":
edition.SetCurrent(edition.Standard)
slog.Info("edition: standard (explicit)")
default:
slog.Warn("unknown GOCLAW_EDITION, using standard", "value", edName)
}
}
// Create core components
msgBus := bus.New()
// V3 domain event bus for consolidation pipeline (episodic → semantic → dreaming)
domainBus := eventbus.NewDomainEventBus(eventbus.Config{
QueueSize: 1000,
WorkerCount: 2,
})
domainBus.Start(context.Background())
defer func() {
if err := domainBus.Drain(10 * time.Second); err != nil {
slog.Warn("domain event bus drain timeout", "error", err)
}
}()
// Create model registry with forward-compat resolvers (shared across all providers)
modelReg := providers.NewInMemoryRegistry()
modelReg.RegisterResolver("anthropic", &providers.AnthropicForwardCompat{})
modelReg.RegisterResolver("openai", &providers.OpenAIForwardCompat{})
// Create provider registry
providerRegistry := providers.NewRegistry(store.TenantIDFromContext)
registerProviders(providerRegistry, cfg, modelReg)
// Resolve workspace (must be absolute for system prompt + file tool path resolution)
workspace := config.ExpandHome(cfg.Agents.Defaults.Workspace)
if !filepath.IsAbs(workspace) {
workspace, _ = filepath.Abs(workspace)
}
os.MkdirAll(workspace, 0755)
// Detect server IPs for output scrubbing (prevents IP leaks via web_fetch, exec, etc.)
// Skip for desktop/lite — localhost-only, no multi-tenant exposure risk
if !edition.Current().IsLimited() {
tools.DetectServerIPs(context.Background())
}
toolsReg, execApprovalMgr, mcpMgr, sandboxMgr, browserMgr, webFetchTool, ttsTool, audioMgr, permPE, toolPE, dataDir, agentCfg := setupToolRegistry(cfg, workspace, providerRegistry)
if browserMgr != nil {
defer browserMgr.Close()
}
if mcpMgr != nil {
defer mcpMgr.Stop()
}
pgStores, traceCollector, snapshotWorker := setupStoresAndTracing(cfg, dataDir, msgBus)
if browserMgr != nil && pgStores != nil && pgStores.BrowserCookies != nil && cfg.Tools.Browser.CookieSyncEnabled {
browserMgr.SetCookieProvider(newStoreBrowserCookieProvider(pgStores.BrowserCookies))
}
// Recover from crashes: flip ghost 'summoning' rows to 'summon_failed'.
// Summon goroutines don't survive process restart; stale DB rows would trap the UI.
if pgStores.Agents != nil {
if n, err := pgStores.Agents.ResetStuckSummoning(context.Background()); err != nil {
slog.Warn("agents.reset_stuck_summoning_failed", "err", err)
} else if n > 0 {
slog.Info("agents.reset_stuck_summoning", "count", n)
}
}
if traceCollector != nil {
defer traceCollector.Stop()
// OTel OTLP export: compiled via build tags. Build with 'go build -tags otel' to enable.
initOTelExporter(context.Background(), cfg, traceCollector)
}
if snapshotWorker != nil {
defer snapshotWorker.Stop()
}
// Redis cache: compiled via build tags. Build with 'go build -tags redis' to enable.
redisClient := initRedisClient(cfg)
defer shutdownRedis(redisClient)
// Register providers from DB (overrides config providers).
if pgStores.Providers != nil {
dbGatewayAddr := loopbackAddr(cfg.Gateway.Host, cfg.Gateway.Port)
registerProvidersFromDB(providerRegistry, pgStores.Providers, pgStores.ConfigSecrets, dbGatewayAddr, cfg.Gateway.Token, pgStores.MCP, cfg, modelReg)
}
slog.Info("model registry initialized", "anthropic_models", len(modelReg.Catalog("anthropic")), "openai_models", len(modelReg.Catalog("openai")))
// Warn if deprecated session scope settings are configured
if cfg.Sessions.Scope != "" && cfg.Sessions.Scope != "per-sender" {
slog.Warn("sessions.scope config is deprecated and ignored — fixed to per-sender", "configured", cfg.Sessions.Scope)
}
if cfg.Sessions.DmScope != "" && cfg.Sessions.DmScope != "per-channel-peer" {
slog.Warn("sessions.dm_scope config is deprecated and ignored — fixed to per-channel-peer", "configured", cfg.Sessions.DmScope)
}
seedSystemConfigs(pgStores.SystemConfigs, pgStores.Tenants, cfg)
// Read back system_configs from DB and overlay onto in-memory config.
if pgStores.SystemConfigs != nil {
if sysConfigs, err := pgStores.SystemConfigs.List(
store.WithTenantID(context.Background(), store.MasterTenantID),
); err == nil && len(sysConfigs) > 0 {
cfg.ApplySystemConfigs(sysConfigs)
slog.Info("system_configs applied to in-memory config", "keys", len(sysConfigs))
}
}
// Re-apply tool rate limiter using DB-overlaid config. setupToolRegistry
// initialised the limiter from the JSON5 default before ApplySystemConfigs
// ran, so DB-driven changes to tools.rate_limit_per_hour were lost. Replace
// the limiter object now that cfg reflects the DB value. Safe: server has
// not started, no in-flight tool calls.
if cfg.Tools.RateLimitPerHour > 0 {
toolsReg.SetRateLimiter(tools.NewToolRateLimiter(cfg.Tools.RateLimitPerHour))
slog.Info("tool rate limiting reapplied from system_configs", "per_hour", cfg.Tools.RateLimitPerHour)
} else {
toolsReg.SetRateLimiter(nil)
}
setupMemoryEmbeddings(pgStores, providerRegistry)
usageCapSvc := usagecaps.NewService(pgStores.UsageCaps, pgStores.Providers)
// Resolve background provider for consolidation + vault enrichment.
// Fallback: background.provider → agent.default_provider → first registered provider.
bgProvider, bgModel := resolveBackgroundProvider(cfg, providerRegistry)
// V3: Wire consolidation pipeline (episodic → semantic → KG → dreaming)
if pgStores.Episodic != nil {
if bgProvider != nil {
var kgExtractor *kg.Extractor
if pgStores.KnowledgeGraph != nil {
kgExtractor = kg.NewExtractor(bgProvider, bgModel, 0)
kgExtractor.SetUsageCapService(usageCapSvc)
}
cleanupConsolidation := consolidation.Register(consolidation.ConsolidationDeps{
EpisodicStore: pgStores.Episodic,
MemoryStore: pgStores.Memory,
KGStore: pgStores.KnowledgeGraph,
SessionStore: pgStores.Sessions,
EventBus: domainBus,
SystemConfigs: pgStores.SystemConfigs,
Registry: providerRegistry,
Extractor: kgExtractor,
AlertDeps: bgalert.AlertDeps{SystemConfigs: pgStores.SystemConfigs, MsgBus: msgBus},
UsageCaps: usageCapSvc,
AgentStore: pgStores.Agents,
})
defer cleanupConsolidation()
slog.Info("consolidation pipeline registered", "provider", bgProvider.Name(), "model", bgModel)
} else {
slog.Warn("consolidation pipeline skipped: no provider available")
}
}
if memorySvc := makeChannelMemoryService(pgStores, domainBus, providerRegistry, usageCapSvc); memorySvc != nil {
cleanupChannelMemory := (&channelmemory.Worker{Service: memorySvc}).Start(context.Background())
defer cleanupChannelMemory()
slog.Info("channel memory extraction worker registered")
}
// V3: Wire vault enrichment worker (async summary + embedding + auto-linking).
// Provider is resolved per-tenant at runtime — no static provider needed.
var enrichProgress *vault.EnrichProgress
var enrichWorker *vault.EnrichWorker
if pgStores.Vault != nil && providerRegistry != nil {
cleanupVaultEnrich, ep, ew := vault.RegisterEnrichWorker(vault.EnrichWorkerDeps{
VaultStore: pgStores.Vault,
SystemConfigs: pgStores.SystemConfigs,
Registry: providerRegistry,
EventBus: domainBus,
MsgBus: msgBus,
TeamStore: pgStores.Teams,
AlertDeps: bgalert.AlertDeps{SystemConfigs: pgStores.SystemConfigs, MsgBus: msgBus},
UsageCaps: usageCapSvc,
})
enrichProgress = ep
enrichWorker = ew
defer cleanupVaultEnrich()
slog.Info("vault enrichment worker registered (per-tenant provider resolution)")
}
loadBootstrapFiles(pgStores, workspace, agentCfg)
// Backfill CAPABILITIES.md for pre-v3 agents that don't have it yet.
if count, err := bootstrap.BackfillCapabilities(context.Background(), pgStores.DB); err != nil {
slog.Warn("bootstrap: capabilities backfill failed", "error", err)
} else if count > 0 {
slog.Info("bootstrap: capabilities backfill complete", "agents", count)
}
if readImage, ok := toolsReg.Get("read_image"); ok {
if t, ok := readImage.(*tools.ReadImageTool); ok {
t.SetUsageCapService(usageCapSvc)
}
}
// Subagent system (secureCLI store wired so subagent ExecTools enforce the gate)
subagentMgr := setupSubagents(providerRegistry, cfg, msgBus, toolsReg, workspace, sandboxMgr, pgStores.SecureCLI, usageCapSvc)
if subagentMgr != nil {
// Wire announce queue for batched subagent result delivery (matching TS debounce pattern).
announceQueue := tools.NewAnnounceQueue(1000, 20, makeDelegateAnnounceCallback(subagentMgr, msgBus))
subagentMgr.SetAnnounceQueue(announceQueue)
if pgStores.SubagentTasks != nil {
subagentMgr.SetTaskStore(pgStores.SubagentTasks)
}
toolsReg.Register(tools.NewSpawnTool(subagentMgr, "default", 0))
slog.Info("subagent system enabled", "tools", []string{"spawn"})
}
skillsLoader, skillSearchTool, globalSkillsDir, bundledSkillsDir, builtinSkillsDir := setupSkillsSystem(cfg, workspace, dataDir, pgStores, toolsReg, providerRegistry, msgBus)
_ = skillSearchTool // used via wireExtras → skillsLoader; kept for type clarity
// Register cron/heartbeat/session/message tools, aliases, allow-paths, store wiring.
heartbeatTool, hasMemory := wireExtraTools(pgStores, toolsReg, msgBus, workspace, dataDir, agentCfg, globalSkillsDir, builtinSkillsDir)
// Register workstation_exec + claude_remote tools (Standard edition only; deny-all until Phase 6).
// cleanupWorkstation stops the activity sink retention goroutine and drains the write buffer.
cleanupWorkstation := wireWorkstationTools(pgStores, toolsReg, domainBus)
defer cleanupWorkstation()
// Create all agents — resolved lazily from database by the managed resolver.
agentRouter := agent.NewRouter()
if traceCollector != nil {
agentRouter.SetTraceCollector(traceCollector)
}
slog.Info("agents will be resolved lazily from database")
// Create gateway server and wire enforcement
server := gateway.NewServer(cfg, msgBus, agentRouter, pgStores.Sessions, toolsReg)
server.SetVersion(Version)
server.SetDB(pgStores.DB)
server.SetPolicyEngine(permPE)
server.SetPairingService(pgStores.Pairing)
server.SetMessageBus(msgBus)
server.SetOAuthHandler(httpapi.NewOAuthHandler(pgStores.Providers, pgStores.ConfigSecrets, providerRegistry, msgBus))
// contextFileInterceptor is created inside wireExtras.
// Declared here so it can be passed to registerAllMethods → AgentsMethods
// for immediate cache invalidation on agents.files.set.
var contextFileInterceptor *tools.ContextFileInterceptor
// Set agent store for tools_invoke context injection + wire extras
if pgStores.Agents != nil {
server.SetAgentStore(pgStores.Agents)
}
// Build OAuth token refresher before wireExtras so the resolver can inject tokens.
var mcpOAuthRefresher mcpbridge.OAuthTokenProvider
if pgStores != nil && pgStores.MCPOAuthTokens != nil {
mcpOAuthRefresher = mcpoauth.NewRefresher(pgStores.MCPOAuthTokens, security.NewSafeClient(15*time.Second))
}
var mcpPool *mcpbridge.Pool
var mediaStore *media.Store
var postTurn tools.PostTurnProcessor
contextFileInterceptor, mcpPool, mediaStore, postTurn = wireExtras(pgStores, agentRouter, providerRegistry, modelReg, msgBus, pgStores.Sessions, toolsReg, toolPE, skillsLoader, hasMemory, traceCollector, workspace, cfg.Gateway.InjectionAction, cfg, sandboxMgr, redisClient, domainBus, usageCapSvc, mcpOAuthRefresher)
if mcpPool != nil {
defer mcpPool.Stop()
}
// Populate shared deps struct used by extracted helper methods.
deps := &gatewayDeps{
cfg: cfg,
server: server,
msgBus: msgBus,
pgStores: pgStores,
providerRegistry: providerRegistry,
agentRouter: agentRouter,
toolsReg: toolsReg,
skillsLoader: skillsLoader,
enrichProgress: enrichProgress,
enrichWorker: enrichWorker,
workspace: workspace,
dataDir: dataDir,
domainBus: domainBus,
usageCapSvc: usageCapSvc,
audioMgr: audioMgr,
}
gatewayAddr := loopbackAddr(cfg.Gateway.Host, cfg.Gateway.Port)
var mcpToolLister httpapi.MCPToolLister
if mcpMgr != nil {
mcpToolLister = mcpMgr
}
httpapi.InitGatewayToken(cfg.Gateway.Token)
mcpbridge.SetAllowedHosts(cfg.Gateway.MCPAllowedHosts) // operator allowlist: trusted MCP hosts exempt from private-IP SSRF block
httpapi.InitGatewayNoAuthFallbackAllowed(config.GatewayNoAuthFallbackAllowed(cfg.Gateway))
exportTokenStore := httpapi.InitExportTokenStore()
defer exportTokenStore.Stop()
agentsH, skillsH, tracesH, mcpH, channelInstancesH, providersH, builtinToolsH, pendingMessagesH, teamEventsH, secureCLIH, secureCLIGrantH, mcpUserCredsH := wireHTTP(pgStores, cfg.Agents.Defaults.Workspace, dataDir, bundledSkillsDir, msgBus, domainBus, toolsReg, providerRegistry, modelReg, permPE.IsOwner, gatewayAddr, mcpToolLister, usageCapSvc, cfg, cfg.Skills)
// Wire dependencies for system prompt preview parity.
if agentsH != nil {
agentsH.SetPreviewDeps(toolsReg, skillsLoader)
var skillAccess store.SkillAccessStore
if pgStores.Skills != nil {
skillAccess, _ = pgStores.Skills.(store.SkillAccessStore)
}
agentsH.SetPreviewStores(pgStores.Teams, pgStores.AgentLinks, skillAccess)
}
// External wake/trigger API
wakeH := httpapi.NewWakeHandler(agentRouter)
if postTurn != nil {
wakeH.SetPostTurnProcessor(postTurn)
}
// MCP OAuth handler — per-server OAuth 2.1 client flows.
var mcpOAuthH *httpapi.MCPOAuthHandler
if pgStores != nil && pgStores.MCP != nil && pgStores.MCPOAuthTokens != nil {
safeHTTPClient := security.NewSafeClient(15 * time.Second)
var oauthRefresher *mcpoauth.Refresher
if r, ok := mcpOAuthRefresher.(*mcpoauth.Refresher); ok {
oauthRefresher = r
}
mcpOAuthH = httpapi.NewMCPOAuthHandler(httpapi.MCPOAuthHandlerDeps{
MCPStore: pgStores.MCP,
OAuthStore: pgStores.MCPOAuthTokens,
Discoverer: mcpoauth.NewDiscoverer(safeHTTPClient),
FlowMgr: mcpoauth.NewFlowManager(safeHTTPClient),
Refresher: oauthRefresher,
EventBus: msgBus,
PublicURL: cfg.Gateway.PublicURL,
Port: cfg.Gateway.Port,
TenantStore: pgStores.Tenants,
})
// Inject OAuth token provider into MCP tools handler so on-demand tool
// discovery can authenticate against OAuth-protected MCP servers.
if mcpH != nil && mcpOAuthRefresher != nil {
mcpH.SetOAuthProvider(mcpOAuthRefresher)
}
// Inject the OAuth token store so the update handler can purge stale tokens
// when a server's URL or OAuth config changes.
if mcpH != nil {
mcpH.SetOAuthStore(pgStores.MCPOAuthTokens)
}
}
// Wire all server.Set*Handler() calls via extracted helper.
deps.wireHTTPHandlersOnServer(
httpHandlers{
agents: agentsH,
skills: skillsH,
traces: tracesH,
mcp: mcpH,
channelInstances: channelInstancesH,
providers: providersH,
builtinTools: builtinToolsH,
pendingMessages: pendingMessagesH,
teamEvents: teamEventsH,
secureCLI: secureCLIH,
secureCLIGrant: secureCLIGrantH,
mcpUserCreds: mcpUserCredsH,
mcpOAuth: mcpOAuthH,
},
wakeH,
mcpPool,
postTurn,
mediaStore,
)
// System backup API — admin + owner only, SSE progress streaming.
server.SetBackupHandler(httpapi.NewBackupHandler(cfg, cfg.Database.PostgresDSN, Version, permPE.IsOwner))
// System restore API — admin + owner only, multipart upload + SSE progress.
server.SetRestoreHandler(httpapi.NewRestoreHandler(cfg, cfg.Database.PostgresDSN, permPE.IsOwner))
// S3 backup integration — admin + owner only.
server.SetBackupS3Handler(httpapi.NewBackupS3Handler(cfg, cfg.Database.PostgresDSN, Version, pgStores.ConfigSecrets, permPE.IsOwner))
// Tenant-scoped backup/restore — owner or tenant admin.
if pgStores.Tenants != nil {
server.SetTenantBackupHandler(httpapi.NewTenantBackupHandler(pgStores.DB, cfg, pgStores.Tenants, Version, permPE.IsOwner))
}
// Register all RPC methods
server.SetLogTee(logTee)
server.SetRuntimeLogsHandler(httpapi.NewRuntimeLogsHandler(logTee))
pairingMethods, heartbeatMethods, chatMethods, cfgPermsMethods := registerAllMethods(server, agentRouter, pgStores.Sessions, pgStores.RunTimeline, pgStores.Cron, pgStores.Pairing, cfg, cfgPath, workspace, dataDir, msgBus, execApprovalMgr, pgStores.Agents, pgStores.Skills, pgStores.ConfigSecrets, pgStores.Teams, contextFileInterceptor, logTee, pgStores.Heartbeats, pgStores.ConfigPermissions, pgStores.SystemConfigs, pgStores.Tenants, pgStores.SkillTenantCfgs, audioMgr, usageCapSvc)
// Phase 3: Agent hooks RPC methods (hooks.list/create/update/delete/toggle/test/history).
if hs, ok := pgStores.Hooks.(hooks.HookStore); ok && hs != nil {
hm := methods.NewHookMethods(hs, edition.Current())
// Reuse dispatcher handlers for dry-run test runner so UI test panel
// exercises the exact code that will run in production.
if sharedHookHandlers != nil {
hm.SetTestRunner(methods.NewDispatcherTestRunner(sharedHookHandlers))
}
hm.Register(server.Router())
slog.Info("registered hooks RPC methods")
}
// Workstations WS methods — Standard edition only.
// Lite (desktop/SQLite) must NOT expose workstation RPC methods.
if edition.Current().Name != "lite" && pgStores.Workstations != nil && pgStores.WorkstationLinks != nil {
wsMethods := methods.NewWorkstationsMethods(pgStores.Workstations, pgStores.WorkstationLinks)
if pgStores.WorkstationPermissions != nil {
wsMethods.SetPermStore(pgStores.WorkstationPermissions)
}
if pgStores.WorkstationActivity != nil {
wsMethods.SetActivityStore(pgStores.WorkstationActivity)
}
wsMethods.Register(server.Router())
slog.Info("registered workstations RPC methods")
}
// Wire post-turn processor for team task dispatch (WS chat.send + HTTP API paths).
if postTurn != nil {
chatMethods.SetPostTurnProcessor(postTurn)
server.SetPostTurnProcessor(postTurn) // HTTP: /v1/chat/completions, /v1/responses
wakeH.SetPostTurnProcessor(postTurn) // HTTP: /v1/agents/{id}/wake
}
// Wire pairing event broadcasts to all WS clients.
pairingMethods.SetBroadcaster(server.BroadcastEvent)
// Wire pairing request callback — works for both PG and SQLite stores.
type pairingRequestNotifier interface {
SetOnRequest(func(code, senderID, channel, chatID string))
}
if ps, ok := pgStores.Pairing.(pairingRequestNotifier); ok {
ps.SetOnRequest(func(code, senderID, channel, chatID string) {
server.BroadcastEvent(*protocol.NewEvent(protocol.EventDevicePairReq, map[string]any{
"code": code, "sender_id": senderID, "channel": channel, "chat_id": chatID,
}))
})
}
// Channel manager
channelMgr := channels.NewManager(msgBus)
deps.channelMgr = channelMgr
// Wire channel member resolver into permission grant paths (WS + HTTP) so
// file_writer grants coming from the Web UI auto-enrich their metadata.
cfgPermsMethods.SetMemberResolver(channelMgr)
if channelInstancesH != nil {
channelInstancesH.SetMemberResolver(channelMgr)
// Setter (not constructor) because wireHTTP runs before channelMgr is
// created — required for handleDelete to invoke ChannelDestroyer on
// Bitrix24 channels (imbot.unregister bot cleanup).
channelInstancesH.SetChannelManager(channelMgr)
}
// Wire channel sender + tenant checker on message tool (now that channelMgr exists)
if t, ok := toolsReg.Get("message"); ok {
if cs, ok := t.(tools.ChannelSenderAware); ok {
cs.SetChannelSender(channelMgr.SendToChannel)
}
if tc, ok := t.(tools.ChannelTenantCheckerAware); ok {
tc.SetChannelTenantChecker(channelMgr.ChannelTenantID)
}
}
// Wire group member lister on list_group_members tool
if t, ok := toolsReg.Get("list_group_members"); ok {
if gl, ok := t.(tools.GroupMemberListerAware); ok {
gl.SetGroupMemberLister(channelMgr.ListGroupMembers)
}
}
// Load channel instances from DB.
var instanceLoader *channels.InstanceLoader
if pgStores.ChannelInstances != nil {
instanceLoader = channels.NewInstanceLoader(pgStores.ChannelInstances, pgStores.Agents, channelMgr, msgBus, pgStores.Pairing)
instanceLoader.SetProviderRegistry(providerRegistry)
instanceLoader.SetPendingCompactionConfig(cfg.Channels.PendingCompaction)
instanceLoader.SetUsageCapService(usageCapSvc)
instanceLoader.RegisterFactory(channels.TypeTelegram, telegram.FactoryWithStoresAndAudio(pgStores.Agents, pgStores.ConfigPermissions, pgStores.Teams, pgStores.SubagentTasks, pgStores.PendingMessages, audioMgr))
instanceLoader.RegisterFactory(channels.TypeDiscord, discord.FactoryWithStoresAndAudio(pgStores.Agents, pgStores.ConfigPermissions, pgStores.PendingMessages, audioMgr))
instanceLoader.RegisterFactory(channels.TypeFeishu, feishu.FactoryWithPendingStoreAndAudio(pgStores.PendingMessages, audioMgr))
instanceLoader.RegisterFactory(channels.TypeZaloOA, zalo.Factory)
instanceLoader.RegisterFactory(channels.TypeZaloPersonal, zalopersonal.FactoryWithPendingStore(pgStores.PendingMessages))
instanceLoader.RegisterFactory(channels.TypeWhatsApp, whatsapp.FactoryWithDBAudio(pgStores.DB, pgStores.PendingMessages, "pgx", audioMgr, pgStores.BuiltinTools))
instanceLoader.RegisterFactory(channels.TypeSlack, slackchannel.FactoryWithPendingStore(pgStores.PendingMessages))
instanceLoader.RegisterFactory(channels.TypeFacebook, facebook.Factory)
instanceLoader.RegisterFactory(channels.TypePancake, pancake.Factory)
// Bitrix24: factory needs the portal store + encKey injected so each
// Channel can resolve its portal on Start(). The encKey here mirrors
// the one used by pg.NewPGStores → NewPGBitrixPortalStore.
bitrixEncKey := os.Getenv("GOCLAW_ENCRYPTION_KEY")
// Use the MCP-aware factory variant so channels that opt into
// lazy per-user credential provisioning (via mcp_server_name +
// mcp_base_url in their instance config) can reach the partner's
// MCPServerStore. The MCP server authenticates each onboard call
// via the caller-supplied Bitrix access_token (the "Bitrix24
// OAuth → existing mcp_user_credentials bridge" — Bitrix-specific
// glue, not a generic MCP architecture pattern) — no shared admin
// secret is required. Channels with none of those set operate
// identically to before — the MCPStore arg is nil-safe inside the
// factory.
instanceLoader.RegisterFactory(channels.TypeBitrix24, bitrix24.FactoryWithPortalStoreAndMCP(pgStores.BitrixPortals, pgStores.MCP, bitrixEncKey))
if err := instanceLoader.LoadAll(context.Background()); err != nil {
slog.Error("failed to load channel instances from DB", "error", err)
}
// Bitrix24 portal management RPC (self-service onboarding).
// Registers bitrix.portals.list/create/get_install_url/delete methods
// on the WS router; install URL is built from the gateway's observed
// public URL via Server.PublicURLSnapshot().
if pgStores.BitrixPortals != nil {
methods.NewBitrixPortalsMethods(
pgStores.BitrixPortals,
pgStores.ChannelInstances,
server.PublicURLSnapshot().Get,
).Register(server.Router())
}
// Warm the shared Bitrix24 router with every portal row so inbound
// webhooks land on the right *Portal even before a channel instance
// is loaded for that portal. Idempotent; no-op on sqlite-lite.
if pgStores.BitrixPortals != nil {
if err := bitrix24.BootstrapPortals(context.Background(), pgStores.BitrixPortals, bitrixEncKey); err != nil {
// Surface the missing-table case loudly so an operator notices
// without having to grep logs — bitrix24 channels silently
// no-op until `goclaw migrate up` runs migration 000058.
if strings.Contains(err.Error(), "bitrix_portals") &&
(strings.Contains(err.Error(), "does not exist") || strings.Contains(err.Error(), "no such table")) {
slog.Warn("bitrix24 bootstrap skipped — bitrix_portals table missing; run `goclaw migrate up` (migration 000068) to enable Bitrix24 channels",
"err", err)
} else {
slog.Warn("bitrix24 bootstrap failed", "err", err)
}
}
}
}
// Register config-based channels as fallback when no DB instances loaded.
registerConfigChannels(cfg, channelMgr, msgBus, pgStores, instanceLoader, audioMgr)
// Register channels/instances/links/teams RPC methods
chInstancesM := wireChannelRPCMethods(server, pgStores, channelMgr, instanceLoader, agentRouter, msgBus, cfg, workspace)
// Bitrix24 orphan-bot cleaner. Fires from channel_instances delete handler
// when the channel is no longer loaded in the Manager (typical scenario:
// admin disabled the channel earlier so InstanceLoader.Reload removed it).
// Without this, deleting a disabled Bitrix24 channel would orphan the bot
// on the portal.
if pgStores.BitrixPortals != nil {
bitrixEncKey := os.Getenv("GOCLAW_ENCRYPTION_KEY")
orphanCleaner := func(ctx context.Context, tenantID uuid.UUID, cfg []byte) error {
return bitrix24.DestroyOrphanBot(ctx, pgStores.BitrixPortals, bitrixEncKey, tenantID, cfg)
}
if channelInstancesH != nil {
channelInstancesH.RegisterOrphanCleaner(channels.TypeBitrix24, orphanCleaner)
}
if chInstancesM != nil {
chInstancesM.RegisterOrphanCleaner(channels.TypeBitrix24, orphanCleaner)
}
}
// Wire channel event subscribers (cache invalidation, pairing, cascade disable)
wireChannelEventSubscribers(msgBus, server, pgStores, channelMgr, instanceLoader, pairingMethods, cfg)
// Audit log subscriber + team task event subscribers.
auditCh := deps.wireAuditSubscriber()
deps.wireEventSubscribers()
// Setup graceful shutdown
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
server.StartUpdateChecker(ctx)
sigCh := make(chan os.Signal, 1)
signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM)
// Skills directory watcher — auto-detect new/removed/modified skills at runtime.
if skillsWatcher, err := skills.NewWatcher(skillsLoader); err != nil {
slog.Warn("skills watcher unavailable", "error", err)
} else {
if err := skillsWatcher.Start(ctx); err != nil {
slog.Warn("skills watcher start failed", "error", err)
} else {
defer skillsWatcher.Stop()
}
}
// Start channels
if err := channelMgr.StartAll(ctx); err != nil {
slog.Error("failed to start channels", "error", err)
}
// Create lane-based scheduler (matching TS CommandLane pattern).
// Must be created before cron setup so cron jobs route through the scheduler.
sched := scheduler.NewScheduler(
scheduler.DefaultLanes(),
scheduler.DefaultQueueConfig(),
makeSchedulerRunFunc(agentRouter, cfg),
)
defer sched.Stop()
// Start cron + heartbeat ticker, wire wake functions and adaptive throttle.
heartbeatTicker := startCronAndHeartbeat(pgStores, server, sched, msgBus, providerRegistry, channelMgr, cfg, heartbeatTool, heartbeatMethods)
// Subscribe to agent events for channel streaming/reaction forwarding.
deps.wireChannelStreamingSubscriber()
// Slow tool notification subscriber — direct outbound when tool exceeds adaptive threshold.
wireSlowToolNotifySubscriber(msgBus)
// Inbound message consumer setup
consumerTeamStore := pgStores.Teams
// Quota checker: enforces per-user/group request limits.
config.MergeChannelGroupQuotas(cfg)
var quotaChecker *channels.QuotaChecker
if cfg.Gateway.Quota != nil && cfg.Gateway.Quota.Enabled {
quotaChecker = channels.NewQuotaChecker(pgStores.DB, *cfg.Gateway.Quota)
defer quotaChecker.Stop()
slog.Info("channel quota enabled",
"default_hour", cfg.Gateway.Quota.Default.Hour,
"default_day", cfg.Gateway.Quota.Default.Day,
"default_week", cfg.Gateway.Quota.Default.Week,
)
}
// Register quota usage RPC.
methods.NewQuotaMethods(quotaChecker, pgStores.DB).Register(server.Router())
// API key management RPC
if pgStores.APIKeys != nil {
methods.NewAPIKeysMethods(pgStores.APIKeys).Register(server.Router())
}
// Tenant management RPC + HTTP
if pgStores.Tenants != nil {
methods.NewTenantsMethods(pgStores.Tenants, msgBus, workspace).Register(server.Router())
server.SetTenantsHandler(httpapi.NewTenantsHandler(pgStores.Tenants, msgBus, workspace))
server.Router().SetTenantStore(pgStores.Tenants)
// Permission cache for tenant membership checks. Store on deps so
// lifecycle shutdown can call Close() to stop the sweep goroutines.
permCache := cache.NewPermissionCache()
deps.permCache = permCache
msgBus.Subscribe("permission-cache", func(e bus.Event) {
if p, ok := e.Payload.(bus.CacheInvalidatePayload); ok {
permCache.HandleInvalidation(p)
}
})
server.Router().SetPermissionCache(permCache)
httpapi.InitTenantStore(pgStores.Tenants, msgBus)
httpapi.InitOwnerIDs(cfg.Gateway.OwnerIDs)
}
// Wire lifecycle: config-reload subscribers, consumer, task recovery, shutdown, server start.
deps.runLifecycle(ctx, cancel, lifecycleDeps{
sched: sched,
heartbeatTicker: heartbeatTicker,
quotaChecker: quotaChecker,
webFetchTool: webFetchTool,
ttsTool: ttsTool,
sandboxMgr: sandboxMgr,
postTurn: postTurn,
subagentMgr: subagentMgr,
consumerTeamStore: consumerTeamStore,
auditCh: auditCh,
sigCh: sigCh,
})
}
// resolveBackgroundProvider picks the LLM provider+model for background workers
// (vault enrichment, consolidation). Fallback chain:
//
// background.provider/model → agent.default_provider/model → first registered provider.
func resolveBackgroundProvider(cfg *config.Config, reg *providers.Registry) (providers.Provider, string) {
try := func(name, model string) (providers.Provider, string, bool) {
if name == "" {
return nil, "", false
}
p, err := reg.GetForTenant(providers.MasterTenantID, name)
if err != nil || p == nil {
return nil, "", false
}
if model == "" {
model = p.DefaultModel()
}
return p, model, true
}
// 1. Explicit background config
if p, m, ok := try(cfg.Gateway.BackgroundProvider, cfg.Gateway.BackgroundModel); ok {
return p, m
}
// 2. Agent default provider
if p, m, ok := try(cfg.Agents.Defaults.Provider, cfg.Agents.Defaults.Model); ok {
return p, m
}
// 3. First registered provider (legacy fallback)
if names := reg.ListForTenant(providers.MasterTenantID); len(names) > 0 {
if p, m, ok := try(names[0], ""); ok {
return p, m
}
}
return nil, ""
}