mirror of
https://github.com/Tencent/WeKnora.git
synced 2026-10-04 06:48:15 +08:00
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:
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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:
|
||||
|
||||
@@ -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)。
|
||||
|
||||
@@ -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 }}
|
||||
@@ -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)
|
||||
# -----------------------------------------------------------------------------
|
||||
|
||||
@@ -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))
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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())
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
|
||||
@@ -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,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()
|
||||
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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()
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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",
|
||||
|
||||
@@ -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]
|
||||
}
|
||||
|
||||
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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))
|
||||
|
||||
@@ -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` 服务独立部署,代码不用改。
|
||||
Reference in New Issue
Block a user