feat(plugin): standalone plugin hosts with Redis discovery (P2 batch 4) (#3725)

* feat(plugin): standalone plugin hosts with Redis discovery
This commit is contained in:
lyingbug
2026-09-26 14:46:26 +08:00
committed by GitHub
parent 728007ada9
commit da4ef68d9d
27 changed files with 1576 additions and 40 deletions
+15
View File
@@ -811,3 +811,18 @@ CONCURRENCY_POOL_SIZE=5
# BROWSERSKILL_INTERNAL_URL=http://10.0.0.12:8080
# 所有副本使用相同的随机密钥(至少 32 字符),从 Secret 注入,不提交到仓库。
# BROWSERSKILL_CLUSTER_SECRET=
# ========== J4. 插件运行(可选) ==========
# 代码插件(runtime.type: host)默认在每个 app 进程内运行。多机或想把 Python 插件隔开时,
# 可以另起 `WeKnora plugin-host`(compose:--profile plugin-host;helm:pluginHost.enabled)。
# plugin-host 与 app 共用数据库、Redis、存储和 SYSTEM_AES_KEY(或 JWT_SECRET),经 Redis 通告自己。
# app 自己运行哪些 kind(binary,python);none 表示全部交给 plugin-host。默认本机能跑的都跑。
# WEKNORA_PLUGIN_EMBEDDED_KINDS=none
# 不在 app 本机运行的插件(remote、plugin-host 上的)回调 Host API 的地址。
# WEKNORA_PLUGIN_HOST_API_URL=http://app:8080
# 仅 plugin-host:运行哪些 kind、监听地址、app 访问它的地址(默认 http://<hostname>:<端口>)。
# WEKNORA_PLUGIN_HOST_KINDS=python
# WEKNORA_PLUGIN_HOST_ADDR=:8081
# WEKNORA_PLUGIN_HOST_URL=http://plugin-host:8081
# Python 插件使用的解释器(默认 PATH 中的 python3)。
# WEKNORA_PLUGIN_PYTHON=python3
+3
View File
@@ -41,6 +41,9 @@ import (
)
func main() {
if code, ran := subcommand(); ran {
os.Exit(code)
}
// Set Gin mode
if os.Getenv("GIN_MODE") == "release" {
gin.SetMode(gin.ReleaseMode)
+30
View File
@@ -0,0 +1,30 @@
package main
import (
"context"
"os"
"os/signal"
"github.com/Tencent/WeKnora/internal/container"
"github.com/Tencent/WeKnora/internal/logger"
)
// runPluginHost is `WeKnora plugin-host`: a standalone plugin host that
// runs host plugins for the app nodes (see container.RunPluginHost).
func runPluginHost() int {
ctx, stop := signal.NotifyContext(context.Background(), shutdownSignals...)
defer stop()
if err := container.RunPluginHost(ctx); err != nil {
logger.Errorf(context.Background(), "[plugin-host] %v", err)
return 1
}
return 0
}
// subcommand runs a subcommand named on the command line, if any.
func subcommand() (code int, ran bool) {
if len(os.Args) > 1 && os.Args[1] == "plugin-host" {
return runPluginHost(), true
}
return 0, false
}
+65
View File
@@ -381,6 +381,11 @@ services:
- OIDC_USER_INFO_MAPPING_EMAIL=${OIDC_USER_INFO_MAPPING_EMAIL:-}
# 飞书云文档解析模式 export:导出为docx,blocks:根据飞书云文档块解析为markdown。默认为:export
- FEISHU_DOCX_PARSE_MODE=${FEISHU_DOCX_PARSE_MODE:-export}
# ========== 插件 ==========
# app 自己运行哪些 kind 的宿主插件(binary,python;none 表示全部交给 plugin-host),默认本机能跑的都跑
- WEKNORA_PLUGIN_EMBEDDED_KINDS=${WEKNORA_PLUGIN_EMBEDDED_KINDS:-}
# 不在本机运行的插件(remote、plugin-host 上的)回调 Host API 的地址;启用 plugin-host 时设为 http://app:8080
- WEKNORA_PLUGIN_HOST_API_URL=${WEKNORA_PLUGIN_HOST_API_URL:-}
depends_on:
redis:
condition: service_started
@@ -394,6 +399,66 @@ services:
extra_hosts:
- "host.docker.internal:host-gateway"
# 独立插件宿主(可选,profile: plugin-host):在单独的容器里运行宿主插件,
# 例如把 Python 插件与 app 隔开。与 app 共用镜像、数据库、Redis 和存储,
# 通过 Redis 向 app 通告自己,app 用 SYSTEM_AES_KEY 派生的密钥签名调用它。
# 启用:docker compose --profile plugin-host up -d,并在 .env 里设
# WEKNORA_PLUGIN_EMBEDDED_KINDS=none(或 binary)
# WEKNORA_PLUGIN_HOST_API_URL=http://app:8080
plugin-host:
image: wechatopenai/weknora-app:${WEKNORA_VERSION:-latest}
container_name: WeKnora-plugin-host
profiles:
- plugin-host
command: ["./WeKnora", "plugin-host"]
volumes:
# STORAGE_TYPE=local 时插件包存在这里
- data-files:/data/files
env_file:
- .env
environment:
- GIN_MODE=${GIN_MODE:-release}
- TZ=${TZ:-Asia/Shanghai}
- LOG_LEVEL=${LOG_LEVEL:-info}
- WEKNORA_PLUGIN_HOST_ADDR=:8081
- WEKNORA_PLUGIN_HOST_URL=http://plugin-host:8081
- WEKNORA_PLUGIN_HOST_KINDS=${WEKNORA_PLUGIN_HOST_KINDS:-}
- WEKNORA_PLUGIN_HOST_API_URL=${WEKNORA_PLUGIN_HOST_API_URL:-http://app:8080}
- DB_DRIVER=${DB_DRIVER:-postgres}
- DB_HOST=${DB_HOST:-postgres}
- DB_PORT=${DB_PORT:-5432}
- DB_USER=${DB_USER:-}
- DB_PASSWORD=${DB_PASSWORD:-}
- DB_NAME=${DB_NAME:-}
- STORAGE_TYPE=${STORAGE_TYPE:-local}
- LOCAL_STORAGE_BASE_DIR=${LOCAL_STORAGE_BASE_DIR:-/data/files}
- MINIO_ENDPOINT=${MINIO_ENDPOINT:-minio:9000}
- MINIO_ACCESS_KEY_ID=${MINIO_ACCESS_KEY_ID:-minioadmin}
- MINIO_SECRET_ACCESS_KEY=${MINIO_SECRET_ACCESS_KEY:-minioadmin}
- REDIS_ADDR=${REDIS_ADDR:-redis:6379}
- REDIS_USERNAME=${REDIS_USERNAME:-}
- REDIS_PASSWORD=${REDIS_PASSWORD:-}
- REDIS_DB=${REDIS_DB:-}
- WEKNORA_REDIS_NAMESPACE=${WEKNORA_REDIS_NAMESPACE:-}
- SYSTEM_AES_KEY=${SYSTEM_AES_KEY:-}
- JWT_SECRET=${JWT_SECRET:-}
- SSRF_WHITELIST=${SSRF_WHITELIST:-}
healthcheck:
test: ["CMD", "curl", "-f", "http://localhost:8081/healthz"]
interval: 30s
timeout: 10s
retries: 3
start_period: 30s
depends_on:
redis:
condition: service_started
postgres:
condition: service_healthy
networks:
- WeKnora-network
restart: unless-stopped
stop_grace_period: 75s
# Sandbox 镜像:仅用于 build/pull,非常驻。Docker 沙箱默认关闭,需
# WEKNORA_SANDBOX_DOCKER_ENABLED=true 并挂载 docker.sock 后,app 才会按会话创建容器。
sandbox:
+9
View File
@@ -126,6 +126,15 @@ spec:
# whatever Engine API this process can reach (often host root).
- name: WEKNORA_SANDBOX_DOCKER_ENABLED
value: {{ .Values.app.env.WEKNORA_SANDBOX_DOCKER_ENABLED | default "false" | quote }}
{{- if .Values.pluginHost.enabled }}
# Host plugins of these kinds still run in the app pods; the rest
# go to the plugin host pods. Plugins there reach the Host API
# through the app Service.
- name: WEKNORA_PLUGIN_EMBEDDED_KINDS
value: {{ .Values.pluginHost.appKinds | quote }}
- name: WEKNORA_PLUGIN_HOST_API_URL
value: "http://app:{{ .Values.app.service.port }}"
{{- end }}
{{- if .Values.neo4j.enabled }}
# Neo4j configuration (for GraphRAG)
# NEO4J_ENABLE 是知识图谱的唯一开关(Go 代码认它,不认 ENABLE_GRAPH_RAG)。
+157
View File
@@ -0,0 +1,157 @@
{{/*
Copyright 2025 Tencent
SPDX-License-Identifier: MIT
WeKnora standalone plugin host (WeKnora plugin-host): runs host plugins for the
app pods. Hosts announce themselves in Redis and app pods call them directly
by pod IP through a gateway signed with a key derived from SYSTEM_AES_KEY, so
no Service is needed.
*/}}
{{- if .Values.pluginHost.enabled }}
apiVersion: apps/v1
kind: Deployment
metadata:
name: {{ include "weknora.fullname" . }}-plugin-host
namespace: {{ .Release.Namespace }}
labels:
{{- include "weknora.componentLabels" (dict "component" "plugin-host" "context" .) | nindent 4 }}
spec:
replicas: {{ .Values.pluginHost.replicaCount }}
selector:
matchLabels:
{{- include "weknora.componentSelectorLabels" (dict "component" "plugin-host" "context" .) | nindent 6 }}
template:
metadata:
labels:
{{- include "weknora.componentSelectorLabels" (dict "component" "plugin-host" "context" .) | nindent 8 }}
spec:
{{- include "weknora.imagePullSecrets" . | nindent 6 }}
serviceAccountName: {{ include "weknora.serviceAccountName" . }}
# The host withdraws from Redis, then lets calls in flight finish.
terminationGracePeriodSeconds: 75
{{- with .Values.pluginHost.podSecurityContext | default .Values.global.podSecurityContext }}
securityContext:
{{- toYaml . | nindent 8 }}
{{- end }}
containers:
- name: plugin-host
image: {{ include "weknora.app.image" . }}
imagePullPolicy: {{ .Values.app.image.pullPolicy }}
args: ["./WeKnora", "plugin-host"]
{{- with .Values.pluginHost.securityContext }}
securityContext:
{{- toYaml . | nindent 12 }}
{{- end }}
ports:
- containerPort: 8081
name: gateway
protocol: TCP
env:
- name: POD_IP
valueFrom:
fieldRef:
fieldPath: status.podIP
- name: WEKNORA_PLUGIN_HOST_ADDR
value: ":8081"
- name: WEKNORA_PLUGIN_HOST_URL
value: "http://$(POD_IP):8081"
- name: WEKNORA_PLUGIN_HOST_KINDS
value: {{ .Values.pluginHost.kinds | quote }}
# Plugins granted Host API scopes call the app through its Service.
- name: WEKNORA_PLUGIN_HOST_API_URL
value: "http://app:{{ .Values.app.service.port }}"
- name: TZ
value: {{ .Values.app.env.TZ | quote }}
# The same database, Redis, storage and secrets as the app.
- name: DB_DRIVER
value: "postgres"
- name: DB_HOST
value: "postgres"
- name: DB_PORT
value: "5432"
- name: DB_USER
valueFrom:
secretKeyRef:
name: {{ include "weknora.secretName" . }}
key: DB_USER
- name: DB_PASSWORD
valueFrom:
secretKeyRef:
name: {{ include "weknora.secretName" . }}
key: DB_PASSWORD
- name: DB_NAME
valueFrom:
secretKeyRef:
name: {{ include "weknora.secretName" . }}
key: DB_NAME
- name: REDIS_ADDR
value: "redis:6379"
- name: REDIS_USERNAME
valueFrom:
secretKeyRef:
name: {{ include "weknora.secretName" . }}
key: REDIS_USERNAME
optional: true
- name: REDIS_PASSWORD
valueFrom:
secretKeyRef:
name: {{ include "weknora.secretName" . }}
key: REDIS_PASSWORD
- name: REDIS_DB
value: "0"
- name: JWT_SECRET
valueFrom:
secretKeyRef:
name: {{ include "weknora.secretName" . }}
key: JWT_SECRET
- name: SYSTEM_AES_KEY
valueFrom:
secretKeyRef:
name: {{ include "weknora.secretName" . }}
key: SYSTEM_AES_KEY
- name: STORAGE_TYPE
value: {{ .Values.app.env.STORAGE_TYPE | quote }}
- name: LOCAL_STORAGE_BASE_DIR
value: {{ .Values.app.env.LOCAL_STORAGE_BASE_DIR | quote }}
{{- with .Values.pluginHost.extraEnv }}
{{- toYaml . | nindent 12 }}
{{- end }}
volumeMounts:
# Plugin packages live in the app's storage; with STORAGE_TYPE=local
# that is this volume, so it must be ReadWriteMany.
- name: data-files
mountPath: /data/files
resources:
{{- toYaml .Values.pluginHost.resources | nindent 12 }}
livenessProbe:
httpGet:
path: /healthz
port: gateway
initialDelaySeconds: 20
periodSeconds: 10
readinessProbe:
httpGet:
path: /healthz
port: gateway
periodSeconds: 5
volumes:
- name: data-files
{{- if .Values.dataFiles.persistence.enabled }}
persistentVolumeClaim:
claimName: {{ .Values.dataFiles.persistence.existingClaim | default (printf "%s-data-files" (include "weknora.fullname" .)) }}
{{- else }}
emptyDir: {}
{{- end }}
{{- with .Values.pluginHost.nodeSelector }}
nodeSelector:
{{- toYaml . | nindent 8 }}
{{- end }}
{{- with .Values.pluginHost.affinity }}
affinity:
{{- toYaml . | nindent 8 }}
{{- end }}
{{- with .Values.pluginHost.tolerations }}
tolerations:
{{- toYaml . | nindent 8 }}
{{- end }}
{{- end }}
+39
View File
@@ -146,6 +146,45 @@ app:
# -- Affinity rules
affinity: {}
# -----------------------------------------------------------------------------
# Plugin host (WeKnora plugin-host)
# -----------------------------------------------------------------------------
# Runs installed host plugins in their own pods instead of inside every app
# pod. Uses the app image, database, Redis and storage. With
# dataFiles.persistence and STORAGE_TYPE=local the data-files PVC must be
# ReadWriteMany, since plugin packages live there.
pluginHost:
# -- Enable standalone plugin hosts
enabled: false
# -- Number of plugin host pods; each runs every plugin of its kinds
replicaCount: 1
# -- Kinds of host plugin the plugin hosts run (binary, python)
kinds: "binary,python"
# -- Kinds the app pods keep running themselves ("none" = all to plugin hosts)
appKinds: "none"
resources:
requests:
cpu: 100m
memory: 256Mi
limits:
cpu: "2"
memory: 2Gi
podSecurityContext: {}
securityContext:
allowPrivilegeEscalation: false
# -- Additional environment variables (e.g. storage credentials)
extraEnv: []
nodeSelector: {}
tolerations: []
affinity: {}
# -----------------------------------------------------------------------------
# Frontend (Web UI)
# -----------------------------------------------------------------------------
+3 -2
View File
@@ -74,7 +74,6 @@ import (
"github.com/Tencent/WeKnora/internal/plugin/activate"
pluginbuiltin "github.com/Tencent/WeKnora/internal/plugin/builtin"
plugindriver "github.com/Tencent/WeKnora/internal/plugin/driver"
pluginhost "github.com/Tencent/WeKnora/internal/plugin/host"
pluginmanifest "github.com/Tencent/WeKnora/internal/plugin/manifest"
"github.com/Tencent/WeKnora/internal/plugin/reconcile"
pluginregistry "github.com/Tencent/WeKnora/internal/plugin/registry"
@@ -168,8 +167,10 @@ func BuildContainer(container *dig.Container) *dig.Container {
must(container.Provide(activate.NewMCPServers))
must(container.Provide(activate.NewSkills))
must(container.Provide(activate.NewModelVendors))
must(container.Provide(pluginhost.NewManager))
must(container.Provide(newPluginHostManager))
must(container.Provide(pluginremote.NewManager))
must(container.Provide(newPluginHostPool))
must(container.Provide(newPluginDelegation))
must(container.Provide(newPluginInvoker))
must(container.Provide(repository.NewPluginKVRepository))
must(container.Provide(newPluginHostAPI))
+152
View File
@@ -0,0 +1,152 @@
package container
import (
"context"
"errors"
"fmt"
"net"
"net/http"
"net/url"
"os"
"strconv"
"strings"
"time"
"github.com/Tencent/WeKnora/internal/application/repository"
"github.com/Tencent/WeKnora/internal/config"
"github.com/Tencent/WeKnora/internal/logger"
"github.com/Tencent/WeKnora/internal/plugin/host"
"github.com/Tencent/WeKnora/internal/plugin/hostpool"
"github.com/Tencent/WeKnora/internal/plugin/manifest"
"github.com/Tencent/WeKnora/internal/plugin/reconcile"
pluginregistry "github.com/Tencent/WeKnora/internal/plugin/registry"
)
// Standalone plugin host settings.
const (
envPluginHostAddr = "WEKNORA_PLUGIN_HOST_ADDR" // listen address, default :8081
envPluginHostURL = "WEKNORA_PLUGIN_HOST_URL" // gateway URL app nodes use
envPluginHostKinds = "WEKNORA_PLUGIN_HOST_KINDS" // kinds to run, default all available
)
// RunPluginHost runs this binary as a standalone plugin host
// (weknora plugin-host) until ctx ends. It shares the app's database,
// object storage and Redis, runs the installed host plugins of the kinds it
// supports, announces them in Redis and serves them to app nodes through a
// gateway signed with the cluster key. It never migrates the database;
// app nodes do.
func RunPluginHost(ctx context.Context) error {
if _, set := os.LookupEnv("AUTO_MIGRATE"); !set {
_ = os.Setenv("AUTO_MIGRATE", "false")
}
key, err := hostpool.ClusterKey()
if err != nil {
return err
}
rdb, err := initRedisClient()
if err != nil {
return err
}
if rdb == nil {
return errors.New("a plugin host needs Redis (REDIS_ADDR) to announce itself to app nodes")
}
cfg, err := config.LoadConfig()
if err != nil {
return fmt.Errorf("load config: %w", err)
}
db, err := initDatabase(cfg)
if err != nil {
return fmt.Errorf("database: %w", err)
}
store, err := newPluginPackageStore(cfg)
if err != nil {
return fmt.Errorf("package storage: %w", err)
}
kinds := host.KindsFromEnv(envPluginHostKinds)
if len(kinds) == 0 {
return fmt.Errorf("%s leaves this plugin host nothing to run", envPluginHostKinds)
}
mgr := host.NewStandaloneManager(kinds)
// Plugins call the Host API on an app node, not through their egress
// proxy.
hostAPI := strings.TrimSpace(os.Getenv("WEKNORA_PLUGIN_HOST_API_URL"))
if u, err := url.Parse(hostAPI); err == nil && u.Hostname() != "" {
mgr.SetDirectHosts(u.Hostname())
}
r := reconcile.New(reconcile.Options{
Repo: repository.NewPluginRepository(db), Store: store, Registry: pluginregistry.New(), Redis: rdb,
Activators: []reconcile.Activator{mgr}, Runtimes: []manifest.RuntimeType{manifest.RuntimeHost},
Accept: func(m *manifest.Manifest) bool { return mgr.Runs(m.Runtime.Kind) }, Role: "plugin-host",
})
mgr.SetReporter(r)
addr := os.Getenv(envPluginHostAddr)
if addr == "" {
addr = ":8081"
}
ln, err := net.Listen("tcp", addr)
if err != nil {
return fmt.Errorf("listen %s: %w", addr, err)
}
advertise, err := advertiseURL(ln.Addr())
if err != nil {
_ = ln.Close()
return err
}
srv := &http.Server{Handler: mgr.Gateway(key), ReadHeaderTimeout: 10 * time.Second}
serveErr := make(chan error, 1)
go func() { serveErr <- srv.Serve(ln) }()
runCtx, stop := context.WithCancel(ctx)
defer stop()
r.Start(runCtx)
hostname, _ := os.Hostname()
announced := make(chan struct{})
go func() {
hostpool.NewAnnouncer(rdb, mgr, r.NodeName(), advertise).Run(runCtx)
close(announced)
}()
logger.Infof(ctx, "[plugin-host] %s runs %s plugins, serving app nodes at %s",
hostname, strings.Join(kinds, ", "), advertise)
select {
case <-ctx.Done():
case err := <-serveErr:
stop()
<-announced
mgr.Close()
return fmt.Errorf("gateway: %w", err)
}
// Withdraw first so app nodes stop picking this host, then let calls in
// flight finish before the plugins stop.
stop()
<-announced
shutdown, cancel := context.WithTimeout(context.Background(), 60*time.Second)
defer cancel()
if err := srv.Shutdown(shutdown); err != nil {
logger.Warnf(ctx, "[plugin-host] shutdown: %v", err)
}
mgr.Close()
return nil
}
// advertiseURL is the gateway URL app nodes use: WEKNORA_PLUGIN_HOST_URL,
// or http://<hostname>:<port>, which suits compose and Kubernetes pod DNS.
func advertiseURL(addr net.Addr) (string, error) {
if u := strings.TrimSpace(os.Getenv(envPluginHostURL)); u != "" {
if _, err := url.Parse(u); err != nil {
return "", fmt.Errorf("%s: %w", envPluginHostURL, err)
}
return strings.TrimSuffix(u, "/"), nil
}
tcp, ok := addr.(*net.TCPAddr)
if !ok {
return "", fmt.Errorf("set %s", envPluginHostURL)
}
name, err := os.Hostname()
if err != nil || name == "" {
return "", fmt.Errorf("no hostname to advertise; set %s", envPluginHostURL)
}
return "http://" + net.JoinHostPort(name, strconv.Itoa(tcp.Port)), nil
}
+90 -21
View File
@@ -15,10 +15,13 @@ import (
"github.com/Tencent/WeKnora/internal/application/service"
"github.com/Tencent/WeKnora/internal/config"
"github.com/Tencent/WeKnora/internal/handler"
"github.com/Tencent/WeKnora/internal/logger"
"github.com/Tencent/WeKnora/internal/plugin/activate"
"github.com/Tencent/WeKnora/internal/plugin/host"
"github.com/Tencent/WeKnora/internal/plugin/hostapi"
"github.com/Tencent/WeKnora/internal/plugin/hostpool"
"github.com/Tencent/WeKnora/internal/plugin/install"
"github.com/Tencent/WeKnora/internal/plugin/manifest"
"github.com/Tencent/WeKnora/internal/plugin/reconcile"
pluginregistry "github.com/Tencent/WeKnora/internal/plugin/registry"
"github.com/Tencent/WeKnora/internal/plugin/remote"
@@ -44,6 +47,7 @@ type pluginActivators struct {
Host *host.Manager
Remote *remote.Manager
Delegation *pluginDelegation
WebSearch *activate.WebSearch
Connectors *activate.Connectors
Parsers *activate.Parsers
@@ -57,27 +61,26 @@ type pluginActivators struct {
// reachable before anything routes calls to it.
func (a pluginActivators) list() []reconcile.Activator {
return []reconcile.Activator{
a.Host, a.Remote, a.WebSearch, a.Connectors, a.Parsers, a.Vendors, a.MCP, a.Skills,
a.Host, a.Remote, a.Delegation, a.WebSearch, a.Connectors, a.Parsers, a.Vendors, a.MCP, a.Skills,
}
}
// newPluginHostAPI serves the Host API and gives calls a way back to it: the
// embedded host's plugins reach this node on loopback, remote plugins only
// through WEKNORA_PLUGIN_HOST_API_URL.
// embedded host's plugins reach this node on loopback, plugins elsewhere
// (remote, on a plugin host) only through WEKNORA_PLUGIN_HOST_API_URL.
func newPluginHostAPI(
cfg *config.Config, iv *activate.Invoker, repo interfaces.PluginKVRepository, cleaner interfaces.ResourceCleaner,
) *hostapi.Handler {
issuer := hostapi.NewIssuerFromEnv()
url := strings.TrimSpace(os.Getenv("WEKNORA_PLUGIN_HOST_API_URL"))
iv.SetPublicHostAPI(url)
if url == "" {
port := 8080
if cfg != nil && cfg.Server != nil && cfg.Server.Port > 0 {
port = cfg.Server.Port
}
url = fmt.Sprintf("http://127.0.0.1:%d", port)
// Plugins on this node always use loopback; the public address is for
// plugins elsewhere (remote, on a plugin host), which reach this node
// the way the deployment routes to it.
port := 8080
if cfg != nil && cfg.Server != nil && cfg.Server.Port > 0 {
port = cfg.Server.Port
}
iv.SetHostAPI(issuer, url)
iv.SetHostAPI(issuer, fmt.Sprintf("http://127.0.0.1:%d", port))
iv.SetPublicHostAPI(strings.TrimSpace(os.Getenv("WEKNORA_PLUGIN_HOST_API_URL")))
kv := hostapi.NewKV(repo)
ctx, cancel := context.WithCancel(context.Background())
kv.StartSweeper(ctx, 10*time.Minute)
@@ -97,25 +100,88 @@ func bindPluginActivators(
skills.SetPluginSkills(a.Skills)
}
// newPluginInvoker reaches code plugins through this node's plugin host or,
// for remote ones, at their registered URLs.
func newPluginInvoker(h *host.Manager, r *remote.Manager) *activate.Invoker {
return activate.NewInvoker(pluginClients{host: h, remote: r})
// newPluginHostManager is this node's embedded plugin host. It runs the
// kinds WEKNORA_PLUGIN_EMBEDDED_KINDS names (default: every kind this
// machine can run; "none" leaves all host plugins to standalone hosts).
func newPluginHostManager() *host.Manager {
m := host.NewManager()
m.SetKinds(host.KindsFromEnv("WEKNORA_PLUGIN_EMBEDDED_KINDS"))
return m
}
// pluginClients finds a code plugin in whichever runtime loaded it.
// newPluginHostPool finds standalone plugin hosts (weknora plugin-host) in
// Redis. It is nil without Redis or a cluster key: host plugins then run on
// this node or nowhere.
func newPluginHostPool(rdb *redis.Client) *hostpool.Pool {
if rdb == nil {
return nil
}
key, err := hostpool.ClusterKey()
if err != nil {
logger.Warnf(context.Background(), "[plugin] standalone plugin hosts are off: %v", err)
return nil
}
return hostpool.NewPool(rdb, key)
}
// newPluginInvoker reaches code plugins through this node's plugin host, a
// standalone plugin host, or a remote plugin's registered URL.
func newPluginInvoker(h *host.Manager, r *remote.Manager, pool *hostpool.Pool) *activate.Invoker {
return activate.NewInvoker(pluginClients{host: h, remote: r, pool: pool})
}
// pluginClients finds a code plugin in whichever runtime serves it.
type pluginClients struct {
host *host.Manager
remote *remote.Manager
pool *hostpool.Pool
}
func (c pluginClients) Client(pluginID string) (*client.Client, error) {
if c.remote.Owns(pluginID) {
return c.remote.Client(pluginID)
func (c pluginClients) Client(ctx context.Context, m *manifest.Manifest) (*client.Client, error) {
switch {
case c.remote.Owns(m.ID):
return c.remote.Client(m.ID)
case c.host.Local(m.ID) || c.pool == nil || m.Runtime.Type != manifest.RuntimeHost:
return c.host.Client(m.ID)
default:
return c.pool.Client(ctx, m.ID, m.Version)
}
return c.host.Client(pluginID)
}
func (c pluginClients) OnThisNode(pluginID string) bool { return c.host.Local(pluginID) }
// pluginDelegation fails a host plugin this node leaves to standalone
// plugin hosts when there are none to leave it to. With hosts configured it
// only logs: they may start later, and calls say which plugin is missing.
type pluginDelegation struct {
host *host.Manager
pool *hostpool.Pool
}
func newPluginDelegation(h *host.Manager, pool *hostpool.Pool) *pluginDelegation {
return &pluginDelegation{host: h, pool: pool}
}
func (d *pluginDelegation) Name() string { return "plugin-hosts" }
func (d *pluginDelegation) Activate(ctx context.Context, l *reconcile.Loaded) error {
m := l.Manifest
if m.Runtime.Type != manifest.RuntimeHost || d.host.Runs(m.Runtime.Kind) {
return nil
}
if d.pool == nil {
return fmt.Errorf("this node does not run %s plugins (WEKNORA_PLUGIN_EMBEDDED_KINDS) and no plugin host "+
"is configured; run weknora plugin-host with Redis and the same SYSTEM_AES_KEY", m.Runtime.Kind)
}
if !d.pool.Runs(ctx, m.ID, m.Version) {
logger.Infof(ctx, "[plugin] %s@%s waits for a plugin host that runs %s plugins",
m.ID, m.Version, m.Runtime.Kind)
}
return nil
}
func (d *pluginDelegation) Deactivate(context.Context, string) error { return nil }
func newMCPServiceRepository(db *gorm.DB, plugins *activate.MCPServers) interfaces.MCPServiceRepository {
return plugins.Repository(repository.NewMCPServiceRepository(db))
}
@@ -146,6 +212,9 @@ func startPluginReconciler(
r *reconcile.Reconciler, hostManager *host.Manager, remoteManager *remote.Manager,
cleaner interfaces.ResourceCleaner,
) {
if kinds := hostManager.Kinds(); len(kinds) > 0 {
logger.Infof(context.Background(), "[plugin] this node runs %s host plugins", strings.Join(kinds, ", "))
}
hostManager.SetReporter(r)
remoteManager.SetReporter(r)
ctx, cancel := context.WithCancel(context.Background())
+72
View File
@@ -0,0 +1,72 @@
package container
import (
"context"
"strings"
"testing"
"github.com/alicebob/miniredis/v2"
"github.com/redis/go-redis/v9"
"github.com/Tencent/WeKnora/internal/plugin/host"
"github.com/Tencent/WeKnora/internal/plugin/hostpool"
"github.com/Tencent/WeKnora/internal/plugin/manifest"
"github.com/Tencent/WeKnora/internal/plugin/reconcile"
"github.com/Tencent/WeKnora/internal/plugin/remote"
"github.com/Tencent/WeKnora/pluginsdk/pluginapi"
)
func pythonPlugin() *reconcile.Loaded {
return &reconcile.Loaded{Manifest: &manifest.Manifest{
ID: "acme.py", Version: "1.0.0",
Runtime: manifest.Runtime{Type: manifest.RuntimeHost, Kind: host.KindPython, Entry: "main.py"},
}}
}
// A node that does not run a kind hands it to plugin hosts; without any it
// says so instead of pretending the plugin loaded.
func TestDelegationNeedsPluginHosts(t *testing.T) {
ctx := context.Background()
h := host.NewManager()
h.SetKinds([]string{host.KindBinary})
err := newPluginDelegation(h, nil).Activate(ctx, pythonPlugin())
if err == nil || !strings.Contains(err.Error(), "no plugin host is configured") {
t.Fatalf("without plugin hosts = %v", err)
}
mr := miniredis.RunT(t)
pool := hostpool.NewPool(redis.NewClient(&redis.Options{Addr: mr.Addr()}), []byte("k"))
if err := newPluginDelegation(h, pool).Activate(ctx, pythonPlugin()); err != nil {
t.Fatalf("with plugin hosts configured, a missing one is not a load failure: %v", err)
}
h.SetKinds([]string{host.KindPython})
if err := newPluginDelegation(h, nil).Activate(ctx, pythonPlugin()); err != nil {
t.Fatalf("a kind the node runs needs no plugin host: %v", err)
}
}
// Calls go to the pool only for host plugins this node does not run.
func TestPluginClientsRouteToThePool(t *testing.T) {
ctx := context.Background()
h := host.NewManager()
h.SetKinds([]string{host.KindBinary})
mr := miniredis.RunT(t)
pool := hostpool.NewPool(redis.NewClient(&redis.Options{Addr: mr.Addr()}), []byte("k"))
clients := pluginClients{host: h, remote: remote.NewManager(), pool: pool}
_, err := clients.Client(ctx, pythonPlugin().Manifest)
if pe, ok := pluginapi.AsError(err); !ok || !strings.Contains(pe.Message, "no plugin host runs acme.py@1.0.0") {
t.Fatalf("a delegated plugin is looked up in the pool, got %v", err)
}
if clients.OnThisNode("acme.py") {
t.Fatal("a delegated plugin is not on this node")
}
clients.pool = nil
_, err = clients.Client(ctx, pythonPlugin().Manifest)
if pe, ok := pluginapi.AsError(err); !ok || !strings.Contains(pe.Message, "not running on this node") {
t.Fatalf("without a pool the node's own host answers, got %v", err)
}
}
+10 -6
View File
@@ -14,9 +14,13 @@ import (
"github.com/Tencent/WeKnora/pluginsdk/pluginapi"
)
// ClientSource reaches running code plugins (the node's plugin host).
// ClientSource reaches running code plugins wherever they run: this node's
// plugin host, a standalone plugin host, a remote service.
type ClientSource interface {
Client(pluginID string) (*client.Client, error)
Client(ctx context.Context, m *manifest.Manifest) (*client.Client, error)
// OnThisNode reports whether a plugin runs on this node, where the
// loopback Host API address reaches WeKnora.
OnThisNode(pluginID string) bool
}
// Invoker is how the remote adapters call a code plugin: it finds the
@@ -29,7 +33,7 @@ type Invoker struct {
tenancy *tenancy.Service
plugins interfaces.PluginRepository
// tokens and hostURL give plugins granted Host API scopes a way back;
// remote plugins, off this node, use publicHostURL.
// plugins off this node (remote, on a plugin host) use publicHostURL.
tokens TokenIssuer
hostURL string
publicHostURL string
@@ -104,7 +108,7 @@ func (iv *Invoker) Envelope(
}
iv.mu.RLock()
t, plugins, tokens, hostURL := iv.tenancy, iv.plugins, iv.tokens, iv.hostURL
if m.Runtime.Type == manifest.RuntimeRemote {
if !iv.clients.OnThisNode(m.ID) {
hostURL = iv.publicHostURL
}
iv.mu.RUnlock()
@@ -160,7 +164,7 @@ func present(pluginID string, err error) error {
func (iv *Invoker) Call(
ctx context.Context, m *manifest.Manifest, path string, instance map[string]any, input, out any,
) error {
c, err := iv.clients.Client(m.ID)
c, err := iv.clients.Client(ctx, m)
if err != nil {
return present(m.ID, err)
}
@@ -176,7 +180,7 @@ func (iv *Invoker) Stream(
ctx context.Context, m *manifest.Manifest, path string, instance map[string]any, input any,
fn func(pluginapi.Event) error,
) ([]byte, error) {
c, err := iv.clients.Client(m.ID)
c, err := iv.clients.Client(ctx, m)
if err != nil {
return nil, present(m.ID, err)
}
+1 -1
View File
@@ -36,7 +36,7 @@ func TestEnvelopeCarriesHostAccessOnlyForGrantedPlugins(t *testing.T) {
}
remote := *granted
remote.Runtime.Type = manifest.RuntimeRemote
remote.ID = "acme.remote"
if env, _ := iv.Envelope(ctx, &remote, nil); env.Context.Host != nil {
t.Fatal("a remote plugin must not be sent to this node's loopback address")
}
+1 -1
View File
@@ -103,7 +103,7 @@ func (e *pluginEngine) FileTypes(bool) []string { return e.fileTypes }
func (e *pluginEngine) PluginID() string { return e.m.ID }
func (e *pluginEngine) DisplayNames() map[string]string { return e.names }
func (e *pluginEngine) CheckAvailable(bool, map[string]string) (bool, string) {
if _, err := e.iv.clients.Client(e.m.ID); err != nil {
if _, err := e.iv.clients.Client(context.Background(), e.m); err != nil {
return false, err.Error()
}
return true, ""
+8 -1
View File
@@ -8,6 +8,7 @@ import (
"testing"
infra_web_search "github.com/Tencent/WeKnora/internal/infrastructure/web_search"
"github.com/Tencent/WeKnora/internal/plugin/manifest"
"github.com/Tencent/WeKnora/internal/plugin/pkg"
"github.com/Tencent/WeKnora/internal/plugin/plugintest"
"github.com/Tencent/WeKnora/internal/plugin/reconcile"
@@ -36,7 +37,13 @@ contributes:
// fakeClients serves one plugin from an in-process SDK handler.
type fakeClients struct{ c *client.Client }
func (f fakeClients) Client(string) (*client.Client, error) { return f.c, nil }
func (f fakeClients) Client(context.Context, *manifest.Manifest) (*client.Client, error) {
return f.c, nil
}
// OnThisNode is true unless the manifest's runtime is remote, like the
// embedded host.
func (f fakeClients) OnThisNode(id string) bool { return id != "acme.remote" }
func searchLoaded(t *testing.T, schema string) *reconcile.Loaded {
t.Helper()
+119
View File
@@ -0,0 +1,119 @@
package host
import (
"bytes"
"encoding/json"
"io"
"net/http"
"strings"
"time"
"github.com/Tencent/WeKnora/pluginsdk/pluginapi"
)
// GatewayPrefix starts the paths a plugin host serves its plugins under:
// /p/{pluginId}/{version}/{protocol path}.
const GatewayPrefix = "/p/"
// maxGatewayBody bounds one relayed request: a parse call carries the whole
// document.
const maxGatewayBody = 512 << 20
// Gateway serves the plugins this host runs to the other WeKnora nodes, at
// /p/{pluginId}/{version}/... Requests must be signed with the cluster key
// (pluginapi.Sign over the body); answers, streams included, are relayed as
// the plugin gives them.
func (m *Manager) Gateway(key []byte) http.Handler {
mux := http.NewServeMux()
mux.HandleFunc("GET /healthz", func(w http.ResponseWriter, _ *http.Request) {
w.WriteHeader(http.StatusOK)
})
mux.HandleFunc(GatewayPrefix, func(w http.ResponseWriter, r *http.Request) {
m.inFlight.Add(1)
defer m.inFlight.Add(-1)
m.relay(w, r, key)
})
return mux
}
func gatewayError(w http.ResponseWriter, e *pluginapi.Error) {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(e.Code.HTTPStatus())
_ = json.NewEncoder(w).Encode(pluginapi.ErrorBody{Error: *e})
}
func (m *Manager) relay(w http.ResponseWriter, r *http.Request, key []byte) {
body, err := io.ReadAll(io.LimitReader(r.Body, maxGatewayBody+1))
if err != nil {
gatewayError(w, pluginapi.Errorf(pluginapi.CodeBadRequest, "read request: %v", err))
return
}
if len(body) > maxGatewayBody {
gatewayError(w, pluginapi.Errorf(pluginapi.CodeBadRequest, "request is over %d bytes", maxGatewayBody))
return
}
err = pluginapi.VerifySignature(key, r.Header.Get(pluginapi.TimestampHeader),
r.Header.Get(pluginapi.SignatureHeader), body, time.Now())
if err != nil {
gatewayError(w, pluginapi.Errorf(pluginapi.CodeUnauthorized, "gateway: %v", err))
return
}
parts := strings.SplitN(strings.TrimPrefix(r.URL.Path, GatewayPrefix), "/", 3)
if len(parts) < 3 || parts[0] == "" || parts[1] == "" {
gatewayError(w, pluginapi.Errorf(pluginapi.CodeNotFound, "no endpoint %s", r.URL.Path))
return
}
id, version, path := parts[0], parts[1], "/"+parts[2]
m.mu.Lock()
p := m.procs[id]
m.mu.Unlock()
if p == nil || p.spec.m.Version != version {
// The caller's view of this host is a heartbeat old; it retries
// elsewhere.
gatewayError(w, pluginapi.Errorf(pluginapi.CodeUnavailable, "this host does not run %s@%s", id, version))
return
}
c, err := p.Client()
if err != nil {
if pe, ok := pluginapi.AsError(err); ok {
gatewayError(w, pe)
} else {
gatewayError(w, pluginapi.Errorf(pluginapi.CodeUnavailable, "%v", err))
}
return
}
var sent []byte
if len(body) > 0 {
sent = body
}
resp, err := c.Raw(r.Context(), r.Method, path, sent, r.Header.Get(pluginapi.RequestIDHeader))
if err != nil {
gatewayError(w, pluginapi.Errorf(pluginapi.CodeUnavailable, "%v", err))
return
}
defer func() { _ = resp.Body.Close() }()
if ct := resp.Header.Get("Content-Type"); ct != "" {
w.Header().Set("Content-Type", ct)
}
w.WriteHeader(resp.StatusCode)
flusher, _ := w.(http.Flusher)
buf := make([]byte, 32<<10)
for {
n, err := resp.Body.Read(buf)
if n > 0 {
if _, werr := w.Write(buf[:n]); werr != nil {
return
}
// Stream events go out as the plugin writes them.
if flusher != nil && bytes.IndexByte(buf[:n], '\n') >= 0 {
flusher.Flush()
}
}
if err != nil {
if flusher != nil {
flusher.Flush()
}
return
}
}
}
+146 -3
View File
@@ -12,7 +12,12 @@ package host
import (
"context"
"fmt"
"os"
"os/exec"
"sort"
"strings"
"sync"
"sync/atomic"
"github.com/Tencent/WeKnora/internal/logger"
"github.com/Tencent/WeKnora/internal/plugin/manifest"
@@ -31,13 +36,101 @@ type StateReporter interface {
// before the activators that route calls to the plugins.
type Manager struct {
reporter StateReporter
inFlight atomic.Int64
mu sync.Mutex
procs map[string]*process
kinds map[string]bool
// standalone hosts have no one to hand other kinds to.
standalone bool
// direct are hosts plugins reach without the egress proxy: the Host API
// when it is not on this machine.
direct []string
}
// NewManager creates an empty host.
func NewManager() *Manager { return &Manager{procs: map[string]*process{}} }
// NewManager creates an empty host that runs every kind this machine can.
func NewManager() *Manager {
m := &Manager{procs: map[string]*process{}}
m.SetKinds(AvailableKinds())
return m
}
// NewStandaloneManager creates the host of a standalone plugin host
// (weknora plugin-host): it runs the given kinds and refuses the others.
func NewStandaloneManager(kinds []string) *Manager {
m := &Manager{procs: map[string]*process{}, standalone: true}
m.SetKinds(kinds)
return m
}
// AvailableKinds are the kinds this machine can run: binaries, and python
// when an interpreter is installed.
func AvailableKinds() []string {
kinds := []string{KindBinary}
if _, err := exec.LookPath(PythonCommand()); err == nil {
kinds = append(kinds, KindPython)
}
return kinds
}
// KindsFromEnv reads a comma-separated list of kinds from an environment
// variable: unset means every available kind, "none" means none.
func KindsFromEnv(name string) []string {
raw, ok := os.LookupEnv(name)
raw = strings.TrimSpace(raw)
if !ok || raw == "" {
return AvailableKinds()
}
if raw == "none" {
return nil
}
var out []string
for _, k := range strings.Split(raw, ",") {
if k = strings.TrimSpace(k); Supported(k) {
out = append(out, k)
}
}
return out
}
// SetKinds limits the kinds this host runs; plugins of other kinds are left
// to plugin hosts elsewhere.
func (m *Manager) SetKinds(kinds []string) {
m.mu.Lock()
defer m.mu.Unlock()
m.kinds = map[string]bool{}
for _, k := range kinds {
m.kinds[k] = true
}
}
// Kinds lists the kinds this host runs.
func (m *Manager) Kinds() []string {
m.mu.Lock()
defer m.mu.Unlock()
out := make([]string, 0, len(m.kinds))
for k := range m.kinds {
out = append(out, k)
}
sort.Strings(out)
return out
}
// Runs reports whether this host runs plugins of a kind.
func (m *Manager) Runs(kind string) bool {
m.mu.Lock()
defer m.mu.Unlock()
return m.kinds[kind]
}
// SetDirectHosts names hosts plugins reach without the egress proxy, such
// as the Host API's when it is on another machine. It applies to processes
// started afterwards.
func (m *Manager) SetDirectHosts(hosts ...string) {
m.mu.Lock()
m.direct = append([]string(nil), hosts...)
m.mu.Unlock()
}
// SetReporter wires health changes to the reconciler.
func (m *Manager) SetReporter(r StateReporter) {
@@ -63,8 +156,18 @@ func (m *Manager) Activate(ctx context.Context, l *reconcile.Loaded) error {
return fmt.Errorf("runtime.kind %q is not supported by this host; binary and python are",
l.Manifest.Runtime.Kind)
}
if !m.Runs(l.Manifest.Runtime.Kind) {
if m.standalone {
return fmt.Errorf("this plugin host does not run %s plugins (it runs %s)",
l.Manifest.Runtime.Kind, strings.Join(m.Kinds(), ", "))
}
return nil // a plugin host elsewhere runs it
}
id := l.Manifest.ID
p, err := startProcess(spec{m: l.Manifest, dir: l.Dir}, func(s State, err error) {
m.mu.Lock()
direct := m.direct
m.mu.Unlock()
p, err := startProcess(spec{m: l.Manifest, dir: l.Dir, direct: direct}, func(s State, err error) {
m.report(id, s, err)
})
if err != nil {
@@ -123,6 +226,46 @@ func (m *Manager) Client(pluginID string) (*client.Client, error) {
return p.Client()
}
// Local reports whether a plugin runs on this host.
func (m *Manager) Local(pluginID string) bool {
m.mu.Lock()
defer m.mu.Unlock()
_, ok := m.procs[pluginID]
return ok
}
// Running is one plugin this host runs.
type Running struct {
ID string `json:"id"`
Version string `json:"version"`
Kind string `json:"kind"`
State State `json:"state"`
}
// Running lists the plugins this host runs and their state.
func (m *Manager) Running() []Running {
m.mu.Lock()
procs := make([]*process, 0, len(m.procs))
for _, p := range m.procs {
procs = append(procs, p)
}
m.mu.Unlock()
out := make([]Running, 0, len(procs))
for _, p := range procs {
p.mu.RLock()
st := p.state
p.mu.RUnlock()
out = append(out, Running{
ID: p.spec.m.ID, Version: p.spec.m.Version, Kind: p.spec.m.Runtime.Kind, State: st,
})
}
sort.Slice(out, func(i, j int) bool { return out[i].ID < out[j].ID })
return out
}
// InFlight is how many gateway calls this host is answering.
func (m *Manager) InFlight() int64 { return m.inFlight.Load() }
// Close stops every plugin, for shutdown.
func (m *Manager) Close() {
m.mu.Lock()
+6 -4
View File
@@ -51,6 +51,8 @@ const (
type spec struct {
m *manifest.Manifest
dir string // extracted package
// direct are hosts reached without the egress proxy.
direct []string
}
// Kinds of host plugin this host runs.
@@ -200,11 +202,11 @@ func (p *process) childEnv(network, socket, token string) []string {
"HTTPS_PROXY=" + p.proxy.URL(),
"http_proxy=" + p.proxy.URL(),
"https_proxy=" + p.proxy.URL(),
// The Host API is on this node's loopback; the proxy is for the
// outside world.
"NO_PROXY=127.0.0.1,localhost,::1",
"no_proxy=127.0.0.1,localhost,::1",
}
// The Host API is on this node's loopback, or on a host named direct;
// the proxy is for the outside world.
noProxy := strings.Join(append([]string{"127.0.0.1", "localhost", "::1"}, p.spec.direct...), ",")
env = append(env, "NO_PROXY="+noProxy, "no_proxy="+noProxy)
for _, k := range []string{"PATH", "LANG", "LC_ALL", "TZ", "SYSTEMROOT", "WINDIR"} {
if v, ok := os.LookupEnv(k); ok {
env = append(env, k+"="+v)
+27
View File
@@ -128,3 +128,30 @@ func TestEntryName(t *testing.T) {
t.Fatal("this host runs binaries and python")
}
}
func TestKindsDecideWhereAPluginRuns(t *testing.T) {
ctx := context.Background()
l := install(t, "1.0.0", "")
embedded := NewManager()
embedded.SetKinds([]string{KindPython})
defer embedded.Close()
if err := embedded.Activate(ctx, l); err != nil || embedded.Local("acme.echo") {
t.Fatalf("an app node hands other kinds on: %v", err)
}
standalone := NewStandaloneManager([]string{KindPython})
defer standalone.Close()
if err := standalone.Activate(ctx, l); err == nil || !strings.Contains(err.Error(), "does not run binary") {
t.Fatalf("a standalone host refuses other kinds, got %v", err)
}
t.Setenv("WEKNORA_PLUGIN_EMBEDDED_KINDS", "none")
if got := KindsFromEnv("WEKNORA_PLUGIN_EMBEDDED_KINDS"); len(got) != 0 {
t.Fatalf("none = %v", got)
}
t.Setenv("WEKNORA_PLUGIN_EMBEDDED_KINDS", "binary, node")
if got := KindsFromEnv("WEKNORA_PLUGIN_EMBEDDED_KINDS"); len(got) != 1 || got[0] != KindBinary {
t.Fatalf("binary, node = %v", got)
}
}
+261
View File
@@ -0,0 +1,261 @@
// Package hostpool connects app nodes to standalone plugin hosts
// (weknora plugin-host) without Kubernetes: each host announces itself and
// the plugins it runs in Redis, and app nodes call the least busy host that
// runs the version they need, through the host's signed gateway.
package hostpool
import (
"context"
"crypto/hmac"
"crypto/sha256"
"encoding/json"
"errors"
"fmt"
"math/rand/v2"
"net/http"
"os"
"runtime"
"strings"
"sync"
"time"
"github.com/redis/go-redis/v9"
"github.com/Tencent/WeKnora/internal/logger"
"github.com/Tencent/WeKnora/internal/plugin/host"
"github.com/Tencent/WeKnora/internal/utils"
"github.com/Tencent/WeKnora/pluginsdk/client"
"github.com/Tencent/WeKnora/pluginsdk/pluginapi"
)
// Heartbeat timing: a host refreshes its record every HeartbeatInterval and
// drops out HeartbeatTTL after its last one.
const (
HeartbeatInterval = 5 * time.Second
HeartbeatTTL = 15 * time.Second
// refreshAfter is how long an app node trusts its view of the hosts.
refreshAfter = 2 * time.Second
)
const keyBase = "weknora:plugin-host:"
func keyPrefix() string {
if ns := strings.TrimSpace(os.Getenv("WEKNORA_REDIS_NAMESPACE")); ns != "" {
return keyBase + ns + ":"
}
return keyBase
}
// ErrNoClusterKey means neither SYSTEM_AES_KEY nor JWT_SECRET is set, so
// app nodes and plugin hosts share no secret to sign calls with.
var ErrNoClusterKey = errors.New("plugin hosts need SYSTEM_AES_KEY or JWT_SECRET, the same on every node")
// ClusterKey is the key app nodes sign gateway calls with, derived from the
// secrets every node already shares.
func ClusterKey() ([]byte, error) {
var secret []byte
switch {
case utils.GetAESKey() != nil:
secret = utils.GetAESKey()
case strings.TrimSpace(os.Getenv("JWT_SECRET")) != "":
secret = []byte(strings.TrimSpace(os.Getenv("JWT_SECRET")))
default:
return nil, ErrNoClusterKey
}
mac := hmac.New(sha256.New, secret)
_, _ = mac.Write([]byte("weknora plugin host gateway"))
return mac.Sum(nil), nil
}
// Info is one plugin host's announcement.
type Info struct {
ID string `json:"id"`
URL string `json:"url"` // gateway base URL, reachable from app nodes
OS string `json:"os"`
// Arch and Kinds say what the host can run.
Arch string `json:"arch"`
Kinds []string `json:"kinds"`
Plugins []host.Running `json:"plugins"`
InFlight int64 `json:"inFlight"`
UpdatedAt time.Time `json:"updatedAt"`
}
// runs reports whether the host serves a plugin version right now.
func (i Info) runs(pluginID, version string) bool {
for _, p := range i.Plugins {
if p.ID == pluginID && p.Version == version && p.State == host.StateReady {
return true
}
}
return false
}
// Announcer keeps a plugin host's record in Redis.
type Announcer struct {
rdb *redis.Client
mgr *host.Manager
id string
url string
}
// NewAnnouncer announces the host at url (its gateway, as app nodes reach
// it) under id.
func NewAnnouncer(rdb *redis.Client, mgr *host.Manager, id, url string) *Announcer {
return &Announcer{rdb: rdb, mgr: mgr, id: id, url: strings.TrimSuffix(url, "/")}
}
func (a *Announcer) info() Info {
return Info{
ID: a.id, URL: a.url, OS: runtime.GOOS, Arch: runtime.GOARCH, Kinds: a.mgr.Kinds(),
Plugins: a.mgr.Running(), InFlight: a.mgr.InFlight(), UpdatedAt: time.Now(),
}
}
// Announce writes the record once.
func (a *Announcer) Announce(ctx context.Context) error {
b, err := json.Marshal(a.info())
if err != nil {
return err
}
return a.rdb.Set(ctx, keyPrefix()+a.id, b, HeartbeatTTL).Err()
}
// Run announces every HeartbeatInterval until ctx ends, then withdraws the
// record so app nodes stop calling at once.
func (a *Announcer) Run(ctx context.Context) {
t := time.NewTicker(HeartbeatInterval)
defer t.Stop()
for {
if err := a.Announce(ctx); err != nil && ctx.Err() == nil {
logger.Warnf(ctx, "[plugin-host] heartbeat: %v", err)
}
select {
case <-ctx.Done():
wctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
_ = a.rdb.Del(wctx, keyPrefix()+a.id).Err()
cancel()
return
case <-t.C:
}
}
}
// Pool is an app node's view of the plugin hosts.
type Pool struct {
rdb *redis.Client
key []byte
http *http.Client
now func() time.Time
mu sync.Mutex
hosts []Info
fetched time.Time
clients map[string]*client.Client
}
// NewPool reads plugin hosts from Redis and signs calls with key.
func NewPool(rdb *redis.Client, key []byte) *Pool {
return &Pool{
rdb: rdb, key: key, now: time.Now, clients: map[string]*client.Client{},
// Calls are bounded by their context; syncs stream for long.
http: &http.Client{Transport: &http.Transport{MaxIdleConnsPerHost: 32, IdleConnTimeout: 90 * time.Second}},
}
}
// Hosts lists the live plugin hosts.
func (p *Pool) Hosts(ctx context.Context) ([]Info, error) {
p.mu.Lock()
defer p.mu.Unlock()
if p.hosts != nil && p.now().Sub(p.fetched) < refreshAfter {
return p.hosts, nil
}
hosts, err := p.fetch(ctx)
if err != nil {
if p.hosts != nil {
logger.Warnf(ctx, "[plugin] read plugin hosts, keeping the last view: %v", err)
return p.hosts, nil
}
return nil, err
}
p.hosts, p.fetched = hosts, p.now()
return hosts, nil
}
func (p *Pool) fetch(ctx context.Context) ([]Info, error) {
var keys []string
iter := p.rdb.Scan(ctx, 0, keyPrefix()+"*", 100).Iterator()
for iter.Next(ctx) {
keys = append(keys, iter.Val())
}
if err := iter.Err(); err != nil {
return nil, err
}
hosts := []Info{}
if len(keys) == 0 {
return hosts, nil
}
vals, err := p.rdb.MGet(ctx, keys...).Result()
if err != nil {
return nil, err
}
for _, v := range vals {
s, ok := v.(string)
if !ok {
continue
}
var info Info
if json.Unmarshal([]byte(s), &info) == nil && info.URL != "" &&
p.now().Sub(info.UpdatedAt) < HeartbeatTTL {
hosts = append(hosts, info)
}
}
return hosts, nil
}
// Client reaches a plugin version on the least busy host that runs it.
func (p *Pool) Client(ctx context.Context, pluginID, version string) (*client.Client, error) {
hosts, err := p.Hosts(ctx)
if err != nil {
return nil, pluginapi.Errorf(pluginapi.CodeUnavailable, "read plugin hosts: %v", err)
}
var best []Info
for _, h := range hosts {
if !h.runs(pluginID, version) {
continue
}
switch {
case len(best) == 0 || h.InFlight < best[0].InFlight:
best = []Info{h}
case h.InFlight == best[0].InFlight:
best = append(best, h)
}
}
if len(best) == 0 {
return nil, pluginapi.Errorf(pluginapi.CodeUnavailable,
"no plugin host runs %s@%s; is weknora plugin-host up?", pluginID, version)
}
h := best[rand.IntN(len(best))] //nolint:gosec // load spreading, not security
base := fmt.Sprintf("%s%s%s/%s", h.URL, host.GatewayPrefix, pluginID, version)
p.mu.Lock()
defer p.mu.Unlock()
c, ok := p.clients[base]
if !ok {
c = client.New(base, p.http, client.Signed(p.key))
p.clients[base] = c
}
return c, nil
}
// Runs reports whether some live host runs a plugin version.
func (p *Pool) Runs(ctx context.Context, pluginID, version string) bool {
hosts, err := p.Hosts(ctx)
if err != nil {
return false
}
for _, h := range hosts {
if h.runs(pluginID, version) {
return true
}
}
return false
}
+179
View File
@@ -0,0 +1,179 @@
package hostpool
import (
"context"
"net/http"
"net/http/httptest"
"os/exec"
"path/filepath"
"runtime"
"testing"
"time"
"github.com/alicebob/miniredis/v2"
"github.com/redis/go-redis/v9"
"github.com/Tencent/WeKnora/internal/plugin/host"
"github.com/Tencent/WeKnora/internal/plugin/manifest"
"github.com/Tencent/WeKnora/internal/plugin/reconcile"
"github.com/Tencent/WeKnora/pluginsdk/client"
"github.com/Tencent/WeKnora/pluginsdk/pluginapi"
)
// echoPackage builds the host package's echo test plugin into an extracted
// package directory.
func echoPackage(t *testing.T) *reconcile.Loaded {
t.Helper()
dir := t.TempDir()
name := "echo"
if runtime.GOOS == "windows" {
name += ".exe"
}
bin := filepath.Join(dir, "bin", runtime.GOOS+"-"+runtime.GOARCH, name)
out, err := exec.Command("go", "build", "-o", bin, "../host/testdata/echoplugin").CombinedOutput()
if err != nil {
t.Fatalf("build echo plugin: %v\n%s", err, out)
}
m := &manifest.Manifest{
SchemaVersion: manifest.SchemaVersion, ID: "acme.echo", Version: "1.0.0", APIVersion: pluginapi.APIVersion,
Name: manifest.Text("Echo", nil), Publisher: manifest.Publisher{ID: "acme"},
Runtime: manifest.Runtime{Type: manifest.RuntimeHost, Kind: host.KindBinary, Entry: "bin/{os}-{arch}/echo"},
Contributes: manifest.Contributions{
manifest.PointWebSearch: {{ID: "echo", Name: manifest.Text("Echo", nil)}},
},
}
return &reconcile.Loaded{Manifest: m, Dir: dir}
}
func search(ctx context.Context, c *client.Client, q string) (string, error) {
var out pluginapi.SearchOutput
err := c.Call(ctx, pluginapi.SearchPath("echo"), pluginapi.Envelope{}, pluginapi.SearchInput{Query: q}, &out)
if err != nil {
return "", err
}
return out.Results[0].Title, nil
}
func TestPoolReachesAPluginThroughAHost(t *testing.T) {
ctx := context.Background()
mr := miniredis.RunT(t)
rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()})
t.Setenv("SYSTEM_AES_KEY", "0123456789abcdef0123456789abcdef")
key, err := ClusterKey()
if err != nil {
t.Fatal(err)
}
// The plugin host: a manager running the plugin behind its gateway.
mgr := host.NewManager()
defer mgr.Close()
if err := mgr.Activate(ctx, echoPackage(t)); err != nil {
t.Fatal(err)
}
gw := httptest.NewServer(mgr.Gateway(key))
defer gw.Close()
a := NewAnnouncer(rdb, mgr, "h1", gw.URL+"/")
if err := a.Announce(ctx); err != nil {
t.Fatal(err)
}
pool := NewPool(rdb, key)
c, err := pool.Client(ctx, "acme.echo", "1.0.0")
if err != nil {
t.Fatal(err)
}
if got, err := search(ctx, c, "hello"); err != nil || got != "hello" {
t.Fatalf("search via host = %q, %v", got, err)
}
if !pool.Runs(ctx, "acme.echo", "1.0.0") || pool.Runs(ctx, "acme.echo", "2.0.0") {
t.Fatal("Runs should match the running version only")
}
if _, err := pool.Client(ctx, "acme.echo", "2.0.0"); !isCode(err, pluginapi.CodeUnavailable) {
t.Fatalf("other version = %v", err)
}
// The gateway refuses calls without the cluster key, and relays a
// version it does not run as unavailable.
unsigned := client.New(gw.URL+"/p/acme.echo/1.0.0", nil, nil)
if _, err := search(ctx, unsigned, "x"); !isCode(err, pluginapi.CodeUnauthorized) {
t.Fatalf("unsigned = %v", err)
}
stale := client.New(gw.URL+"/p/acme.echo/0.9.0", nil, client.Signed(key))
if _, err := search(ctx, stale, "x"); !isCode(err, pluginapi.CodeUnavailable) {
t.Fatalf("stale version = %v", err)
}
resp, err := http.Get(gw.URL + "/healthz")
if err != nil || resp.StatusCode != http.StatusOK {
t.Fatalf("healthz = %v, %v", resp, err)
}
_ = resp.Body.Close()
}
func TestHostsAgeOutAndWithdraw(t *testing.T) {
ctx := context.Background()
mr := miniredis.RunT(t)
rdb := redis.NewClient(&redis.Options{Addr: mr.Addr()})
mgr := host.NewManager()
pool := NewPool(rdb, []byte("k"))
now := time.Now()
pool.now = func() time.Time { return now }
runCtx, stop := context.WithCancel(ctx)
done := make(chan struct{})
go func() { NewAnnouncer(rdb, mgr, "h1", "http://h1:8081").Run(runCtx); close(done) }()
waitFor(t, func() bool { return mr.Exists(keyPrefix() + "h1") })
hosts, err := pool.Hosts(ctx)
if err != nil || len(hosts) != 1 || hosts[0].URL != "http://h1:8081" || len(hosts[0].Kinds) == 0 {
t.Fatalf("hosts = %+v, %v", hosts, err)
}
// A host that stops announcing drops out when its record expires.
stop()
<-done
if mr.Exists(keyPrefix() + "h1") {
t.Fatal("a stopped host must withdraw its record")
}
now = now.Add(refreshAfter)
if hosts, _ := pool.Hosts(ctx); len(hosts) != 0 {
t.Fatalf("hosts after withdrawal = %+v", hosts)
}
// A record that outlived its heartbeat (clock skew, a stuck host) is
// ignored even if Redis still has it.
_ = NewAnnouncer(rdb, mgr, "h2", "http://h2:8081").Announce(ctx)
now = now.Add(refreshAfter + HeartbeatTTL)
if hosts, _ := pool.Hosts(ctx); len(hosts) != 0 {
t.Fatalf("stale hosts = %+v", hosts)
}
}
func TestClusterKey(t *testing.T) {
t.Setenv("SYSTEM_AES_KEY", "")
t.Setenv("JWT_SECRET", "")
if _, err := ClusterKey(); err == nil {
t.Fatal("no shared secret must be an error")
}
t.Setenv("JWT_SECRET", "s")
a, _ := ClusterKey()
t.Setenv("JWT_SECRET", "t")
b, _ := ClusterKey()
if len(a) != 32 || string(a) == string(b) {
t.Fatal("the key derives from the secret")
}
}
func isCode(err error, code pluginapi.ErrorCode) bool {
pe, ok := pluginapi.AsError(err)
return ok && pe.Code == code
}
func waitFor(t *testing.T, cond func() bool) {
t.Helper()
deadline := time.Now().Add(5 * time.Second)
for !cond() {
if time.Now().After(deadline) {
t.Fatal("timed out")
}
time.Sleep(10 * time.Millisecond)
}
}
+3
View File
@@ -5,6 +5,7 @@ import (
"archive/zip"
"bytes"
"context"
"encoding/json"
"fmt"
"sort"
"sync"
@@ -195,8 +196,10 @@ func Install(t testing.TB, repo *MemRepo, store *MemStore, data []byte, state st
t.Fatal(err)
}
uri, _ := store.Put(ctx, p.Digest, data)
manifestJSON, _ := json.Marshal(p.Manifest)
_ = repo.SaveVersion(ctx, &types.PluginVersion{
PluginID: p.Manifest.ID, Version: p.Manifest.Version, Digest: p.Digest, PackageURI: uri,
Manifest: types.JSON(manifestJSON),
})
_ = repo.SavePlugin(ctx, &types.InstalledPlugin{
ID: p.Manifest.ID, ActiveVersion: p.Manifest.Version, DesiredState: state, Runtime: "declarative",
+3
View File
@@ -37,6 +37,9 @@ func (r *Reconciler) NodeName() string {
if host == "" {
host = "node"
}
if r.role != "" {
host = r.role + ":" + host
}
return host + "/" + r.instanceID[:8]
}
+38 -1
View File
@@ -105,6 +105,9 @@ type Reconciler struct {
activators []Activator
instanceID string
interval time.Duration
runtimes map[string]bool
accept func(*manifest.Manifest) bool
role string
mu sync.Mutex // serializes passes
loaded map[string]*Loaded
@@ -149,6 +152,14 @@ type Options struct {
Redis *redis.Client // nil on single-node deployments
Activators []Activator
Interval time.Duration
// Runtimes limits the node to plugins of these runtimes (a standalone
// plugin host loads host plugins only); empty means all.
Runtimes []manifest.RuntimeType
// Accept further limits the node by the active version's manifest (a
// plugin host that runs python plugins only); nil accepts all.
Accept func(*manifest.Manifest) bool
// Role names what the node is in status reports, e.g. "plugin-host".
Role string
}
// New creates a Reconciler.
@@ -159,9 +170,17 @@ func New(o Options) *Reconciler {
if o.CacheDir == "" {
o.CacheDir = DefaultCacheDir()
}
var runtimes map[string]bool
if len(o.Runtimes) > 0 {
runtimes = map[string]bool{}
for _, rt := range o.Runtimes {
runtimes[string(rt)] = true
}
}
return &Reconciler{
repo: o.Repo, store: o.Store, registry: o.Registry, cacheDir: o.CacheDir, rdb: o.Redis,
activators: o.Activators, instanceID: uuid.NewString(), interval: o.Interval,
runtimes: runtimes, accept: o.Accept, role: o.Role,
loaded: map[string]*Loaded{}, digests: map[string]string{}, retries: map[string]retry{},
status: map[string]Status{}, now: time.Now,
}
@@ -188,7 +207,8 @@ func (r *Reconciler) Reconcile(ctx context.Context) error {
var errs []error
desired := map[string]bool{}
for _, row := range rows {
if row.DesiredState != types.PluginStateEnabled {
if row.DesiredState != types.PluginStateEnabled || (r.runtimes != nil && !r.runtimes[row.Runtime]) ||
!r.accepts(ctx, row) {
continue
}
desired[row.ID] = true
@@ -215,6 +235,23 @@ func (r *Reconciler) Reconcile(ctx context.Context) error {
return errors.Join(errs...)
}
// accepts applies Options.Accept to a plugin's active version. A version
// that cannot be read is accepted, so ensure reports why.
func (r *Reconciler) accepts(ctx context.Context, row types.InstalledPlugin) bool {
if r.accept == nil {
return true
}
v, err := r.repo.GetVersion(ctx, row.ID, row.ActiveVersion)
if err != nil || v == nil {
return true
}
var m manifest.Manifest
if json.Unmarshal(v.Manifest, &m) != nil {
return true
}
return r.accept(&m)
}
// runtimeTarget identifies where a plugin runs beyond its package: a new
// remote URL or secret reloads the plugin like a new version would.
func runtimeTarget(row types.InstalledPlugin) string {
@@ -226,3 +226,35 @@ func TestReconcileReloadsOnRuntimeTargetChange(t *testing.T) {
t.Fatalf("Loaded carries secret %q", got)
}
}
// A standalone plugin host loads host plugins only.
func TestRuntimesLimitWhatANodeLoads(t *testing.T) {
ctx := context.Background()
repo, store, reg, act := plugintest.NewMemRepo(), &plugintest.MemStore{}, registry.New(), &recorder{}
r := New(Options{
Repo: repo, Store: store, Registry: reg, CacheDir: t.TempDir(), Activators: []Activator{act},
Runtimes: []manifest.RuntimeType{manifest.RuntimeHost}, Role: "plugin-host",
})
plugintest.Install(t, repo, store, plugintest.KitPackage(t, "1.0.0"), types.PluginStateEnabled)
if err := r.Reconcile(ctx); err != nil {
t.Fatal(err)
}
if len(act.calls) != 0 || len(r.Loaded()) != 0 {
t.Fatalf("a declarative plugin was loaded: %v", act.calls)
}
if !strings.HasPrefix(r.NodeName(), "plugin-host:") {
t.Fatalf("node name = %s", r.NodeName())
}
// Accept filters on the manifest: nothing is loaded, nothing reported.
r = New(Options{
Repo: repo, Store: store, Registry: registry.New(), CacheDir: t.TempDir(), Activators: []Activator{act},
Accept: func(m *manifest.Manifest) bool { return m.Version != "1.0.0" },
})
if err := r.Reconcile(ctx); err != nil || len(r.Loaded()) != 0 {
t.Fatalf("an unaccepted plugin was loaded: %v", err)
}
if _, ok := r.Status("acme.kit"); ok {
t.Fatal("an unaccepted plugin has no status here")
}
}
+7
View File
@@ -126,6 +126,13 @@ func (c *Client) do(ctx context.Context, method, path string, body []byte, reque
return resp, nil
}
// Raw sends one request as the client would and returns the answer
// unread, whatever its status: for gateways that relay a plugin's answers,
// streams included. The caller closes the body.
func (c *Client) Raw(ctx context.Context, method, path string, body []byte, requestID string) (*http.Response, error) {
return c.do(ctx, method, path, body, requestID)
}
// decodeError turns a non-2xx answer into a protocol error.
func decodeError(resp *http.Response) error {
b, _ := io.ReadAll(io.LimitReader(resp.Body, 64<<10))
+100
View File
@@ -0,0 +1,100 @@
# 插件
插件是 WeKnora 扩展能力的打包与分发单元:模型厂商、数据源、联网搜索、文档解析、技能、MCP 服务等都可以由插件提供。内置能力本身也以内置插件的形式登记,和第三方插件共用一套目录与开关。
两类角色分工如下:
| | 系统管理员 | 空间管理员 |
| --- | --- | --- |
| 做什么 | 安装、升级、回滚、卸载插件,填写平台配置 | 在本空间启用或停用插件,填写空间配置 |
| 入口 | 「设置 → 插件管理」 | 「设置 → 插件」 |
安装后的插件对所有空间可见,但**默认停用**,由各空间管理员自行启用。停用插件后,它的集成不再出现在类型列表中、不能新建,已有实例照常工作。
## 插件的运行方式
插件包是一个 `.wkp` 文件(根目录带 `plugin.yaml` 的 zip)。`runtime.type` 决定插件代码在哪里运行:
| 运行方式 | 说明 | 适用 |
| --- | --- | --- |
| `declarative` | 没有代码,由 WeKnora 解释清单 | 技能、模型厂商定义、远程 MCP |
| `host` | WeKnora 的插件宿主拉起的子进程,`kind` 为 `binary` 或 `python` | 自托管的代码插件 |
| `remote` | 插件作者自行部署的 HTTP 服务,安装时登记地址 | SaaS 类插件、独立团队维护的服务 |
代码插件与 WeKnora 之间走统一的扩展协议 v1(HTTP + JSON,流式同步用 NDJSON)。
## 安装插件
1. 以系统管理员身份打开「设置 → 插件管理」,点击「安装插件」。
2. 上传 `.wkp` 或填写下载地址。WeKnora 先解析插件包,展示它提供的能力、申请的权限(可访问的外部域名、Host API 权限等)和包摘要。
3. 确认后安装。安装请求带着审阅时的摘要,下载地址在此期间换了内容会被拒绝。
同一插件再次安装更高版本即升级,旧版本保留,可在详情中回滚。
### 远程插件
`runtime.type: remote` 的插件在安装时还需填写服务地址:
- 服务在内网时,需把地址加入 `SSRF_WHITELIST`。
- 安装后 WeKnora 生成一个签名密钥,**只显示这一次**。把它配置为服务的 `WEKNORA_PLUGIN_SECRET` 环境变量,服务据此确认请求来自 WeKnora。
- 插件详情中可以修改服务地址或轮换密钥。轮换后,服务换上新密钥之前的调用都会失败。
- 服务版本必须与安装的插件包一致,否则插件显示为异常、调用被拒绝。升级时同时升级服务与插件包。
## 运行代码插件
### 内嵌宿主(默认)
默认情况下,每个 app 进程自带插件宿主,在本机运行已安装的 `host` 插件:
- 插件进程只拿到 `WEKNORA_PLUGIN_*` 环境变量,看不到数据库密码等 WeKnora 密钥。
- 出网流量经宿主的出口代理,只放行插件在 `permissions.egress` 中声明的域名,内网地址一律拒绝。
- 宿主检查运行的进程与安装的包一致,健康检查失败时按退避重启,并在插件详情中显示为「异常」。
- Python 插件使用本机的 `python3`(可用 `WEKNORA_PLUGIN_PYTHON` 指定解释器)。Docker app 镜像已带 Python。
### 独立插件宿主
插件较多、希望把插件与 app 隔开,或 app 节点不便运行 Python 时,可以单独运行 `WeKnora plugin-host`:
- 它与 app 使用同一镜像,共用数据库、Redis、对象存储和 `SYSTEM_AES_KEY`(或 `JWT_SECRET`)。
- 启动后每 5 秒在 Redis 中通告自己运行的插件。app 按版本挑选最空闲的宿主,用由 `SYSTEM_AES_KEY` 派生的密钥签名调用它。
- 宿主停止时先撤回通告,再等进行中的调用结束。
- 插件详情的「节点状态」中,`plugin-host:` 开头的就是独立宿主。
docker compose 启用方式:在 `.env` 中设置
```bash
WEKNORA_PLUGIN_EMBEDDED_KINDS=none
WEKNORA_PLUGIN_HOST_API_URL=http://app:8080
```
然后执行:
```bash
docker compose --profile plugin-host up -d
```
Helm 设置 `pluginHost.enabled=true` 即可。使用本地存储(`STORAGE_TYPE=local`)时,插件包存放在 data-files 卷中,该卷需支持多 Pod 读写(ReadWriteMany)。
相关环境变量:
| 变量 | 作用于 | 说明 |
| --- | --- | --- |
| `WEKNORA_PLUGIN_EMBEDDED_KINDS` | app | app 自己运行的 kind(`binary`、`python`,逗号分隔);`none` 表示全部交给独立宿主。默认本机能跑的都跑 |
| `WEKNORA_PLUGIN_HOST_API_URL` | app、plugin-host | 不在 app 本机运行的插件(远程插件、独立宿主上的插件)回调 Host API 的地址 |
| `WEKNORA_PLUGIN_HOST_KINDS` | plugin-host | 宿主运行的 kind,默认本机能跑的都跑 |
| `WEKNORA_PLUGIN_HOST_ADDR` | plugin-host | 监听地址,默认 `:8081` |
| `WEKNORA_PLUGIN_HOST_URL` | plugin-host | app 访问该宿主的地址,默认 `http://<主机名>:<端口>`;Helm 中为 Pod IP |
| `WEKNORA_PLUGIN_PYTHON` | 两者 | Python 插件的解释器,默认 `python3` |
独立宿主需要 Redis,且所有节点的 `SYSTEM_AES_KEY`(或 `JWT_SECRET`)必须一致。app 不运行某个 kind、又没有配置独立宿主时,该 kind 的插件在插件详情中显示为加载失败,并说明原因。
## 开发插件
- **Go**:[pluginsdk](https://github.com/Tencent/WeKnora/tree/main/pluginsdk),含协议定义、SDK、客户端和一致性测试工具 `weknora-plugin-conformance`。
- **Python**:[pluginsdk/python](https://github.com/Tencent/WeKnora/tree/main/pluginsdk/python),Python 3.9+,只依赖标准库。
- **示例**:[examples/plugins](https://github.com/Tencent/WeKnora/tree/main/examples/plugins):
- `rss`:数据源连接器,Go;
- `subtitles`:文档解析器,Go,使用 Host API;
- `notebooks`:Jupyter 笔记本解析器,Python。
同一个插件既可以打包成 `host` 插件由 WeKnora 运行,也可以作为 `remote` 服务独立部署,代码不用改。