11 Commits

Author SHA1 Message Date
8313120d2a chore(version): release/v1.3.x 路牌推到 1.3.12(含驻留 inputch 残留修复 + Lua events 修复) 2026-09-13 20:34:09 +08:00
5f63f6ec6c fix(lua): events.subscribe 改用内部 Subscribe + 订阅生命周期(修死锁/use-after-close)
上一版 Lua 对齐引入的 sdk.events.subscribe 有两个真问题,本提交修掉:

1) 用了公共 SDK 的 Events(),但本内核从未注入 event subscriber
   (SetEventSubscriber 全仓无调用点),拿到永远是 nil ⇒ subscribe 只会
   返回 "events unavailable"。改用内部 SDK 的 s.Subscribe——内置插件走的就是
   这条路径(cli/webui/skillmgr 全用它)。

2) 自死锁:subscribe 会在 Lua 的 plugin.start(sdk) 回调里被调用,而
   luaPlugin.Start 正持有 p.mu;原实现在 subscribe 里再 lock p.mu 追加 subs,
   不可重入 ⇒ 测试实测 30s 超时。改用独立的 subsMu。

3) use-after-close:Stop 会 Close LState,但事件订阅此前无人取消,残留回调
   再触发就会碰已关的 L。现在:Stop 先(不持 p.mu,避免与 Bus.Publish
   锁序反转)取 subsMu 取消全部订阅,再置 closed 并关 L;事件回调持 p.mu 后
   先查 closed,已进入等锁的旧回调会直接返回。

4) plugin_mgr 访问补 nil 保护(部分单测构造的 SDK 不含 pluginMgr)。

回归:TestLuaEventsSubscribeAndStopCleanup——订阅后 Publish 命中、Stop 后
再 Publish 不 panic。全套 Lua 测试在 -race 下通过。
2026-09-13 20:33:25 +08:00
a13504be38 fix(resident): 销毁驻留子时注销其入站 inputch(child/<id>)—— 修登记表脏数据累积
根因:residentInboundChannel 在 create 时把 child/<id> 登记进共享登记表
(Plugin=resident, Owner=父),但 teardownResident 只把划入的 inputch
(如 timer)归还为未分配,从未注销这条入站登记。于是每次 create/destroy
都在登记表里留下一条脏记录,且随次数单调累积。

实测(HomeAgent 侧,HΔ-Kernel v1.3.10): 后
 仍列出 child/<id>,归属 main;HomeAgent 没有任何
工具能单独注销 inputch,只能重启 homed 清掉。

修法:teardownResident 里用纯函数 inboundChannelName 算出名字并 Unregister。
不能复用 residentInboundChannel——它有重新登记的副作用。
该路径同时覆盖 destroy / reclaim / StopResidents(父退出)。

测试:TestResident_LifecycleAndNoOrphans 增加两条断言——销毁后与父退出后
child/<id> 都必须从登记表消失。
2026-09-13 20:19:47 +08:00
e75859668b chore(version): release/v1.3.x 路牌推到 1.3.11(本版内容:Lua 语义对齐 SDK 1.3.0) 2026-09-13 20:04:47 +08:00
b6a66c57fe feat(lua): Lua 插件桥全量对齐 SDK 1.3.0(媒体/注入标志位/优先级/事件/通道注销)
内核 Lua 桥(internal/plugin/lua_plugin.go)此前停在 v0.8.0 时代能力面,
1.1/1.2/1.3 新增能力只在 Go 侧存在,而 PLUGIN_DEV.md 宣称『能力完全对齐』。
本补丁把 Lua 侧补齐到与公开 SDK 1.3.0 对齐:

- 1.1 媒体:memory.commit 支持 sentence_text/media_digests;
  doc.insert_with_media + attachments;text_memory.append attachments;
  set_tool_blocks / inject_input_media(_sync) / inject_interrupt_media。
- 1.2 注入语义:inject_input_sync(_opts)、六个 *_opts 变体
  (no_memory/context_policy/cleaner_name/priority);
  ToolDef/ChannelDef 解析 context_policy。
- 1.3 优先级与动态通道:priority 常量透传;unregister_output_channel。
- StageContext 暴露 reasoning_content/context_msgs/token_usage/memory/extra/errors。
- 新增 sdk.events.subscribe 与 sdk.plugin_mgr.*。
- sdk.lua mock 同步(单一事实源在 SDK 仓 sdk/lua/sdk.lua,内核副本由
  third_party/homeagent-sdk/scripts/sync-lua-sdk.sh 同步)。

契约测试(lua_surface_test.go):
- 守住内核内嵌 mock 与 SDK 仓事实源一致;
- 守住 mock 承诺的每个函数都有运行时 RawSetString 绑定;
- 覆盖 opts/media/attachments 解析与 context_policy 透传。

文档:中英 PLUGIN_DEV.md 的 Lua API 表补齐并改为『对齐至 SDK 1.3.0』。
2026-09-13 20:04:31 +08:00
1b49365d46 chore(version): release/v1.3.x 路牌推到 1.3.10(本版内容:输出次数不再受限的提示词 + type 缺省) 2026-09-13 16:04:11 +08:00
17ea7fd5f0 fix(prompt): 去掉"每轮只能发一次 output_send"的凭空限制;type 缺省即 text
用户现场指出:**qq 插件的输出通道判据太严了**(那条判据在插件侧,已单独修:
`output_send__qq` 不再受"当前会话身份"限制)。同时内核提示词里还有一条**同类的凭空限制**:

  「每轮对话**通常只需调用一次** output_send__{通道名} 即可完成回复。
    仅在内容确实超过单条消息长度上限(如 >4000 字)时才拆分为多条」

可设计上输出是 agent 的**主动调用**:收到一次输入后,可以往**任意(已授权的)通道**
发**任意多次**(分段播报、先回执后结论、同时通知多个通道都合法)。这句话会让模型
自己收起合理的多次输出 —— 而且它不是任何机制的要求,只是当初为压 output-loop 写的
措辞(真正的防环机制是"回执只回 ok、不回传富结果",那条保留)。

改法:
- 提示词改为明确授权:**输出次数与目标通道由你自己决定**,没有「一轮只能发一次」的限制;
  只保留两条真话:单条长度上限(超长拆完整段落)、别反复重发**完全相同**的内容。
- `output_send__*` 的 `type` 参数改为**可选**(缺省 text):判据该拦的是"不知道发什么",
  不是"没写众所周知的默认值"——此前缺 type 会直接失败并让模型重试一次。

判据 3 条(新增 `output_rules_test.go`):提示词不得含输出次数限制且必须显式授权 /
省略 type 时按 text 发送成功且 schema 的 required 只有 payload / 空 payload 仍被拦。
2026-09-13 16:04:11 +08:00
ffcfaf46e2 chore(version): release/v1.3.x 路牌推到 1.3.9(本版内容:驻留子轮次计数) 2026-09-13 15:41:03 +08:00
cd88b2dfe5 fix(resident): 子的「轮次」一直显示 0 —— info() 根本没填 Rounds
现象(用户线上联调实录 + 我复验):父侧 `resident_agents` 列出 `输入ch=[timer] 轮次=0
处理表=2` —— **处理表已有两条记录,轮次却是 0**,自相矛盾,容易被读成"子没干活"。

根因:`residentChild.info()` 构造 `ResidentInfo` 时**从来没有填过 Rounds 字段**
(结构体里有这个字段,于是永远输出零值),不是计数漏加。

改法:`Rounds = 已执行轮次数`(调度器执行计数,单调不减)。新增 `Agent.roundsExecuted()`
并写明为什么**不能**用 inputch 处理表条数当轮次:那张表记的是"当前上下文窗口内"的轮次,
压缩会清空(设计 §8.3)—— 用它会让父看到轮次倒退。

判据:inputch 路由测试里补一条断言 —— 子处理完输入后 `info().Rounds > 0`。
2026-09-13 15:41:02 +08:00
1d46c6c0f6 chore(version): release/v1.3.x 路牌推到 1.3.8(本版内容:inputch 归属路由) 2026-09-13 15:35:54 +08:00
4707b05498 fix(scheduler): inputch 划给子后输入只流向子 —— 补上「进内核之前」的输入路由
用户指出的语义(设计稿 §4.1 早已写明):
**inputch 是可分配资源**,「路由发生在**进内核之前**」—— 划给某个 agent 后,
该通道的输入**只流向那个 agent**;outputch 不同,授权是**非独占**的,
父依旧可以通过它发送内容。

而代码里 inputch 划拨只做了**登记**,没有做**路由**:
- 插件注入输入的 io 是**根 agent 的**(`cmd/homed` 里 `pluginReg.SetIOManager(iom)`);
- 唯一消费输入的是「该 io 自己的调度器」(`scheduler.go` 读 `a.io.InputChan()`);
- `ChannelRegistry.Assign` 只把 Owner 写进登记表,**没有任何转发动作**。

⇒ 现场表现(用户线上联调):子挂 `inputch=[timer]`,**timer 的输入却打在父身上**
(日志 `[agent] interrupt from timer/timer`),子侧 `轮次=0` 永远不动。
登记表里的 Owner 于是沦为标签。

改法(按 §4.1 把路由放回"进内核之前"):
- `IOManager` 增加 `InputRouter`(`SetInputRouter`),并把**五处直接入队**收口到
  `deliverInput`:`InjectInput` / `InjectInputSync` / `InjectInputTo` /
  `InjectInputSyncTo` / `InjectInterrupt`(排队与中断两条路都过路由)。
- 内核注入路由器 `Agent.routeInputByOwner`:查 inputch 的 Owner —— 归自己/未分配 ⇒
  本内核处理;归自己的某个驻留子 ⇒ `DeliverRouted` 交给它(**不再进父的队列**);
  归一个不存在的 agent ⇒ **不吞输入**,父兜底 + 留痕(吞掉输入比多处理一条更糟)。
- `DeliverRouted` 是"已路由"的投递口,不再二次路由(避免成环)。
- 同步输入的 `ResponseCh` 随事件一起走 ⇒ 回答由持有者写回同一回程(§4.3)。

判据(新增 6 条):
- io 层:被接管时排队/中断都**不入本内核队列**(且中断确实经过路由)/ 放行与未设
  路由器时与历史行为一致 / `DeliverRouted` 不再触发路由
- 内核层:划给子的 inputch 输入进**子**(子 Executed>0)且**父 Enqueued 不变** /
  归属到不存在的 agent 时父兜底(不吞)/ 未分配的 inputch 仍归父
2026-09-13 15:35:54 +08:00
21 changed files with 1581 additions and 72 deletions

1
.gitignore vendored
View File

@ -31,6 +31,7 @@ cmd/gui/dist/
third_party/homeagent-sdk/bin/
third_party/homeagent-sdk/tools/
third_party/homeagent-sdk/package/
third_party/homeagent-sdk/scripts/
third_party/homeagent-sdk/.gitignore
third_party/homeagent-sdk/README*
third_party/homeagent-sdk/example/

View File

@ -599,23 +599,26 @@ When running inside the kernel, `sdk.*` global variables are injected by the Go
### Lua SDK API
The `sdk.*` API of Lua plugins is fully aligned with external plugins (toolchain-built `plugin.bin` subprocesses): registration functions raise a Lua error on failure; data functions uniformly return `(result, err)` with `err == nil` on success. Subsystems not wired by the core (e.g. SocialAPI) return empty values instead of errors.
The `sdk.*` API of Lua plugins is aligned with external plugins (toolchain-built `plugin.bin` subprocesses) up to **SDK 1.3.0** (requires kernel **1.4.0+**, also backfilled by the Lua-alignment patch `v1.3.11`): registration functions raise a Lua error on failure; data functions uniformly return `(result, err)` with `err == nil` on success. Subsystems not wired by the core (e.g. SocialAPI) return empty values instead of errors.
> Historical note: the 1.11.3 media / inject-flags / priority capabilities were long available only on the Go side and were silently missing on the Lua side. They are now fully aligned, guarded by the contract test in `internal/plugin/lua_surface_test.go` (every function promised by the mock has a runtime binding).
**Registration**
| Function | Description |
|----------|-------------|
| `sdk.log(level, msg)` | Log output |
| `sdk.register_tool(name, def, handler)` | Register tool; `def` supports `description`, `parameters`, `no_memory`, `cleaner` |
| `sdk.register_tool(name, def, handler)` | Register tool; `def` supports `description`, `parameters`, `no_memory`, `context_policy` (`"none"`/`"prune"`), `cleaner` |
| `sdk.register_stage(stage, handler, scope)` | Register stage hook; `scope` is `nil`/`"global"` (default) or `"own_tools"` (fires only for `before_toolcall`/`after_toolcall` when the tool belongs to this plugin) |
| `sdk.register_api(name)` | Register API |
| `sdk.register_output_channel(name, caps, desc, def, handler)` | Register output channel; `def` supports `no_memory`, `cleaner` |
| `sdk.register_output_channel(name, caps, desc, def, handler)` | Register output channel; `def` supports `no_memory`, `context_policy`, `cleaner` |
| `sdk.register_input_channel(name, def)` | Register input channel; `def` as above |
| `sdk.unregister_output_channel(name)` | Unregister an output channel (for resource-bound channels, e.g. remote devices); returns `(nil, err)` |
| `sdk.set_auto_restart(enabled)` | Auto-restart the plugin after a crash |
**Stage hook context**
Stage handlers receive the full context (same as external plugins): `raw_message`, `user_id`, `group_id`, `phase`, `llm_text`, `final_text`, `no_memory`, `response` (when responded), `tool_calls`, `tool_results`.
Stage handlers receive the full context (same as external plugins): `raw_message`, `user_id`, `group_id`, `phase`, `llm_text`, `reasoning_content`, `final_text`, `no_memory`, `context_msgs`, `token_usage`, `memory`, `extra`, `errors`, `response` (when responded), `tool_calls`, `tool_results`.
**Stage writeback**: the `ctx` table passed to the handler is a reference — mutating writable fields inside the handler syncs back to the core `StageContext` (aligned with subprocess external-plugin capability):
@ -644,20 +647,37 @@ Writable fields: `raw_message`, `llm_text`, `final_text`, `user_id`, `group_id`,
| `sdk.inject_text(source, channel, text)` | Deliver text message |
| `sdk.inject_interrupt(source, channel, text)` | Interrupt delivery |
| `sdk.inject_text_no_memory(source, channel, text)` | Deliver without memory computation |
| `sdk.inject_text_opts` / `sdk.inject_interrupt_opts(source, channel, text, opts)` | Delivery with flags; `opts = { no_memory=bool, context_policy="none"|"prune", cleaner_name=string, priority="L1".."L3" }` |
| `sdk.inject_input_sync(source, channel, text)` | Inject synchronously and wait for this turn's reply; returns `(reply, err)`, reply is nil when there is none |
| `sdk.inject_input_sync_opts(source, channel, text, opts)` | Same, with flags |
| `sdk.inject_input_media(source, channel, text, blocks)` | Inject text + multimodal content blocks |
| `sdk.inject_input_media_opts(source, channel, text, blocks, opts)` | Same, with flags |
| `sdk.inject_input_media_sync` / `..._sync_opts(...)` | Synchronous media injection; returns `(reply, err)` |
| `sdk.inject_interrupt_media(source, channel, text, blocks)` | Interrupt delivery with media |
| `sdk.inject_interrupt_media_opts(source, channel, text, blocks, opts)` | Same, with flags |
| `sdk.set_tool_blocks(blocks)` | Set multimodal blocks carried by the next tool message (lets the model see images / hear audio) |
Each `blocks` item: `{ type="text", text="..." }`, `{ type="image_url", image_url={ url="...", detail="high" } }`, or `{ type="audio_url", audio_url={ url="..." } }`. An absent `opts` is the zero value (recorded in memory + no pruning), equivalent to the three-argument form.
**Data APIs (aligned with subprocess external plugins, all return `(result, err)`)**
| Sub-table | Functions |
|-----------|-----------|
| `sdk.memory.*` | `recall(query, depth)`, `commit({triples})`, `introspect()`, `merge(source, target)`, `purge(criteria, hard)` |
| `sdk.doc.*` | `query(text, top_k)`, `insert({id,title,content})`, `remove(id)`, `stats()` |
| `sdk.memory.*` | `recall(query, depth)`, `commit({triples})` (triple supports `subject/relation/object/confidence/subject_type/object_type/sentence_text/media_digests`), `introspect()`, `merge(source, target)`, `purge(criteria, hard)` |
| `sdk.doc.*` | `query(text, top_k)`, `insert({id,title,content})`, `insert_with_media(doc, attachments)`, `remove(id)`, `stats()` |
| `sdk.knowledge.*` | `search(query, limit)`, `add(tag, content)`, `list()` |
| `sdk.text_memory.*` | `append({role,content,timestamp,channel})` |
| `sdk.text_memory.*` | `append({role,content,timestamp,channel,attachments})` |
| `sdk.llm.*` | `list_sources()`, `set_source(name)`, `current_source()` |
| `sdk.social.*` (read-only) | `get_person(name)`, `get_network(name, depth)`, `get_trait(name, trait)`, `get_relations(name)`, `list_persons()` |
| `sdk.events.*` | `subscribe(event_type, handler)` → returns an unsubscribe function; handler receives `{type,source,timestamp,payload}` |
| `sdk.plugin_mgr.*` | `reload_one(name)`, `list_loaded()`, `is_disabled(name)` |
| `sdk.json.*` | `encode(val)`, `decode(str)` |
| `sdk.http.*` | `get(url)`, `post(url, body, content_type)` |
Each `attachments` item: `{ digest=, mime=, name=, data=<base64> }`; with `data` it is new content (stored in the content-addressed store), with only `digest` it references existing content.
> The `sdk.events.subscribe` callback runs on the kernel's event-publishing goroutine, and Lua is single-state + mutex-guarded — **do only lightweight forwarding inside the callback; never block**, or every call of this plugin will stall.
---
<img src="../../assets/branding/mascot-xiaozhai.webp" width="20" style="border-radius:50%;vertical-align:middle"> :

View File

@ -592,23 +592,26 @@ lua main.lua
### Lua SDK API
Lua 插件的 `sdk.*` API 与外部插件(工具链编译的 `plugin.bin` 子进程)能力完全对齐:注册类函数调用即时报错(抛 Lua error数据类函数统一返回 `(result, err)``err` 为 nil 表示成功。核心未装配的子系统(如 SocialAPI返回空值而非报错。
Lua 插件的 `sdk.*` API 与外部插件(工具链编译的 `plugin.bin` 子进程)能力对齐**SDK 1.3.0**(需内核 **1.4.0+**,也在 `v1.3.11` 的 Lua 对齐补丁中回填):注册类函数调用即时报错(抛 Lua error数据类函数统一返回 `(result, err)``err` 为 nil 表示成功。核心未装配的子系统(如 SocialAPI返回空值而非报错。
> 历史提醒1.11.3 的媒体/注入标志位/优先级能力曾长期只在 Go 侧Lua 侧静默缺失。现已全量对齐,并由 `internal/plugin/lua_surface_test.go` 的契约测试守住「mock 承诺的每个函数都有运行时绑定」。
**注册类**
| 函数 | 说明 |
|------|------|
| `sdk.log(level, msg)` | 日志输出 |
| `sdk.register_tool(name, def, handler)` | 注册工具;`def` 支持 `description``parameters``no_memory``cleaner` |
| `sdk.register_tool(name, def, handler)` | 注册工具;`def` 支持 `description``parameters``no_memory``context_policy``"none"`/`"prune"`)、`cleaner` |
| `sdk.register_stage(stage, handler, scope)` | 注册阶段钩子;`scope``nil`/`"global"`(默认)或 `"own_tools"`(仅 `before_toolcall`/`after_toolcall` 且工具属于本插件时触发) |
| `sdk.register_api(name)` | 注册 API |
| `sdk.register_output_channel(name, caps, desc, def, handler)` | 注册输出通道;`def` 支持 `no_memory``cleaner` |
| `sdk.register_output_channel(name, caps, desc, def, handler)` | 注册输出通道;`def` 支持 `no_memory``context_policy``cleaner` |
| `sdk.register_input_channel(name, def)` | 注册输入通道;`def` 同上 |
| `sdk.unregister_output_channel(name)` | 注销输出通道(随资源生灭的动态通道,如远程设备);返回 `(nil, err)` |
| `sdk.set_auto_restart(enabled)` | 崩溃时内核自动拉起插件 |
**阶段钩子上下文**
`register_stage` 的 handler 收到完整上下文(与外部插件一致):`raw_message``user_id``group_id``phase``llm_text``final_text``no_memory``response`(已响应时)、`tool_calls``tool_results`
`register_stage` 的 handler 收到完整上下文(与外部插件一致):`raw_message``user_id``group_id``phase``llm_text``reasoning_content``final_text``no_memory``context_msgs``token_usage``memory``extra``errors``response`(已响应时)、`tool_calls``tool_results`
**Stage 写回**handler 收到的 `ctx` 是引用 table——在 handler 内直接修改可写回字段并同步至内核 `StageContext`(与子进程外部插件能力对齐):
@ -637,20 +640,37 @@ end)
| `sdk.inject_text(source, channel, text)` | 投递文本消息 |
| `sdk.inject_interrupt(source, channel, text)` | 中断投递 |
| `sdk.inject_text_no_memory(source, channel, text)` | 免记忆投递 |
| `sdk.inject_text_opts` / `sdk.inject_interrupt_opts(source, channel, text, opts)` | 带标志位投递;`opts = { no_memory=bool, context_policy="none"|"prune", cleaner_name=string, priority="L1".."L3" }` |
| `sdk.inject_input_sync(source, channel, text)` | 同步注入并等本轮回复;返回 `(reply, err)`,无回复时 reply 为 nil |
| `sdk.inject_input_sync_opts(source, channel, text, opts)` | 同上带标志位 |
| `sdk.inject_input_media(source, channel, text, blocks)` | 注入文本 + 多模态内容块 |
| `sdk.inject_input_media_opts(source, channel, text, blocks, opts)` | 同上带标志位 |
| `sdk.inject_input_media_sync` / `..._sync_opts(...)` | 带媒体的同步注入;返回 `(reply, err)` |
| `sdk.inject_interrupt_media(source, channel, text, blocks)` | 带媒体的中断注入 |
| `sdk.inject_interrupt_media_opts(source, channel, text, blocks, opts)` | 同上带标志位 |
| `sdk.set_tool_blocks(blocks)` | 设置下一轮 tool message 携带的多模态内容块(模型据此看图/听音频) |
`blocks` 每项形如:`{ type="text", text="..." }``{ type="image_url", image_url={ url="...", detail="high" } }``{ type="audio_url", audio_url={ url="..." } }``opts` 缺省即零值(记入记忆 + 不裁剪),与三参数版本等价。
**数据类(与子进程外部插件对齐,均返回 `(result, err)`**
| 子表 | 函数 |
|------|------|
| `sdk.memory.*` | `recall(query, depth)``commit({triples})``introspect()``merge(source, target)``purge(criteria, hard)` |
| `sdk.doc.*` | `query(text, top_k)``insert({id,title,content})``remove(id)``stats()` |
| `sdk.memory.*` | `recall(query, depth)``commit({triples})`triple 支持 `subject/relation/object/confidence/subject_type/object_type/sentence_text/media_digests``introspect()``merge(source, target)``purge(criteria, hard)` |
| `sdk.doc.*` | `query(text, top_k)``insert({id,title,content})``insert_with_media(doc, attachments)``remove(id)``stats()` |
| `sdk.knowledge.*` | `search(query, limit)``add(tag, content)``list()` |
| `sdk.text_memory.*` | `append({role,content,timestamp,channel})` |
| `sdk.text_memory.*` | `append({role,content,timestamp,channel,attachments})` |
| `sdk.llm.*` | `list_sources()``set_source(name)``current_source()` |
| `sdk.social.*`(只读) | `get_person(name)``get_network(name, depth)``get_trait(name, trait)``get_relations(name)``list_persons()` |
| `sdk.events.*` | `subscribe(event_type, handler)` → 返回取消订阅函数handler 收到 `{type,source,timestamp,payload}` |
| `sdk.plugin_mgr.*` | `reload_one(name)``list_loaded()``is_disabled(name)` |
| `sdk.json.*` | `encode(val)``decode(str)` |
| `sdk.http.*` | `get(url)``post(url, body, content_type)` |
`attachments` 每项:`{ digest=, mime=, name=, data=<base64> }`;带 `data` 是新内容(落进内容寻址存储),只带 `digest` 是引用已有内容。
> `sdk.events.subscribe` 的回调在内核事件发布 goroutine 上执行,且 Lua 是单状态 + 互斥锁——**回调内只做轻量转发,不可阻塞**,否则会卡死本插件的全部调用。
---
<img src="../../assets/branding/mascot-xiaozhai.webp" width="20" style="border-radius:50%;vertical-align:middle"> :

View File

@ -325,7 +325,7 @@ func New(cfg AgentConfig) *Agent {
}
}
return &Agent{
a := &Agent{
id: cfg.ID,
startTime: time.Now(),
provider: cfg.Provider,
@ -376,6 +376,14 @@ func New(cfg AgentConfig) *Agent {
noMergeMarkers: make(map[string]int),
lastInput: make(map[string]time.Time),
}
// 输入路由inputch 是可分配资源,划给某个 agent 后输入**只**流向那个 agent
// (设计 §4.1「路由发生在进内核之前」。io 层不认识 agent所以在这里把路由器
// 注入进去:插件注入输入时先问它,被别的 agent 接管就不再进本内核队列。
if a.io != nil {
a.io.SetInputRouter(a.routeInputByOwner)
}
return a
}
// SetSkillIndexProvider 注入技能索引提供者skillmgr 插件加载后由 main 接线)。

View File

@ -0,0 +1,44 @@
package core
import (
"log"
agentIO "gitcode.com/JianFeeeee/HomeAgent/internal/agent/io"
)
// routeInputByOwner 实现**输入路由**inputch 是最基本的输入路由单位,
// 划给某个 agent 之后,该通道的输入**只流向那个 agent**,本内核看不到它
// docs/zh/resident-subagent-design.md §4.1:「路由发生在进内核之前」)。
//
// 为什么必须在进内核之前做:插件注入输入的收口是**父**的 IOManager
// `cmd/homed` 里 pluginReg 拿到的就是它),而父的内核是该 io 唯一的消费者。
// 如果不按归属路由,登记表里的 Owner 就只是个标签 —— 现场表现正是如此:
// 子挂着 `inputch=[timer]`timer 的输入却打在父身上,子的轮次永远是 0。
//
// 返回 true = 本次注入已被"持有该 inputch 的 agent"接管,本内核不再处理。
//
// 已知边界:路由只在本 agent 的**直接**驻留子里找。若孙辈的 inputch 由子划拨,
// 而插件注入打在根 io 上,根解析不到那个 owner ⇒ 兜底给根处理(有日志)。
// 这一层要等"孙辈 + 根可见的 agent 表"再收口,此处不静默丢输入。
func (a *Agent) routeInputByOwner(evt *agentIO.InputEvent, isInterrupt bool) bool {
if evt == nil || evt.OutputChannel == "" || a.io == nil {
return false
}
entry, ok := a.io.LookupInputChannel(evt.OutputChannel)
if !ok || entry.Owner == "" || entry.Owner == string(a.id) {
return false // 未分配 / 归自己 ⇒ 本内核处理
}
a.residentMu.Lock()
rc := a.residents[entry.Owner]
a.residentMu.Unlock()
if rc == nil || rc.agent == nil || rc.agent.io == nil {
// 归属到一个不存在(或已销毁、登记表尚未归还)的 agent
// **不吞输入** —— 由本内核兜底处理并留痕。吞掉一条输入比多处理一条更糟:
// 用户会看到"消息发出去了却没人理",而日志里什么都没有。
log.Printf("[route] inputch %s 归属 %s 无对应 agent输入由 %s 兜底", evt.OutputChannel, entry.Owner, a.id)
return false
}
rc.agent.io.DeliverRouted(evt, isInterrupt)
return true
}

View File

@ -0,0 +1,85 @@
package core
import (
"path/filepath"
"testing"
agentIO "gitcode.com/JianFeeeee/HomeAgent/internal/agent/io"
)
// 输入路由是**独占**的inputch 划给子之后,该通道的输入只流向子,父不再收到。
//
// 现场缺陷(用户线上联调实录):子挂着 inputch=[timer]timer 的输入却打在父身上
// (日志 `[agent] interrupt from timer/timer`),子的轮次永远是 0 —— 因为
// `Assign` 只把 Owner 写进登记表,注入路径根本没有按归属路由。
func TestResident_InputchRoutingIsExclusive(t *testing.T) {
parent, _, dir := newRootForResidents(t)
reg := parent.io.ChannelRegistry()
if err := reg.Register(agentIO.InputChannel{Name: "sub/in", Plugin: "sub"}); err != nil {
t.Fatal(err)
}
spawnTestResident(t, parent, dir, "r-route", "sub/in")
child := parent.residents["r-route"].agent
before := parent.DumpScheduler().Stats.Enqueued
// 插件往"已划给子"的 inputch 投输入
parent.io.InjectTextTo("plugin-sub", "sub/in", "去查一下这个")
waitFor(t, "子处理了划给它的输入", func() bool {
if child.DumpScheduler().Stats.Executed > 0 {
return true
}
return parent.residents["r-route"].info().TableSize > 0
})
if got := parent.DumpScheduler().Stats.Enqueued; got != before {
t.Fatalf("划给子的 inputch父不应再入队before=%d after=%d", before, got)
}
// 轮次必须真的涨:此前 info() 根本没填 Rounds ⇒ 父永远读到 0
// (现场:子处理表已有 2 条,轮次却显示 0被误判成"子没干活")。
if got := parent.residents["r-route"].info().Rounds; got <= 0 {
t.Fatalf("子处理的轮次应 > 0实际 %d", got)
}
}
// 归属到一个不存在(或已销毁)的 agent 时**不吞输入**:父兜底处理。
// 吞掉一条输入比多处理一条更糟 —— 用户会看到"消息发出去了却没人理",日志里什么都没有。
func TestResident_InputchRoutingFallsBackWhenOwnerMissing(t *testing.T) {
parent, _, dir := newRootForResidents(t)
reg := parent.io.ChannelRegistry()
if err := reg.Register(agentIO.InputChannel{Name: "ghost/in", Plugin: "ghost"}); err != nil {
t.Fatal(err)
}
// 故意划给一个不存在的 agent id
if err := reg.Assign("ghost/in", "no-such-agent", 0); err != nil {
t.Fatal(err)
}
_ = dir
before := parent.DumpScheduler().Stats.Enqueued
parent.io.InjectTextTo("plugin-ghost", "ghost/in", "兜底测试")
waitFor(t, "父兜底处理了无人认领的输入", func() bool {
return parent.DumpScheduler().Stats.Enqueued > before
})
}
// 未划拨的 inputchOwner 为空)仍然由父处理 —— 路由不能把默认路径也改掉。
func TestResident_UnassignedInputchStaysWithParent(t *testing.T) {
parent, _, dir := newRootForResidents(t)
reg := parent.io.ChannelRegistry()
if err := reg.Register(agentIO.InputChannel{Name: "own/in", Plugin: "own"}); err != nil {
t.Fatal(err)
}
// 造一个子在跑,确保路由逻辑是"有子存在"的情形
spawnTestResident(t, parent, dir, "r-other", "own/in")
_ = filepath.Join(dir, "residents")
// 把通道退还给父(未分配)
if err := reg.Assign("own/in", "", 0); err != nil {
t.Fatal(err)
}
before := parent.DumpScheduler().Stats.Enqueued
parent.io.InjectTextTo("plugin-own", "own/in", "还是我的")
waitFor(t, "未分配的 inputch 仍由父处理", func() bool {
return parent.DumpScheduler().Stats.Enqueued > before
})
}

View File

@ -13,8 +13,13 @@ func (a *Agent) executeOutputSendTool(tc agentAPI.ToolCall) string {
channel := strings.TrimPrefix(tc.Name, "output_send__")
payload, _ := tc.Arguments["payload"].(string)
rawType, _ := tc.Arguments["type"].(string)
if channel == "" || payload == "" || rawType == "" {
return "工具名称格式: output_send__{channel}payload 和 type 不能为空"
if channel == "" || payload == "" {
return "工具名称格式: output_send__{channel}payload 不能为空"
}
// type 缺省按 text 处理:绝大多数输出就是文本,让模型为"省略一个默认值"付一次
// 失败重试没有意义(判据该拦的是"不知道发什么",不是"没写众所周知的默认值")。
if rawType == "" {
rawType = "text"
}
// 授权闸(纵深防御):模型可能凭名字直接调未授权的输出门。
if !a.IsOutputAllowed(channel) {

View File

@ -0,0 +1,100 @@
package core
import (
"strings"
"testing"
agentAPI "gitcode.com/JianFeeeee/HomeAgent/internal/agent/api"
)
// 设计口径:输出是 agent 的**主动调用** —— 收到一次输入后,可以往任意(已授权的)
// 通道发**任意多次**(分段播报、先回执后结论、同时通知多个通道都合法)。
//
// 这条判据钉住的是"提示词里不得出现输出次数限制"。此前 `tooldefs.go` 里写着
// 「每轮对话通常只需调用一次 output_send__{通道名} 即可完成回复」—— 一条凭空的限制,
// 会让模型自己收起合理的多次输出(用户现场指出)。
func TestSystemPromptDoesNotRestrictOutputCount(t *testing.T) {
parent, _, _ := newRootForResidents(t)
defer parent.Stop()
prompt := parent.buildSystemPrompt("", "你好")
banned := []string{
"只需调用一次",
"只能调用一次",
"通常只需调用",
"不要重复发送",
}
for _, b := range banned {
if strings.Contains(prompt, b) {
t.Fatalf("系统提示词里仍有输出次数限制 %q —— 设计上次数不限", b)
}
}
if !strings.Contains(prompt, "输出次数与目标通道由你自己决定") {
t.Fatal("系统提示词应明确「输出次数与目标通道由你自己决定」")
}
if !strings.Contains(prompt, "没有任何「一轮只能发一次」的限制") {
t.Fatal("系统提示词应显式否认「一轮只能发一次」")
}
}
// 输出工具的 type 可省略,缺省按 text 处理(判据该拦的是"不知道发什么"
// 不是"没写众所周知的默认值")。
func TestOutputSendTypeDefaultsToText(t *testing.T) {
parent, _, _ := newRootForResidents(t)
defer parent.Stop()
dev := &outputTestDevice{name: "fakeout"}
if err := parent.io.RegisterDevice(dev); err != nil {
t.Fatal(err)
}
out := parent.executeOutputSendTool(agentAPI.ToolCall{
Name: "output_send__fakeout",
Arguments: map[string]interface{}{"payload": "只给 payload不给 type"},
})
if out != "ok" {
t.Fatalf("省略 type 时应默认 text 并发送成功,得到 %q", out)
}
if len(dev.sent) != 1 {
t.Fatalf("通道应收到 1 次输出,得到 %d", len(dev.sent))
}
if args, _ := dev.sent[0]["args"].(map[string]interface{}); args["type"] != "text" {
t.Fatalf("缺省类型应为 text实际 %v", args["type"])
}
// 工具 schemarequired 只应含 payload
var found bool
for _, td := range parent.buildToolDefs() {
entry, _ := td.(map[string]interface{})
fn, _ := entry["function"].(map[string]interface{})
if n, _ := fn["name"].(string); n != "output_send__fakeout" {
continue
}
found = true
params, _ := fn["parameters"].(map[string]interface{})
req, _ := params["required"].([]string)
if len(req) != 1 || req[0] != "payload" {
t.Fatalf("output_send 的 required 应只有 payload实际 %v", req)
}
}
if !found {
t.Fatal("未生成 output_send__fakeout 工具")
}
}
// 空 payload 仍应被拦(这条判据是对的:不知道发什么不能放过)。
func TestOutputSendStillRequiresPayload(t *testing.T) {
parent, _, _ := newRootForResidents(t)
defer parent.Stop()
if err := parent.io.RegisterDevice(&outputTestDevice{name: "fakeout"}); err != nil {
t.Fatal(err)
}
out := parent.executeOutputSendTool(agentAPI.ToolCall{
Name: "output_send__fakeout",
Arguments: map[string]interface{}{"type": "text"},
})
if !strings.Contains(out, "payload 不能为空") {
t.Fatalf("空 payload 应被拦,得到 %q", out)
}
}

View File

@ -198,8 +198,13 @@ func (a *Agent) SpawnResident(opts ResidentOptions) (ResidentInfo, error) {
func (a *Agent) residentParentSource() string { return "parent/" + string(a.id) }
// residentInboundChannel 是"父接收某个子的消息"的 inputch 名(登记进登记表可见)。
//
// inboundChannelName 只拼名字,不产生副作用(注销路径要用它算出同一个名字,
// 不能再去调 residentInboundChannel——那会顺手把刚摘掉的登记又写回去
func inboundChannelName(childID string) string { return "child/" + childID }
func (a *Agent) residentInboundChannel(childID string) string {
ch := "child/" + childID
ch := inboundChannelName(childID)
if reg := a.io.ChannelRegistry(); reg != nil {
// 归属父自己:它是父的入站 inputch。
_ = reg.Register(agentIO.InputChannel{Name: ch, Plugin: "resident", Owner: string(a.id)})
@ -246,6 +251,14 @@ func (a *Agent) teardownResident(rc *residentChild) {
for _, ch := range rc.inputChs {
_ = reg.Assign(ch, "", 0)
}
// 注销"父接收该子消息"的入站 inputchchild/<id>)。
//
// 它由 residentInboundChannel 在 create 时登记Owner=父),销毁时必须
// 一并摘掉:登记表是共享的、按 name 全局唯一,残留会随 create/destroy
// 次数单调累积脏数据。实测destroy 后 child/<id> 仍挂在根 agent 名下,
// 而外部没有任何工具能单独注销 inputch只能重启 homed 清。
// 注意用纯函数算名字,不要再走 residentInboundChannel会重新登记
reg.Unregister(inboundChannelName(rc.id))
}
}
@ -504,6 +517,13 @@ func (rc *residentChild) info() ResidentInfo {
ID: rc.id, State: state, InputChs: append([]string(nil), rc.inputChs...),
AllowedOutputs: append([]string(nil), rc.allowed...),
ContextFull: full, CreatedAt: rc.createdAt, TableSize: len(table),
// Rounds = 子**已执行的轮次数**(调度器的执行计数,单调不减)。
//
// 此前这里根本没填这个字段 ⇒ 父看到的永远是 `轮次=0`,与"处理表已有 N 条"
// 自相矛盾(现场:子明明处理了两轮,父读到 rounds=0误判成"子没干活")。
// 注意它**不等于** len(table):处理表记的是"当前上下文窗口内"的轮次,
// 压缩会清空§8.3),所以窗口内的条数会被重置,而轮次总数不会。
Rounds: rc.agent.roundsExecuted(),
}
if len(table) > 0 {
info.Table = table

View File

@ -142,6 +142,11 @@ func TestResident_LifecycleAndNoOrphans(t *testing.T) {
if ch, _ := reg.Lookup("sub/in"); ch.Owner != "" {
t.Fatalf("销毁后 inputch 应回到未分配:%+v", ch)
}
// 父的入站 inputchchild/<id>)必须在销毁时一并注销,否则登记表残留脏数据——
// HomeAgent 实测destroy 后 child/<id> 仍挂在根 agent 名下,且无工具可单独注销。
if inbound, ok := reg.Lookup("child/child-1"); ok {
t.Fatalf("销毁后父的入站 inputch 应被注销,实际残留:%+v", inbound)
}
if err := parent.DestroyResident("child-1"); err == nil {
t.Fatal("重复销毁应报错")
}
@ -159,6 +164,9 @@ func TestResident_LifecycleAndNoOrphans(t *testing.T) {
if _, err := osStat(filepath.Join(dir, "residents", id)); err == nil {
t.Fatalf("子 %s 的 temp 目录应被丢弃", id)
}
if _, ok := reg.Lookup(inboundChannelName(id)); ok {
t.Fatalf("父退出后子 %s 的入站 inputch 应被注销", id)
}
}
}

View File

@ -745,6 +745,18 @@ func newKernelInterruptTask(evt *agentIO.InputEvent) *Task {
}
// DumpScheduler 返回调度器的原子快照(供状态页/测试断言)。
// roundsExecuted 返回本 agent 已执行的轮次数(供驻留子状态面展示)。
//
// 一轮 = 一次被执行的输入(排队与中断都算)。为什么不用 inputch 处理表的条数:
// 那张表记的是"当前上下文窗口内"的轮次,压缩会清空(设计 §8.3)——
// 拿它当轮次会让父看到轮次倒退。
func (a *Agent) roundsExecuted() int {
if a.sched == nil {
return 0
}
return int(a.DumpScheduler().Stats.Executed)
}
func (a *Agent) DumpScheduler() SchedulerSnapshot {
if a.sched == nil {
return SchedulerSnapshot{}

View File

@ -102,7 +102,14 @@ func (a *Agent) buildSystemPrompt(memContext string, userInput string) string {
prompt += "- 同步通道webui / cli / 终端):直接返回纯文本,内核会把文本交给等待方显示,无需调用工具。\n"
prompt += "- 异步通道qq / wechat / 群聊等):返回纯文本**【不会】**自动送达用户,必须调用 output_send__{通道名} 工具(注意 meta 里带上正确的 user_id 或 group_id才能真正把消息发出去。\n"
prompt += "- 不确定当前通道的发送方式时,先用 output_send__{通道名}_help 查看该通道的 meta 格式和 type 枚举,再决定。\n"
prompt += "- 每轮对话**通常只需调用一次** output_send__{通道名} 即可完成回复。仅在内容确实超过单条消息长度上限(如 >4000 字)时才拆分为多条;拆分时每条应是完整段落,不要碎片化。\n"
// ❗这里**不得**限制"一轮只能发一次"。设计上输出是 agent 的**主动调用**
// 收到一次输入后,可以往**任意(已授权的)通道**发**任意多次**(分段播报、
// 先回执后结论、同时通知多个通道都合法)。此前这里写着"每轮对话通常只需调用
// 一次 output_send"——那是一条**凭空的限制**,会让模型自己收起合理的多次输出。
// 真正需要提醒的只有两件事:单条长度上限(超长拆成完整段落)与"别反复重发
// 完全相同的内容"(自律,不是判据)。
prompt += "- **输出次数与目标通道由你自己决定**:一次输入可以对同一通道发多条(先回执后结论、分步播报、分段长文),也可以同时发到多个通道(例如同时通知 webui 与 qq。**没有任何「一轮只能发一次」的限制。**\n"
prompt += "- 输出时只需注意两点:单条消息的长度上限(超长就拆成完整段落,不要碎片化);别反复重发**完全相同**的内容(那是浪费,不是限制)。\n"
prompt += "- 需要多步执行的长任务:**必须先**向当前对话通道发一条确认消息告诉用户已收到(异步通道用输出门工具,同步通道直接返回文本),**然后再**执行具体排查工具。确认消息不代表任务完成,发出后仍需继续执行实际工具并最终汇报结果。\n"
prompt += "- 用户从其他渠道发来「在哪里/怎么样了」这类追问时,先回忆上次任务的通道与上下文,再回同一通道。"
@ -624,7 +631,7 @@ func (a *Agent) buildToolDefs() []interface{} {
"type": "function",
"function": map[string]interface{}{
"name": "output_send__" + ch.Name,
"description": desc + "。能力: " + capStr + "。payload 为消息载荷meta 为 JSON 发送元数据type 为载荷类型。用 _help 查看 meta 格式 type 枚举。",
"description": desc + "。能力: " + capStr + "。payload 为消息载荷type 默认 text可省略meta 为 JSON 发送元数据。用 _help 查看 meta 格式 type 枚举。",
"parameters": map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{
@ -638,10 +645,10 @@ func (a *Agent) buildToolDefs() []interface{} {
},
"type": map[string]interface{}{
"type": "string",
"description": "载荷类型,用 channel._help 查看支持的枚举值",
"description": "载荷类型,默认 text其它枚举用 channel._help 查看",
},
},
"required": []string{"payload", "type"},
"required": []string{"payload"},
},
},
})

View File

@ -118,6 +118,15 @@ type IOManager struct {
// 回退只解决"看得见",不解决"能不能用"。
parent *IOManager
// inputRouter 决定一条输入是否被"别的 agent"接管(返回 true = 已接管)。
//
// 为什么放在 ioinputch 是**最基本的输入路由单位**,而**路由发生在进内核之前**
// docs/zh/resident-subagent-design.md §4.1)。插件注入输入的收口就在这里,
// 所以路由必须在这里生效 —— inputch 划给某个 agent 后,输入**只流向那个 agent**
// 本内核根本看不到它。io 层不认识 agent路由器由内核注入
// (见 core.Agent.routeInputByOwner
inputRouter InputRouter
// toolBlocks插件工具注入多模态内容块process.go 在下一条 tool message 时消费。
// 用 interface{}[] 避免 import api.ContentBlock 导致的循环依赖。
toolBlocksMu sync.Mutex
@ -134,6 +143,47 @@ func NewIOManager() *IOManager {
}
}
// InputRouter 是输入路由器的签名。
//
// evt 待投递的输入事件OutputChannel 即它的 inputch
// isInterrupt 该输入是中断还是排队(两者都要按归属路由)
// 返回 true = 已被别的 agent 接管,本内核不再处理
type InputRouter func(evt *InputEvent, isInterrupt bool) bool
// SetInputRouter 注入输入路由器nil = 不路由,行为与以前完全一致)。
func (m *IOManager) SetInputRouter(r InputRouter) {
m.mu.Lock()
m.inputRouter = r
m.mu.Unlock()
}
// deliverInput 是**本内核**接收一条外部输入的收口:先按 inputch 归属路由,
// 被别的 agent 接管就不进本内核队列(划给子的 inputch父不再收到 —— 这是「划拨」
// 的语义,不是"父也顺便看一眼")。
func (m *IOManager) deliverInput(evt *InputEvent, isInterrupt bool) {
m.mu.RLock()
router := m.inputRouter
m.mu.RUnlock()
if router != nil && router(evt, isInterrupt) {
return
}
m.pushLocal(evt, isInterrupt)
}
// DeliverRouted 把**已被路由**的事件放进本内核队列(不再二次路由)。
// 由路由器实现调用:父把输入交给持有该 inputch 的子。
func (m *IOManager) DeliverRouted(evt *InputEvent, isInterrupt bool) {
m.pushLocal(evt, isInterrupt)
}
func (m *IOManager) pushLocal(evt *InputEvent, isInterrupt bool) {
if isInterrupt {
m.interruptCh <- evt
return
}
m.inputCh <- evt
}
// SetParentIO 设置上级 IOManagernil 表示无上级,行为与以前完全一致)。
// 见 parent 字段的说明:用于驻留子继承父的输出通道/设备视图。
func (m *IOManager) SetParentIO(p *IOManager) {
@ -232,50 +282,52 @@ func (m *IOManager) StopAll() {
}
func (m *IOManager) InjectInput(source string, eventType string, payload map[string]interface{}) {
m.inputCh <- &InputEvent{
m.deliverInput(&InputEvent{
RequestID: m.nextRequestID(),
Source: source,
Type: eventType,
Payload: payload,
OutputChannel: source,
}
}, false)
}
func (m *IOManager) InjectInputSync(source string, eventType string, payload map[string]interface{}) *OutputEvent {
ch := make(chan *OutputEvent, 1)
m.inputCh <- &InputEvent{
m.deliverInput(&InputEvent{
RequestID: m.nextRequestID(),
Source: source,
Type: eventType,
Payload: payload,
ResponseCh: ch,
OutputChannel: source,
}
}, false)
// 被路由走时,回答由持有该 inputch 的 agent 写进同一个 ResponseCh
//§4.3:同步输入的回程是事前定好的)——所以这里照常等待。
return <-ch
}
// InjectInputTo 注入输入事件并指定输出通道
func (m *IOManager) InjectInputTo(source, outputChannel, eventType string, payload map[string]interface{}) {
m.inputCh <- &InputEvent{
m.deliverInput(&InputEvent{
RequestID: m.nextRequestID(),
Source: source,
Type: eventType,
Payload: payload,
OutputChannel: outputChannel,
}
}, false)
}
// InjectInputSyncTo 注入输入事件(同步等待)并指定输出通道
func (m *IOManager) InjectInputSyncTo(source, outputChannel, eventType string, payload map[string]interface{}) *OutputEvent {
ch := make(chan *OutputEvent, 1)
m.inputCh <- &InputEvent{
m.deliverInput(&InputEvent{
RequestID: m.nextRequestID(),
Source: source,
Type: eventType,
Payload: payload,
ResponseCh: ch,
OutputChannel: outputChannel,
}
}, false)
return <-ch
}
@ -370,13 +422,13 @@ func (m *IOManager) InjectInterrupt(source, channel string, payload map[string]i
payload = map[string]interface{}{}
}
evtType, _ := payload["type"].(string)
m.interruptCh <- &InputEvent{
m.deliverInput(&InputEvent{
RequestID: m.nextRequestID(),
Source: source,
Type: evtType,
Payload: payload,
OutputChannel: channel,
}
}, true)
}
func (m *IOManager) InjectInterruptText(source, channel, text string) {

View File

@ -0,0 +1,74 @@
package io
import "testing"
// 输入路由:路由器说"已被别的 agent 接管"时,事件**不得**进本内核队列。
//
// 语义(设计 §4.1inputch 是可分配资源,划给某个 agent 后输入只流向它 ——
// "父也顺便看到一份"是错的。
func TestInputRouter_TakesOverExclusively(t *testing.T) {
m := NewIOManager()
var got []*InputEvent
var sawInterrupt bool
m.SetInputRouter(func(evt *InputEvent, isInterrupt bool) bool {
got = append(got, evt)
sawInterrupt = sawInterrupt || isInterrupt
return true // 全部接管
})
m.InjectInputTo("plugin-x", "sub/in", "text", map[string]interface{}{"content": "a"})
m.InjectInterruptTextOpts("plugin-x", "sub/in", "b", InjectOptions{})
if len(got) < 2 {
t.Fatalf("路由器应被调用(含中断路径),实际 %d 次", len(got))
}
if n := len(m.InputChan()); n != 0 {
t.Fatalf("被接管的排队输入不得进本内核队列,实际 %d 条", n)
}
if !sawInterrupt {
t.Fatal("中断注入也必须经过路由(否则中断会绕过 inputch 归属直投父)")
}
if got[0].OutputChannel != "sub/in" {
t.Fatalf("路由器应拿到事件的 inputch得到 %q", got[0].OutputChannel)
}
}
// 路由器放行(返回 false或未设置时行为与以前完全一致。
func TestInputRouter_PassthroughKeepsOldBehaviour(t *testing.T) {
m := NewIOManager()
calls := 0
m.SetInputRouter(func(evt *InputEvent, isInterrupt bool) bool { calls++; return false })
m.InjectInputTo("plugin-x", "sub/in", "text", map[string]interface{}{"content": "a"})
if calls != 1 {
t.Fatalf("路由器应被调用一次,实际 %d", calls)
}
if n := len(m.InputChan()); n != 1 {
t.Fatalf("放行的输入应进本内核队列,实际 %d 条", n)
}
if _, ok := <-m.InputChan(); !ok {
t.Fatal("队列应可读")
}
// 未设路由器:直接入队(历史行为)
m2 := NewIOManager()
m2.InjectInterruptText("plugin-x", "sub/in", "c")
if n := len(m2.InputInterruptChan()); n != 1 {
t.Fatalf("未设路由器时中断应直接入队,实际 %d 条", n)
}
}
// DeliverRouted 是不再二次路由的投递口(路由器实现把事件交给持有者)。
func TestDeliverRouted_SkipsSecondRouting(t *testing.T) {
m := NewIOManager()
routerCalls := 0
m.SetInputRouter(func(evt *InputEvent, isInterrupt bool) bool { routerCalls++; return true })
m.DeliverRouted(&InputEvent{OutputChannel: "sub/in"}, false)
if routerCalls != 0 {
t.Fatalf("DeliverRouted 不应再触发路由(会成环),实际 %d 次", routerCalls)
}
if n := len(m.InputChan()); n != 1 {
t.Fatalf("应已入队,实际 %d 条", n)
}
}

View File

@ -69,6 +69,65 @@ function sdk.inject_text_no_memory(source, channel, text)
print("[lua-plugin] inject_text_no_memory: " .. tostring(source))
end
-- !impl
-- opts: { no_memory=bool, context_policy="none"|"prune", cleaner_name=string, priority="L1".."L3" }
-- 零值/缺省 = 记入记忆 + 不裁剪(与三参数版本等价)。
function sdk.inject_text_opts(source, channel, text, opts)
print("[lua-plugin] inject_text_opts: " .. tostring(source))
end
-- !impl
function sdk.inject_interrupt_opts(source, channel, text, opts)
print("[lua-plugin] inject_interrupt_opts: " .. tostring(source))
end
-- !impl
-- 同步注入:等待本轮回复 -> (reply, err);无回复时 reply 为 nil。
function sdk.inject_input_sync(source, channel, text) return nil, nil end
-- !impl
function sdk.inject_input_sync_opts(source, channel, text, opts) return nil, nil end
-- !impl
-- blocks: ContentBlock 数组,见 sdk.inject_input_media。
-- 设置下一轮 tool message 携带的多模态内容块(模型据此看图/听音频)。
function sdk.set_tool_blocks(blocks)
print("[lua-plugin] set_tool_blocks: " .. tostring(blocks and #blocks or 0))
end
-- !impl
-- blocks 每项:{ type="text", text="..." }
-- | { type="image_url", image_url={ url="...", detail="high" } }
-- | { type="audio_url", audio_url={ url="..." } }
function sdk.inject_input_media(source, channel, text, blocks)
print("[lua-plugin] inject_input_media: " .. tostring(source))
end
-- !impl
function sdk.inject_input_media_opts(source, channel, text, blocks, opts)
print("[lua-plugin] inject_input_media_opts: " .. tostring(source))
end
-- !impl
function sdk.inject_input_media_sync(source, channel, text, blocks) return nil, nil end
-- !impl
function sdk.inject_input_media_sync_opts(source, channel, text, blocks, opts) return nil, nil end
-- !impl
function sdk.inject_interrupt_media(source, channel, text, blocks)
print("[lua-plugin] inject_interrupt_media: " .. tostring(source))
end
-- !impl
function sdk.inject_interrupt_media_opts(source, channel, text, blocks, opts)
print("[lua-plugin] inject_interrupt_media_opts: " .. tostring(source))
end
-- !impl
-- 注销输出通道(随资源生灭的动态通道,如远程设备)。返回 (nil, err)。
function sdk.unregister_output_channel(name) return nil, nil end
-- !impl
-- enabled: true/false崩溃时内核自动拉起
function sdk.set_auto_restart(enabled)
@ -101,6 +160,9 @@ function sdk.doc.query(text, top_k) return {} end
-- doc: { id=, title=, content= }
function sdk.doc.insert(doc) return nil end
-- !impl
-- attachments 每项:{ digest=, mime=, name=, data=<base64> }
function sdk.doc.insert_with_media(doc, attachments) return nil end
-- !impl
function sdk.doc.remove(id) return nil end
-- !impl
function sdk.doc.stats() return {} end
@ -174,6 +236,24 @@ function sdk.settings.dump() return {} end
-- !impl
function sdk.settings.plugins() return {} end
-- ============ events只读订阅 ============
-- !impl
-- subscribe(event_type, handler) -> unsubscribe()
-- handler 收到 { type=, source=, timestamp=, payload= }
-- 回调在其内核事件发布 goroutine 上执行只做轻量转发不可阻塞Lua 单状态 + 互斥锁)。
sdk.events = {}
function sdk.events.subscribe(event_type, handler)
print("[lua-plugin] events.subscribe: " .. tostring(event_type))
return function() end
end
-- ============ plugin_mgr ============
-- !impl
sdk.plugin_mgr = {}
function sdk.plugin_mgr.reload_one(name) return nil end
function sdk.plugin_mgr.list_loaded() return {} end
function sdk.plugin_mgr.is_disabled(name) return false end
-- json utils (pure Lua)
sdk.json = {}

View File

@ -28,13 +28,24 @@ var (
//
// ❗main 上此值始终是**下一个未发布中版本**,不随 patch 发布变动
//(见 docs/git-branching.md §2.1);已发布的版本号看对应的 release/vX.Y.x 与 tag。
// 1.3.10:去掉提示词里"每轮只能发一次 output_send"的凭空限制type 缺省即 text。
// 1.3.9:驻留子的「轮次」不再是恒 0info() 此前没填 Rounds
// 1.3.8inputch 划给子后输入只流向子(补上"进内核之前"的输入路由)。
// 1.3.7:驻留子继承父的输出通道(此前子侧 childIO 空壳 ⇒ 子不会发消息)。
// 1.3.6:人格文本不再在播种时固化版本 + 存量实例一次性去版本化(生产实例
// 曾自报 v1.0.3);系统提示词支持 {{kernel_version}} 等占位符。
// 1.3.5:系统提示词(人格卡)支持版本占位符 —— 人格卡是配置项,写死版本号
// 会随发版说谎(线上写 v1.0.3、内核 1.3.xagent 就自报 1.0.3)。
// 支持 {{kernel_version}} / {{kernel_commit}} / {{sdk_version}}。
Version = "1.3.7"
// 1.3.11Lua 插件桥全量对齐 SDK 1.3.0——把 1.1/1.2/1.3 新增的媒体、
// 注入标志位、中断优先级、事件订阅、动态输出通道注销补进 Lua 侧
// (此前只在 Go 侧存在而文档宣称“完全对齐”)。公开 Go SDK 接口
// 零变更,故 SDK 保持 1.3.0。
// 1.3.12:修 1.3.11 引入的两个真问题 —— ① 驻留子销毁后残留入站 inputch
// child/<id>(改用纯函数名并 Unregister覆盖 destroy/reclaim/StopResidents
// ② sdk.events.subscribe 用了从未注入的公共 Events(),且订阅生命周期
// 管理会自死锁/use-after-close改用内部 Subscribe + 独立 subsMu + Stop 取消)。
Version = "1.3.12"
// Commit 是构建时的 Git commit hash。
Commit = "unknown"

View File

@ -1,6 +1,7 @@
package plugin
import (
"encoding/base64"
"encoding/json"
"fmt"
"io"
@ -9,9 +10,10 @@ import (
"strings"
"sync"
lua "github.com/yuin/gopher-lua"
agentEvents "gitcode.com/JianFeeeee/HomeAgent/internal/events"
luaSDK "gitcode.com/JianFeeeee/HomeAgent/internal/lua/sdk"
sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk"
lua "github.com/yuin/gopher-lua"
)
type toolReg struct {
@ -40,7 +42,16 @@ type luaPlugin struct {
stages map[sdk.Stage]*stageReg
outputChs map[string]*outputChReg
inputDefs map[string]sdk.ChannelDef
mu sync.Mutex
// subs 是本插件注册的事件订阅取消函数Stop 时兜底取消,
// 避免 L 已 Close 后残留回调被触发use-after-close
// 用独立的 subsMu 而非 musubscribe 会在 Lua 的 start 回调里被调,
// 而 Start 正持着 mu —— 用 mu 就是不可重入的自死锁。
subs []func()
subsMu sync.Mutex
// closed 在 Stop 里置位(持 mu事件回调持 mu 后先查它,
// 防止“回调已通过取消订阅检查、但等锁期间 L 被 Close”的竞态。
closed bool
mu sync.Mutex
}
func newLuaPlugin(luaPath, name string) (*luaPlugin, error) {
@ -287,6 +298,19 @@ func replaceSDKReal(L *lua.LState, t *lua.LTable, plg *luaPlugin, s *sdk.PluginS
return 0
}))
// pushReply 统一同步注入的返回约定:非空回复返回 (reply, nil)
// 无回复返回 (nil, nil),与数据类 API 的 (result, err) 约定一致。
pushReply := func(reply string) int {
if reply == "" {
L.Push(lua.LNil)
L.Push(lua.LNil)
return 2
}
L.Push(lua.LString(reply))
L.Push(lua.LNil)
return 2
}
t.RawSetString("inject_text", L.NewFunction(func(L *lua.LState) int {
s.InjectText(L.CheckString(1), L.CheckString(2), L.CheckString(3))
return 0
@ -300,6 +324,67 @@ func replaceSDKReal(L *lua.LState, t *lua.LTable, plg *luaPlugin, s *sdk.PluginS
return 0
}))
// ---- 1.2.0 注入标志位no_memory / context_policy / cleaner_name / priority----
t.RawSetString("inject_text_opts", L.NewFunction(func(L *lua.LState) int {
s.InjectTextOpts(L.CheckString(1), L.CheckString(2), L.CheckString(3), parseInjectOptions(L, 4))
return 0
}))
t.RawSetString("inject_interrupt_opts", L.NewFunction(func(L *lua.LState) int {
s.InjectInterruptTextOpts(L.CheckString(1), L.CheckString(2), L.CheckString(3), parseInjectOptions(L, 4))
return 0
}))
// ---- 同步注入:注入后等待本轮回复,返回 (reply, err) ----
// 注意:内置 SDK 的同名 InjectInputSync 是 (eventType, payload) 形态并遮蔽了
// 公共 SDK 的三参文本版本,故这里显式走 PluginSDK 的公共方法。
t.RawSetString("inject_input_sync", L.NewFunction(func(L *lua.LState) int {
return pushReply(s.PluginSDK.InjectInputSync(L.CheckString(1), L.CheckString(2), L.CheckString(3)))
}))
t.RawSetString("inject_input_sync_opts", L.NewFunction(func(L *lua.LState) int {
return pushReply(s.InjectInputSyncOpts(L.CheckString(1), L.CheckString(2), L.CheckString(3), parseInjectOptions(L, 4)))
}))
// ---- 多模态注入1.1.0):内容块随下一次 LLM 请求送达 ----
t.RawSetString("set_tool_blocks", L.NewFunction(func(L *lua.LState) int {
s.SetToolBlocks(luaToContentBlocks(L, 1))
return 0
}))
t.RawSetString("inject_input_media", L.NewFunction(func(L *lua.LState) int {
s.InjectInputMedia(L.CheckString(1), L.CheckString(2), L.CheckString(3), luaToContentBlocks(L, 4))
return 0
}))
t.RawSetString("inject_input_media_opts", L.NewFunction(func(L *lua.LState) int {
s.InjectInputMediaOpts(L.CheckString(1), L.CheckString(2), L.CheckString(3), luaToContentBlocks(L, 4), parseInjectOptions(L, 5))
return 0
}))
t.RawSetString("inject_input_media_sync", L.NewFunction(func(L *lua.LState) int {
return pushReply(s.InjectInputMediaSync(L.CheckString(1), L.CheckString(2), L.CheckString(3), luaToContentBlocks(L, 4)))
}))
t.RawSetString("inject_input_media_sync_opts", L.NewFunction(func(L *lua.LState) int {
return pushReply(s.InjectInputMediaSyncOpts(L.CheckString(1), L.CheckString(2), L.CheckString(3), luaToContentBlocks(L, 4), parseInjectOptions(L, 5)))
}))
t.RawSetString("inject_interrupt_media", L.NewFunction(func(L *lua.LState) int {
s.InjectInterruptMedia(L.CheckString(1), L.CheckString(2), L.CheckString(3), luaToContentBlocks(L, 4))
return 0
}))
t.RawSetString("inject_interrupt_media_opts", L.NewFunction(func(L *lua.LState) int {
s.InjectInterruptMediaOpts(L.CheckString(1), L.CheckString(2), L.CheckString(3), luaToContentBlocks(L, 4), parseInjectOptions(L, 5))
return 0
}))
// ---- 1.3.0 动态输出通道注销:随资源生灭的通道(如远程设备)必须能注销,
// 否则 output_list_channels 会一直列着死通道骗模型。 ----
t.RawSetString("unregister_output_channel", L.NewFunction(func(L *lua.LState) int {
if err := s.UnregisterOutputChannel(L.CheckString(1)); err != nil {
L.Push(lua.LNil)
L.Push(lua.LString(err.Error()))
return 2
}
L.Push(lua.LNil)
L.Push(lua.LNil)
return 2
}))
// ---- 数据类 API与 C ABI 外部插件面完全对齐)----
// 约定:结果型返回 (result, err)void 型返回 (nil, err),成功时 err 为 nil。
@ -373,13 +458,21 @@ func replaceSDKReal(L *lua.LState, t *lua.LTable, plg *luaPlugin, s *sdk.PluginS
if tbl := L.OptTable(1, nil); tbl != nil {
tbl.ForEach(func(_, v lua.LValue) {
if t2, ok := v.(*lua.LTable); ok {
// sentence_text / media_digests 是媒体绑定链的必经环节:
// 媒体引用挂在句子上,漏掉这两个字段会让图片永远绑不上记忆。
var digests []string
if mt, ok := t2.RawGetString("media_digests").(*lua.LTable); ok {
mt.ForEach(func(_, e lua.LValue) { digests = append(digests, e.String()) })
}
triples = append(triples, sdk.Triple{
Subject: t2.RawGetString("subject").String(),
Relation: t2.RawGetString("relation").String(),
Object: t2.RawGetString("object").String(),
Confidence: float64(lua.LVAsNumber(t2.RawGetString("confidence"))),
SubjectType: t2.RawGetString("subject_type").String(),
ObjectType: t2.RawGetString("object_type").String(),
Subject: t2.RawGetString("subject").String(),
Relation: t2.RawGetString("relation").String(),
Object: t2.RawGetString("object").String(),
Confidence: float64(lua.LVAsNumber(t2.RawGetString("confidence"))),
SubjectType: t2.RawGetString("subject_type").String(),
ObjectType: t2.RawGetString("object_type").String(),
SentenceText: t2.RawGetString("sentence_text").String(),
MediaDigests: digests,
})
}
})
@ -442,12 +535,17 @@ func replaceSDKReal(L *lua.LState, t *lua.LTable, plg *luaPlugin, s *sdk.PluginS
}))
docTbl.RawSetString("insert", L.NewFunction(func(L *lua.LState) int {
if dm := s.DocMemory(); dm != nil {
tbl := L.CheckTable(1)
if err := dm.Insert(&sdk.Doc{
ID: tbl.RawGetString("id").String(),
Title: tbl.RawGetString("title").String(),
Content: tbl.RawGetString("content").String(),
}); err != nil {
if err := dm.Insert(docFromLua(L.CheckTable(1))); err != nil {
return pushErr(err)
}
}
return pushNil()
}))
// insert_with_media1.1.0):文档直接持有媒体块,文档向量融合其原生向量,
// 图片按自己的向量被召回,不依赖任何生成的描述文本。
docTbl.RawSetString("insert_with_media", L.NewFunction(func(L *lua.LState) int {
if dm := s.DocMemory(); dm != nil {
if err := dm.InsertWithMedia(docFromLua(L.CheckTable(1)), luaToAttachments(L, L.Get(2))); err != nil {
return pushErr(err)
}
}
@ -503,10 +601,11 @@ func replaceSDKReal(L *lua.LState, t *lua.LTable, plg *luaPlugin, s *sdk.PluginS
if tmem := s.TextMemory(); tmem != nil {
tbl := L.CheckTable(1)
if err := tmem.Append(sdk.TextEvent{
Role: tbl.RawGetString("role").String(),
Content: tbl.RawGetString("content").String(),
Timestamp: int64(lua.LVAsNumber(tbl.RawGetString("timestamp"))),
Channel: tbl.RawGetString("channel").String(),
Role: tbl.RawGetString("role").String(),
Content: tbl.RawGetString("content").String(),
Timestamp: int64(lua.LVAsNumber(tbl.RawGetString("timestamp"))),
Channel: tbl.RawGetString("channel").String(),
Attachments: luaToAttachments(L, tbl.RawGetString("attachments")),
}); err != nil {
return pushErr(err)
}
@ -700,6 +799,76 @@ func replaceSDKReal(L *lua.LState, t *lua.LTable, plg *luaPlugin, s *sdk.PluginS
}
return pushVal([]interface{}{})
}))
// ---- sdk.events.*(只读事件订阅)----
//
// 用内部 SDK 的 Subscribe内置插件用的是同一条路径
// 不用公共 SDK 的 Events()——那个 subscriber 在本内核里从未被注入
// SetEventSubscriber 无调用点),拿到的永远是 nil。
//
// 回调用内核事件发布 goroutine 上执行必须只做轻量转发Lua 单状态 + 互斥锁);
// 阻塞会卡死本插件的全部调用。返回一个取消订阅函数,并在 Stop 时兜底取消
// (否则插件停掉/重载后 L 已 Close残留回调再触发就是 use-after-close
evTbl := subTable("events")
evTbl.RawSetString("subscribe", L.NewFunction(func(L *lua.LState) int {
eventType := L.CheckString(1)
fn := L.CheckFunction(2)
unsub := s.Subscribe(agentEvents.EventType(eventType), func(evt *agentEvents.Event) {
plg.mu.Lock()
defer plg.mu.Unlock()
if plg.closed {
return
}
L2 := plg.L
tbl := L2.NewTable()
tbl.RawSetString("type", lua.LString(string(evt.Type)))
tbl.RawSetString("source", lua.LString(evt.Source))
tbl.RawSetString("timestamp", lua.LNumber(evt.Timestamp))
tbl.RawSetString("payload", goValueToLua(L2, evt.Payload))
L2.Push(fn)
L2.Push(tbl)
if err := L2.PCall(1, 0, nil); err != nil {
fmt.Printf("[lua-plugin/%s] event handler error: %v\n", plg.name, err)
}
})
plg.subsMu.Lock()
plg.subs = append(plg.subs, unsub)
plg.subsMu.Unlock()
L.Push(L.NewFunction(func(L *lua.LState) int {
unsub() // 事件总线的取消订阅是幂等的(重复调用只会匹配不到)
return 0
}))
L.Push(lua.LNil)
return 2
}))
// ---- sdk.plugin_mgr.*(插件管理,与外部插件的 PluginMgrAPI 对齐)----
// PluginMgr 可能未装配(如部分单测的 SDK 构造),此时返回"不可用"而不是 panic。
pmTbl := subTable("plugin_mgr")
pmTbl.RawSetString("reload_one", L.NewFunction(func(L *lua.LState) int {
pm := s.PluginMgr()
if pm == nil {
return pushErr(fmt.Errorf("plugin manager unavailable"))
}
if err := pm.ReloadOne(L.CheckString(1)); err != nil {
return pushErr(err)
}
return pushNil()
}))
pmTbl.RawSetString("list_loaded", L.NewFunction(func(L *lua.LState) int {
pm := s.PluginMgr()
if pm == nil {
return pushList([]interface{}{})
}
return pushList(pm.ListLoadedPlugins())
}))
pmTbl.RawSetString("is_disabled", L.NewFunction(func(L *lua.LState) int {
pm := s.PluginMgr()
if pm == nil {
return pushVal(false)
}
return pushVal(pm.IsPluginDisabled(L.CheckString(1)))
}))
}
func makeToolHandler(plg *luaPlugin, name string, fn *lua.LFunction) sdk.ToolHandler {
@ -724,13 +893,29 @@ func makeStageHandler(plg *luaPlugin, stage sdk.Stage, fn *lua.LFunction) sdk.St
defer plg.mu.Unlock()
L := plg.L
ctx := map[string]interface{}{
"raw_message": sc.RawMessage,
"user_id": sc.UserID,
"group_id": sc.GroupID,
"phase": string(sc.Phase),
"llm_text": sc.LLMText,
"final_text": sc.FinalText,
"no_memory": sc.NoMemory,
"raw_message": sc.RawMessage,
"user_id": sc.UserID,
"group_id": sc.GroupID,
"phase": string(sc.Phase),
"llm_text": sc.LLMText,
"reasoning_content": sc.ReasoningContent,
"final_text": sc.FinalText,
"no_memory": sc.NoMemory,
}
if len(sc.ContextMsgs) > 0 {
ctx["context_msgs"] = jsonToIface(sc.ContextMsgs)
}
if len(sc.TokenUsage) > 0 {
ctx["token_usage"] = jsonToIface(sc.TokenUsage)
}
if len(sc.Memory) > 0 {
ctx["memory"] = jsonToIface(sc.Memory)
}
if len(sc.Extra) > 0 {
ctx["extra"] = jsonToIface(sc.Extra)
}
if len(sc.Errors) > 0 {
ctx["errors"] = jsonToIface(sc.Errors)
}
if sc.Response != nil {
ctx["response"] = *sc.Response
@ -815,6 +1000,7 @@ func parseToolDef(L *lua.LState, defTbl *lua.LTable, plg *luaPlugin, name string
if v := defTbl.RawGetString("no_memory"); v != nil {
goDef.NoMemory = lua.LVAsBool(v)
}
goDef.ContextPolicy = defTbl.RawGetString("context_policy").String()
if v := defTbl.RawGetString("cleaner"); v != nil && v.Type() == lua.LTFunction {
goDef.Cleaner = makeLuaCleaner(plg, v.(*lua.LFunction))
}
@ -834,12 +1020,104 @@ func parseChannelDef(L *lua.LState, defTbl *lua.LTable, plg *luaPlugin) sdk.Chan
if v := defTbl.RawGetString("no_memory"); v != nil {
chDef.NoMemory = lua.LVAsBool(v)
}
chDef.ContextPolicy = defTbl.RawGetString("context_policy").String()
if v := defTbl.RawGetString("cleaner"); v != nil && v.Type() == lua.LTFunction {
chDef.Cleaner = makeLuaCleaner(plg, v.(*lua.LFunction))
}
return chDef
}
// parseInjectOptions 解析 Lua 侧 options table 为 SDK InjectOptions。
// 支持的键no_memory(bool)、context_policy(string)、cleaner_name(string)、priority(string)。
// 缺省/非表等价于零值(记入记忆 + 不裁剪),与旧的三参数注入完全等价。
func parseInjectOptions(L *lua.LState, idx int) sdk.InjectOptions {
opts := sdk.InjectOptions{}
tbl, ok := L.Get(idx).(*lua.LTable)
if !ok {
return opts
}
if v := tbl.RawGetString("no_memory"); v != nil {
opts.NoMemory = lua.LVAsBool(v)
}
opts.ContextPolicy = tbl.RawGetString("context_policy").String()
opts.CleanerName = tbl.RawGetString("cleaner_name").String()
opts.Priority = tbl.RawGetString("priority").String()
return opts
}
// luaToContentBlocks 把 Lua 的 blocks 数组解析为 SDK ContentBlock。
// 每项形如:
//
// { type = "text", text = "..." }
// { type = "image_url", image_url = { url = "...", detail = "high" } }
// { type = "audio_url", audio_url = { url = "..." } }
func luaToContentBlocks(L *lua.LState, idx int) []sdk.ContentBlock {
tbl, ok := L.Get(idx).(*lua.LTable)
if !ok {
return nil
}
var blocks []sdk.ContentBlock
tbl.ForEach(func(_, v lua.LValue) {
bt, ok := v.(*lua.LTable)
if !ok {
return
}
b := sdk.ContentBlock{
Type: bt.RawGetString("type").String(),
Text: bt.RawGetString("text").String(),
}
if iu, ok := bt.RawGetString("image_url").(*lua.LTable); ok {
b.ImageURL = &sdk.ImageURL{
URL: iu.RawGetString("url").String(),
Detail: iu.RawGetString("detail").String(),
}
}
if au, ok := bt.RawGetString("audio_url").(*lua.LTable); ok {
b.AudioURL = &sdk.AudioURL{URL: au.RawGetString("url").String()}
}
blocks = append(blocks, b)
})
return blocks
}
// luaToAttachments 把 Lua 附件数组解析为 SDK MediaAttachment。
// 每项:{ digest=, mime=, name=, data=<base64 字符串> }。
// 带 data 的是新内容(内核落进内容寻址存储),只带 digest 的是引用已有内容。
// base64 解码失败时忽略 data不整单失败——坏附件不应阻断一条记忆写入。
func luaToAttachments(L *lua.LState, val lua.LValue) []sdk.MediaAttachment {
tbl, ok := val.(*lua.LTable)
if !ok {
return nil
}
var out []sdk.MediaAttachment
tbl.ForEach(func(_, v lua.LValue) {
at, ok := v.(*lua.LTable)
if !ok {
return
}
a := sdk.MediaAttachment{
Digest: at.RawGetString("digest").String(),
MIME: at.RawGetString("mime").String(),
Name: at.RawGetString("name").String(),
}
if s := at.RawGetString("data").String(); s != "" {
if b, err := base64.StdEncoding.DecodeString(s); err == nil {
a.Data = b
}
}
out = append(out, a)
})
return out
}
func docFromLua(tbl *lua.LTable) *sdk.Doc {
return &sdk.Doc{
ID: tbl.RawGetString("id").String(),
Title: tbl.RawGetString("title").String(),
Content: tbl.RawGetString("content").String(),
}
}
// jsonToIface 通过 JSON 往返把任意 Go 值转换为 JSON 兼容的 interface{} 树。
func jsonToIface(v interface{}) interface{} {
b, err := json.Marshal(v)
@ -959,8 +1237,21 @@ func (p *luaPlugin) Start(s *sdk.PluginSDK) error {
}
func (p *luaPlugin) Stop() error {
// ① 先取消事件订阅。**不持 p.mu**Bus.Publish 持总线锁回调 handler
// 而 handler 要 p.mu若此处持 p.mu 再取总线锁,就是锁序反转死锁。
p.subsMu.Lock()
subs := p.subs
p.subs = nil
p.subsMu.Unlock()
for _, unsub := range subs {
unsub()
}
// ② 置 closed 并关 L。置位在持锁下完成已进入但等锁的 event 回调
// 拿到锁后会先看到 closed 而直接返回,不会碰已关的 L。
p.mu.Lock()
defer p.mu.Unlock()
p.closed = true
if p.tbl != nil {
fn := p.tbl.RawGetString("stop")

View File

@ -5,9 +5,10 @@ import (
"path/filepath"
"testing"
lua "github.com/yuin/gopher-lua"
internalConfig "gitcode.com/JianFeeeee/HomeAgent/internal/config"
"gitcode.com/JianFeeeee/HomeAgent/internal/events"
sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk"
lua "github.com/yuin/gopher-lua"
)
func TestTryLoadLua_Basic(t *testing.T) {
@ -597,3 +598,66 @@ return plugin
t.Errorf("llm_text writeback: got %q, want %q", sc2.LLMText, "模型输出[尾部标记]")
}
}
// TestLuaEventsSubscribeAndStopCleanup 覆盖 sdk.events.subscribe
// 1. 订阅真的能收到内核事件(走内部 SDK 的 Subscribe不是永远为 nil 的公共 Events()
// 2. Stop 会取消订阅,之后 Publish 不得再触碰已 Close 的 LState。
func TestLuaEventsSubscribeAndStopCleanup(t *testing.T) {
dir := t.TempDir()
os.WriteFile(filepath.Join(dir, "plugin.json"), []byte(`{"name":"evlua","entry":"main.lua"}`), 0644)
os.WriteFile(filepath.Join(dir, "main.lua"), []byte(`
local plugin = { name = "evlua" }
function plugin.start(sdk)
_G.hits = 0
local unsub, err = sdk.events.subscribe("agent_output", function(evt)
_G.hits = _G.hits + 1
_G.last_type = evt.type
_G.last_source = evt.source
end)
_G.sub_err = err
_G.unsub_type = type(unsub)
end
function plugin.stop() end
return plugin
`), 0644)
plg, err := tryLoadLua(dir, "evlua", nil)
if err != nil {
t.Fatalf("tryLoadLua failed: %v", err)
}
lp := plg.(*luaPlugin)
bus := events.NewBus()
reg := internalConfig.NewConfigRegistry("")
sett := sdk.NewSettings("evlua", reg)
s := sdk.New("evlua", sdk.SDKConfig{EventBus: bus, Settings: sett})
if err := plg.Start(s); err != nil {
t.Fatalf("Start failed: %v", err)
}
L := lp.L
if errStr := L.GetGlobal("sub_err").String(); errStr != "nil" {
t.Fatalf("subscribe returned error: %s", errStr)
}
if got := L.GetGlobal("unsub_type").String(); got != "function" {
t.Fatalf("subscribe should return an unsubscribe function, got %s", got)
}
bus.Publish(&events.Event{Type: events.EventAgentOutput, Source: "test-src"})
if hits := int(lua.LVAsNumber(L.GetGlobal("hits"))); hits != 1 {
t.Fatalf("event handler hits = %d, want 1", hits)
}
if got := L.GetGlobal("last_source").String(); got != "test-src" {
t.Fatalf("event source = %q, want test-src", got)
}
// Stop 取消订阅 + 关 L此后再 Publish 不得 panic / use-after-close。
if err := plg.Stop(); err != nil {
t.Fatalf("Stop failed: %v", err)
}
bus.Publish(&events.Event{Type: events.EventAgentOutput, Source: "after-stop"})
}

View File

@ -0,0 +1,188 @@
package plugin
import (
"os"
"path/filepath"
"regexp"
"strings"
"testing"
luaSDK "gitcode.com/JianFeeeee/HomeAgent/internal/lua/sdk"
sdk "gitcode.com/JianFeeeee/HomeAgent/internal/sdk"
lua "github.com/yuin/gopher-lua"
)
// TestLuaSDKMockSingleSource 守住「Lua mock 只有一份事实源」。
//
// 三份 sdk.lua内核内嵌 / 工具链模板 / 项目副本)历史上各自漂移过,
// 表现为「mock 里有的 API内核运行时是 nil」这类静默失配。
// 事实源是 SDK 仓的 sdk/lua/sdk.lua内核副本由
// third_party/homeagent-sdk/scripts/sync-lua-sdk.sh 同步。
func TestLuaSDKMockSingleSource(t *testing.T) {
canonical, err := os.ReadFile(filepath.Join("..", "..", "third_party", "homeagent-sdk", "sdk", "lua", "sdk.lua"))
if err != nil {
t.Skipf("SDK repo canonical sdk.lua not available: %v", err)
}
if string(canonical) != luaSDK.SDKSource {
t.Fatal("内核内嵌 sdk.lua 与 SDK 仓 sdk/lua/sdk.lua 不一致;" +
"请跑 third_party/homeagent-sdk/scripts/sync-lua-sdk.sh")
}
}
var luaMockFuncRe = regexp.MustCompile(`(?m)^function sdk\.([A-Za-z0-9_.]+)\s*\(`)
// TestLuaBridgeCoversMock 守住「mock 承诺的每个函数,运行时都有绑定」。
//
// 只查 mock → 运行时这一向mock 定义了但没绑定,插件会先看到 mock 能调、
// 之后内核里是 nil或反向的假象。反向运行时多出未文档化的函数无害。
func TestLuaBridgeCoversMock(t *testing.T) {
bridge, err := os.ReadFile("lua_plugin.go")
if err != nil {
t.Fatalf("read lua_plugin.go: %v", err)
}
src := string(bridge)
// 纯 Lua 实现,不经内核绑定。
exempt := map[string]bool{"json.encode": true, "json.decode": true}
matches := luaMockFuncRe.FindAllStringSubmatch(luaSDK.SDKSource, -1)
if len(matches) < 40 {
t.Fatalf("parsed only %d sdk.* functions from mock; parser likely broken", len(matches))
}
for _, m := range matches {
full := m[1]
if exempt[full] {
continue
}
leaf := full
if i := strings.LastIndex(full, "."); i >= 0 {
leaf = full[i+1:]
}
if !strings.Contains(src, `RawSetString("`+leaf+`"`) {
t.Errorf("sdk.%s: mock 有定义,但 lua_plugin.go 没有 RawSetString(%q) 绑定", full, leaf)
}
}
}
func luaTableFrom(t *testing.T, code string) (*lua.LState, *lua.LTable) {
t.Helper()
L := lua.NewState()
if err := L.DoString("return " + code); err != nil {
L.Close()
t.Fatalf("eval lua table: %v", err)
}
tbl, ok := L.Get(-1).(*lua.LTable)
if !ok {
L.Close()
t.Fatalf("expected a table from %q", code)
}
L.Pop(1)
return L, tbl
}
func TestLuaParseInjectOptions(t *testing.T) {
L, tbl := luaTableFrom(t, `{
no_memory = true,
context_policy = "prune",
cleaner_name = "sanitize",
priority = "L2",
}`)
defer L.Close()
L.Push(tbl)
got := parseInjectOptions(L, 1)
L.Pop(1)
if !got.NoMemory {
t.Error("NoMemory should be true")
}
if got.ContextPolicy != sdk.ContextPolicyPrune {
t.Errorf("ContextPolicy = %q, want prune", got.ContextPolicy)
}
if got.CleanerName != "sanitize" {
t.Errorf("CleanerName = %q, want sanitize", got.CleanerName)
}
if got.Priority != sdk.PriorityL2 {
t.Errorf("Priority = %q, want L2", got.Priority)
}
// 缺省/非表 = 零值(记入记忆 + 不裁剪),与旧三参数注入等价。
if z := parseInjectOptions(L, 99); z != (sdk.InjectOptions{}) {
t.Errorf("missing opts should be zero value, got %#v", z)
}
}
func TestLuaContentBlocksParse(t *testing.T) {
L, tbl := luaTableFrom(t, `{
{ type = "text", text = "看图" },
{ type = "image_url", image_url = { url = "data:image/png;base64,AAAA", detail = "high" } },
{ type = "audio_url", audio_url = { url = "https://x/a.mp3" } },
}`)
defer L.Close()
L.Push(tbl)
blocks := luaToContentBlocks(L, 1)
L.Pop(1)
if len(blocks) != 3 {
t.Fatalf("got %d blocks, want 3", len(blocks))
}
if blocks[0].Type != "text" || blocks[0].Text != "看图" {
t.Errorf("block[0] = %#v", blocks[0])
}
if blocks[1].ImageURL == nil || blocks[1].ImageURL.URL != "data:image/png;base64,AAAA" || blocks[1].ImageURL.Detail != "high" {
t.Errorf("block[1] image_url = %#v", blocks[1].ImageURL)
}
if blocks[2].AudioURL == nil || blocks[2].AudioURL.URL != "https://x/a.mp3" {
t.Errorf("block[2] audio_url = %#v", blocks[2].AudioURL)
}
}
func TestLuaAttachmentsBase64(t *testing.T) {
// "hello" 的 base64 是 aGVsbG8=
L, tbl := luaTableFrom(t, `{
{ digest = "sha256:abc", mime = "image/png", name = "a.png" },
{ mime = "image/jpeg", data = "aGVsbG8=" },
{ mime = "image/png", data = "!!!not-base64!!!" },
}`)
defer L.Close()
atts := luaToAttachments(L, tbl)
if len(atts) != 3 {
t.Fatalf("got %d attachments, want 3", len(atts))
}
if atts[0].Digest != "sha256:abc" || atts[0].MIME != "image/png" || atts[0].Name != "a.png" || atts[0].Data != nil {
t.Errorf("att[0] = %#v", atts[0])
}
if string(atts[1].Data) != "hello" {
t.Errorf("att[1] data = %q, want hello", string(atts[1].Data))
}
// 坏 base64 只丢 data不整单失败——坏附件不应阻断记忆写入。
if atts[2].Data != nil {
t.Errorf("att[2] bad base64 should be dropped, got %q", string(atts[2].Data))
}
}
func TestLuaDefinitionsCarryContextPolicy(t *testing.T) {
L, tbl := luaTableFrom(t, `{
description = "t",
no_memory = true,
context_policy = "prune",
}`)
defer L.Close()
plg := &luaPlugin{name: "cp"}
def := parseToolDef(L, tbl, plg, "t")
if def.ContextPolicy != sdk.ContextPolicyPrune {
t.Errorf("ToolDef.ContextPolicy = %q, want prune", def.ContextPolicy)
}
if !def.NoMemory {
t.Error("ToolDef.NoMemory should be true")
}
chDef := parseChannelDef(L, tbl, plg)
if chDef.ContextPolicy != sdk.ContextPolicyPrune {
t.Errorf("ChannelDef.ContextPolicy = %q, want prune", chDef.ContextPolicy)
}
}

View File

@ -15,7 +15,7 @@ import (
type PluginType string
const (
PluginTypeSKILL PluginType = "skill"
PluginTypeSKILL PluginType = "skill"
)
type IOConfig struct {
@ -78,18 +78,28 @@ func LoadSKILL(path string) (*SKILLPlugin, error) {
metaFile := filepath.Join(path, "skill.json")
if data, err := os.ReadFile(metaFile); err == nil {
var meta struct {
Name string `json:"name"`
Description string `json:"description"`
Version string `json:"version"`
Author string `json:"author"`
Name string `json:"name"`
Description string `json:"description"`
Version string `json:"version"`
Author string `json:"author"`
IO *IOConfig `json:"io,omitempty"`
}
if err := json.Unmarshal(data, &meta); err == nil {
if meta.Name != "" { p.name = meta.Name }
if meta.Description != "" { p.description = meta.Description }
if meta.Version != "" { p.version = meta.Version }
if meta.Author != "" { p.author = meta.Author }
if meta.IO != nil { p.ioConfig = meta.IO }
if meta.Name != "" {
p.name = meta.Name
}
if meta.Description != "" {
p.description = meta.Description
}
if meta.Version != "" {
p.version = meta.Version
}
if meta.Author != "" {
p.author = meta.Author
}
if meta.IO != nil {
p.ioConfig = meta.IO
}
}
}
} else if filepath.Ext(path) == ".md" {
@ -198,7 +208,9 @@ func extractToolDefs(content string) []ToolDef {
inCodeBlock = !inCodeBlock
continue
}
if inCodeBlock { continue }
if inCodeBlock {
continue
}
if strings.HasPrefix(trimmed, "## ") && !strings.HasPrefix(trimmed, "### ") {
if currentTool != nil && currentTool.Name != "" {
@ -232,7 +244,9 @@ func extractToolDefs(content string) []ToolDef {
continue
}
if currentTool == nil || currentTool.Name == "" { continue }
if currentTool == nil || currentTool.Name == "" {
continue
}
if currentTool.Description == "" && trimmed != "" &&
!strings.HasPrefix(trimmed, "- ") && !strings.HasPrefix(trimmed, "#") {
@ -280,7 +294,9 @@ func isNonToolSection(name string) bool {
func extractIOConfig(content string) *IOConfig {
ioType := extractField(content, "io_type")
if ioType == "" { return nil }
if ioType == "" {
return nil
}
cfg := &IOConfig{
Type: ioType,
InputRoute: extractField(content, "io_input_route"),

View File

@ -0,0 +1,403 @@
-- HomeAgent Lua Plugin SDK
-- Interface contract between Lua plugins and HomeAgent kernel.
-- !impl functions are replaced by Go implementations at runtime.
-- Standalone/debug: pure Lua mock implementations are used.
-- Usage: local sdk = require("sdk")
sdk = {}
-- !impl
-- level: "debug" | "info" | "warn" | "error"
function sdk.log(level, msg)
print("[lua-plugin] " .. tostring(level) .. ": " .. tostring(msg))
end
-- !impl
-- def: { description="...", parameters={...}, no_memory=true/false, cleaner=function(text)->text }
-- handler: function(args) -> result
function sdk.register_tool(name, def, handler)
print("[lua-plugin] register_tool: " .. tostring(name))
end
-- !impl
-- stage: "on_input" | "pre_action" | "post_action" | ...
-- scope: nil/"global" (默认) | "own_tools"(仅 before_toolcall/after_toolcall 且工具属于本插件时触发)
function sdk.register_stage(stage, handler, scope)
print("[lua-plugin] register_stage: " .. tostring(stage) .. " scope=" .. tostring(scope))
end
-- !impl
function sdk.register_api(name)
print("[lua-plugin] register_api: " .. tostring(name))
end
-- !impl
-- def: { no_memory=true/false, cleaner=function(text)->text }
-- handler: function(args) -> result
function sdk.register_output_channel(name, caps, desc, def, handler)
print("[lua-plugin] register_output_channel: " .. tostring(name))
end
-- !impl
-- def: { no_memory=true/false, cleaner=function(text)->text }
function sdk.register_input_channel(name, def)
print("[lua-plugin] register_input_channel: " .. tostring(name))
end
-- !impl
function sdk.get_setting(key)
return nil
end
-- !impl
function sdk.set_setting(key, value)
print("[lua-plugin] set_setting: " .. tostring(key))
end
-- !impl
function sdk.inject_text(source, channel, text)
print("[lua-plugin] inject_text: " .. tostring(source) .. "/" .. tostring(channel))
end
-- !impl
function sdk.inject_interrupt(source, channel, text)
print("[lua-plugin] inject_interrupt: " .. tostring(source))
end
-- !impl
function sdk.inject_text_no_memory(source, channel, text)
print("[lua-plugin] inject_text_no_memory: " .. tostring(source))
end
-- !impl
-- opts: { no_memory=bool, context_policy="none"|"prune", cleaner_name=string, priority="L1".."L3" }
-- 零值/缺省 = 记入记忆 + 不裁剪(与三参数版本等价)。
function sdk.inject_text_opts(source, channel, text, opts)
print("[lua-plugin] inject_text_opts: " .. tostring(source))
end
-- !impl
function sdk.inject_interrupt_opts(source, channel, text, opts)
print("[lua-plugin] inject_interrupt_opts: " .. tostring(source))
end
-- !impl
-- 同步注入:等待本轮回复 -> (reply, err);无回复时 reply 为 nil。
function sdk.inject_input_sync(source, channel, text) return nil, nil end
-- !impl
function sdk.inject_input_sync_opts(source, channel, text, opts) return nil, nil end
-- !impl
-- blocks: ContentBlock 数组,见 sdk.inject_input_media。
-- 设置下一轮 tool message 携带的多模态内容块(模型据此看图/听音频)。
function sdk.set_tool_blocks(blocks)
print("[lua-plugin] set_tool_blocks: " .. tostring(blocks and #blocks or 0))
end
-- !impl
-- blocks 每项:{ type="text", text="..." }
-- | { type="image_url", image_url={ url="...", detail="high" } }
-- | { type="audio_url", audio_url={ url="..." } }
function sdk.inject_input_media(source, channel, text, blocks)
print("[lua-plugin] inject_input_media: " .. tostring(source))
end
-- !impl
function sdk.inject_input_media_opts(source, channel, text, blocks, opts)
print("[lua-plugin] inject_input_media_opts: " .. tostring(source))
end
-- !impl
function sdk.inject_input_media_sync(source, channel, text, blocks) return nil, nil end
-- !impl
function sdk.inject_input_media_sync_opts(source, channel, text, blocks, opts) return nil, nil end
-- !impl
function sdk.inject_interrupt_media(source, channel, text, blocks)
print("[lua-plugin] inject_interrupt_media: " .. tostring(source))
end
-- !impl
function sdk.inject_interrupt_media_opts(source, channel, text, blocks, opts)
print("[lua-plugin] inject_interrupt_media_opts: " .. tostring(source))
end
-- !impl
-- 注销输出通道(随资源生灭的动态通道,如远程设备)。返回 (nil, err)。
function sdk.unregister_output_channel(name) return nil, nil end
-- !impl
-- enabled: true/false崩溃时内核自动拉起
function sdk.set_auto_restart(enabled)
print("[lua-plugin] set_auto_restart: " .. tostring(enabled))
end
-- ============ graph memory ============
-- !impl
sdk.memory = {}
-- !impl
-- query: string, depth: number -> {entities={...}, relations={...}}
function sdk.memory.recall(query, depth) return {entities={}, relations={}} end
-- !impl
-- triples: { {subject=, relation=, object=, [confidence=], [sentence_text=]} } -> err
function sdk.memory.commit(triples) return nil end
-- !impl
function sdk.memory.introspect() return {} end
-- !impl
function sdk.memory.merge(source, target) return 0 end
-- !impl
-- criteria: {key=value}, hard: boolean
function sdk.memory.purge(criteria, hard) return 0 end
-- ============ document memory ============
-- !impl
sdk.doc = {}
-- !impl
function sdk.doc.query(text, top_k) return {} end
-- !impl
-- doc: { id=, title=, content= }
function sdk.doc.insert(doc) return nil end
-- !impl
-- attachments 每项:{ digest=, mime=, name=, data=<base64> }
function sdk.doc.insert_with_media(doc, attachments) return nil end
-- !impl
function sdk.doc.remove(id) return nil end
-- !impl
function sdk.doc.stats() return {} end
-- ============ knowledge ============
-- !impl
sdk.knowledge = {}
-- !impl
function sdk.knowledge.search(query, limit) return {} end
-- !impl
function sdk.knowledge.add(tag, content) return nil end
-- !impl
function sdk.knowledge.list() return {} end
-- ============ text memory ============
-- !impl
sdk.text_memory = {}
-- !impl
-- evt: { timestamp=, role=, content=, channel= }
function sdk.text_memory.append(evt) return nil end
-- ============ llm ============
-- !impl
sdk.llm = {}
-- !impl
function sdk.llm.list_sources() return {} end
-- !impl
function sdk.llm.set_source(name) return nil end
-- !impl
function sdk.llm.current_source() return nil end
-- ============ social (只读) ============
-- !impl
sdk.social = {}
-- !impl
function sdk.social.get_person(name) return {} end
-- !impl
function sdk.social.get_network(name, depth) return {} end
-- !impl
function sdk.social.get_trait(name, trait) return {value=nil, found=false} end
-- !impl
function sdk.social.get_relations(name) return {} end
-- !impl
function sdk.social.list_persons() return {} end
-- ============ settings (作用域变体) ============
-- !impl
sdk.settings = {}
-- !impl
function sdk.settings.get_core(key) return nil end
-- !impl
function sdk.settings.set_core(key, value) return nil end
-- !impl
function sdk.settings.list_core(prefix) return {} end
-- !impl
function sdk.settings.get_plugin(plugin, key) return nil end
-- !impl
function sdk.settings.set_plugin(plugin, key, value) return nil end
-- !impl
function sdk.settings.list_plugin(plugin, prefix) return {} end
-- !impl
function sdk.settings.list(prefix) return {} end
-- !impl
-- def: { key=, type=, display_name=, description=, category=, options=, default=,
-- min=, max=, step=, required=, secret= }
function sdk.settings.register_def(def) return nil end
-- !impl
function sdk.settings.defs(prefix) return {} end
-- !impl
function sdk.settings.dump() return {} end
-- !impl
function sdk.settings.plugins() return {} end
-- ============ events只读订阅 ============
-- !impl
-- subscribe(event_type, handler) -> unsubscribe()
-- handler 收到 { type=, source=, timestamp=, payload= }
-- 回调在其内核事件发布 goroutine 上执行只做轻量转发不可阻塞Lua 单状态 + 互斥锁)。
sdk.events = {}
function sdk.events.subscribe(event_type, handler)
print("[lua-plugin] events.subscribe: " .. tostring(event_type))
return function() end
end
-- ============ plugin_mgr ============
-- !impl
sdk.plugin_mgr = {}
function sdk.plugin_mgr.reload_one(name) return nil end
function sdk.plugin_mgr.list_loaded() return {} end
function sdk.plugin_mgr.is_disabled(name) return false end
-- json utils (pure Lua)
sdk.json = {}
function sdk.json.encode(val)
local ok, result = pcall(function()
local function _encode(v)
local t = type(v)
if t == "string" then
local s = v:gsub('\\', '\\\\'):gsub('"', '\\"'):gsub('\n', '\\n'):gsub('\r', '\\r'):gsub('\t', '\\t')
return '"' .. s .. '"'
elseif t == "number" then
return tostring(v)
elseif t == "boolean" then
return tostring(v)
elseif t == "table" then
local keys = {}
local is_array = true
local maxn = 0
for k in pairs(v) do
keys[#keys + 1] = k
if type(k) ~= "number" or k < 1 or k ~= math.floor(k) then
is_array = false
end
if type(k) == "number" and k > maxn then maxn = k end
end
if is_array and #keys >= maxn then
local parts = {}
for i = 1, maxn do
parts[#parts + 1] = _encode(v[i])
end
return "[" .. table.concat(parts, ",") .. "]"
else
local parts = {}
for _, k in ipairs(keys) do
parts[#parts + 1] = _encode(tostring(k)) .. ":" .. _encode(v[k])
end
return "{" .. table.concat(parts, ",") .. "}"
end
else
return "null"
end
end
return _encode(val)
end)
if ok then return result end
return "null"
end
function sdk.json.decode(str)
local ok, result = pcall(function()
local pos, _end = 1, #str
local function skip()
while pos <= _end and str:sub(pos, pos):match("%s") do pos = pos + 1 end
end
local function parse()
skip()
if pos > _end then return nil end
local c = str:sub(pos, pos)
if c == '"' then
local s = {}
pos = pos + 1
while pos <= _end do
local ch = str:sub(pos, pos)
if ch == '"' then
pos = pos + 1
return table.concat(s)
elseif ch == '\\' then
pos = pos + 1
local n = str:sub(pos, pos)
if n == '"' then s[#s+1] = '"'
elseif n == '\\' then s[#s+1] = '\\'
elseif n == '/' then s[#s+1] = '/'
elseif n == 'b' then s[#s+1] = '\b'
elseif n == 'f' then s[#s+1] = '\f'
elseif n == 'n' then s[#s+1] = '\n'
elseif n == 'r' then s[#s+1] = '\r'
elseif n == 't' then s[#s+1] = '\t'
elseif n == 'u' then
local hex = str:sub(pos+1, pos+4)
pos = pos + 4
s[#s+1] = utf8 and utf8.char(tonumber(hex, 16)) or '?'
end
pos = pos + 1
else
s[#s+1] = ch
pos = pos + 1
end
end
return table.concat(s)
elseif c == 't' then pos = pos + 4; return true
elseif c == 'f' then pos = pos + 5; return false
elseif c == 'n' then pos = pos + 4; return nil
elseif c == '{' then
pos = pos + 1; skip()
local t = {}
if str:sub(pos, pos) == '}' then pos = pos + 1; return t end
while true do
skip(); local k = parse(); skip()
if str:sub(pos, pos) == ':' then pos = pos + 1 end
skip(); t[k] = parse(); skip()
local sep = str:sub(pos, pos)
if sep == '}' then pos = pos + 1; return t end
if sep == ',' then pos = pos + 1 end
end
elseif c == '[' then
pos = pos + 1; skip()
local t = {}
if str:sub(pos, pos) == ']' then pos = pos + 1; return t end
local idx = 1
while true do
skip(); t[idx] = parse(); idx = idx + 1; skip()
local sep = str:sub(pos, pos)
if sep == ']' then pos = pos + 1; return t end
if sep == ',' then pos = pos + 1 end
end
else
local s, e = str:find('^[-%d%.eE]+', pos)
if s then
local num = tonumber(str:sub(s, e))
pos = e + 1
return num
end
return nil
end
end
return parse()
end)
if ok then return result end
return nil
end
-- http utils
sdk.http = {}
-- !impl
function sdk.http.get(url)
print("[lua-plugin] http.get: " .. tostring(url))
return {status=200, body='{"mock":true}', headers={}}
end
-- !impl
function sdk.http.post(url, body, content_type)
print("[lua-plugin] http.post: " .. tostring(url))
return {status=200, body='{"mock":true}', headers={}}
end
return sdk