210 lines
6.6 KiB
Go
210 lines
6.6 KiB
Go
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)
|
||
}
|