Files
omnichannel-configserver-mcp/internal/configserver/manager.go
T
2026-10-07 20:13:23 +07:00

210 lines
6.6 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package configserver
import (
"context"
"errors"
"fmt"
"log/slog"
"os"
"path/filepath"
"strconv"
"sync"
"time"
"git.totmin.ru/en2zmax/forge-toolkit/configreload"
)
// tenantFileName — имя файла конфига внутри каталога OMNI_CONFIG_DIR
// (например, при постраничной/per-call передаче конфига).
const tenantFileName = "omnichannel-mcp.json"
// Manager — одиночный потокобезопасный диспетчер процесса (ARCHITECTURE.md §2).
// Кэширует per-путь загрузчики конфига (live-reload по контент-хэшу), сессии
// (отдельно rw/ro) и семафоры на мутации. Конкурентный доступ агентов безопасен.
type Manager struct {
staticPath string
mu sync.Mutex
loaders map[string]*configreload.Loader[*Config]
sessions map[string]*Session
sems map[string]chan struct{}
}
// NewManager создаёт Manager. staticPath — путь из `-config` (основной режим);
// пусто — путь ищется в OMNI_CONFIG_DIR.
func NewManager(staticPath string) *Manager {
return &Manager{
staticPath: staticPath,
loaders: map[string]*configreload.Loader[*Config]{},
sessions: map[string]*Session{},
sems: map[string]chan struct{}{},
}
}
// ResolvePath выбирает путь конфига по приоритету: явный путь вызова →
// каталог из окружения (OMNI_CONFIG_DIR, затем legacy FORGE_TENANT_CONFIG) →
// статический `-config`. Если переменная задаёт каталог, в нём берётся
// tenantFileName.
func (m *Manager) ResolvePath(tenantPath string) string {
if tenantPath != "" {
return tenantPath
}
for _, env := range []string{"OMNI_CONFIG_DIR", "FORGE_TENANT_CONFIG"} {
if dir := os.Getenv(env); dir != "" {
return filepath.Join(dir, tenantFileName)
}
}
return m.staticPath
}
// Config загружает конфиг (live-reload) и возвращает его вместе с путём.
// Отсутствие файла — fail-closed (ErrNoConfig); битый файл с last-good —
// работаем со старым значением и логируем.
func (m *Manager) Config(tenantPath string) (*Config, string, error) {
path := m.ResolvePath(tenantPath)
if path == "" {
return nil, "", ErrNoConfig
}
loader := m.loaderFor(path)
cfg, err := loader.Get()
if err != nil {
if errors.Is(err, configreload.ErrNotFound) {
return nil, path, fmt.Errorf("%w (файл %s)", ErrNoConfig, path)
}
if cfg == nil {
return nil, path, err
}
// last-good: продолжаем на прежнем конфиге, но сигналим в лог.
slog.Warn("конфиг изменён к невалидному, работаем на последнем рабочем", "path", path, "error", err)
}
return cfg, path, nil
}
// Open собирает фасад API для вызова: резолвит конфиг и сервер, проверяет
// read_only, поднимает сессию и (для мутаций) берёт семафор конкурентности.
// Close обязателен.
func (m *Manager) Open(ctx context.Context, tenantPath, alias string, write bool) (*API, error) {
cfg, path, err := m.Config(tenantPath)
if err != nil {
return nil, err
}
srv, err := cfg.Server(alias)
if err != nil {
return nil, err
}
if write && cfg.ReadOnly {
return nil, ErrReadOnly
}
// Read-инструменты используют read-only креды, если они заданы.
sess, err := m.session(path, srv, !write)
if err != nil {
return nil, err
}
api := &API{
CfgPath: path,
Srv: *srv,
Config: cfg,
Sess: sess,
Write: write,
}
if write {
sem := m.semaphore(path, srv.Alias, srv.MaxConcurrentMutations)
select {
case sem <- struct{}{}:
api.sem, api.held = sem, true
case <-ctx.Done():
return nil, ctx.Err()
}
}
return api, nil
}
// loaderFor возвращает (и кэширует) live-reload загрузчик по пути.
func (m *Manager) loaderFor(path string) *configreload.Loader[*Config] {
m.mu.Lock()
defer m.mu.Unlock()
if l, ok := m.loaders[path]; ok {
return l
}
l := configreload.New(path, func([]byte) (*Config, error) { return LoadFile(path) })
m.loaders[path] = l
return l
}
// session возвращает (и кэширует) сессию по (путь, алиас, rw/ro).
func (m *Manager) session(path string, srv *Server, readonly bool) (*Session, error) {
key := path + "|" + srv.Alias + "|" + strconv.FormatBool(readonly)
m.mu.Lock()
defer m.mu.Unlock()
if s, ok := m.sessions[key]; ok {
return s, nil
}
s, err := NewSession(srv, readonly)
if err != nil {
return nil, err
}
m.sessions[key] = s
return s, nil
}
// semaphore возвращает (и кэширует) семафор конкурентных мутаций на сервер.
func (m *Manager) semaphore(path, alias string, n int) chan struct{} {
if n < 1 {
n = 1
}
key := path + "|" + alias
m.mu.Lock()
defer m.mu.Unlock()
ch, ok := m.sems[key]
if !ok || cap(ch) != n {
ch = make(chan struct{}, n)
m.sems[key] = ch
}
return ch
}
// Close закрывает все сессии (graceful shutdown).
func (m *Manager) Close() {
m.mu.Lock()
defer m.mu.Unlock()
m.sessions = nil
m.sems = nil
}
// API — фасад одного вызова: конфиг сервера + авторизованная сессия + лимиты.
// Методы API (см. api_*.go) соответствуют эндпоинтам config_server v1.1.0.
type API struct {
CfgPath string
Srv Server
Config *Config
Sess *Session
Write bool
sem chan struct{}
held bool
}
// Close освобождает семафор мутаций (если он был взят). Идемпотентен.
func (a *API) Close() {
if a.held {
<-a.sem
a.held = false
}
}
// Alias возвращает алиас сервера (для сообщений).
func (a *API) Alias() string { return a.Srv.Alias }
// withTimeout накладывает таймаут сервера на вызов.
func (a *API) withTimeout(ctx context.Context) (context.Context, context.CancelFunc) {
return context.WithTimeout(ctx, time.Duration(a.Srv.TimeoutSec)*time.Second)
}
// withUploadTimeout — увеличенный таймаут для загрузки фронта.
func (a *API) withUploadTimeout(ctx context.Context) (context.Context, context.CancelFunc) {
return context.WithTimeout(ctx, time.Duration(a.Srv.UploadTimeoutSec)*time.Second)
}