mirror of
https://gitcode.com/JianFeeeee/ModelRouter.git
synced 2026-10-03 23:54:06 +00:00
生产故障:pi 客户端报 `model "AUTO" returned a completed response with no content`, 重试延迟 8 秒。实测 AUTO 20 次有 2 次返回空 content,全部是 claude-opus-4-8。 根因(直连上游抓包确认):思考型模型先吐 reasoning_content,max_tokens 小到 思考阶段就把预算用完时,上游返回 200 / finish_reason=length,28 个 chunk 全是 reasoning_content、content 一片空白。runTier 只看 err == nil 就当成功返回, 客户端拿到一个空响应。 修复: - resultIsEmpty:非流式路径把「无 content、无 tool_calls、无 image」的响应当作 slot 失败继续降级。注意 ReasoningContent 不算内容——客户端要的是文本,为 另一个模型的思考阶段扣住请求比降级更糟。 - peekStream:流式路径在出现首个真实内容前缓冲 reasoning 前导,流结束仍无内容 则回报空结果,让 chainDrive 换 slot。缓冲只覆盖思考前导,拿到内容后立即 转发。tool_calls delta 算内容,agent 回合不会被误判。 - 3 条判据 + 3 个变异(恒 false / 恒 true / reasoning 算内容)全部被捕获。 同时修三个 WebUI 布局缺陷(都靠截图而非 DOM 断言发现): - 插件侧栏项只渲染图标没有标题:btn.innerHTML 只塞 pluginIconHTML(pg.icon), 与原生页的「图标 + <span>标题</span>」不一致,侧栏是一排无名图标。 - 计费维度表 8 列挤在 465px 卡片里:table-layout:fixed 把每列压到 62px, 23/80 个单元格溢出、数字互相重叠。改为 6 列(token 细分合并为 「输入(新鲜+缓存)」,细分进 title)+ table-layout:auto,实测 0/60 溢出。 - 数字列 word-break:break-all 让每个字符独占一行(USD 0.56 竖排成 U/S/D), 改 nowrap + 容器横向滚动。
1585 lines
54 KiB
Go
1585 lines
54 KiB
Go
package lua
|
|
|
|
// Plugin runtime.
|
|
//
|
|
// A plugin is a single .lua file, loaded from its own directory, that extends
|
|
// the gateway in two ways:
|
|
//
|
|
// 1. Hooks: it registers callbacks on the request pipeline's stages
|
|
// (request_start, response_end, …). A hook receives a JSON table and
|
|
// returns either nil (no opinion) or a JSON object.
|
|
// 2. UI: at boot it returns HTML/CSS/JS fragments that the kernel injects
|
|
// into the WebUI — either as a whole new page, or as an extra element on
|
|
// an existing page.
|
|
//
|
|
// Design notes that are load-bearing (each one cost something to learn):
|
|
//
|
|
// - Plugins run in a SEPARATE VM from adapters, and each plugin gets its own
|
|
// elastic pool, exactly like an adapter. Sharing one state would let a
|
|
// plugin's globals corrupt an adapter's protocol translation (or vice
|
|
// versa), and a plugin is third-party code while an adapter is core.
|
|
//
|
|
// - A plugin that errors must NEVER break request forwarding. Hook calls are
|
|
// wrapped so a plugin error is logged and the original value is returned
|
|
// unchanged. A broken plugin is a missing feature, not an outage — the same
|
|
// reason transform_stream_chunk swallows errors today.
|
|
//
|
|
// - The billing plugin therefore cannot be trusted to be the source of truth
|
|
// for anything the gateway must enforce. It reads usage off the hook
|
|
// payload and accumulates in its own state; the gateway's own quota
|
|
// accounting (internal/gateway/stats.go) stays authoritative for limits.
|
|
// Two accounting paths that disagree is worse than one that is slightly
|
|
// less featureful, so the split is explicit and documented.
|
|
|
|
import (
|
|
"encoding/json"
|
|
"fmt"
|
|
"os"
|
|
"path/filepath"
|
|
"sort"
|
|
"strings"
|
|
"sync"
|
|
"sync/atomic"
|
|
"time"
|
|
|
|
golua "github.com/aarzilli/golua/lua"
|
|
)
|
|
|
|
// Stage identifies a point in the request pipeline. Plugins may register a
|
|
// function for any stage; unknown stages are ignored at call time so an old
|
|
// plugin survives a gateway that grew new stages.
|
|
type Stage string
|
|
|
|
// The pipeline stages. The order here is the order they fire in; it is the
|
|
// contract plugins are written against.
|
|
const (
|
|
// StageRequestStart fires after the gateway has parsed and authorized a
|
|
// request but BEFORE any upstream slot is chosen. payload:
|
|
// stage, type ("chat"|"stream"|"image"), model (as requested),
|
|
// key (masked gateway key id), role, source (empty), stream (bool),
|
|
// messages_count, tools_count, ts (unix seconds).
|
|
StageRequestStart Stage = "request_start"
|
|
|
|
// StageRouted fires once a (source, model) slot has been chosen and before
|
|
// the upstream call. payload adds: source, model (the resolved one),
|
|
// tier (AUTO tier, -1 for the direct path), stream.
|
|
StageRouted Stage = "routed"
|
|
|
|
// StageChainStep fires ONCE PER STEP of an AUTO chain walk, and only on the
|
|
// AUTO path (a direct request has no chain and therefore emits nothing).
|
|
//
|
|
// This exists because StageRouted cannot express degradation: it fires once,
|
|
// after the walk, with the slot that finally won. "tier 1 was cooling so we
|
|
// dropped to tier 3" and "tier 1 served it" were indistinguishable. That
|
|
// distinction is the whole point of a priority chain, and it is what an
|
|
// operator debugging "why did my expensive model not get used" needs.
|
|
//
|
|
// payload:
|
|
// kind "tier_skip" | "slot_fail" | "tier_busy" | "selected"
|
|
// tier the AUTO tier this step belongs to (1 = highest priority)
|
|
// source / model set for slot_fail and selected
|
|
// reason human-readable cause, for tier_skip and tier_busy
|
|
// error the underlying error text, for slot_fail
|
|
// attempt 1-based slot attempt within this walk
|
|
//
|
|
// Ordering: every step precedes StageRouted, and the "selected" step is the
|
|
// last one. A plugin accumulating the walk therefore has the full picture
|
|
// by the time request_end arrives.
|
|
//
|
|
// These events are OBSERVATION ONLY — see the accounting note in
|
|
// docs/plugins.md: nothing here feeds back into routing, cooldown or quota.
|
|
StageChainStep Stage = "chain_step"
|
|
|
|
// StageRequestEnd fires exactly once per request, after the client response
|
|
// has been produced (or after a failure was recorded). payload adds:
|
|
// source, model, ok, status, latency_ms, first_byte_ms, prompt_tokens,
|
|
// completion_tokens, cache_hit_tokens, cache_miss_tokens, image_count,
|
|
// error ("" when ok).
|
|
//
|
|
// This is the stage a billing plugin should read: it carries the final
|
|
// accounting for the request, including the upstream's own usage numbers.
|
|
StageRequestEnd Stage = "request_end"
|
|
)
|
|
|
|
// AllStages is the firing order, used by the docs and by the hook listing.
|
|
var AllStages = []Stage{
|
|
StageRequestStart,
|
|
StageChainStep,
|
|
StageRouted,
|
|
StageRequestEnd,
|
|
}
|
|
|
|
// UIExtension is what a plugin contributes to the WebUI at boot.
|
|
type UIExtension struct {
|
|
// Page is a whole new sidebar entry + pane. Requires PageID and Title.
|
|
// The kernel renders Page's HTML into a pane whose id is "tab-"+PageID and
|
|
// adds a sidebar button with data-tab="<PageID>".
|
|
Page *UIPage `json:"page,omitempty"`
|
|
// Pages is the multi-page form of the same thing. A plugin with several
|
|
// distinct screens (billing totals vs. the price rules that produced
|
|
// them) would otherwise have to cram both into one pane behind tabs, or
|
|
// smuggle the second one in as a hidden element. Both make the sidebar
|
|
// lie about what the plugin contributes.
|
|
//
|
|
// Page and Pages merge: a plugin may use either or both.
|
|
Pages []*UIPage `json:"pages,omitempty"`
|
|
// Elements are snippets injected into EXISTING pages, keyed by target page
|
|
// id (e.g. "status", "keys"). Order within a target is plugin load order.
|
|
Elements []UIElement `json:"elements,omitempty"`
|
|
}
|
|
|
|
// UIPage is a plugin-provided page.
|
|
type UIPage struct {
|
|
PageID string `json:"page_id"` // kebab-case; becomes data-tab and #tab-<id>
|
|
Title string `json:"title"` // sidebar label
|
|
Icon string `json:"icon"` // optional inline SVG or short glyph
|
|
Order int `json:"order"` // sidebar sort key (default 100)
|
|
// Mount is the page body. It may contain <script> and <style>; the kernel
|
|
// executes scripts AFTER injecting the HTML so the DOM exists, and exposes
|
|
// `pluginAPI` to them (see docs/plugins.md).
|
|
Mount string `json:"mount"`
|
|
}
|
|
|
|
// UIElement is a snippet injected into an existing page.
|
|
type UIElement struct {
|
|
// Target is the id of the pane to inject into: "status", "chat", "keys",
|
|
// "sort", "sources" or "adapters".
|
|
Target string `json:"target"`
|
|
// Anchor is where in the target pane: "top", "bottom" or "before:<sel>" /
|
|
// "after:<sel>" for a CSS selector. Empty = "bottom".
|
|
Anchor string `json:"anchor,omitempty"`
|
|
Order int `json:"order,omitempty"` // sort key within the target
|
|
Mount string `json:"mount"`
|
|
}
|
|
|
|
// PluginInfo is the manifest a plugin declares.
|
|
type PluginInfo struct {
|
|
Name string `json:"name"`
|
|
Version string `json:"version"`
|
|
Description string `json:"description,omitempty"`
|
|
Author string `json:"author,omitempty"`
|
|
}
|
|
|
|
// Plugin is one loaded plugin.
|
|
type Plugin struct {
|
|
Info PluginInfo
|
|
Hooks map[Stage]string // stage -> exported function name
|
|
UI *UIExtension
|
|
// LoadError is non-empty when the plugin failed to compile or register. Such
|
|
// a plugin is listed in the UI with its error but is never called.
|
|
LoadError string
|
|
// Disabled marks a plugin the operator switched off. It stays on disk and
|
|
// keeps its name, hooks and UI declared (so the UI can show it greyed out
|
|
// and report what it WOULD contribute), but Fire never calls it and its
|
|
// UI extension is excluded from the inject payload.
|
|
//
|
|
// Disabling is deliberately NOT deleting: a plugin that breaks a live
|
|
// gateway is often one line away from being fixed, and an operator needs a
|
|
// way to take it out of the request path without losing it. That is the same
|
|
// reasoning as a systemd unit being masked rather than removed.
|
|
Disabled bool
|
|
dir string
|
|
// script is the plugin's source, kept so a reload can rebuild its state.
|
|
script string
|
|
// Builtin marks a plugin that shipped with the gateway. The UI shows it as
|
|
// such so an operator can tell "example I can read" from "mine", and a
|
|
// delete of a builtin is allowed but re-seeds on a fresh plugin dir.
|
|
Builtin bool
|
|
|
|
// state is the plugin's SINGLE authoritative Lua state, guarded by mu.
|
|
//
|
|
// A plugin deliberately does NOT use the adapter's elastic pool. Adapters
|
|
// are per-worker stateless transforms, so N independent states are correct
|
|
// (and necessary) for them. A plugin's hook, however, typically ACCUMULATES
|
|
// into `plugin.state` — the billing plugin's totals live there — so two
|
|
// states would mean two divergent sets of totals, and whichever worker a
|
|
// hook happened to get would see a different number. That was a real bug
|
|
// found by a test: SetState wrote prices to one worker, the hook then ran on
|
|
// another and priced everything at zero.
|
|
//
|
|
// A single state is sufficient because every entry point (Fire, State,
|
|
// SetState) serializes on mu, and a plugin hook is a short synchronous call.
|
|
// A hook that blocks for seconds would stall every other plugin's hook,
|
|
// which is the real cost — documented as a rule in docs/plugins.md.
|
|
mu sync.Mutex
|
|
state *worker
|
|
|
|
// persisted marks this plugin's state as having been restored from disk.
|
|
// It gates the first save: without it, a freshly loaded plugin whose state
|
|
// is still the compiled-in default would immediately overwrite the file it
|
|
// was supposed to inherit from.
|
|
persisted bool
|
|
}
|
|
|
|
// Plugins is the loaded plugin set, owned by the VM.
|
|
type Plugins struct {
|
|
mu sync.RWMutex
|
|
vm *VM
|
|
dir string // plugin directory; "" disables plugin loading entirely
|
|
plugins []*Plugin // load order; the index is the stable plugin id
|
|
// ui caches the merged UI extensions so the boot payload is computed once.
|
|
ui atomic.Pointer[UIExtension]
|
|
// stageFuncs is the precomputed stage -> []hookCall, in plugin load order.
|
|
stageFuncs map[Stage][]hookCall
|
|
// hookErr records per-stage plugin failures so a silently broken plugin is
|
|
// visible in /api/status rather than merely missing.
|
|
hookErr *hookErrors
|
|
// logf, when set, receives persistence warnings. It is a field rather than a
|
|
// direct log call so the plugin package stays free of a logging dependency
|
|
// and tests can capture the warnings.
|
|
logf func(format string, args ...interface{})
|
|
// stateDirDisabled turns persistence off. Used by tests that assert state
|
|
// starts empty, and by an embedder that has no writable plugin dir.
|
|
stateDirDisabled bool
|
|
// save coalesces state writes into one background goroutine.
|
|
//
|
|
// A hook must NOT write synchronously: the billing plugin mutates state on
|
|
// every single request, and serializing the whole state per request would put
|
|
// a file write (plus a full JSON encode) on the hot path — measured at
|
|
// 14.6us per hook call already, a write would dominate it. Instead hooks mark
|
|
// the plugin dirty and this saver flushes, so N requests between two flushes
|
|
// cost one write.
|
|
save *stateSaver
|
|
}
|
|
|
|
// stateSaver coalesces state writes: one goroutine, a minimum interval between
|
|
// flushes, and a dirty set. Under load the number of writes is bounded by the
|
|
// timer rather than by the request rate.
|
|
type stateSaver struct {
|
|
mu sync.Mutex
|
|
ps *Plugins
|
|
dirty map[string]bool
|
|
wake chan struct{}
|
|
stop chan struct{}
|
|
done chan struct{}
|
|
// wakeDelay is how long a burst suppresses the next tick. A separate field
|
|
// (not the constant 2s) so tests can shorten it: Close() waits for the saver
|
|
// goroutine, so a test with the production interval pays that interval on
|
|
// every teardown.
|
|
wakeDelay time.Duration
|
|
once sync.Once
|
|
// interval is the minimum gap between flushes.
|
|
interval time.Duration
|
|
// writes counts completed file writes; tests read it to assert that N
|
|
// mutations did NOT become N writes.
|
|
writes atomic.Int64
|
|
}
|
|
|
|
// setWakeDelayForTest shortens the debounce so tests do not pay it on teardown.
|
|
func (s *stateSaver) setWakeDelayForTest(d time.Duration) { s.wakeDelay = d }
|
|
|
|
func newStateSaver(ps *Plugins) *stateSaver {
|
|
s := &stateSaver{
|
|
ps: ps,
|
|
dirty: map[string]bool{},
|
|
wake: make(chan struct{}, 1),
|
|
stop: make(chan struct{}),
|
|
done: make(chan struct{}),
|
|
interval: 2 * time.Second,
|
|
wakeDelay: 2 * time.Second,
|
|
}
|
|
go s.loop()
|
|
return s
|
|
}
|
|
|
|
// markDirty records that a plugin's state changed. Never blocks: the channel is
|
|
// buffered, and a full buffer means a flush is already pending, which is exactly
|
|
// the state we want.
|
|
func (s *stateSaver) markDirty(name string) {
|
|
if s == nil {
|
|
return
|
|
}
|
|
s.mu.Lock()
|
|
s.dirty[name] = true
|
|
s.mu.Unlock()
|
|
select {
|
|
case s.wake <- struct{}{}:
|
|
default:
|
|
}
|
|
}
|
|
|
|
func (s *stateSaver) loop() {
|
|
defer close(s.done)
|
|
t := time.NewTicker(s.interval)
|
|
defer t.Stop()
|
|
for {
|
|
select {
|
|
case <-s.stop:
|
|
// Final flush so a clean shutdown does not lose the tail — losing
|
|
// the last interval of spend is the same bug in a smaller window.
|
|
s.flush()
|
|
return
|
|
case <-t.C:
|
|
s.flush()
|
|
case <-s.wake:
|
|
// Debounce a burst of marks into one write.
|
|
time.Sleep(s.wakeDelay)
|
|
s.flush()
|
|
}
|
|
}
|
|
}
|
|
|
|
// flush writes every dirty plugin's state.
|
|
func (s *stateSaver) flush() {
|
|
if s == nil {
|
|
return
|
|
}
|
|
s.mu.Lock()
|
|
if len(s.dirty) == 0 {
|
|
s.mu.Unlock()
|
|
return
|
|
}
|
|
names := make([]string, 0, len(s.dirty))
|
|
for n := range s.dirty {
|
|
names = append(names, n)
|
|
}
|
|
s.dirty = map[string]bool{}
|
|
s.mu.Unlock()
|
|
|
|
// Snapshot + write here, off the request path. This is the ONLY place the
|
|
// plugin's Lua tables are walked for persistence.
|
|
for _, n := range names {
|
|
s.ps.persistByName(n)
|
|
s.writes.Add(1)
|
|
}
|
|
}
|
|
|
|
// Close stops the saver AFTER its final flush completes.
|
|
//
|
|
// It must WAIT for that flush, not just signal it. Signalling and returning
|
|
// leaves the write to a goroutine that the caller is about to tear down
|
|
// (vm.Stop() frees the plugin states; the test process is exiting), so the
|
|
// final state is silently lost — which is the very defect persistence was added
|
|
// to fix. Close is on the shutdown path, where a few milliseconds is free.
|
|
func (s *stateSaver) Close() {
|
|
if s == nil {
|
|
return
|
|
}
|
|
s.once.Do(func() { close(s.stop) })
|
|
<-s.done
|
|
}
|
|
|
|
type hookCall struct {
|
|
pluginIdx int
|
|
plugin string
|
|
fn string
|
|
}
|
|
|
|
// hookErrors counts hook failures per stage, surfaced in /api/status so a
|
|
// silently broken plugin is visible instead of just missing.
|
|
type hookErrors struct {
|
|
mu sync.Mutex
|
|
counts map[Stage]int
|
|
last map[Stage]string
|
|
}
|
|
|
|
func newHookErrors() *hookErrors {
|
|
return &hookErrors{counts: map[Stage]int{}, last: map[Stage]string{}}
|
|
}
|
|
|
|
func (h *hookErrors) note(s Stage, msg string) {
|
|
h.mu.Lock()
|
|
h.counts[s]++
|
|
if len(msg) > 200 {
|
|
msg = msg[:200]
|
|
}
|
|
h.last[s] = msg
|
|
h.mu.Unlock()
|
|
}
|
|
|
|
func (h *hookErrors) snapshot() map[string]map[string]interface{} {
|
|
h.mu.Lock()
|
|
defer h.mu.Unlock()
|
|
out := map[string]map[string]interface{}{}
|
|
for s, n := range h.counts {
|
|
if n == 0 {
|
|
continue
|
|
}
|
|
out[string(s)] = map[string]interface{}{"count": n, "last_error": h.last[s]}
|
|
}
|
|
return out
|
|
}
|
|
|
|
// pluginGlobal is the Lua global the plugin's returned table is stored under,
|
|
// mirroring adapterGlobal for adapters.
|
|
const pluginGlobal = "__llmsproxy_plugin"
|
|
|
|
// NewPlugins creates the plugin registry for a VM. dir is the plugin directory;
|
|
// a missing directory is not an error (plugins are optional).
|
|
func NewPlugins(vm *VM, dir string) *Plugins {
|
|
ps := &Plugins{
|
|
vm: vm,
|
|
stageFuncs: map[Stage][]hookCall{},
|
|
hookErr: newHookErrors(),
|
|
dir: dir,
|
|
}
|
|
// The saver is started even with no dir: a gateway can be given a plugin dir
|
|
// later, and a saver that only exists when dir != "" would silently never
|
|
// flush. markDirtyLocked and flush both no-op when the dir is empty.
|
|
ps.save = newStateSaver(ps)
|
|
return ps
|
|
}
|
|
|
|
// Close stops the background state saver, flushing once more first.
|
|
//
|
|
// The gateway MUST call this on shutdown. Without it the last flush interval's
|
|
// accumulation is lost, which is the same "restart loses the books" defect this
|
|
// persistence exists to fix, just scoped to a few seconds instead of the whole
|
|
// process lifetime.
|
|
// DisableStatePersistence turns persistence off (tests/benchmarks only).
|
|
func (ps *Plugins) DisableStatePersistence() { ps.stateDirDisabled = true }
|
|
|
|
func (ps *Plugins) Close() {
|
|
if ps == nil {
|
|
return
|
|
}
|
|
ps.save.Close()
|
|
}
|
|
|
|
// LoadDir loads every .lua file in dir as a plugin. Files are loaded in
|
|
// lexical order so a plugin's UI order is stable across restarts.
|
|
//
|
|
// A plugin that fails to compile or register is NOT fatal: it is kept with its
|
|
// LoadError so the UI can show it, and it is never called. This is the same
|
|
// posture as adapters, except adapters are core while plugins are not.
|
|
func (ps *Plugins) LoadDir() error {
|
|
if ps.dir == "" {
|
|
return nil
|
|
}
|
|
if _, err := os.Stat(ps.dir); os.IsNotExist(err) {
|
|
return nil
|
|
}
|
|
entries, err := os.ReadDir(ps.dir)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
var names []string
|
|
for _, e := range entries {
|
|
if e.IsDir() || !strings.HasSuffix(e.Name(), ".lua") {
|
|
continue
|
|
}
|
|
names = append(names, e.Name())
|
|
}
|
|
sort.Strings(names)
|
|
for _, n := range names {
|
|
path := filepath.Join(ps.dir, n)
|
|
code, err := os.ReadFile(path)
|
|
if err != nil {
|
|
continue
|
|
}
|
|
if err := ps.LoadSource(strings.TrimSuffix(n, ".lua"), string(code)); err != nil {
|
|
// LoadSource records the error on the plugin itself; keep going so
|
|
// one bad plugin does not stop the others from loading.
|
|
continue
|
|
}
|
|
}
|
|
ps.rebuild()
|
|
return nil
|
|
}
|
|
|
|
// SeedBundled writes the plugins shipped with the gateway into dir when the
|
|
// directory does not exist yet, mirroring the adapter seeding rule: once the
|
|
// directory exists it is authoritative, so deleting or editing a shipped plugin
|
|
// is a real action that survives restarts.
|
|
func (ps *Plugins) SeedBundled() error {
|
|
if ps.dir == "" {
|
|
return nil
|
|
}
|
|
if _, err := os.Stat(ps.dir); err == nil {
|
|
return nil
|
|
} else if !os.IsNotExist(err) {
|
|
return err
|
|
}
|
|
return writeBundledPlugins(ps.dir)
|
|
}
|
|
|
|
// bootPluginState compiles one plugin into a fresh Lua state and stores its
|
|
// returned table under pluginGlobal. It is the plugin counterpart of
|
|
// adapterPool.boot, minus the pooling: a plugin keeps exactly one state.
|
|
func bootPluginState(code, name string) (*worker, error) {
|
|
L := golua.NewState()
|
|
L.OpenLibs()
|
|
setupGlobals(L)
|
|
if err := L.DoString(code); err != nil {
|
|
L.Close()
|
|
return nil, fmt.Errorf("compile plugin %s: %w", name, err)
|
|
}
|
|
if L.Type(-1) != golua.LUA_TTABLE {
|
|
L.Close()
|
|
return nil, fmt.Errorf("plugin %s must return a table", name)
|
|
}
|
|
// SetGlobal POPS, so this is the only global it can set (see vm.go boot()).
|
|
L.SetGlobal(pluginGlobal)
|
|
L.SetTop(0)
|
|
return &worker{L: L}, nil
|
|
}
|
|
|
|
// LoadSource loads one plugin from source text. name is the plugin id (the
|
|
// file's base name). It returns an error only for conditions the caller should
|
|
// see; a plugin that merely registers nothing is not an error.
|
|
func (ps *Plugins) LoadSource(name, code string) error {
|
|
if name == "" {
|
|
return fmt.Errorf("plugin name required")
|
|
}
|
|
if err := os.MkdirAll(ps.dir, 0755); err != nil && ps.dir != "" {
|
|
return err
|
|
}
|
|
p := &Plugin{
|
|
Info: PluginInfo{Name: name},
|
|
Hooks: map[Stage]string{},
|
|
dir: ps.dir,
|
|
}
|
|
|
|
// A plugin runs in its OWN Lua state, separate from every adapter and every
|
|
// other plugin, so an error in one cannot corrupt another. The returned
|
|
// table is stored under pluginGlobal, exactly like adapterGlobal.
|
|
//
|
|
// One state, held for the plugin's lifetime (see Plugin.state): a hook that
|
|
// accumulates into plugin.state would otherwise split its totals across
|
|
// whichever worker it happened to run on.
|
|
w, err := bootPluginState(code, name)
|
|
if err != nil {
|
|
p.LoadError = err.Error()
|
|
ps.append(p)
|
|
ps.rebuild()
|
|
return err
|
|
}
|
|
p.script = code
|
|
p.state = w
|
|
// Builtin = the shipped source is byte-identical to what this file was
|
|
// loaded from. That is stronger than "the name matches a bundled plugin": an
|
|
// operator who EDITED billing.lua must not be told their copy is builtin,
|
|
// or the UI would offer to overwrite their changes.
|
|
if orig, err := ReadBundledPlugin(name); err == nil && orig == code {
|
|
p.Builtin = true
|
|
}
|
|
|
|
// ---- manifest ----
|
|
w.L.GetGlobal(pluginGlobal)
|
|
if !w.L.IsNil(-1) {
|
|
w.L.GetField(-1, "name")
|
|
if s := w.L.ToString(-1); s != "" {
|
|
p.Info.Name = s
|
|
}
|
|
w.L.SetTop(-2)
|
|
w.L.GetField(-1, "version")
|
|
if s := w.L.ToString(-1); s != "" {
|
|
p.Info.Version = s
|
|
}
|
|
w.L.SetTop(-2)
|
|
w.L.GetField(-1, "description")
|
|
if s := w.L.ToString(-1); s != "" {
|
|
p.Info.Description = s
|
|
}
|
|
w.L.SetTop(-2)
|
|
w.L.GetField(-1, "author")
|
|
if s := w.L.ToString(-1); s != "" {
|
|
p.Info.Author = s
|
|
}
|
|
w.L.SetTop(-2)
|
|
}
|
|
w.L.SetTop(0)
|
|
|
|
// ---- hooks ----
|
|
for _, st := range AllStages {
|
|
if fn, ok := pluginHookName(w.L, string(st)); ok {
|
|
p.Hooks[st] = fn
|
|
}
|
|
}
|
|
|
|
// ---- UI ----
|
|
if ui, ok := readPluginUI(w.L); ok {
|
|
p.UI = ui
|
|
}
|
|
|
|
// Restore persisted state AFTER the plugin compiled and registered its
|
|
// hooks, so the restore overwrites the compiled-in defaults instead of being
|
|
// overwritten by them. Doing it earlier would mean a restart resets the
|
|
// totals back to whatever the .lua source initialises them to.
|
|
ps.restore(p)
|
|
|
|
ps.append(p)
|
|
// Rebuild here rather than only in LoadDir: LoadSource is also the single-
|
|
// plugin entry point (the WebUI upload path), and a caller that loads one
|
|
// plugin and immediately fires a stage must not silently get nothing.
|
|
ps.rebuild()
|
|
return nil
|
|
}
|
|
|
|
// pluginHookName returns the exported function name a plugin registered for a
|
|
// stage. A plugin registers either `hooks = {request_end = "on_end"}` or a
|
|
// direct `request_end = function(...) end` on the returned table; both forms are
|
|
// accepted because the table form keeps the manifest tidy while the direct form
|
|
// is shorter for a single-hook plugin.
|
|
//
|
|
// For the anonymous-function forms the value is re-keyed onto the plugin table
|
|
// under a synthetic name so the hot path can fetch every hook by name.
|
|
//
|
|
// STACK DISCIPLINE (this binding aborts the PROCESS on a bad index — SIGABRT,
|
|
// not a Go panic, so nothing can recover it):
|
|
//
|
|
// - GetField/SetField take the table by ABSOLUTE index, so the plugin table's
|
|
// index must be re-read after every SetTop, since SetTop(0) invalidates it.
|
|
// - Therefore each form re-pushes the plugin table and re-reads its index,
|
|
// instead of caching one index across a reset. Getting this wrong was a
|
|
// real crash found by running the test, not by reading the code.
|
|
func pluginHookName(L *golua.State, stage string) (string, bool) {
|
|
L.SetTop(0)
|
|
L.GetGlobal(pluginGlobal)
|
|
if L.IsNil(-1) {
|
|
L.SetTop(0)
|
|
return "", false
|
|
}
|
|
plug := L.GetTop()
|
|
|
|
// form 1: hooks = { request_end = "fn" } (or = function)
|
|
L.GetField(plug, "hooks")
|
|
if L.Type(-1) != golua.LUA_TNIL {
|
|
hooksIdx := L.GetTop()
|
|
L.GetField(hooksIdx, stage)
|
|
switch L.Type(-1) {
|
|
case golua.LUA_TSTRING:
|
|
name := L.ToString(-1)
|
|
L.SetTop(0)
|
|
if name != "" {
|
|
return name, true
|
|
}
|
|
return "", false
|
|
case golua.LUA_TFUNCTION:
|
|
// An anonymous function: key it onto the plugin table under a stable
|
|
// per-stage name so invoke() can fetch it like any other hook.
|
|
name := "__hook_" + stage
|
|
L.SetField(plug, name)
|
|
L.SetTop(0)
|
|
return name, true
|
|
}
|
|
L.SetTop(0)
|
|
}
|
|
// form 2: request_end = function(...) end directly on the table.
|
|
// Re-push and re-read the index: SetTop(0) above invalidated `plug`.
|
|
L.GetGlobal(pluginGlobal)
|
|
if L.IsNil(-1) {
|
|
L.SetTop(0)
|
|
return "", false
|
|
}
|
|
plug = L.GetTop()
|
|
L.GetField(plug, stage)
|
|
if L.Type(-1) == golua.LUA_TFUNCTION {
|
|
name := "__hook_" + stage
|
|
L.SetField(plug, name)
|
|
L.SetTop(0)
|
|
return name, true
|
|
}
|
|
L.SetTop(0)
|
|
return "", false
|
|
}
|
|
|
|
// readPluginUI reads the optional ui extension block.
|
|
func readPluginUI(L *golua.State) (*UIExtension, bool) {
|
|
L.GetGlobal(pluginGlobal)
|
|
if L.IsNil(-1) {
|
|
return nil, false
|
|
}
|
|
L.GetField(-1, "ui")
|
|
if L.Type(-1) == golua.LUA_TNIL {
|
|
L.SetTop(0)
|
|
return nil, false
|
|
}
|
|
var ui UIExtension
|
|
if err := luaToJSON(L, -1, &ui); err != nil {
|
|
L.SetTop(0)
|
|
return nil, false
|
|
}
|
|
L.SetTop(0)
|
|
if ui.Page == nil && len(ui.Pages) == 0 && len(ui.Elements) == 0 {
|
|
return nil, false
|
|
}
|
|
return &ui, true
|
|
}
|
|
|
|
func (ps *Plugins) append(p *Plugin) {
|
|
ps.mu.Lock()
|
|
ps.plugins = append(ps.plugins, p)
|
|
ps.mu.Unlock()
|
|
}
|
|
|
|
// rebuild recomputes the stage dispatch table and the merged UI payload. It is
|
|
// called after any load so the hot path (Fire) is a slice walk with no map
|
|
// lookups or locking beyond one RLock.
|
|
func (ps *Plugins) rebuild() {
|
|
ps.mu.Lock()
|
|
defer ps.mu.Unlock()
|
|
stageFuncs := map[Stage][]hookCall{}
|
|
for i, p := range ps.plugins {
|
|
// Disabled plugins are excluded from BOTH the dispatch table and the
|
|
// merged UI below. Including them in the UI would render a page whose
|
|
// refresh calls go to a plugin that is never consulted.
|
|
if p.LoadError != "" || p.Disabled {
|
|
continue
|
|
}
|
|
for _, st := range AllStages {
|
|
if fn, ok := p.Hooks[st]; ok {
|
|
stageFuncs[st] = append(stageFuncs[st], hookCall{pluginIdx: i, plugin: p.Info.Name, fn: fn})
|
|
}
|
|
}
|
|
}
|
|
ps.stageFuncs = stageFuncs
|
|
|
|
merged := &UIExtension{}
|
|
seenPage := map[string]bool{}
|
|
for _, p := range ps.plugins {
|
|
if p.LoadError != "" || p.Disabled || p.UI == nil {
|
|
continue
|
|
}
|
|
// The single `page` field is folded into the same list as `pages`.
|
|
// Assigning it to its own slot meant the LAST plugin to declare a page
|
|
// silently replaced every earlier one — a second plugin contributing a
|
|
// page erased the first from the sidebar with no error. One list, one
|
|
// rule: first writer wins per page_id (a duplicate id would collide in
|
|
// the DOM, and the plugin that loaded first is the better answer than
|
|
// whichever happened to load last).
|
|
var declared []*UIPage
|
|
if p.UI.Page != nil {
|
|
declared = append(declared, p.UI.Page)
|
|
}
|
|
declared = append(declared, p.UI.Pages...)
|
|
for _, pg := range declared {
|
|
if pg == nil || seenPage[pg.PageID] {
|
|
continue
|
|
}
|
|
seenPage[pg.PageID] = true
|
|
merged.Pages = append(merged.Pages, pg)
|
|
}
|
|
merged.Elements = append(merged.Elements, p.UI.Elements...)
|
|
}
|
|
// Stable order by declared Order then page id, so the sidebar does not
|
|
// reshuffle when a plugin is reloaded.
|
|
sort.SliceStable(merged.Pages, func(i, j int) bool {
|
|
if merged.Pages[i].Order != merged.Pages[j].Order {
|
|
return merged.Pages[i].Order < merged.Pages[j].Order
|
|
}
|
|
return merged.Pages[i].PageID < merged.Pages[j].PageID
|
|
})
|
|
ps.ui.Store(merged)
|
|
}
|
|
|
|
// DiskEntry is one .lua file in the plugin directory, whether or not it loaded.
|
|
// The management UI lists these rather than only the loaded set, so an operator
|
|
// can see (and fix) a plugin that failed to compile instead of finding it
|
|
// missing from the list.
|
|
type DiskEntry struct {
|
|
Name string `json:"name"`
|
|
Path string `json:"path"`
|
|
Size int64 `json:"size"`
|
|
Loaded bool `json:"loaded"`
|
|
Disabled bool `json:"disabled"`
|
|
Builtin bool `json:"builtin"`
|
|
Error string `json:"error,omitempty"`
|
|
// Version/Description are read from the loaded plugin when available.
|
|
Version string `json:"version,omitempty"`
|
|
Description string `json:"description,omitempty"`
|
|
Hooks int `json:"hooks"`
|
|
}
|
|
|
|
// OnDisk lists every .lua file in the plugin directory with its load state.
|
|
func (ps *Plugins) OnDisk() []DiskEntry {
|
|
out := []DiskEntry{}
|
|
if ps.dir == "" {
|
|
return out
|
|
}
|
|
entries, err := os.ReadDir(ps.dir)
|
|
if err != nil {
|
|
return out
|
|
}
|
|
for _, e := range entries {
|
|
if e.IsDir() || !strings.HasSuffix(e.Name(), ".lua") {
|
|
continue
|
|
}
|
|
name := strings.TrimSuffix(e.Name(), ".lua")
|
|
info, _ := e.Info()
|
|
de := DiskEntry{Name: name, Path: e.Name()}
|
|
if info != nil {
|
|
de.Size = info.Size()
|
|
}
|
|
ps.mu.RLock()
|
|
if p := ps.findLocked(name); p != nil {
|
|
de.Loaded = p.LoadError == ""
|
|
de.Disabled = p.Disabled
|
|
de.Builtin = p.Builtin
|
|
de.Error = p.LoadError
|
|
de.Version = p.Info.Version
|
|
de.Description = p.Info.Description
|
|
de.Hooks = len(p.Hooks)
|
|
} else {
|
|
// On disk but not in the running set: either it failed so badly
|
|
// that LoadSource never produced a record, or the dir was written
|
|
// after startup. Mark it by comparing with the bundled source.
|
|
if code, err := os.ReadFile(filepath.Join(ps.dir, e.Name())); err == nil {
|
|
if orig, err := ReadBundledPlugin(name); err == nil && orig == string(code) {
|
|
de.Builtin = true
|
|
}
|
|
}
|
|
}
|
|
ps.mu.RUnlock()
|
|
out = append(out, de)
|
|
}
|
|
sort.Slice(out, func(i, j int) bool {
|
|
// builtins first, then alphabetical: the example plugin an operator
|
|
// is most likely to want to read should not be buried under whatever
|
|
// they installed most recently.
|
|
if out[i].Builtin != out[j].Builtin {
|
|
return out[i].Builtin
|
|
}
|
|
return out[i].Name < out[j].Name
|
|
})
|
|
return out
|
|
}
|
|
|
|
// Count returns how many plugins loaded (including ones with LoadError).
|
|
func (ps *Plugins) Count() int {
|
|
ps.mu.RLock()
|
|
defer ps.mu.RUnlock()
|
|
return len(ps.plugins)
|
|
}
|
|
|
|
// List returns a JSON-friendly view of every loaded plugin, for the status API.
|
|
func (ps *Plugins) List() []map[string]interface{} {
|
|
ps.mu.RLock()
|
|
defer ps.mu.RUnlock()
|
|
out := make([]map[string]interface{}, 0, len(ps.plugins))
|
|
for _, p := range ps.plugins {
|
|
stages := make([]string, 0, len(p.Hooks))
|
|
for _, st := range AllStages {
|
|
if _, ok := p.Hooks[st]; ok {
|
|
stages = append(stages, string(st))
|
|
}
|
|
}
|
|
row := map[string]interface{}{
|
|
"name": p.Info.Name,
|
|
"version": p.Info.Version,
|
|
"description": p.Info.Description,
|
|
"author": p.Info.Author,
|
|
"hooks": stages,
|
|
"loaded": p.LoadError == "",
|
|
"disabled": p.Disabled,
|
|
"builtin": p.Builtin,
|
|
}
|
|
if p.LoadError != "" {
|
|
row["error"] = p.LoadError
|
|
}
|
|
if p.UI != nil {
|
|
ui := map[string]interface{}{}
|
|
// Every contributed page, listed under `pages`, INCLUDING the single
|
|
// `page` one. Two keys for the same thing looks redundant until you
|
|
// try to render the count: the management panel asked for `pages`,
|
|
// got an empty list for a plugin that plainly has a page, and showed
|
|
// "contributes 0 pages" next to a working screen. `page` stays as the
|
|
// backwards-compatible single-page accessor.
|
|
ids := make([]string, 0, len(p.UI.Pages)+1)
|
|
if p.UI.Page != nil && p.UI.Page.PageID != "" {
|
|
ids = append(ids, p.UI.Page.PageID)
|
|
}
|
|
for _, pg := range p.UI.Pages {
|
|
if pg != nil && pg.PageID != "" {
|
|
ids = append(ids, pg.PageID)
|
|
}
|
|
}
|
|
if len(ids) > 0 {
|
|
ui["pages"] = ids
|
|
if p.UI.Page != nil {
|
|
ui["page"] = p.UI.Page.PageID
|
|
}
|
|
}
|
|
if len(p.UI.Elements) > 0 {
|
|
ui["elements"] = len(p.UI.Elements)
|
|
}
|
|
row["ui"] = ui
|
|
}
|
|
out = append(out, row)
|
|
}
|
|
return out
|
|
}
|
|
|
|
// HookErrors returns per-stage hook failure counts (empty when all is well).
|
|
func (ps *Plugins) HookErrors() map[string]map[string]interface{} {
|
|
return ps.hookErr.snapshot()
|
|
}
|
|
|
|
// UI returns the merged UI extensions to inject into the WebUI.
|
|
func (ps *Plugins) UI() *UIExtension {
|
|
if u := ps.ui.Load(); u != nil {
|
|
return u
|
|
}
|
|
return &UIExtension{}
|
|
}
|
|
|
|
// State returns a plugin's own published state. A plugin publishes it by
|
|
// setting `plugin.state = {...}` inside its hook; that is the only way a hook's
|
|
// numbers reach the WebUI, because request_end is the LAST pipeline stage and
|
|
// has no downstream consumer to hand a return value to.
|
|
//
|
|
// Returns nil when the plugin does not exist or has not published anything.
|
|
func (ps *Plugins) State(name string) interface{} {
|
|
ps.mu.RLock()
|
|
p := ps.find(name)
|
|
ps.mu.RUnlock()
|
|
if p == nil {
|
|
return nil
|
|
}
|
|
p.mu.Lock()
|
|
w := p.state
|
|
if w == nil {
|
|
p.mu.Unlock()
|
|
return nil
|
|
}
|
|
defer p.mu.Unlock()
|
|
L := w.L
|
|
L.SetTop(0)
|
|
defer L.SetTop(0)
|
|
L.GetGlobal(pluginGlobal)
|
|
if L.IsNil(-1) {
|
|
return nil
|
|
}
|
|
plug := L.GetTop()
|
|
L.GetField(plug, "state")
|
|
if L.Type(-1) == golua.LUA_TNIL {
|
|
return nil
|
|
}
|
|
var out interface{}
|
|
if err := luaToJSON(L, -1, &out); err != nil {
|
|
return nil
|
|
}
|
|
return out
|
|
}
|
|
|
|
// Prices returns the plugin's current price table (the `prices` field, kept
|
|
// separate from `state`).
|
|
//
|
|
// Exposed because the billing rules editor has to show what is actually in
|
|
// effect. Reconstructing it from config alone would be wrong the moment a
|
|
// price was injected directly through SetState or the two drifted, and an
|
|
// editor that displays a different table from the one billing uses is worse
|
|
// than no editor at all.
|
|
func (ps *Plugins) Prices(name string) interface{} {
|
|
ps.mu.RLock()
|
|
p := ps.find(name)
|
|
ps.mu.RUnlock()
|
|
if p == nil {
|
|
return nil
|
|
}
|
|
p.mu.Lock()
|
|
defer p.mu.Unlock()
|
|
if p.state == nil || p.state.L == nil {
|
|
return nil
|
|
}
|
|
return snapshotField(p.state.L, "prices")
|
|
}
|
|
|
|
// SetState replaces a plugin's published state (admin API). It is how a
|
|
// configuration change (a new price for a model) reaches the plugin without
|
|
// reloading it.
|
|
//
|
|
// CONTRACT: a `prices` key in the payload is applied to the plugin's SEPARATE
|
|
// `prices` field and STRIPPED from `state`. That split is deliberate: prices are
|
|
// configuration while state is accumulated history, and a single replaceable
|
|
// field would make a price update wipe the totals (or make the totals carry a
|
|
// stale price table). A plugin that keeps its config elsewhere can ignore the
|
|
// convention and read the whole payload from `state` instead.
|
|
func (ps *Plugins) SetState(name string, state interface{}) error {
|
|
ps.mu.RLock()
|
|
p := ps.find(name)
|
|
ps.mu.RUnlock()
|
|
if p == nil {
|
|
return fmt.Errorf("plugin %s not loaded", name)
|
|
}
|
|
p.mu.Lock()
|
|
w := p.state
|
|
if w == nil {
|
|
p.mu.Unlock()
|
|
return fmt.Errorf("plugin %s has no state", name)
|
|
}
|
|
defer p.mu.Unlock()
|
|
L := w.L
|
|
L.SetTop(0)
|
|
defer L.SetTop(0)
|
|
L.GetGlobal(pluginGlobal)
|
|
if L.IsNil(-1) {
|
|
return fmt.Errorf("plugin %s has no table", name)
|
|
}
|
|
plug := L.GetTop()
|
|
|
|
// Pull `prices` out of the payload before storing the rest as state.
|
|
body := state
|
|
if m, ok := state.(map[string]interface{}); ok {
|
|
if prices, has := m["prices"]; has {
|
|
pushGoValue(L, prices)
|
|
L.SetField(plug, "prices")
|
|
rest := make(map[string]interface{}, len(m))
|
|
for k, v := range m {
|
|
if k != "prices" {
|
|
rest[k] = v
|
|
}
|
|
}
|
|
if len(rest) == 0 {
|
|
// A prices-ONLY payload is a configuration change, not a state
|
|
// reset. Leaving `state` untouched is what makes repricing safe:
|
|
// replacing it with an empty table would silently erase every
|
|
// accumulated total, so the next request would start from zero
|
|
// and the dashboard would show a sudden drop in spend.
|
|
//
|
|
// The prices still have to be persisted here: returning before
|
|
// the write below would leave the new price table in memory
|
|
// only, and a restart would silently revert to the old prices
|
|
// while the operator believed the change took effect.
|
|
var curState, curPrices interface{}
|
|
curState = snapshotField(L, "state")
|
|
curPrices = snapshotField(L, "prices")
|
|
if ps.dir != "" && !ps.stateDirDisabled && curState != nil {
|
|
writeJSONAtomic(ps.stateFile(name), map[string]interface{}{
|
|
"version": 1,
|
|
"state": curState,
|
|
"prices": curPrices,
|
|
})
|
|
}
|
|
return nil
|
|
}
|
|
body = rest
|
|
}
|
|
}
|
|
|
|
pushGoValue(L, body)
|
|
L.SetField(plug, "state")
|
|
L.SetTop(0)
|
|
// A state replacement is exactly the kind of thing an operator restarts the
|
|
// gateway for, so it must survive the restart. The snapshot is taken from
|
|
// the values just written — the caller already holds p.mu, so calling
|
|
// persistNow here would deadlock on that same non-reentrant mutex.
|
|
writeJSONAtomic(ps.stateFile(name), map[string]interface{}{
|
|
"version": 1,
|
|
"state": body,
|
|
"prices": readPricesFromPayload(state),
|
|
})
|
|
return nil
|
|
}
|
|
|
|
// readPricesFromPayload recovers the prices table a caller passed to SetState,
|
|
// which SetState stores on the plugin (not in state). Used only for the
|
|
// persistence record, so that a restart restores configuration and history to
|
|
// the two fields they belong in rather than collapsing them.
|
|
func readPricesFromPayload(state interface{}) interface{} {
|
|
if m, ok := state.(map[string]interface{}); ok {
|
|
return m["prices"]
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// markDirty snapshots a plugin's state and hands it to the saver.
|
|
//
|
|
// It is called from the hook path, so it must not block on file I/O — that is
|
|
// the saver's job. It DOES walk the Lua tables, because that has to happen
|
|
// while the caller still holds p.mu and the VM is guaranteed alive; the flush
|
|
// (Snapshotting in the flush goroutine instead was a use-after-free: vm.Stop()
|
|
// frees the Lua states. See persistByName for why the walk is safe here.)
|
|
//
|
|
// Cost is one JSON conversion per request, which the billing plugin would pay
|
|
// anyway inside its own hook.
|
|
func (ps *Plugins) markDirtyLocked(p *Plugin) {
|
|
if ps == nil || p == nil || ps.stateDirDisabled || p.state == nil || ps.dir == "" {
|
|
return
|
|
}
|
|
// NOTE: this only records THAT the state changed. It deliberately does not
|
|
// snapshot it.
|
|
//
|
|
// Snapshotting here cost 264us per request — the snapshot walks the plugin's
|
|
// Lua tables through luaValueToGo + json.Marshal + json.Unmarshal (three
|
|
// conversions), and the billing plugin does it on every single request. That
|
|
// turned a 14.6us hook into a 385us one and put the gateway's hot path
|
|
// behind a bookkeeping step. The saver pulls the snapshot on its own schedule
|
|
// instead, so N requests between flushes cost ONE snapshot.
|
|
ps.save.markDirty(p.Info.Name)
|
|
}
|
|
|
|
// persistByName snapshots one plugin and writes its state file.
|
|
//
|
|
// It is called from the saver goroutine, which means it walks a Lua state that
|
|
// vm.Stop() may be about to free. That was a real crash: the first version had
|
|
// the saver read the tables itself and the process died with SIGSEGV inside
|
|
// golua once shutdown raced a flush.
|
|
//
|
|
// Two things make it safe now:
|
|
// - Plugins.Close() stops the saver AND waits for its final flush BEFORE the
|
|
// caller stops the VM (see Plugins.Close), so no flush can be in flight when
|
|
// the Lua states go away.
|
|
// - p.mu is held while the tables are read, so a hook cannot be mutating them
|
|
// underneath. A hook that arrives after the VM stopped cannot run either,
|
|
// because the request path is already torn down by then.
|
|
func (ps *Plugins) persistByName(name string) {
|
|
if ps == nil || ps.stateDirDisabled || ps.dir == "" {
|
|
return
|
|
}
|
|
ps.mu.RLock()
|
|
p := ps.find(name)
|
|
ps.mu.RUnlock()
|
|
if p == nil || p.state == nil {
|
|
return
|
|
}
|
|
ps.persistNow(p)
|
|
}
|
|
|
|
// ---------- state persistence ----------
|
|
//
|
|
// A plugin's state lives in the Lua VM, which dies with the process. For the
|
|
// billing plugin that means a restart silently zeroes every total — measured on
|
|
// this gateway: 218,241 requests and 26.2M prompt tokens gone after one
|
|
// `systemctl restart`. For a REPORTING plugin whose whole purpose is the number
|
|
// it accumulates, that is not a rounding error, it is the feature not working.
|
|
//
|
|
// Three rules shape this:
|
|
//
|
|
// 1. `prices` (configuration) and `state` (accumulated history) are persisted
|
|
// to SEPARATE files. Restoring them together would let a price edit look
|
|
// like a state reset, or a state restore resurrect stale prices — SetState
|
|
// already keeps them apart in memory, and disk has to agree.
|
|
// 2. Writes are atomic (temp file + rename). A crash mid-write must leave the
|
|
// previous state readable, not a truncated JSON file that fails to parse on
|
|
// the next boot and loses the total anyway.
|
|
// 3. A corrupt or unreadable state file is a WARNING, never a startup error.
|
|
// Forwarding must not depend on a plugin's bookkeeping surviving.
|
|
|
|
// stateFile is the per-plugin state path, kept beside the plugin source so an
|
|
// operator can find (and delete) it next to the plugin it belongs to.
|
|
func (ps *Plugins) stateFile(name string) string {
|
|
return filepath.Join(ps.dir, "."+name+".state.json")
|
|
}
|
|
|
|
// persistNow writes a plugin's state and prices to disk synchronously. Callers
|
|
// on the hot path must use ps.markDirtyLocked instead; this is the flush worker and
|
|
// the admin-state path, where durability matters more than latency.
|
|
//
|
|
// It is called after every
|
|
// state mutation, so it must be cheap enough not to matter: the billing plugin
|
|
// mutates on every request, and writing the whole state per request would put a
|
|
// file write on the hot path.
|
|
func (ps *Plugins) persistNow(p *Plugin) {
|
|
if ps == nil || p == nil || ps.stateDirDisabled || p.state == nil || ps.dir == "" {
|
|
return
|
|
}
|
|
p.mu.Lock()
|
|
state, prices := readStateAndPrices(p.state.L)
|
|
p.mu.Unlock()
|
|
if state == nil {
|
|
return
|
|
}
|
|
writeJSONAtomic(ps.stateFile(p.Info.Name), map[string]interface{}{
|
|
"version": 1,
|
|
"state": state,
|
|
"prices": prices,
|
|
})
|
|
}
|
|
|
|
// readStateAndPrices pulls both tables out of the Lua state. The caller holds
|
|
// p.mu. Returns nils when the plugin table is missing, which happens for a
|
|
// plugin that failed to compile.
|
|
//
|
|
// Each field is read by REBUILDING the stack from the global, because the
|
|
// conversion helper (luaToJSON) ends with L.SetTop(0) on every path — it treats
|
|
// the whole stack as its own. Reusing a saved index across two conversions
|
|
// addresses a slot that no longer exists, and this binding aborts the process
|
|
// (SIGABRT) instead of reporting a bad index. The first version did exactly
|
|
// that and crashed inside the Lua C layer on every hook call.
|
|
func readStateAndPrices(L *golua.State) (state, prices interface{}) {
|
|
state = snapshotField(L, "state")
|
|
prices = snapshotField(L, "prices")
|
|
L.SetTop(0)
|
|
return state, prices
|
|
}
|
|
|
|
// snapshotField converts plugin.<field> into Go values, leaving the stack clean.
|
|
// Caller holds p.mu; caller is responsible for the Lua state being alive.
|
|
func snapshotField(L *golua.State, field string) interface{} {
|
|
L.SetTop(0)
|
|
defer L.SetTop(0)
|
|
L.GetGlobal(pluginGlobal)
|
|
if L.IsNil(-1) {
|
|
return nil
|
|
}
|
|
L.GetField(L.GetTop(), field)
|
|
if L.Type(-1) != golua.LUA_TTABLE {
|
|
return nil
|
|
}
|
|
var out interface{}
|
|
if err := luaToJSON(L, -1, &out); err != nil {
|
|
return nil
|
|
}
|
|
return out
|
|
}
|
|
|
|
// writeJSONAtomic writes v to path via a temp file + rename, so a reader (or a
|
|
// crash) never observes a half-written file.
|
|
func writeJSONAtomic(path string, v interface{}) {
|
|
b, err := json.MarshalIndent(v, "", " ")
|
|
if err != nil {
|
|
return
|
|
}
|
|
dir := filepath.Dir(path)
|
|
if err := os.MkdirAll(dir, 0755); err != nil {
|
|
return
|
|
}
|
|
tmp, err := os.CreateTemp(dir, filepath.Base(path)+".tmp*")
|
|
if err != nil {
|
|
return
|
|
}
|
|
tmpName := tmp.Name()
|
|
if _, err := tmp.Write(b); err != nil {
|
|
tmp.Close()
|
|
os.Remove(tmpName)
|
|
return
|
|
}
|
|
if err := tmp.Close(); err != nil {
|
|
os.Remove(tmpName)
|
|
return
|
|
}
|
|
if err := os.Rename(tmpName, path); err != nil {
|
|
os.Remove(tmpName)
|
|
}
|
|
}
|
|
|
|
// restore re-applies a persisted state (and prices) onto a freshly compiled
|
|
// plugin. Called once per plugin after it compiles and registers its hooks.
|
|
//
|
|
// Ordering matters: this runs AFTER compilation so the plugin's own defaults
|
|
// exist, and it OVERWRITES them, so a restart continues from the saved totals
|
|
// rather than from the values baked into the .lua source.
|
|
func (ps *Plugins) restore(p *Plugin) {
|
|
if ps == nil || p == nil || p.state == nil || ps.stateDirDisabled {
|
|
return
|
|
}
|
|
b, err := os.ReadFile(ps.stateFile(p.Info.Name))
|
|
if err != nil {
|
|
return // no saved state yet: the compiled-in default stands
|
|
}
|
|
var saved struct {
|
|
Version int `json:"version"`
|
|
State interface{} `json:"state"`
|
|
Prices interface{} `json:"prices"`
|
|
}
|
|
if err := json.Unmarshal(b, &saved); err != nil {
|
|
// A corrupt file must not stop the gateway: the plugin keeps its
|
|
// compiled-in defaults and the operator sees a warning in the log.
|
|
ps.logf("plugin %s: ignoring unreadable state file %s: %v",
|
|
p.Info.Name, ps.stateFile(p.Info.Name), err)
|
|
return
|
|
}
|
|
p.mu.Lock()
|
|
defer p.mu.Unlock()
|
|
L := p.state.L
|
|
L.SetTop(0)
|
|
defer L.SetTop(0)
|
|
L.GetGlobal(pluginGlobal)
|
|
if L.IsNil(-1) {
|
|
return
|
|
}
|
|
plug := L.GetTop()
|
|
if saved.Prices != nil {
|
|
pushGoValue(L, saved.Prices)
|
|
L.SetField(plug, "prices")
|
|
}
|
|
if saved.State != nil {
|
|
pushGoValue(L, saved.State)
|
|
L.SetField(plug, "state")
|
|
}
|
|
p.persisted = true
|
|
}
|
|
|
|
// SetEnabled turns a plugin's dispatch on or off without touching its file.
|
|
//
|
|
// The state is on the Plugin record (not derived from disk) so a disable survives
|
|
// as long as the process lives and is trivially re-enabled; it deliberately does
|
|
// NOT persist across restarts, because a "disable" that silently outlives the
|
|
// operator's intent is its own surprise. An operator who wants it permanent
|
|
// moves the file out of the plugin dir.
|
|
func (ps *Plugins) SetEnabled(name string, enabled bool) error {
|
|
ps.mu.Lock()
|
|
p := ps.findLocked(name)
|
|
ps.mu.Unlock()
|
|
if p == nil {
|
|
return fmt.Errorf("plugin %s not loaded", name)
|
|
}
|
|
if p.LoadError != "" {
|
|
return fmt.Errorf("plugin %s failed to load (%s); fix the file before enabling it", name, p.LoadError)
|
|
}
|
|
p.mu.Lock()
|
|
p.Disabled = !enabled
|
|
p.mu.Unlock()
|
|
ps.rebuild()
|
|
return nil
|
|
}
|
|
|
|
// Enabled reports whether a plugin is currently dispatching.
|
|
func (ps *Plugins) Enabled(name string) bool {
|
|
ps.mu.RLock()
|
|
p := ps.findLocked(name)
|
|
ps.mu.RUnlock()
|
|
return p != nil && p.LoadError == "" && !p.Disabled
|
|
}
|
|
|
|
// Unload removes a plugin from the running set. Its states are closed so the
|
|
// memory goes back; a subsequent LoadSource with the same name works again.
|
|
func (ps *Plugins) Unload(name string) error {
|
|
ps.mu.Lock()
|
|
idx := -1
|
|
for i, p := range ps.plugins {
|
|
if p.Info.Name == name || strings.HasSuffix(filepath.Base(p.dir), name+".lua") {
|
|
idx = i
|
|
break
|
|
}
|
|
}
|
|
if idx < 0 {
|
|
ps.mu.Unlock()
|
|
return fmt.Errorf("plugin %s not loaded", name)
|
|
}
|
|
p := ps.plugins[idx]
|
|
ps.plugins = append(ps.plugins[:idx], ps.plugins[idx+1:]...)
|
|
ps.mu.Unlock()
|
|
p.mu.Lock()
|
|
if p.state != nil && p.state.L != nil {
|
|
p.state.L.Close()
|
|
p.state = nil
|
|
}
|
|
p.mu.Unlock()
|
|
ps.rebuild()
|
|
return nil
|
|
}
|
|
|
|
// find returns a loaded plugin by declared name, taking the read lock.
|
|
func (ps *Plugins) find(name string) *Plugin {
|
|
ps.mu.RLock()
|
|
defer ps.mu.RUnlock()
|
|
return ps.findLocked(name)
|
|
}
|
|
|
|
// findLocked returns a loaded plugin by declared name. The caller holds ps.mu.
|
|
func (ps *Plugins) findLocked(name string) *Plugin {
|
|
for _, p := range ps.plugins {
|
|
if p.Info.Name == name {
|
|
return p
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Fire calls every plugin registered for a stage, in plugin load order.
|
|
//
|
|
// A plugin may MUTATE the payload by returning a JSON object: any keys it
|
|
// returns are merged into the payload for the next hook and returned to the
|
|
// caller. Returning nil or an empty table means "no opinion". This lets a
|
|
// plugin add fields (the billing plugin adds `cost_usd`) without the gateway
|
|
// having to know about them.
|
|
//
|
|
// Errors are contained: a plugin that throws is logged against the stage and
|
|
// skipped. Forwarding never depends on plugin health.
|
|
func (ps *Plugins) Fire(stage Stage, payload map[string]interface{}) map[string]interface{} {
|
|
ps.mu.RLock()
|
|
calls := ps.stageFuncs[stage]
|
|
ps.mu.RUnlock()
|
|
if len(calls) == 0 {
|
|
return payload
|
|
}
|
|
// Resolve the plugins and drop the lock BEFORE running any hook.
|
|
//
|
|
// A hook can install, disable or remove a plugin (the admin API is reachable
|
|
// from a hook that has the key), and those paths take ps.mu for write. Taking
|
|
// a per-iteration RLock like the single-plugin case did works but serializes
|
|
// on a shared cache line under load; a snapshot plus one lookup is cheaper and
|
|
// — more importantly — it means the parallel section below never holds ps.mu,
|
|
// so a hook cannot deadlock against a concurrent reload.
|
|
type target struct {
|
|
p *Plugin
|
|
hc hookCall
|
|
}
|
|
targets := make([]target, 0, len(calls))
|
|
for _, hc := range calls {
|
|
ps.mu.RLock()
|
|
p := ps.plugins[hc.pluginIdx]
|
|
ps.mu.RUnlock()
|
|
if p == nil || p.LoadError != "" {
|
|
continue
|
|
}
|
|
targets = append(targets, target{p: p, hc: hc})
|
|
}
|
|
if len(targets) == 0 {
|
|
return payload
|
|
}
|
|
|
|
// One plugin is the overwhelmingly common case (a gateway with the bundled
|
|
// billing plugin has exactly one). Paying for goroutines there would be pure
|
|
// overhead, so it takes the direct path.
|
|
if len(targets) == 1 {
|
|
if out := ps.runHook(stage, targets[0].p, targets[0].hc.fn, payload); len(out) > 0 {
|
|
for k, v := range out {
|
|
payload[k] = v
|
|
}
|
|
}
|
|
return payload
|
|
}
|
|
|
|
// Parallel across plugins. Each hook receives the SAME payload snapshot, and
|
|
// the returned tables are merged afterwards IN PLUGIN LOAD ORDER, so the
|
|
// documented contract ("a returned table's keys are merged into the payload")
|
|
// still holds deterministically.
|
|
//
|
|
// What changes: a hook no longer sees the keys another hook just added. The
|
|
// sequential behaviour made that possible, and docs/plugins.md described it
|
|
// ("the payload is passed to the next plugin unchanged"). It was never used —
|
|
// the bundled billing plugin returns nil on every stage, with a comment saying
|
|
// nobody downstream reads it — but it IS a contract change and is called out
|
|
// there rather than left as a surprise.
|
|
//
|
|
// Why this is safe for the per-plugin lock: each plugin has its own Lua state
|
|
// and its own p.mu, so two plugins never touch the same state. What is shared
|
|
// is the payload, and it is only READ here — merging happens after every hook
|
|
// has returned, on the caller's goroutine.
|
|
outs := make([]map[string]interface{}, len(targets))
|
|
var wg sync.WaitGroup
|
|
for i := range targets {
|
|
wg.Add(1)
|
|
go func(i int) {
|
|
defer wg.Done()
|
|
// A hook that panics would otherwise take the whole process with it,
|
|
// and in the sequential version a panic could not escape Fire either.
|
|
// recover() here restores that: the plugin is skipped and noted.
|
|
defer func() {
|
|
if rec := recover(); rec != nil {
|
|
ps.hookErr.note(stage, fmt.Sprintf("%s: panic in hook: %v", targets[i].hc.plugin, rec))
|
|
}
|
|
}()
|
|
outs[i] = ps.runHook(stage, targets[i].p, targets[i].hc.fn, payload)
|
|
}(i)
|
|
}
|
|
wg.Wait()
|
|
|
|
// Merge in load order so the result does not depend on goroutine scheduling.
|
|
// A later plugin's value wins on a key collision, exactly as it did when the
|
|
// hooks ran one after another.
|
|
for _, out := range outs {
|
|
for k, v := range out {
|
|
payload[k] = v
|
|
}
|
|
}
|
|
return payload
|
|
}
|
|
|
|
// runHook invokes one hook and folds its returned table into payload. It
|
|
// contains the plugin's error handling so both the sequential and the parallel
|
|
// path behave identically on failure.
|
|
func (ps *Plugins) runHook(stage Stage, p *Plugin, fn string, payload map[string]interface{}) map[string]interface{} {
|
|
out, err := ps.invoke(p, fn, payload)
|
|
if err != nil {
|
|
ps.hookErr.note(stage, p.Info.Name+": "+err.Error())
|
|
return nil
|
|
}
|
|
return out
|
|
}
|
|
|
|
// invoke runs one plugin hook on that plugin's own state, under its pool's
|
|
// concurrency cap. The plugin's returned table is re-fetched each call because
|
|
// the pool is elastic: a plugin may have several states, and the hook function
|
|
// lives in each.
|
|
func (ps *Plugins) invoke(p *Plugin, fn string, payload map[string]interface{}) (map[string]interface{}, error) {
|
|
p.mu.Lock()
|
|
w := p.state
|
|
if w == nil {
|
|
p.mu.Unlock()
|
|
return nil, fmt.Errorf("no state for plugin %s", p.Info.Name)
|
|
}
|
|
defer p.mu.Unlock()
|
|
L := w.L
|
|
L.SetTop(0)
|
|
defer L.SetTop(0)
|
|
|
|
L.GetGlobal(pluginGlobal)
|
|
if L.IsNil(-1) {
|
|
return nil, fmt.Errorf("plugin table missing")
|
|
}
|
|
plug := L.GetTop() // absolute, so nothing below shifts
|
|
// The hook is called THROUGH a pcall guard (see hookGuardSrc) so a script
|
|
// error returns as values instead of raising into golua, whose error path
|
|
// SIGSEGVs the process. Desired stack before the payload push:
|
|
// [pluginGlobal, guard, hookfn]. Build it in that order — GetGlobal(guard)
|
|
// then GetField(hookfn) lands the function exactly above the guard, with no
|
|
// Remove/Insert juggling. (The first version reordered with Remove/Insert
|
|
// and ended up calling pluginGlobal as if it were the hook, producing
|
|
// "attempt to call a table value" and silently zeroed every total.)
|
|
L.GetGlobal(hookGuardName)
|
|
if L.IsNil(-1) {
|
|
// Guard absent (a state built before this existed): drop the nil so the
|
|
// stack is [pluginGlobal, hookfn] and the hook is called directly.
|
|
// Correct plugins still work; only their errors stop being survivable.
|
|
L.Pop(1)
|
|
}
|
|
L.GetField(plug, fn)
|
|
if !L.IsFunction(-1) {
|
|
L.SetTop(0)
|
|
return nil, fmt.Errorf("hook %s missing", fn)
|
|
}
|
|
// The hook is called with a DECODED table, not the raw JSON string: the
|
|
// adapter protocol passes JSON text to its transforms (they decode it
|
|
// themselves), but a plugin hook receives a table so it can read
|
|
// payload.model directly. Passing the string made every hook fail with
|
|
// "attempt to index local 'payload' (a string value)".
|
|
//
|
|
// Decoding is done WITHOUT JSON. json.Marshal + json.Unmarshal here cost
|
|
// 9.6us of the 14.6us hook — 67% of the call spent re-deriving a tree that
|
|
// pushGoValue walks natively anyway. plainForLua does the same conversion
|
|
// by type-switching, and falls back to JSON only for a type it does not
|
|
// model, so an exotic payload still arrives instead of vanishing.
|
|
pushGoValue(L, plainForLua(payload))
|
|
if err := L.Call(2, 2); err != nil {
|
|
return nil, err
|
|
}
|
|
// The guard returns TWO values: the hook's value, then an error string
|
|
// (nil on success). Drop the error slot unconditionally so the value is
|
|
// left on top. Peeking instead of popping was the first bug here: on a
|
|
// SUCCESSFUL call the top is nil, so "the error is non-nil" is false, the
|
|
// pop was skipped, and the code then read the (always-present) nil error as
|
|
// if it were the hook's answer. Every plugin's returned table silently
|
|
// became "no opinion" while the tests that only checked totals kept
|
|
// passing.
|
|
errMsg := ""
|
|
if L.GetTop() >= 2 {
|
|
if !L.IsNil(-1) {
|
|
errMsg = L.ToString(-1)
|
|
}
|
|
L.Pop(1) // the error slot, nil or not
|
|
}
|
|
if errMsg != "" {
|
|
// A hook that raised partway may already have mutated plugin.state, so
|
|
// mark it dirty before bailing: the accounting it managed to do is
|
|
// still real and should be persisted.
|
|
ps.markDirtyLocked(p)
|
|
return nil, fmt.Errorf("plugin %s: %s", p.Info.Name, errMsg)
|
|
}
|
|
// Read the return value FIRST, then snapshot state for persistence.
|
|
//
|
|
// The order is load-bearing. markDirtyLocked walks the Lua tables and resets
|
|
// the stack, so calling it before the return value is read destroyed the
|
|
// hook's answer — every plugin that returned a table silently became a
|
|
// plugin that "had no opinion". Reading first costs nothing and keeps the
|
|
// documented merge contract intact.
|
|
var out map[string]interface{}
|
|
if L.GetTop() < 1 || L.IsNil(-1) {
|
|
ps.markDirtyLocked(p)
|
|
return nil, nil
|
|
}
|
|
if err := luaToJSON(L, -1, &out); err != nil {
|
|
ps.markDirtyLocked(p)
|
|
return nil, nil // not a table: treat as "no opinion"
|
|
}
|
|
// A hook that ran at all may have mutated plugin.state, whether or not it
|
|
// returned anything. Marking here (not only on a returned table) is what
|
|
// makes an accumulating plugin like billing durable: its totals change on
|
|
// every call and it returns nil every time.
|
|
ps.markDirtyLocked(p)
|
|
return out, nil
|
|
}
|