Files
ModelRouter/internal/lua/plugins.go
JianFeeeee 42764bc99e feat(plugin): AUTO 调度轨迹可见(chain_step stage)
被问"还有 auto 调度相关 stage 呢?"问出来的真实缺口。

## 问题
chainDrive 只返回 (resp, src, model, err),调用方只知道**最终哪个槽位赢了**。
遍历过程中算出来又丢掉的东西——哪些档被跳过、为什么跳过、哪些槽位硬失败、
哪档全忙——一律不可见。ChainErr 里其实有这些,但**只在全部失败时**才填,
而它是 error 返回值不是记录。于是:

    "tier 1 冷却所以降级到 tier 3"  ==  "tier 1 正常接单"

对插件而言 tier 只是个常量 -2("resolved by the chain"),信息量为零。而这
恰恰是优先级链存在的全部理由,也是"我那个贵模型为什么没被用"的答案。

## 做法(scheduler 侧零新依赖)
新增 TraceEvent / TraceSink,chainDrive 多一个可选 sink 参数:

  - TraceEvent 是本包的普通 struct,sink 是 func 参数 ⇒ **不新增 import**,
    scheduler 仍然可独立测试
  - sink 为 nil 时每次 emit 只多一次 nil 判断;没有插件的网关在 AUTO 热路径上
    零开销(gateway 的 chainTraceSink 直接返回 nil)
  - 事件是纯观测:scheduler 不基于它做任何分支,gateway 也不把它喂回路由/
    冷却/配额

四种 kind:tier_skip / slot_fail / tier_busy / selected,selected 每次成功
遍历恰好一次且是最后一步。顺序保证所有 step 在 routed 之前。

## 暴露给插件
新增 chain_step stage(逐个步骤),并在 request_end 载荷里加三个便于做报表的
字段:chain_walk(上限 12 步,防审计记录膨胀)、degraded、tier_served。

## ★ 计费口径(我按推荐的做,已写进文档,需要你确认)
**按实际服务的模型计费**:降级到 tier 3 仍按 tier 3 的价算,轨迹只作观测。
理由与 §7.5 的边界一致——插件只报表不执法,两套口径混在一起会引出"降级该不该
多收钱"这种无法从代码判断的争议。若要改成"按本该用的档计价",需要在 models
价目里允许按 tier 定价,这我没做,因为那是个产品决策。

## 计费插件同步消费
by_tier_served / skip_reasons / degraded_reqs 三个新维度。skip_reasons 的等待
时长做了归一(`no free slot within <wait>`),否则 busy-wait 文案一变就多一行。
降级次数在 request_end 里计而不是在 chain_step 里计:一次降级的请求要走多步,
按步计会重复计数。

## 判据(346 个测试全绿,新增 15 个)
  scheduler  6 个:正常路径只发一个 selected / 跳档+降级可见 / 硬失败与跳档
                严格区分(不可混为一谈,否则抖动上游看起来像空闲上游)/
                nil sink 安全 / 全失败时轨迹与 ChainErr 并存且不互相破坏 /
                空链不发事件
  gateway    1 个端到端:tier 1 全 500 → 插件收到 slot_fail(tier 1) +
                selected(tier 2),request_end 的 tier_served=2 且 degraded=true
  lua        2 个:降级计数与按实际模型计价 / 跳过原因归一聚合
  lua        1 个:chain_step 是真 stage 且顺序正确

3 个变异都红:去掉 slot_fail(3 个判据红)/ 去掉 tier_skip(1 个)/
去掉 degraded 字段(1 个)。
2026-10-02 01:03:39 +08:00

835 lines
26 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"
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"`
// 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
dir string
// script is the plugin's source, kept so a reload can rebuild its state.
script string
// 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
}
// 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
}
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 {
return &Plugins{
vm: vm,
stageFuncs: map[Stage][]hookCall{},
hookErr: newHookErrors(),
dir: dir,
}
}
// 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
// ---- 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
}
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.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 {
if p.LoadError != "" {
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{}
for _, p := range ps.plugins {
if p.LoadError != "" || p.UI == nil {
continue
}
if p.UI.Page != nil {
merged.Page = p.UI.Page
}
merged.Elements = append(merged.Elements, p.UI.Elements...)
}
ps.ui.Store(merged)
}
// 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 == "",
}
if p.LoadError != "" {
row["error"] = p.LoadError
}
if p.UI != nil {
ui := map[string]interface{}{}
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
}
// 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.
L.SetTop(0)
return nil
}
body = rest
}
}
pushGoValue(L, body)
L.SetField(plug, "state")
L.SetTop(0)
return nil
}
// 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. Caller holds ps.mu.
func (ps *Plugins) find(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
}
for _, hc := range calls {
ps.mu.RLock()
p := ps.plugins[hc.pluginIdx]
ps.mu.RUnlock()
if p == nil || p.LoadError != "" {
continue
}
out, err := ps.invoke(p, hc.fn, payload)
if err != nil {
ps.hookErr.note(stage, p.Info.Name+": "+err.Error())
continue
}
if len(out) > 0 {
for k, v := range out {
payload[k] = v
}
}
}
return payload
}
// 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
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)".
raw, err := json.Marshal(payload)
if err != nil {
return nil, err
}
var decoded interface{}
if err := json.Unmarshal(raw, &decoded); err != nil {
return nil, err
}
pushGoValue(L, decoded)
// Call takes NO function index: it invokes whatever sits directly below the
// nargs values it just pushed. Passing an index here is a compile-time no-op
// in this binding and the call lands on the argument instead
// ("attempt to call a table value").
if err := L.Call(1, 1); err != nil {
return nil, err
}
if L.GetTop() < 1 || L.IsNil(-1) {
return nil, nil
}
var out map[string]interface{}
if err := luaToJSON(L, -1, &out); err != nil {
return nil, nil // not a table: treat as "no opinion"
}
return out, nil
}