Files
homeagent-sdk/example/rss/plugin.go
JianFeeeee ba49dfda44 feat(sdk): 注入行为的记忆/裁剪标志位(纯追加)+ plugindev 退出码修复 + 示例 hmap 随发版
## 1. 注入标志位(公开 API 纯追加,无签名变更)

给注入行为补上工具早已有的两类声明位,并让通道定义也带上:
- `InjectOptions{NoMemory, ContextPolicy, CleanerName}`
- `IOInjector` 新增六个 `*Opts` 变体(排队/中断/同步 × 纯文本/带媒体)
- `ChannelDef.ContextPolicy`,并**补上 JSON tag**(Cleaner 标 `json:"-"`)

零值 InjectOptions 与旧的三参数方法完全等价(记入记忆 + 不裁剪),
存量插件不需要改一行、也不需要重编;旧方法保留为转发到零值 opts 的语法糖。

三条设计要点:
- **默认不裁剪**:裁剪会归档丢弃低相关事件,必须显式声明(ContextPolicy=prune)。
- **中断也允许声明 prune**(已确认):中断同样携带内容进上下文。
- `CleanerName`:注入的 source 未必是注册过的输入通道名,而注入内容常带
  ANSI/JSON 包装;允许显式指定用哪个已注册 cleaner 清洗。

顺带修掉一个易静默丢字段的坑:`ChannelDef` 原来没有 JSON tag,只能手写字段
白名单跨进程传(`{"NoMemory": ...}`),新增字段会被丢掉。现在模板整体传 `def`。

## 2. 示例调用点统一写明意图
rss/memo/calendar/qq 的中断注入显式 `NoMemory: true`(行为等价,写清语义)。

## 3. plugindev 出错却 exit 0(真缺陷)
`buildBundle`/`buildTarget` 遇错只 Printf 后 return,`cmdBuild` 返回 void,
于是**构建失败也退 0**。实测中一个示例的 windows 目标编译失败,批量脚本却报
「17/17 全绿」,并因此少产出 16 个 .hmap。现在累计 `buildFailed` 并以非零退出。

## 4. 平台策略:插件目标去掉 windows
homed 已放弃 Windows 原生(见核心仓 cmd/homed/platform_windows.go:插件体系依赖
fd 继承 + 统一共享内存区的段内偏移解引用,Windows 句柄模型无法表达),
插件只运行在 homed 能跑的平台上,故 `allBundleTargets` 去掉 windows,
并对 windows 目标给出**可执行的报错**(指引 WSL2),而不是让它死在一句
`undefined: attachUnifiedShm` 上。

## 5. 发版附带各示例插件的 .hmap
新增 `package/build-examples.sh` 并接入 `package/build.sh`(组件 all|plugindev|examples):
- 用**刚构建出来的**那把工具链构建示例,保证与本次发版同源
- 逐平台 `--no-bundle --target <os/arch>`(bundle 会连 windows 一起编)
- 判成功同时看**退出码 + 产物存在**
- 全部产物齐了才 `sha256sum`(边打边算会漏掉后生成的包)
- 有任一失败即整体失败,不生成 SHA256SUMS

## 6. 版本
SDK 仍为 1.2.0(main 是下一个未发布中版本);1.2.0 条目补记本次新增接口,
并注明新标志位需核心 1.2.0+(旧核心会忽略这些字段,不报错但不生效)。
2026-09-11 20:29:30 +08:00

501 lines
11 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

package main
import (
"encoding/json"
"fmt"
"net/http"
"os"
"path/filepath"
"sort"
"strconv"
"strings"
"sync"
"time"
sdk "gitcode.com/JianFeeeee/homeagent-sdk/sdk"
"github.com/mmcdole/gofeed"
)
const injectDedupWindow = 5 * time.Minute
type FeedSub struct {
URL string `json:"url"`
Title string `json:"title"`
AddedAt string `json:"added_at"`
Interval int `json:"interval"`
}
type Plugin struct {
name string
sdk *sdk.PluginSDK
client *http.Client
fp *gofeed.Parser
dataDir string
mu sync.RWMutex
feeds []FeedSub
seenGUIDs map[string]bool
injected map[string]time.Time
stopCh chan struct{}
stopOnce sync.Once
wg sync.WaitGroup
pollTicker *time.Ticker
}
func NewPluginFactory(name string, config map[string]interface{}) (sdk.Plugin, error) {
return &Plugin{name: name}, nil
}
func (p *Plugin) Name() string { return p.name }
func getSetting[T string | int64 | float64](s sdk.SettingsAPI, key string, fallback T) T {
v, err := s.Get(key)
if err != nil || v == nil {
return fallback
}
switch any(fallback).(type) {
case string:
if sv, ok := v.(string); ok {
return any(sv).(T)
}
case int64:
switch val := v.(type) {
case float64:
return any(int64(val)).(T)
case string:
if n, err := strconv.ParseInt(val, 10, 64); err == nil {
return any(n).(T)
}
}
case float64:
switch val := v.(type) {
case float64:
return any(val).(T)
case string:
if n, err := strconv.ParseFloat(val, 64); err == nil {
return any(n).(T)
}
}
}
return fallback
}
func readArg(args map[string]interface{}, key string) string {
if v, ok := args[key]; ok && v != nil {
if s, ok := v.(string); ok {
return s
}
}
return ""
}
func readArgInt(args map[string]interface{}, key string, fallback int) int {
if v, ok := args[key]; ok && v != nil {
switch n := v.(type) {
case float64:
return int(n)
case int64:
return int(n)
}
}
return fallback
}
func (p *Plugin) Start(s *sdk.PluginSDK) error {
s.SetAutoRestart(true)
p.sdk = s
p.client = &http.Client{Timeout: 30 * time.Second}
p.fp = gofeed.NewParser()
p.stopCh = make(chan struct{})
p.seenGUIDs = make(map[string]bool)
p.injected = make(map[string]time.Time)
p.feeds = []FeedSub{}
dataDirVal, err := s.Settings().GetCore("core.daemon.data_dir")
if err != nil || dataDirVal == "" {
dataDirVal = "."
}
p.dataDir = filepath.Join(fmt.Sprint(dataDirVal), "rss")
if err := os.MkdirAll(p.dataDir, 0755); err != nil {
fmt.Printf("[%s] mkdir %s: %v\n", p.name, p.dataDir, err)
}
p.loadData()
// 卸载(删除)时清理订阅数据目录;重载不触发
s.RegisterOnRemoveHandler(p.cleanupData)
s.Settings().RegisterDef(sdk.ConfigDef{
Key: "poll_interval", Default: "30", Type: "string",
DisplayName: "Poll Interval", Description: "Default polling interval in minutes (default: 30)",
Category: "rss",
})
tp := p.name + "_"
s.RegisterTool(tp+"subscribe", sdk.ToolDef{
Name: tp + "subscribe", Description: "Subscribe to an RSS/Atom feed URL",
NoMemory: true,
Parameters: map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{
"url": map[string]interface{}{"type": "string", "description": "Feed URL"},
"interval": map[string]interface{}{"type": "integer", "description": "Poll interval in minutes (default: 30, minimum: 5)"},
},
"required": []string{"url"},
},
}, p.handleSubscribe)
s.RegisterTool(tp+"unsubscribe", sdk.ToolDef{
Name: tp + "unsubscribe", Description: "Unsubscribe from a feed",
NoMemory: true,
Parameters: map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{
"url": map[string]interface{}{"type": "string", "description": "Feed URL to unsubscribe"},
},
"required": []string{"url"},
},
}, p.handleUnsubscribe)
s.RegisterTool(tp+"list", sdk.ToolDef{
Name: tp + "list", Description: "List all subscribed feeds",
Parameters: map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{},
},
}, p.handleList)
s.RegisterTool(tp+"check_now", sdk.ToolDef{
Name: tp + "check_now", Description: "Manually check all feeds for new articles now",
NoMemory: true,
Parameters: map[string]interface{}{
"type": "object",
"properties": map[string]interface{}{},
},
}, p.handleCheckNow)
pollMin := int(getSetting(s.Settings(), "poll_interval", int64(30)))
if pollMin < 5 {
pollMin = 5
}
p.pollTicker = time.NewTicker(time.Duration(pollMin) * time.Minute)
p.wg.Add(1)
go p.pollLoop()
fmt.Printf("[%s] started (%d feeds, poll every %dm)\n", p.name, len(p.feeds), pollMin)
return nil
}
func (p *Plugin) Stop() error {
p.stopOnce.Do(func() { close(p.stopCh) })
p.pollTicker.Stop()
p.wg.Wait()
p.saveData()
fmt.Printf("[%s] stopped\n", p.name)
return nil
}
func (p *Plugin) pollLoop() {
defer p.wg.Done()
p.checkAllFeeds()
for {
select {
case <-p.pollTicker.C:
p.checkAllFeeds()
case <-p.stopCh:
return
}
}
}
func (p *Plugin) checkAllFeeds() {
p.mu.RLock()
feeds := make([]FeedSub, len(p.feeds))
copy(feeds, p.feeds)
p.mu.RUnlock()
for _, feed := range feeds {
select {
case <-p.stopCh:
return
default:
}
p.checkFeed(feed)
}
}
func (p *Plugin) checkFeed(sub FeedSub) {
parsed, err := p.fp.ParseURL(sub.URL)
if err != nil {
return
}
title := parsed.Title
if title == "" {
title = sub.URL
}
var newArticles []*gofeed.Item
for _, item := range parsed.Items {
guid := item.GUID
if guid == "" {
guid = item.Link
}
if guid == "" {
continue
}
guid = sub.URL + "|" + guid
p.mu.RLock()
seen := p.seenGUIDs[guid]
p.mu.RUnlock()
if !seen {
newArticles = append(newArticles, item)
}
}
if len(newArticles) == 0 {
return
}
now := time.Now()
toInject := make([]*gofeed.Item, 0, len(newArticles))
p.mu.Lock()
for _, item := range newArticles {
guid := item.GUID
if guid == "" {
guid = item.Link
}
if guid == "" {
continue
}
key := sub.URL + "|" + guid
if t, ok := p.injected[key]; ok && now.Sub(t) < injectDedupWindow {
continue
}
p.injected[key] = now
p.seenGUIDs[key] = true
toInject = append(toInject, item)
}
p.mu.Unlock()
if len(toInject) == 0 {
return
}
var lines []string
lines = append(lines, fmt.Sprintf("📡 %s (%s) — %d 篇新文章:", title, sub.URL, len(toInject)))
for _, item := range toInject {
pubDate := ""
if item.PublishedParsed != nil {
pubDate = item.PublishedParsed.Format("01-02 15:04")
}
line := fmt.Sprintf(" • %s", item.Title)
if pubDate != "" {
line += fmt.Sprintf(" [%s]", pubDate)
}
if item.Link != "" {
line += "\n " + item.Link
}
lines = append(lines, line)
}
// 中断注入是「系统通知」NoMemory 写明意图:这类提醒不参与记忆计算,
// 原文仍进上下文(模型当轮看得到)。
p.sdk.InjectInterruptTextOpts("rss", "rss", strings.Join(lines, "\n"),
sdk.InjectOptions{NoMemory: true})
p.saveData()
}
func (p *Plugin) handleSubscribe(args map[string]interface{}) (interface{}, error) {
url := readArg(args, "url")
if url == "" {
return map[string]interface{}{"isError": true, "content": "URL is required"}, nil
}
p.mu.RLock()
for _, f := range p.feeds {
if f.URL == url {
p.mu.RUnlock()
return map[string]interface{}{"isError": true, "content": "Already subscribed to: " + url}, nil
}
}
p.mu.RUnlock()
interval := readArgInt(args, "interval", 30)
if interval < 5 {
interval = 5
}
parsed, err := p.fp.ParseURL(url)
if err != nil {
return map[string]interface{}{"isError": true, "content": "Failed to parse feed: " + err.Error()}, nil
}
feedTitle := parsed.Title
if feedTitle == "" {
feedTitle = url
}
sub := FeedSub{
URL: url,
Title: feedTitle,
AddedAt: time.Now().Format("2006-01-02 15:04"),
Interval: interval,
}
guidCount := 0
p.mu.Lock()
for _, item := range parsed.Items {
guid := item.GUID
if guid == "" {
guid = item.Link
}
if guid == "" {
continue
}
p.seenGUIDs[url+"|"+guid] = true
guidCount++
}
p.mu.Unlock()
p.mu.Lock()
p.feeds = append(p.feeds, sub)
p.mu.Unlock()
p.saveData()
return map[string]interface{}{
"content": fmt.Sprintf("Subscribed to: %s\nTitle: %s\nArticles found: %d\nPoll interval: %d min", url, feedTitle, guidCount, interval),
}, nil
}
func (p *Plugin) handleUnsubscribe(args map[string]interface{}) (interface{}, error) {
url := readArg(args, "url")
if url == "" {
return map[string]interface{}{"isError": true, "content": "URL is required"}, nil
}
p.mu.Lock()
found := false
for i, f := range p.feeds {
if f.URL == url {
p.feeds = append(p.feeds[:i], p.feeds[i+1:]...)
found = true
break
}
}
if !found {
p.mu.Unlock()
return map[string]interface{}{"isError": true, "content": "Not subscribed to: " + url}, nil
}
for guid := range p.seenGUIDs {
if strings.HasPrefix(guid, url+"|") {
delete(p.seenGUIDs, guid)
}
}
p.mu.Unlock()
p.saveData()
return map[string]interface{}{"content": "Unsubscribed: " + url}, nil
}
func (p *Plugin) handleList(args map[string]interface{}) (interface{}, error) {
p.mu.RLock()
defer p.mu.RUnlock()
if len(p.feeds) == 0 {
return map[string]interface{}{"content": "No subscriptions. Use rss_subscribe to add one."}, nil
}
sort.Slice(p.feeds, func(i, j int) bool {
return p.feeds[i].Title < p.feeds[j].Title
})
var lines []string
lines = append(lines, fmt.Sprintf("📡 Subscriptions (%d):", len(p.feeds)))
for _, f := range p.feeds {
lines = append(lines, fmt.Sprintf(" • %s\n %s (every %dm, added %s)", f.Title, f.URL, f.Interval, f.AddedAt))
}
return map[string]interface{}{"content": strings.Join(lines, "\n")}, nil
}
func (p *Plugin) handleCheckNow(args map[string]interface{}) (interface{}, error) {
select {
case <-p.stopCh:
return map[string]interface{}{"isError": true, "content": "plugin is stopping"}, nil
default:
}
p.wg.Add(1)
go func() {
defer p.wg.Done()
p.checkAllFeeds()
}()
return map[string]interface{}{"content": "Checking all feeds for updates..."}, nil
}
func (p *Plugin) dataFile() string {
return filepath.Join(p.dataDir, "feeds.json")
}
func (p *Plugin) loadData() {
b, err := os.ReadFile(p.dataFile())
if err != nil {
return
}
var data struct {
Feeds []FeedSub `json:"feeds"`
SeenGUIDs map[string]bool `json:"seen"`
}
if json.Unmarshal(b, &data) != nil {
return
}
if data.Feeds != nil {
p.feeds = data.Feeds
}
if data.SeenGUIDs != nil {
p.seenGUIDs = data.SeenGUIDs
}
}
func (p *Plugin) saveData() {
p.mu.RLock()
defer p.mu.RUnlock()
data := struct {
Feeds []FeedSub `json:"feeds"`
SeenGUIDs map[string]bool `json:"seen"`
}{
Feeds: p.feeds,
SeenGUIDs: p.seenGUIDs,
}
b, _ := json.MarshalIndent(data, "", " ")
atomicWriteJSON(p.dataFile(), b)
}
// cleanupData 卸载时清理订阅数据目录feeds.json 等)
func (p *Plugin) cleanupData() {
p.mu.Lock()
defer p.mu.Unlock()
if p.dataDir == "" {
return
}
for _, f := range []string{"feeds.json"} {
path := filepath.Join(p.dataDir, f)
if err := os.Remove(path); err != nil && !os.IsNotExist(err) {
fmt.Printf("[%s] onRemove cleanup %s: %v\n", p.name, path, err)
}
}
}
// atomicWriteJSON 原子写 JSON先写临时文件再 rename避免进程崩溃截断数据文件。
func atomicWriteJSON(path string, data []byte) error {
tmp := path + ".tmp"
if err := os.WriteFile(tmp, data, 0644); err != nil {
return err
}
return os.Rename(tmp, path)
}