mirror of
https://github.com/prosolis/gogobee.git
synced 2026-09-14 19:01:09 +00:00
Compare commits
7
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1e00d2585f | ||
|
|
b1e6937c0c | ||
|
|
ae5e10435a | ||
|
|
77dde5d133 | ||
|
|
583616f9d0 | ||
|
|
3f9c338e67 | ||
|
|
afa49b2d4c |
@@ -13,7 +13,7 @@ require (
|
||||
github.com/robfig/cron/v3 v3.0.1
|
||||
github.com/rs/zerolog v1.35.1
|
||||
github.com/skip2/go-qrcode v0.0.0-20200617195104-da1b6568686e
|
||||
golang.org/x/image v0.40.0
|
||||
golang.org/x/image v0.45.0
|
||||
maunium.net/go/mautrix v0.28.1
|
||||
modernc.org/sqlite v1.50.1
|
||||
)
|
||||
@@ -40,8 +40,8 @@ require (
|
||||
golang.org/x/crypto v0.53.0 // indirect
|
||||
golang.org/x/exp v0.0.0-20260611194520-c48552f49976 // indirect
|
||||
golang.org/x/net v0.56.0 // indirect
|
||||
golang.org/x/sys v0.46.0 // indirect
|
||||
golang.org/x/text v0.38.0 // indirect
|
||||
golang.org/x/sys v0.47.0 // indirect
|
||||
golang.org/x/text v0.41.0 // indirect
|
||||
gopkg.in/yaml.v3 v3.0.1 // indirect
|
||||
modernc.org/libc v1.72.3 // indirect
|
||||
modernc.org/mathutil v1.7.1 // indirect
|
||||
|
||||
@@ -87,15 +87,15 @@ golang.org/x/crypto v0.53.0 h1:QZ4Muo8THX6CizN2vPPd5fBGHyogrdK9fG4wLPFUsto=
|
||||
golang.org/x/crypto v0.53.0/go.mod h1:DNLU434OwVakk9PzuwV8w62mAJpRJL3vsgcfp4Qnsio=
|
||||
golang.org/x/exp v0.0.0-20260611194520-c48552f49976 h1:X8Hz2ImujgbmetVuW+w2YkyZChE3cBpZi2P158rTG9M=
|
||||
golang.org/x/exp v0.0.0-20260611194520-c48552f49976/go.mod h1:vnf4pv9iKZXY58sQE1L86zmNWJ4159e1RkcWiLCkeEY=
|
||||
golang.org/x/image v0.40.0 h1:Tw4GyDXMo+daZN1znreBRC3VayR1aLFUyUEOLUdW1a8=
|
||||
golang.org/x/image v0.40.0/go.mod h1:uIc348UZMSvS5Z65CVZ7iDPaNobNFEPeJ4kbqTOszmA=
|
||||
golang.org/x/image v0.45.0 h1:FMb1nTbH5H9vF55SriQHgFw5GnNL9Jg6L25BwXKzhB0=
|
||||
golang.org/x/image v0.45.0/go.mod h1:n62x/7RqlwXDvGsSU4u6IUTUf6KghUZ9Bt7cG/T9Fx4=
|
||||
golang.org/x/mod v0.6.0-dev.0.20220419223038-86c51ed26bb4/go.mod h1:jJ57K6gSWd91VN4djpZkiMVwK6gcyfeH4XE8wZrZaV4=
|
||||
golang.org/x/mod v0.8.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs=
|
||||
golang.org/x/mod v0.12.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs=
|
||||
golang.org/x/mod v0.15.0/go.mod h1:hTbmBsO62+eylJbnUtE2MGJUyE7QWk4xUqPFrRgJ+7c=
|
||||
golang.org/x/mod v0.17.0/go.mod h1:hTbmBsO62+eylJbnUtE2MGJUyE7QWk4xUqPFrRgJ+7c=
|
||||
golang.org/x/mod v0.37.0 h1:vF1DjpVEshcIqoEaauuHebaLk1O1forxjxBaVn884JQ=
|
||||
golang.org/x/mod v0.37.0/go.mod h1:m8S8VeM9r4dzDwjrKO0a1sZP3YjeMamRRlD+fmR2Q/0=
|
||||
golang.org/x/mod v0.38.0 h1:MECBjubtXD7yj4HrhIUcywNaGeNVUdfVnxmPajOk4yk=
|
||||
golang.org/x/mod v0.38.0/go.mod h1:V6Xz0pq8TQ3dGqVQ1FVHuelZpAL0uNhSkk9ogYP3c40=
|
||||
golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s=
|
||||
golang.org/x/net v0.0.0-20210226172049-e18ecbb05110/go.mod h1:m0MpNAwzfU5UDzcl9v0D8zg8gWTRqZa9RBIspLL5mdg=
|
||||
golang.org/x/net v0.0.0-20220722155237-a158d28d115b/go.mod h1:XRhObCWvk6IyKnWLug+ECip1KBveYUHfp+8e9klMJ9c=
|
||||
@@ -114,8 +114,8 @@ golang.org/x/sync v0.3.0/go.mod h1:FU7BRWz2tNW+3quACPkgCx/L+uEAv1htQ0V83Z9Rj+Y=
|
||||
golang.org/x/sync v0.6.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk=
|
||||
golang.org/x/sync v0.7.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk=
|
||||
golang.org/x/sync v0.10.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk=
|
||||
golang.org/x/sync v0.21.0 h1:HLII4xRRTtCRkxYp4HNFF0Js/Og6q2i++KXbg0gHCwM=
|
||||
golang.org/x/sync v0.21.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0=
|
||||
golang.org/x/sync v0.22.0 h1:SZjpbeLmrCk4xhRSZFNZW5gFUeCeFgjekvI/+gfScek=
|
||||
golang.org/x/sync v0.22.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0=
|
||||
golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
|
||||
golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
|
||||
golang.org/x/sys v0.0.0-20210615035016-665e8c7367d1/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||
@@ -128,8 +128,8 @@ golang.org/x/sys v0.12.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||
golang.org/x/sys v0.17.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
|
||||
golang.org/x/sys v0.20.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
|
||||
golang.org/x/sys v0.28.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
|
||||
golang.org/x/sys v0.46.0 h1:noSf2Fq6F8DBgS+LysIkx7rIExoNHJsxOAtPp4rthXw=
|
||||
golang.org/x/sys v0.46.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
|
||||
golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs=
|
||||
golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
|
||||
golang.org/x/telemetry v0.0.0-20240228155512-f48c80bd79b2/go.mod h1:TeRTkGYfJXctD9OcfyVLyj2J3IxLnKwHJR8f4D8a3YE=
|
||||
golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo=
|
||||
golang.org/x/term v0.0.0-20210927222741-03fcf44c2211/go.mod h1:jbD1KX2456YbFQfuXm/mYQcufACuNUgVhRMnK/tPxf8=
|
||||
@@ -148,16 +148,16 @@ golang.org/x/text v0.13.0/go.mod h1:TvPlkZtksWOMsz7fbANvkp4WM8x/WCo/om8BMLbz+aE=
|
||||
golang.org/x/text v0.14.0/go.mod h1:18ZOQIKpY8NJVqYksKHtTdi31H5itFRjB5/qKTNYzSU=
|
||||
golang.org/x/text v0.15.0/go.mod h1:18ZOQIKpY8NJVqYksKHtTdi31H5itFRjB5/qKTNYzSU=
|
||||
golang.org/x/text v0.21.0/go.mod h1:4IBbMaMmOPCJ8SecivzSH54+73PCFmPWxNTLm+vZkEQ=
|
||||
golang.org/x/text v0.38.0 h1:sXmwo9DwP3OK9EZ7PqAdaooSGozfl/3a6/xJcbzPRhE=
|
||||
golang.org/x/text v0.38.0/go.mod h1:YXZt3QhHUKYT53r2lLKFIVi6Ao1jdzrTR/KQ09qyxF4=
|
||||
golang.org/x/text v0.41.0 h1:vz/seA0lnX87Othu2f/0L24RcgrXD9/YFTSuGjj3rH8=
|
||||
golang.org/x/text v0.41.0/go.mod h1:jvf1O8ajNzZqhSrQBPbutR/EB83Cc0CFrezNQIwbb5M=
|
||||
golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ=
|
||||
golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
|
||||
golang.org/x/tools v0.1.12/go.mod h1:hNGJHUnrk76NpqgfD5Aqm5Crs+Hm0VOH/i9J2+nxYbc=
|
||||
golang.org/x/tools v0.6.0/go.mod h1:Xwgl3UAJ/d3gWutnCtw505GrjyAbvKui8lOU390QaIU=
|
||||
golang.org/x/tools v0.13.0/go.mod h1:HvlwmtVNQAhOuCjW7xxvovg8wbNq7LwfXh/k7wXUl58=
|
||||
golang.org/x/tools v0.21.1-0.20240508182429-e35e4ccd0d2d/go.mod h1:aiJjzUbINMkxbQROHiO6hDPo2LHcIPhhQsa9DLh0yGk=
|
||||
golang.org/x/tools v0.46.0 h1:7jTurBkPZu4moS/Uy4OQT1M+QBlsj3wejyZwsT8Z7rk=
|
||||
golang.org/x/tools v0.46.0/go.mod h1:FrD85F8l+NWL+9XWBSyVSHO6Ne4jutsfIFba7AWQ5Ys=
|
||||
golang.org/x/tools v0.48.0 h1:3+hClM1aLL5mjMKm5ovokw9epgRXPuu2tILgismM6RE=
|
||||
golang.org/x/tools v0.48.0/go.mod h1:08xX0orndb/F7jJxGDicx061tyd5pcMto75YMAXr6lk=
|
||||
golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
|
||||
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405 h1:yhCVgyC4o1eVCa2tZl7eS0r+SDo693bJlVdllGtEeKM=
|
||||
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
|
||||
|
||||
@@ -1255,6 +1255,16 @@ CREATE TABLE IF NOT EXISTS presence (
|
||||
updated_at INTEGER DEFAULT (unixepoch())
|
||||
);
|
||||
|
||||
-- DM rooms. The bot's m.direct account data is not a reliable store for an
|
||||
-- appservice user (no /sync, and nothing writes it back), so the mapping lives
|
||||
-- here. Without it every restart lost the cache and the bot created a fresh DM
|
||||
-- room per user.
|
||||
CREATE TABLE IF NOT EXISTS dm_rooms (
|
||||
user_id TEXT PRIMARY KEY,
|
||||
room_id TEXT NOT NULL,
|
||||
updated_at INTEGER DEFAULT (unixepoch())
|
||||
);
|
||||
|
||||
-- Markov
|
||||
CREATE TABLE IF NOT EXISTS markov_corpus (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
|
||||
@@ -0,0 +1,161 @@
|
||||
// Package llm wraps the local inference endpoint behind a backend-agnostic
|
||||
// interface. Two concrete backends — Ollama (native /api/generate) and vLLM
|
||||
// (OpenAI-compatible /v1/chat/completions) — implement Client; plugin code
|
||||
// calls the interface only and never knows which one is active.
|
||||
//
|
||||
// Deliberately not routed through internal/safehttp: that client blocks
|
||||
// RFC1918 and loopback destinations to defend against SSRF from feed-supplied
|
||||
// URLs, and the inference endpoint is precisely such a destination. The URL
|
||||
// here comes from our own config, never from user input.
|
||||
package llm
|
||||
|
||||
import (
|
||||
"context"
|
||||
"os"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Request is the backend-neutral generation request. Zero-valued fields fall
|
||||
// back to backend defaults.
|
||||
type Request struct {
|
||||
// Prompt is a single raw instruction. The vLLM backend wraps it as one
|
||||
// user message so the model's chat template still applies; sending it to
|
||||
// /v1/completions instead would bypass the template and degrade an
|
||||
// instruction-tuned model badly.
|
||||
Prompt string
|
||||
// System is an optional system message. Empty means none, which keeps the
|
||||
// single-message shape the majority of callers use.
|
||||
System string
|
||||
// NumCtx is the per-request context window. Ollama honours it directly;
|
||||
// vLLM fixes the window server-side at launch (--max-model-len), so this
|
||||
// is ignored there rather than silently misapplied.
|
||||
NumCtx int
|
||||
// MaxTokens caps the completion length. 0 means the backend default.
|
||||
MaxTokens int
|
||||
// Temperature is passed through when non-zero.
|
||||
Temperature float64
|
||||
// Timeout overrides the client's default per-request budget.
|
||||
Timeout time.Duration
|
||||
}
|
||||
|
||||
// Client is the single surface plugin code depends on.
|
||||
type Client interface {
|
||||
// Generate returns the full completion in one shot. Reasoning blocks are
|
||||
// stripped before returning (see StripThink) — every caller in this repo
|
||||
// wants the visible answer, not the chain of thought.
|
||||
Generate(ctx context.Context, req Request) (string, error)
|
||||
// Model reports the configured model id, for logging and /botinfo.
|
||||
Model() string
|
||||
// Ping reports the model ids the backend is currently serving. Used by
|
||||
// /botinfo for a liveness line; the two backends expose this on different
|
||||
// paths (/api/tags vs /v1/models), which is exactly the sort of difference
|
||||
// this interface exists to hide.
|
||||
Ping(ctx context.Context) ([]string, error)
|
||||
// Backend reports "ollama" or "vllm", for logging and /botinfo.
|
||||
Backend() string
|
||||
}
|
||||
|
||||
// Config selects and configures a backend.
|
||||
type Config struct {
|
||||
Backend string // "ollama" | "vllm"
|
||||
Endpoint string
|
||||
Model string
|
||||
Timeout time.Duration
|
||||
}
|
||||
|
||||
// DefaultTimeout matches the budget the pre-refactor callOllama used.
|
||||
const DefaultTimeout = 120 * time.Second
|
||||
|
||||
// ConfigFromEnv reads backend settings, preferring the new LLM_* names and
|
||||
// falling back to the legacy OLLAMA_* pair so an existing deployment keeps
|
||||
// working untouched after this refactor.
|
||||
func ConfigFromEnv() Config {
|
||||
backend := strings.ToLower(strings.TrimSpace(os.Getenv("LLM_BACKEND")))
|
||||
if backend == "" {
|
||||
backend = "ollama"
|
||||
}
|
||||
|
||||
endpoint := firstNonEmpty(os.Getenv("LLM_ENDPOINT"), os.Getenv("OLLAMA_HOST"))
|
||||
model := firstNonEmpty(os.Getenv("LLM_MODEL"), os.Getenv("OLLAMA_MODEL"))
|
||||
|
||||
timeout := DefaultTimeout
|
||||
if d, err := time.ParseDuration(os.Getenv("LLM_TIMEOUT")); err == nil && d > 0 {
|
||||
timeout = d
|
||||
}
|
||||
|
||||
return Config{Backend: backend, Endpoint: endpoint, Model: model, Timeout: timeout}
|
||||
}
|
||||
|
||||
// Configured reports whether enough config is present to talk to a backend.
|
||||
// Plugins check this to stay dormant rather than erroring on every invocation,
|
||||
// which is what the old `if ollamaHost == "" || ollamaModel == ""` guards did.
|
||||
func (c Config) Configured() bool {
|
||||
return c.Endpoint != "" && c.Model != ""
|
||||
}
|
||||
|
||||
// New builds the client for cfg.Backend. An unrecognised backend falls back to
|
||||
// Ollama, which is what every existing deployment runs.
|
||||
func New(cfg Config) Client {
|
||||
if cfg.Timeout <= 0 {
|
||||
cfg.Timeout = DefaultTimeout
|
||||
}
|
||||
base := backend{
|
||||
endpoint: strings.TrimRight(cfg.Endpoint, "/"),
|
||||
model: cfg.Model,
|
||||
timeout: cfg.Timeout,
|
||||
}
|
||||
switch cfg.Backend {
|
||||
case "vllm":
|
||||
return &VLLMClient{base}
|
||||
default:
|
||||
return &OllamaClient{base}
|
||||
}
|
||||
}
|
||||
|
||||
// backend holds the fields shared by both concrete clients.
|
||||
type backend struct {
|
||||
endpoint string
|
||||
model string
|
||||
timeout time.Duration
|
||||
}
|
||||
|
||||
func (b backend) Model() string { return b.model }
|
||||
|
||||
// timeoutFor lets a single call widen or narrow the client default. The two
|
||||
// dispatch-voice callers rely on this: a dispatch is authored on a game
|
||||
// chokepoint and must not stall it, while a run summary rides a background
|
||||
// ticker and can afford a bigger model.
|
||||
func (b backend) timeoutFor(req Request) time.Duration {
|
||||
if req.Timeout > 0 {
|
||||
return req.Timeout
|
||||
}
|
||||
return b.timeout
|
||||
}
|
||||
|
||||
// StripThink removes a leading <think>...</think> reasoning block, which Qwen
|
||||
// models emit even when thinking is disabled by some backends. Callers that
|
||||
// parse JSON out of the completion depend on this running first.
|
||||
func StripThink(s string) string {
|
||||
for {
|
||||
i := strings.Index(s, "<think>")
|
||||
if i < 0 {
|
||||
break
|
||||
}
|
||||
j := strings.Index(s, "</think>")
|
||||
if j < 0 || j < i {
|
||||
break
|
||||
}
|
||||
s = s[:i] + s[j+len("</think>"):]
|
||||
}
|
||||
return strings.TrimSpace(s)
|
||||
}
|
||||
|
||||
func firstNonEmpty(vals ...string) string {
|
||||
for _, v := range vals {
|
||||
if v = strings.TrimSpace(v); v != "" {
|
||||
return v
|
||||
}
|
||||
}
|
||||
return ""
|
||||
}
|
||||
@@ -0,0 +1,103 @@
|
||||
package llm
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
)
|
||||
|
||||
// OllamaClient talks to Ollama's native /api/generate endpoint.
|
||||
type OllamaClient struct {
|
||||
backend
|
||||
}
|
||||
|
||||
func (c *OllamaClient) Backend() string { return "ollama" }
|
||||
|
||||
type ollamaOptions struct {
|
||||
NumCtx int `json:"num_ctx,omitempty"`
|
||||
NumPredict int `json:"num_predict,omitempty"`
|
||||
Temperature float64 `json:"temperature,omitempty"`
|
||||
}
|
||||
|
||||
type ollamaRequest struct {
|
||||
Model string `json:"model"`
|
||||
Prompt string `json:"prompt"`
|
||||
System string `json:"system,omitempty"`
|
||||
Stream bool `json:"stream"`
|
||||
Think bool `json:"think"`
|
||||
Options ollamaOptions `json:"options,omitempty"`
|
||||
}
|
||||
|
||||
// Generate posts a single non-streaming generation and returns the completion.
|
||||
func (c *OllamaClient) Generate(ctx context.Context, req Request) (string, error) {
|
||||
ctx, cancel := context.WithTimeout(ctx, c.timeoutFor(req))
|
||||
defer cancel()
|
||||
|
||||
body, err := json.Marshal(ollamaRequest{
|
||||
Model: c.model,
|
||||
Prompt: req.Prompt,
|
||||
System: req.System,
|
||||
Stream: false,
|
||||
Think: false,
|
||||
Options: ollamaOptions{
|
||||
NumCtx: req.NumCtx,
|
||||
NumPredict: req.MaxTokens,
|
||||
Temperature: req.Temperature,
|
||||
},
|
||||
})
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("ollama: marshal payload: %w", err)
|
||||
}
|
||||
|
||||
httpReq, err := http.NewRequestWithContext(ctx, http.MethodPost,
|
||||
c.endpoint+"/api/generate", bytes.NewReader(body))
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("ollama: build request: %w", err)
|
||||
}
|
||||
httpReq.Header.Set("Content-Type", "application/json")
|
||||
|
||||
resp, err := http.DefaultClient.Do(httpReq)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("ollama request: %w", err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
respBody, err := io.ReadAll(resp.Body)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("ollama: read response: %w", err)
|
||||
}
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return "", fmt.Errorf("ollama HTTP %d: %s", resp.StatusCode, string(respBody))
|
||||
}
|
||||
|
||||
var result struct {
|
||||
Response string `json:"response"`
|
||||
}
|
||||
if err := json.Unmarshal(respBody, &result); err != nil {
|
||||
return "", fmt.Errorf("ollama: parse response: %w", err)
|
||||
}
|
||||
return StripThink(result.Response), nil
|
||||
}
|
||||
|
||||
// Ping lists locally installed models via Ollama's native /api/tags.
|
||||
func (c *OllamaClient) Ping(ctx context.Context) ([]string, error) {
|
||||
ctx, cancel := context.WithTimeout(ctx, pingTimeout)
|
||||
defer cancel()
|
||||
|
||||
var out struct {
|
||||
Models []struct {
|
||||
Name string `json:"name"`
|
||||
} `json:"models"`
|
||||
}
|
||||
if err := getJSON(ctx, c.endpoint+"/api/tags", &out); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
names := make([]string, 0, len(out.Models))
|
||||
for _, m := range out.Models {
|
||||
names = append(names, m.Name)
|
||||
}
|
||||
return names, nil
|
||||
}
|
||||
@@ -0,0 +1,29 @@
|
||||
package llm
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"time"
|
||||
)
|
||||
|
||||
// pingTimeout keeps a liveness probe short — /botinfo renders synchronously and
|
||||
// a hung endpoint must not hold the reply.
|
||||
const pingTimeout = 5 * time.Second
|
||||
|
||||
func getJSON(ctx context.Context, url string, out any) error {
|
||||
req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
resp, err := http.DefaultClient.Do(req)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return fmt.Errorf("HTTP %d", resp.StatusCode)
|
||||
}
|
||||
return json.NewDecoder(resp.Body).Decode(out)
|
||||
}
|
||||
@@ -0,0 +1,116 @@
|
||||
package llm
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
)
|
||||
|
||||
// VLLMClient talks to an OpenAI-compatible /v1/chat/completions endpoint.
|
||||
type VLLMClient struct {
|
||||
backend
|
||||
}
|
||||
|
||||
func (c *VLLMClient) Backend() string { return "vllm" }
|
||||
|
||||
type vllmMessage struct {
|
||||
Role string `json:"role"`
|
||||
Content string `json:"content"`
|
||||
}
|
||||
|
||||
type vllmRequest struct {
|
||||
Model string `json:"model"`
|
||||
Messages []vllmMessage `json:"messages"`
|
||||
MaxTokens int `json:"max_tokens,omitempty"`
|
||||
Temperature float64 `json:"temperature,omitempty"`
|
||||
Stream bool `json:"stream"`
|
||||
// ChatTemplateKwargs is a vLLM extension to the OpenAI schema. It is how
|
||||
// Qwen3-family reasoning is switched off; the Ollama backend spells the
|
||||
// same intent as its native "think": false.
|
||||
ChatTemplateKwargs map[string]any `json:"chat_template_kwargs,omitempty"`
|
||||
}
|
||||
|
||||
// Generate posts a single non-streaming completion and returns the message
|
||||
// content. The raw prompt is sent as one user message so the server-side chat
|
||||
// template still wraps it.
|
||||
func (c *VLLMClient) Generate(ctx context.Context, req Request) (string, error) {
|
||||
ctx, cancel := context.WithTimeout(ctx, c.timeoutFor(req))
|
||||
defer cancel()
|
||||
|
||||
msgs := make([]vllmMessage, 0, 2)
|
||||
if req.System != "" {
|
||||
msgs = append(msgs, vllmMessage{Role: "system", Content: req.System})
|
||||
}
|
||||
msgs = append(msgs, vllmMessage{Role: "user", Content: req.Prompt})
|
||||
|
||||
body, err := json.Marshal(vllmRequest{
|
||||
Model: c.model,
|
||||
Messages: msgs,
|
||||
MaxTokens: req.MaxTokens,
|
||||
Temperature: req.Temperature,
|
||||
Stream: false,
|
||||
ChatTemplateKwargs: map[string]any{"enable_thinking": false},
|
||||
})
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("vllm: marshal payload: %w", err)
|
||||
}
|
||||
|
||||
httpReq, err := http.NewRequestWithContext(ctx, http.MethodPost,
|
||||
c.endpoint+"/v1/chat/completions", bytes.NewReader(body))
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("vllm: build request: %w", err)
|
||||
}
|
||||
httpReq.Header.Set("Content-Type", "application/json")
|
||||
|
||||
resp, err := http.DefaultClient.Do(httpReq)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("vllm request: %w", err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
respBody, err := io.ReadAll(resp.Body)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("vllm: read response: %w", err)
|
||||
}
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return "", fmt.Errorf("vllm HTTP %d: %s", resp.StatusCode, string(respBody))
|
||||
}
|
||||
|
||||
var result struct {
|
||||
Choices []struct {
|
||||
Message struct {
|
||||
Content string `json:"content"`
|
||||
} `json:"message"`
|
||||
} `json:"choices"`
|
||||
}
|
||||
if err := json.Unmarshal(respBody, &result); err != nil {
|
||||
return "", fmt.Errorf("vllm: parse response: %w", err)
|
||||
}
|
||||
if len(result.Choices) == 0 {
|
||||
return "", fmt.Errorf("vllm: empty choices in response")
|
||||
}
|
||||
return StripThink(result.Choices[0].Message.Content), nil
|
||||
}
|
||||
|
||||
// Ping lists served models via the OpenAI-compatible /v1/models.
|
||||
func (c *VLLMClient) Ping(ctx context.Context) ([]string, error) {
|
||||
ctx, cancel := context.WithTimeout(ctx, pingTimeout)
|
||||
defer cancel()
|
||||
|
||||
var out struct {
|
||||
Data []struct {
|
||||
ID string `json:"id"`
|
||||
} `json:"data"`
|
||||
}
|
||||
if err := getJSON(ctx, c.endpoint+"/v1/models", &out); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
names := make([]string, 0, len(out.Data))
|
||||
for _, m := range out.Data {
|
||||
names = append(names, m.ID)
|
||||
}
|
||||
return names, nil
|
||||
}
|
||||
@@ -1374,6 +1374,7 @@ func (p *AdventurePlugin) sendTreasureDiscoveryDM(userID id.UserID, char *Advent
|
||||
"{treasure_name}": def.Name,
|
||||
"{bonus_desc}": def.InventoryDesc,
|
||||
"{location}": loc.Name,
|
||||
"{location_mid}": advLocationMidSentence(loc.Name),
|
||||
})
|
||||
|
||||
p.SendDM(userID, text)
|
||||
@@ -1400,8 +1401,9 @@ func (p *AdventurePlugin) announceTreasureToRoom(char *AdventureCharacter, def *
|
||||
}
|
||||
displayName, _ := loadDisplayName(char.UserID)
|
||||
announce := advSubstituteFlavor(def.RoomAnnounce, map[string]string{
|
||||
"{name}": displayName,
|
||||
"{location}": loc.Name,
|
||||
"{name}": displayName,
|
||||
"{location}": loc.Name,
|
||||
"{location_mid}": advLocationMidSentence(loc.Name),
|
||||
})
|
||||
p.SendMessage(id.RoomID(gr), announce)
|
||||
}
|
||||
|
||||
@@ -232,7 +232,7 @@ var TreasureDiscovery = map[int][]string{
|
||||
// No exclamation marks. No enthusiasm. Just weight.
|
||||
5: {
|
||||
"The {treasure_name}.\n\n" +
|
||||
"You found it in the Abyssal Maw. It found you in the Abyssal Maw. " +
|
||||
"You found it in {location_mid}. It found you in {location_mid}. " +
|
||||
"The distinction matters less at depth.\n" +
|
||||
"BONUS: {bonus_desc}.\n\n" +
|
||||
"It came with you. Some things don't come with you. This one did. " +
|
||||
@@ -240,15 +240,15 @@ var TreasureDiscovery = map[int][]string{
|
||||
|
||||
"The {treasure_name} is in your inventory.\n\n" +
|
||||
"BONUS: {bonus_desc}.\n\n" +
|
||||
"This is a Tier 5 rare. The Abyssal Maw doesn't give these up. " +
|
||||
"The Abyssal Maw gave this one up. " +
|
||||
"This is a Tier 5 rare. Nothing in {location_mid} gives these up. " +
|
||||
"Something in {location_mid} gave this one up. " +
|
||||
"That's a sentence worth sitting with.",
|
||||
|
||||
"You have the {treasure_name}.\n\n" +
|
||||
"BONUS: {bonus_desc}.\n\n" +
|
||||
"It was in the Abyssal Maw. The things that were between you and it " +
|
||||
"It was in {location_mid}. The things that were between you and it " +
|
||||
"are not between anything and anything anymore. " +
|
||||
"The Maw is noting this. So is the item.",
|
||||
"The place is noting this. So is the item.",
|
||||
|
||||
"The {treasure_name}.\n\n" +
|
||||
"It's warm. It was warm when you found it, which it shouldn't be, " +
|
||||
@@ -262,14 +262,14 @@ var TreasureDiscovery = map[int][]string{
|
||||
"The {treasure_name}. {bonus_desc}.\n\n" +
|
||||
"You have it. " +
|
||||
"Keep it somewhere it won't be lost. " +
|
||||
"The Abyssal Maw does not give second chances.",
|
||||
"There are no second chances in {location_mid}.",
|
||||
|
||||
"The {treasure_name} is yours.\n\n" +
|
||||
"BONUS: {bonus_desc}.\n\n" +
|
||||
"The Abyssal Maw had it. You have it now. " +
|
||||
"The Maw is aware of the transfer. " +
|
||||
"The Maw is considering its position. " +
|
||||
"You should not be in the Maw when it finishes considering.",
|
||||
"It sat in {location_mid} until you took it. You have it now. " +
|
||||
"The place is aware of the transfer. " +
|
||||
"The place is considering its position. " +
|
||||
"You should not be in {location_mid} when it finishes considering.",
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
@@ -70,6 +70,18 @@ func advClearFlavorHistory() {
|
||||
})
|
||||
}
|
||||
|
||||
// advLocationMidSentence renders a location display name for use after a
|
||||
// preposition: "The Abyssal Maw" becomes "the Abyssal Maw" so a line reads
|
||||
// "found in the Abyssal Maw" rather than "found in The Abyssal Maw".
|
||||
// Templates opt in with {location_mid}; {location} keeps the display form for
|
||||
// sentence-initial and "The {location}" positions.
|
||||
func advLocationMidSentence(name string) string {
|
||||
if strings.HasPrefix(name, "The ") {
|
||||
return "the " + strings.TrimPrefix(name, "The ")
|
||||
}
|
||||
return name
|
||||
}
|
||||
|
||||
// advSubstituteFlavor replaces {var} placeholders in a flavor text string.
|
||||
func advSubstituteFlavor(template string, vars map[string]string) string {
|
||||
pairs := make([]string, 0, len(vars)*2)
|
||||
|
||||
@@ -46,6 +46,16 @@ var advTreasureDropRates = map[int]float64{
|
||||
3: 0.008,
|
||||
4: 0.004,
|
||||
5: 0.0015,
|
||||
// Tier 6 (Mythic post-game) breaks the downward curve on purpose. The rate
|
||||
// per roll falls tier over tier because the lower tiers are ground daily;
|
||||
// a Mythic run is gated behind L18 and both T5 bosses and happens rarely,
|
||||
// so the same declining rate would mean a postgame player effectively never
|
||||
// sees a treasure. Slightly above T5 keeps the per-run odds in the band the
|
||||
// other tiers land in.
|
||||
//
|
||||
// NOTE: this rate is inert until advAllTreasures gains a tier 6 pool — see
|
||||
// the TODO there. A tier with a rate but no pool drops nothing.
|
||||
6: 0.002,
|
||||
}
|
||||
|
||||
const advMaxTreasures = 3
|
||||
@@ -213,7 +223,7 @@ var advAllTreasures = map[int][]AdvTreasureDef{
|
||||
{Type: "success_chance", Value: 10},
|
||||
},
|
||||
InventoryDesc: "[THUNDERFURY, BLESSED BLADE OF THE WINDSEEKER]. Yes you got it. +12 Combat.",
|
||||
RoomAnnounce: "⚡ Did {name} get Thunderfury? {name} got Thunderfury. [THUNDERFURY, BLESSED BLADE OF THE WINDSEEKER] has been found in {location}.",
|
||||
RoomAnnounce: "⚡ Did {name} get Thunderfury? {name} got Thunderfury. [THUNDERFURY, BLESSED BLADE OF THE WINDSEEKER] has been found in {location_mid}.",
|
||||
},
|
||||
{
|
||||
Key: "ocarina", Name: "The Ocarina (Cracked, Still Plays)", Tier: 4,
|
||||
@@ -221,6 +231,13 @@ var advAllTreasures = map[int][]AdvTreasureDef{
|
||||
InventoryDesc: "The Ocarina (Cracked). Three songs. +10 all skills. Do not play the third one.",
|
||||
},
|
||||
},
|
||||
// TODO(t6-treasures): there is no tier 6 pool yet, so Mythic zones drop no
|
||||
// treasure at all — advTreasureDropRates has a tier 6 rate waiting for it.
|
||||
// What's missing is the content: four Mythic treasures with bonuses above
|
||||
// the T5 line, InventoryDesc + RoomAnnounce strings for each (RoomAnnounce
|
||||
// must use {location_mid}, never a literal zone name), and a TIER 6 register
|
||||
// in TreasureDiscovery. Write them against internal/flavor/VOICE_CANON.md;
|
||||
// the T5 pool below is the tonal floor to clear, not the ceiling.
|
||||
5: {
|
||||
{
|
||||
Key: "shard_of_unnamed", Name: "Shard of the Unnamed", Tier: 5,
|
||||
@@ -230,13 +247,13 @@ var advAllTreasures = map[int][]AdvTreasureDef{
|
||||
{Type: "death_chance", Value: -5},
|
||||
},
|
||||
InventoryDesc: "Shard of the Unnamed. +15 Combat, +10% XP, -5% death.",
|
||||
RoomAnnounce: "🔴 {name} has recovered the Shard of the Unnamed from the Abyssal Maw. The server feels different.",
|
||||
RoomAnnounce: "🔴 {name} has recovered the Shard of the Unnamed from {location_mid}. The server feels different.",
|
||||
},
|
||||
{
|
||||
Key: "cartographers_final_map", Name: "The Cartographer's Final Map", Tier: 5,
|
||||
Bonuses: []advTreasureBonusDef{{Type: "all_skills", Value: 12}},
|
||||
InventoryDesc: "The Cartographer's Final Map. Updates on its own. +12 all skills, full map.",
|
||||
RoomAnnounce: "🔴 {name} has found the Cartographer's Final Map in the Abyssal Maw. It has their name on it. It always did.",
|
||||
RoomAnnounce: "🔴 {name} has found the Cartographer's Final Map in {location_mid}. It has their name on it. It always did.",
|
||||
},
|
||||
{
|
||||
Key: "triforce_shard", Name: "The Triforce Shard (One Third of Something Larger)", Tier: 5,
|
||||
@@ -246,7 +263,7 @@ var advAllTreasures = map[int][]AdvTreasureDef{
|
||||
// Note: +15 to chosen skill is v2 interactive
|
||||
},
|
||||
InventoryDesc: "Triforce Shard (×1/3). Warm. Waiting. +5 all skills, -8% death.",
|
||||
RoomAnnounce: "🔺 {name} has recovered a Triforce Shard from the Abyssal Maw. One third of something. The other two thirds are somewhere. Probably.",
|
||||
RoomAnnounce: "🔺 {name} has recovered a Triforce Shard from {location_mid}. One third of something. The other two thirds are somewhere. Probably.",
|
||||
},
|
||||
{
|
||||
Key: "the_corridor", Name: "The Corridor (You Know the One)", Tier: 5,
|
||||
@@ -255,7 +272,7 @@ var advAllTreasures = map[int][]AdvTreasureDef{
|
||||
{Type: "special_monthly_death_bypass", Value: 1}, // v2
|
||||
},
|
||||
InventoryDesc: "The Corridor. Folded. Don't look back. +12 all skills, monthly death bypass.",
|
||||
RoomAnnounce: "🔴 {name} found The Corridor in the Abyssal Maw. They know the one. So does it.",
|
||||
RoomAnnounce: "🔴 {name} found The Corridor in {location_mid}. They know the one. So does it.",
|
||||
},
|
||||
},
|
||||
}
|
||||
@@ -285,7 +302,12 @@ func rollAdvTreasureDropDetailed(tier int, userID id.UserID, chatLevel int, weig
|
||||
|
||||
pool, ok := advAllTreasures[tier]
|
||||
if !ok || len(pool) == 0 {
|
||||
return nil, roll, rate
|
||||
// A tier that has a rate but no pool (tier 6, until the TODO above is
|
||||
// filled) cannot drop anything. Report a zero rate rather than the
|
||||
// configured one: the caller turns a close roll into a "just missed"
|
||||
// DM, and telling a postgame player they nearly won a treasure that
|
||||
// cannot be won is worse than staying quiet.
|
||||
return nil, 0, 0
|
||||
}
|
||||
|
||||
// Pick random treasure
|
||||
|
||||
@@ -66,6 +66,36 @@ func TestAdvTreasureDropDetailed_ForcedRollGrantsTreasure(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestAdvTreasureTier6_PoolStillMissing pins the known tier 6 gap: Mythic zones
|
||||
// carry a drop rate but no treasure pool, so a forced roll yields nothing and
|
||||
// the player is told nothing either (a "just missed" DM for an unwinnable
|
||||
// treasure would be a lie).
|
||||
//
|
||||
// This test is the reminder. Writing advAllTreasures[6] — see TODO(t6-treasures)
|
||||
// in adventure_treasure.go — turns it red, and the fix is to delete it and
|
||||
// assert the tier 6 drop instead, the way ForcedRollGrantsTreasure does for
|
||||
// tier 1.
|
||||
func TestAdvTreasureTier6_PoolStillMissing(t *testing.T) {
|
||||
if err := db.Init(t.TempDir()); err != nil {
|
||||
t.Fatalf("db.Init: %v", err)
|
||||
}
|
||||
if _, ok := advAllTreasures[6]; ok {
|
||||
t.Fatal("a tier 6 treasure pool now exists — drop this test and assert the drop instead")
|
||||
}
|
||||
// The rate is configured and waiting, so the gap is the pool alone.
|
||||
if advTreasureDropRates[6] == 0 {
|
||||
t.Error("tier 6 drop rate went missing; a Mythic pool would be unreachable")
|
||||
}
|
||||
drop, roll, rate := rollAdvTreasureDropDetailed(6, "@mythic:example.org", 0, 1000)
|
||||
if drop != nil {
|
||||
t.Fatalf("tier 6 produced a drop from an empty pool: %+v", drop.Def)
|
||||
}
|
||||
// Zeroed, not merely no-drop: this is what keeps the near-miss DM quiet.
|
||||
if roll != 0 || rate != 0 {
|
||||
t.Errorf("empty pool reported roll=%v rate=%v, want 0/0 so no near-miss fires", roll, rate)
|
||||
}
|
||||
}
|
||||
|
||||
// The treasure and masterwork systems predate zones and speak AdvLocation.
|
||||
func TestAdvLocForZone(t *testing.T) {
|
||||
loc := advLocForZone(ZoneGoblinWarrens)
|
||||
|
||||
+12
-43
@@ -1,12 +1,9 @@
|
||||
package plugin
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"os"
|
||||
"strings"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
@@ -121,10 +118,8 @@ func (p *BotInfoPlugin) handleBotInfo(ctx MessageContext) error {
|
||||
sb.WriteString(fmt.Sprintf("Active reminders: %d\n", activeReminders))
|
||||
|
||||
// LLM status
|
||||
ollamaHost := os.Getenv("OLLAMA_HOST")
|
||||
if ollamaHost != "" {
|
||||
llmStatus := p.checkLLMStatus(ollamaHost)
|
||||
sb.WriteString(fmt.Sprintf("LLM status: %s\n", llmStatus))
|
||||
if llmConfigured() {
|
||||
sb.WriteString(fmt.Sprintf("LLM status: %s\n", p.checkLLMStatus()))
|
||||
} else {
|
||||
sb.WriteString("LLM status: not configured\n")
|
||||
}
|
||||
@@ -159,42 +154,16 @@ func (p *BotInfoPlugin) handleBotInfo(ctx MessageContext) error {
|
||||
return p.SendReply(ctx.RoomID, ctx.EventID, sb.String())
|
||||
}
|
||||
|
||||
func (p *BotInfoPlugin) checkLLMStatus(ollamaHost string) string {
|
||||
client := &http.Client{Timeout: 5 * time.Second}
|
||||
apiURL := strings.TrimRight(ollamaHost, "/") + "/api/tags"
|
||||
|
||||
resp, err := client.Get(apiURL)
|
||||
// checkLLMStatus reports backend liveness for /botinfo. The endpoint it probes
|
||||
// differs per backend, which the llm package hides behind Ping.
|
||||
func (p *BotInfoPlugin) checkLLMStatus() string {
|
||||
c := llmClient()
|
||||
models, err := c.Ping(context.Background())
|
||||
if err != nil {
|
||||
return fmt.Sprintf("offline (%s)", err.Error())
|
||||
return fmt.Sprintf("offline (%s: %s)", c.Backend(), err.Error())
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
if resp.StatusCode != 200 {
|
||||
return fmt.Sprintf("error (HTTP %d)", resp.StatusCode)
|
||||
if len(models) == 0 {
|
||||
return fmt.Sprintf("online (%s, no models loaded)", c.Backend())
|
||||
}
|
||||
|
||||
body, err := io.ReadAll(resp.Body)
|
||||
if err != nil {
|
||||
return "online (could not read response)"
|
||||
}
|
||||
|
||||
var result struct {
|
||||
Models []struct {
|
||||
Name string `json:"name"`
|
||||
} `json:"models"`
|
||||
}
|
||||
if err := json.Unmarshal(body, &result); err != nil {
|
||||
return "online (could not parse response)"
|
||||
}
|
||||
|
||||
modelNames := make([]string, 0, len(result.Models))
|
||||
for _, m := range result.Models {
|
||||
modelNames = append(modelNames, m.Name)
|
||||
}
|
||||
|
||||
if len(modelNames) == 0 {
|
||||
return "online (no models loaded)"
|
||||
}
|
||||
|
||||
return fmt.Sprintf("online (%d models: %s)", len(modelNames), strings.Join(modelNames, ", "))
|
||||
return fmt.Sprintf("online (%s, %d models: %s)", c.Backend(), len(models), strings.Join(models, ", "))
|
||||
}
|
||||
|
||||
@@ -350,6 +350,12 @@ func scanExpeditionRows(rows *sql.Rows) ([]*Expedition, error) {
|
||||
// A double-fire on the same expedition is a no-op.
|
||||
func (p *AdventurePlugin) deliverBriefing(e *Expedition, now time.Time) error {
|
||||
priorBriefing := e.LastBriefingAt
|
||||
// Everything logged since the previous briefing is what the overnight
|
||||
// digest reports. Captured before the CAS below clobbers the column.
|
||||
digestSince := e.StartDate
|
||||
if priorBriefing != nil {
|
||||
digestSince = *priorBriefing
|
||||
}
|
||||
threshold := time.Date(now.Year(), now.Month(), now.Day(),
|
||||
expeditionBriefingHour, 0, 0, 0, time.UTC)
|
||||
res, err := db.Get().Exec(`
|
||||
@@ -373,7 +379,7 @@ func (p *AdventurePlugin) deliverBriefing(e *Expedition, now time.Time) error {
|
||||
// DM (rollover happened recently) or force-fires processNightCamp
|
||||
// itself (safety net for stalled autopilots).
|
||||
if isEventAnchored(e) {
|
||||
return p.deliverBriefingEventAnchored(e, priorBriefing)
|
||||
return p.deliverBriefingEventAnchored(e, priorBriefing, digestSince)
|
||||
}
|
||||
|
||||
burn, err := p.nightRolloverBurn(e)
|
||||
@@ -400,6 +406,9 @@ func (p *AdventurePlugin) deliverBriefing(e *Expedition, now time.Time) error {
|
||||
|
||||
line := pickMorningBriefing(e.CurrentDay)
|
||||
body := renderMorningBriefing(e, line, burn)
|
||||
// The single daily message: fold in what the now-silent recap, night
|
||||
// check and ambient events recorded since the last briefing.
|
||||
body = appendOvernightDigest(body, e.ID, digestSince)
|
||||
if sl := p.shadowBriefingLine(e); sl != "" {
|
||||
body += "\n" + sl + "\n"
|
||||
}
|
||||
@@ -413,7 +422,13 @@ func (p *AdventurePlugin) deliverBriefing(e *Expedition, now time.Time) error {
|
||||
body += "\n" + ml
|
||||
}
|
||||
|
||||
p.fanOutExpeditionDM(e, body, p.briefingPetPrefix)
|
||||
p.fanOutExpeditionDM(e, body, p.briefingPerReader)
|
||||
// N1/A6 anchor, relocated here from the retired night-camp digest DM.
|
||||
// Only on a run still under way: an expedition that just ended does not
|
||||
// want a mid-day event landing on top of the emergence.
|
||||
if e.Status == ExpeditionStatusActive {
|
||||
p.fireDigestEventAnchor(e)
|
||||
}
|
||||
// Emergence seam: a briefing-time forced extraction (starvation / abyss
|
||||
// collapse) surfaces the players alive — roll pet arrival. Combat/patrol
|
||||
// deaths never reach deliverBriefing (the row is already abandoned), so an
|
||||
@@ -489,7 +504,9 @@ func (p *AdventurePlugin) maybeDeliverDeferredBriefing(uid id.UserID, now time.T
|
||||
//
|
||||
// priorBriefing is the last_briefing_at value as of entry into deliverBriefing
|
||||
// (before the CAS clobbered it). nil means day-1 or genuinely never rolled.
|
||||
func (p *AdventurePlugin) deliverBriefingEventAnchored(e *Expedition, priorBriefing *time.Time) error {
|
||||
// digestSince is the same instant collapsed to a non-nil cutoff (start date on
|
||||
// day 1) — the window the overnight digest reports.
|
||||
func (p *AdventurePlugin) deliverBriefingEventAnchored(e *Expedition, priorBriefing *time.Time, digestSince time.Time) error {
|
||||
now := time.Now().UTC()
|
||||
var since time.Duration
|
||||
if priorBriefing != nil {
|
||||
@@ -516,6 +533,7 @@ func (p *AdventurePlugin) deliverBriefingEventAnchored(e *Expedition, priorBrief
|
||||
|
||||
line := pickMorningBriefing(e.CurrentDay)
|
||||
body := renderMorningBriefing(e, line, burn)
|
||||
body = appendOvernightDigest(body, e.ID, digestSince)
|
||||
if sl := p.shadowBriefingLine(e); sl != "" {
|
||||
body += "\n" + sl + "\n"
|
||||
}
|
||||
@@ -529,7 +547,13 @@ func (p *AdventurePlugin) deliverBriefingEventAnchored(e *Expedition, priorBrief
|
||||
body += "\n" + ml
|
||||
}
|
||||
|
||||
p.fanOutExpeditionDM(e, body, p.briefingPetPrefix)
|
||||
p.fanOutExpeditionDM(e, body, p.briefingPerReader)
|
||||
// N1/A6 anchor, relocated here from the retired night-camp digest DM.
|
||||
// Only on a run still under way: an expedition that just ended does not
|
||||
// want a mid-day event landing on top of the emergence.
|
||||
if e.Status == ExpeditionStatusActive {
|
||||
p.fireDigestEventAnchor(e)
|
||||
}
|
||||
if forced && e.Status == ExpeditionStatusAbandoned {
|
||||
for _, uid := range expeditionAudience(e) {
|
||||
p.maybeRollPetArrivalOnEmerge(uid)
|
||||
@@ -566,9 +590,10 @@ func (p *AdventurePlugin) deliverRecap(e *Expedition, now time.Time) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// E2b: night phase wandering check fires before the recap so its
|
||||
// outcome is part of today's log when the recap renders.
|
||||
var night *NightCheck
|
||||
// E2b: night phase wandering check. Still fires on the 21:00 clock and
|
||||
// still writes its own "night" log entry — processNightCheck owns that —
|
||||
// which is how the outcome reaches the next morning's digest now that
|
||||
// the recap itself no longer sends anything.
|
||||
if e.Camp != nil && e.Camp.Active {
|
||||
c, _ := LoadDnDCharacter(id.UserID(e.UserID))
|
||||
var charClass DnDClass
|
||||
@@ -579,7 +604,6 @@ func (p *AdventurePlugin) deliverRecap(e *Expedition, now time.Time) error {
|
||||
if err := processNightCheck(e, nc); err != nil {
|
||||
slog.Warn("expedition: night check", "expedition", e.ID, "err", err)
|
||||
}
|
||||
night = &nc
|
||||
// §7.4: Feywild double-day fires an extra wandering check.
|
||||
if e.ZoneID == ZoneFeywildCrossing {
|
||||
if today, _ := e.RegionState["feywild_today"].(string); today == string(FeywildDistortionDouble) {
|
||||
@@ -601,12 +625,11 @@ func (p *AdventurePlugin) deliverRecap(e *Expedition, now time.Time) error {
|
||||
return err
|
||||
}
|
||||
line := pickEveningRecap(e, dayEntries)
|
||||
body := renderEveningRecap(e, line, dayEntries)
|
||||
if night != nil {
|
||||
body += "\n" + renderNightCheck(*night)
|
||||
}
|
||||
|
||||
p.fanOutExpeditionDM(e, body, nil)
|
||||
// Once-a-day cadence: the night check above has already run and already
|
||||
// written its own "night" log entry, so the morning digest reports the
|
||||
// outcome. Nothing is sent here. The recap entry below still lands so
|
||||
// the site keeps a day boundary to render against.
|
||||
if err := appendExpeditionLog(e.ID, e.CurrentDay, "recap",
|
||||
fmt.Sprintf("evening recap — %d log entries today", len(dayEntries)), line); err != nil {
|
||||
return err
|
||||
|
||||
@@ -186,13 +186,10 @@ func (p *AdventurePlugin) deliverAmbient(e *Expedition, now time.Time) error {
|
||||
}
|
||||
|
||||
footer := p.applyAmbientEffect(e, ev)
|
||||
body := renderAmbientDM(e, ev, line, footer)
|
||||
|
||||
if uid := id.UserID(e.UserID); uid != "" {
|
||||
if err := p.SendDM(uid, body); err != nil {
|
||||
slog.Warn("expedition: ambient DM", "user", uid, "err", err)
|
||||
}
|
||||
}
|
||||
// Once-a-day cadence: the effect above still lands on schedule, but the
|
||||
// DM does not. The log entry below is what the next morning's briefing
|
||||
// digest reads back, and the site renders it as it happens.
|
||||
summary := fmt.Sprintf("ambient: %s", ev.Kind)
|
||||
if footer != "" {
|
||||
summary += " — " + footer
|
||||
|
||||
@@ -304,15 +304,11 @@ func (p *AdventurePlugin) tryAutoRun(e *Expedition, now time.Time) error {
|
||||
// every other quiet path stays silent until something interactive fires.
|
||||
if body, ok := buildAutoRunDM(e.ID, r, campBlock, campDecision); ok {
|
||||
p.fanOutExpeditionDM(e, body, nil)
|
||||
// N1/A6 — the end-of-day digest is the primary mid-day event anchor.
|
||||
// The anchor is a per-player roll against a per-player daily slot, so
|
||||
// each member rolls their own; a party does not share one event.
|
||||
if campDecision.Night {
|
||||
for _, member := range expeditionAudience(e) {
|
||||
p.maybeFireAnchoredEvent(member, advEventChanceDigest)
|
||||
}
|
||||
}
|
||||
}
|
||||
// N1/A6's digest event anchor has moved to the 06:00 briefing along with
|
||||
// the digest itself — see deliverBriefing. The anchor's whole premise is
|
||||
// firing at a moment the player is demonstrably reading a DM, and the
|
||||
// night camp no longer sends one.
|
||||
|
||||
// Emergence seam: a run-complete reached by the background ticker is
|
||||
// still a live emergence — roll pet arrival. See maybeRollPetArrivalOnEmerge.
|
||||
@@ -336,8 +332,10 @@ func (p *AdventurePlugin) tryAutoRun(e *Expedition, now time.Time) error {
|
||||
// Surface rules:
|
||||
// - stopFork / stopEnded / stopComplete → render the walk DM. These
|
||||
// are the interactive / climax beats and stay their own messages.
|
||||
// - Night camp pitched → render the EoD digest +
|
||||
// camp block. Walk stream is dropped (the digest summarizes the day).
|
||||
// - Night camp pitched → silent. The once-a-day cadence
|
||||
// (2026-07-26) retired the EoD digest DM: the camp writes its own
|
||||
// `rest` log entry and the day it summarised is already in the log,
|
||||
// so the 06:00 briefing reads the whole thing back the next morning.
|
||||
// - Boss-safety camp pitched → short hold notice + camp
|
||||
// block; walk stream dropped (compact bail was deliberate).
|
||||
// - Anything else → silent.
|
||||
@@ -354,24 +352,9 @@ func buildAutoRunDM(expID string, r autopilotWalkResult, camp string, dec autoCa
|
||||
return "", false
|
||||
}
|
||||
if dec.Night {
|
||||
// EoD digest. The camp pitch already bumped current_day in
|
||||
// nightRolloverBurn, so the day-that-just-ended is CurrentDay-1.
|
||||
// digest is the day rollup, then the camp block lays out the rest.
|
||||
fresh, ferr := getExpedition(expID)
|
||||
prevDay := 0
|
||||
if ferr == nil && fresh != nil {
|
||||
prevDay = fresh.CurrentDay - 1
|
||||
}
|
||||
digest := ""
|
||||
if prevDay > 0 {
|
||||
digest = renderEndOfDayDigest(expID, prevDay)
|
||||
}
|
||||
if digest == "" {
|
||||
// No structured day yet — fall back to a thin header so the
|
||||
// camp block isn't dropped on the player without context.
|
||||
digest = "🌙 *The day winds down.*\n\n"
|
||||
}
|
||||
return digest + camp, true
|
||||
// Silent: the morning briefing is the one daily message now, and it
|
||||
// renders this same day out of the expedition log.
|
||||
return "", false
|
||||
}
|
||||
if dec.Reason == "boss-safety hold — resting before re-engaging" {
|
||||
return "⏸ *Holding before the boss — pitching a rest camp.*\n" + camp, true
|
||||
|
||||
@@ -0,0 +1,248 @@
|
||||
package plugin
|
||||
|
||||
// Once-a-day cadence (2026-07-26).
|
||||
//
|
||||
// Adventure's web feed is now the place to watch a run move minute to minute,
|
||||
// so the bot no longer narrates every beat into Matrix. The three per-player
|
||||
// DM sources that fired on a clock — the 06:00 briefing, the 21:00 recap, and
|
||||
// the 6-hourly ambient event — collapse into a single morning message.
|
||||
//
|
||||
// The rule that keeps this honest: only the *messaging* goes quiet. Every
|
||||
// mechanical effect still fires on exactly the schedule it always did. The
|
||||
// ambient ticker still applies its ±SU nudges, the recap still runs the night
|
||||
// wandering check and its threat bump, the briefing still burns supply and
|
||||
// rolls the day. What changes is that ambient and recap now write their
|
||||
// outcome to the expedition log and stop there; the next morning's briefing
|
||||
// reads that log back and reports it.
|
||||
//
|
||||
// Interrupt-driven DMs are deliberately untouched. A fork needs a human, a
|
||||
// death and a run-completion are terminal, and a mischief hit or a rival
|
||||
// challenge is somebody else acting on you. Those still arrive when they
|
||||
// happen; batching them to the next morning would either strand a decision
|
||||
// behind an 8h auto-pick or report a finished story.
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"os"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"gogobee/internal/peteclient"
|
||||
|
||||
"maunium.net/go/mautrix/id"
|
||||
)
|
||||
|
||||
const (
|
||||
// defaultPeteSiteURL — Pete's public site. Distinct from PETE_INGEST_URL,
|
||||
// which is the Headscale-only ingest endpoint and is not reachable by a
|
||||
// player clicking a link in a DM.
|
||||
defaultPeteSiteURL = "https://news.parodia.dev"
|
||||
|
||||
// digestMaxLines — how many prior-day log lines the morning digest
|
||||
// carries before it defers to the site. The cap is the whole point of
|
||||
// the change: the digest is a teaser for the feed, not a transcript.
|
||||
digestMaxLines = 8
|
||||
|
||||
// digestScanLimit — how far back the digest reads before it gives up on
|
||||
// finding the window's oldest entry. A day of autopilot ticks, ambient
|
||||
// beats and room events runs to a few dozen rows; this is slack, not a
|
||||
// budget.
|
||||
digestScanLimit = 500
|
||||
)
|
||||
|
||||
// peteSiteURL returns the public base URL for Pete's site, without a
|
||||
// trailing slash. Overridable so a dev instance can point its links at
|
||||
// a local Pete instead of prod.
|
||||
func peteSiteURL() string {
|
||||
if v := strings.TrimRight(os.Getenv("PETE_PUBLIC_URL"), "/"); v != "" {
|
||||
return v
|
||||
}
|
||||
return defaultPeteSiteURL
|
||||
}
|
||||
|
||||
// adventureFeedURL is the general Adventure feed: everyone's activity.
|
||||
func adventureFeedURL() string {
|
||||
return peteSiteURL() + "/adventure"
|
||||
}
|
||||
|
||||
// adventureWhoURL is the reader's own adventurer page. Keyed by the same
|
||||
// salted, one-way roster token the board already publishes (see
|
||||
// pete_roster.go), so a DM link and a board link resolve to one page and
|
||||
// neither one leaks a Matrix handle.
|
||||
func adventureWhoURL(uid id.UserID) string {
|
||||
if uid == "" {
|
||||
return adventureFeedURL()
|
||||
}
|
||||
return peteSiteURL() + "/adventure/who/" + eventToken(uid, "roster")
|
||||
}
|
||||
|
||||
// digestSiteFooter appends the reader's own site link. Per-reader rather
|
||||
// than per-expedition: a party shares a briefing body but each member's
|
||||
// link goes to their own sheet.
|
||||
//
|
||||
// Gated on the Pete seam: peteclient.Enabled() is what starts the roster
|
||||
// ticker, and the roster push is what creates the page this link points at.
|
||||
// With the seam off (a dev instance, or a deploy without an ingest token)
|
||||
// every daily DM would otherwise carry a guaranteed 404.
|
||||
func digestSiteFooter(uid id.UserID, body string) string {
|
||||
if !peteclient.Enabled() {
|
||||
return body
|
||||
}
|
||||
return body + "\n\n🔗 _Watch it live: " + adventureWhoURL(uid) + "_"
|
||||
}
|
||||
|
||||
// briefingPerReader is the per-reader decorator for the one daily message:
|
||||
// the reader's own pet event on the front, the reader's own site link on the
|
||||
// back. Both are per-member, so a party's shared briefing body still reaches
|
||||
// each player personalised at both ends.
|
||||
func (p *AdventurePlugin) briefingPerReader(uid id.UserID, body string) string {
|
||||
return digestSiteFooter(uid, p.briefingPetPrefix(uid, body))
|
||||
}
|
||||
|
||||
// fireDigestEventAnchor rolls N1/A6's digest-anchored mid-day event for each
|
||||
// member. It used to hang off the autopilot's night-camp digest DM; that DM is
|
||||
// gone, so it moved here — the briefing is now the message the player is
|
||||
// demonstrably reading, which is the whole premise of an anchored roll.
|
||||
//
|
||||
// Still a per-player roll against a per-player daily slot: a party does not
|
||||
// share one event.
|
||||
func (p *AdventurePlugin) fireDigestEventAnchor(e *Expedition) {
|
||||
for _, member := range expeditionAudience(e) {
|
||||
p.maybeFireAnchoredEvent(member, advEventChanceDigest)
|
||||
}
|
||||
}
|
||||
|
||||
// appendOvernightDigest folds everything that happened since the previous
|
||||
// briefing into a briefing body. A log read failure is non-fatal: the briefing
|
||||
// is the player's only daily message now, so a missing digest block must never
|
||||
// cost them the whole DM.
|
||||
//
|
||||
// The window is a timestamp, not a day number, because the two disagree on
|
||||
// every event-anchored expedition: the autopilot's night camp rolls
|
||||
// current_day at camp time, so by 06:00 the day that just ended is already
|
||||
// current_day-1 — and on a night the autopilot never camped, it isn't. A
|
||||
// since-last-briefing window reports each entry exactly once either way.
|
||||
func appendOvernightDigest(body, expID string, since time.Time) string {
|
||||
entries, err := logEntriesSince(expID, since)
|
||||
if err != nil {
|
||||
slog.Warn("expedition: digest entries", "expedition", expID, "err", err)
|
||||
return body
|
||||
}
|
||||
digest := renderOvernightDigest(entries)
|
||||
if digest == "" {
|
||||
return body
|
||||
}
|
||||
return body + "\n" + digest
|
||||
}
|
||||
|
||||
// logEntriesSince returns an expedition's log entries stamped at or after
|
||||
// `since`, oldest first.
|
||||
//
|
||||
// The cutoff is applied in Go rather than in the WHERE clause on purpose:
|
||||
// dnd_expedition_log.timestamp is a DATETIME column filled by SQLite's own
|
||||
// CURRENT_TIMESTAMP, and comparing it against a bound parameter goes through
|
||||
// numeric affinity and does not reliably answer the question. Scanning the
|
||||
// column into a time.Time does.
|
||||
func logEntriesSince(expID string, since time.Time) ([]ExpeditionEntry, error) {
|
||||
recent, err := recentExpeditionLog(expID, digestScanLimit)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
// CURRENT_TIMESTAMP has one-second resolution while the cutoff we are
|
||||
// handed (a briefing stamp, or the start date on day 1) carries
|
||||
// sub-second precision. Floor it, or an entry written in the same second
|
||||
// as the previous briefing falls out of both windows and is never
|
||||
// reported. Re-reporting inside that one second is the safe direction.
|
||||
since = since.Truncate(time.Second)
|
||||
out := make([]ExpeditionEntry, 0, len(recent))
|
||||
for i := len(recent) - 1; i >= 0; i-- { // recent is newest-first
|
||||
if recent[i].Timestamp.Before(since) {
|
||||
continue
|
||||
}
|
||||
out = append(out, recent[i])
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// digestSkipTypes — log entry types the morning digest never echoes.
|
||||
// `briefing` and `recap` are the frame itself, and the free-narration
|
||||
// types are the per-room prose the site renders in full.
|
||||
var digestSkipTypes = map[string]bool{
|
||||
"briefing": true,
|
||||
"recap": true,
|
||||
"narrative": true,
|
||||
"transit": true,
|
||||
"action": true,
|
||||
"journal": true,
|
||||
}
|
||||
|
||||
// renderOvernightDigest condenses one expedition-day of log entries into the
|
||||
// "here is what you missed" block that opens the morning briefing. Walks
|
||||
// collapse to a count; everything notable keeps its own summary line, capped
|
||||
// at digestMaxLines with an explicit overflow note so a truncated digest
|
||||
// never reads as a complete one.
|
||||
//
|
||||
// Returns "" when there is nothing worth reporting — a day with only walks
|
||||
// and narration gets no block at all rather than an empty header.
|
||||
func renderOvernightDigest(entries []ExpeditionEntry) string {
|
||||
var rooms int
|
||||
var lines []string
|
||||
for _, en := range entries {
|
||||
if en.Type == "walk" {
|
||||
rooms += walkEntryRooms(en.Summary)
|
||||
continue
|
||||
}
|
||||
if digestSkipTypes[en.Type] {
|
||||
continue
|
||||
}
|
||||
s := strings.TrimSpace(en.Summary)
|
||||
if s == "" {
|
||||
continue
|
||||
}
|
||||
lines = append(lines, s)
|
||||
}
|
||||
if rooms == 0 && len(lines) == 0 {
|
||||
return ""
|
||||
}
|
||||
|
||||
var b strings.Builder
|
||||
b.WriteString("📜 **Since yesterday**\n")
|
||||
if rooms > 0 {
|
||||
b.WriteString(fmt.Sprintf("• walked %s\n", pluralRooms(rooms)))
|
||||
}
|
||||
shown := lines
|
||||
overflow := 0
|
||||
if len(shown) > digestMaxLines {
|
||||
overflow = len(shown) - digestMaxLines
|
||||
shown = shown[:digestMaxLines]
|
||||
}
|
||||
for _, l := range shown {
|
||||
b.WriteString("• " + l + "\n")
|
||||
}
|
||||
if overflow > 0 {
|
||||
b.WriteString(fmt.Sprintf("• _...and %d more, on the site._\n", overflow))
|
||||
}
|
||||
return b.String()
|
||||
}
|
||||
|
||||
// walkEntryRooms reads the room count back out of an auto-walk log summary
|
||||
// ("auto-walk: 3 room(s)"). One `walk` entry is one background tick, and a
|
||||
// tick covers as many rooms as the autopilot got through — counting entries
|
||||
// would report a 12-room day as a 3-room one. Unparseable summaries count as
|
||||
// a single room rather than vanishing.
|
||||
func walkEntryRooms(summary string) int {
|
||||
var n int
|
||||
if _, err := fmt.Sscanf(strings.TrimSpace(summary), "auto-walk: %d room", &n); err == nil && n > 0 {
|
||||
return n
|
||||
}
|
||||
return 1
|
||||
}
|
||||
|
||||
// pluralRooms renders a room count with the right noun.
|
||||
func pluralRooms(n int) string {
|
||||
if n == 1 {
|
||||
return "1 room"
|
||||
}
|
||||
return fmt.Sprintf("%d rooms", n)
|
||||
}
|
||||
@@ -0,0 +1,253 @@
|
||||
package plugin
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"maunium.net/go/mautrix/id"
|
||||
)
|
||||
|
||||
// Coverage for the once-a-day cadence (2026-07-26). The behavioural tests
|
||||
// below are the point of the file: before this change nothing asserted that
|
||||
// the recap and ambient paths sent a DM, so nothing would have caught them
|
||||
// silently continuing to send — or, worse, silently dropping the mechanical
|
||||
// effect along with the message.
|
||||
|
||||
func TestRenderOvernightDigest_CollapsesWalksAndSkipsFrameTypes(t *testing.T) {
|
||||
entries := []ExpeditionEntry{
|
||||
// One `walk` entry is one autopilot tick, and a tick can cover
|
||||
// several rooms — the digest reports rooms, not ticks.
|
||||
{Type: "walk", Summary: "auto-walk: 2 room(s)"},
|
||||
{Type: "walk", Summary: "auto-walk: 3 room(s)"},
|
||||
{Type: "briefing", Summary: "morning briefing — 1.0 SU consumed overnight"},
|
||||
{Type: "narrative", Summary: "the corridor bends left"},
|
||||
{Type: "ambient", Summary: "ambient: pack_rat — Supplies -0.5"},
|
||||
{Type: "night", Summary: "Signs of passage near camp; no encounter."},
|
||||
{Type: "recap", Summary: "evening recap — 6 log entries today"},
|
||||
}
|
||||
got := renderOvernightDigest(entries)
|
||||
|
||||
if !strings.Contains(got, "walked 5 rooms") {
|
||||
t.Errorf("walks not collapsed to a count:\n%s", got)
|
||||
}
|
||||
if !strings.Contains(got, "ambient: pack_rat") {
|
||||
t.Errorf("ambient entry missing from digest:\n%s", got)
|
||||
}
|
||||
if !strings.Contains(got, "Signs of passage") {
|
||||
t.Errorf("night check missing from digest:\n%s", got)
|
||||
}
|
||||
for _, unwanted := range []string{"morning briefing", "evening recap", "corridor bends"} {
|
||||
if strings.Contains(got, unwanted) {
|
||||
t.Errorf("digest echoed frame/narration entry %q:\n%s", unwanted, got)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestRenderOvernightDigest_CapsAndReportsOverflow(t *testing.T) {
|
||||
var entries []ExpeditionEntry
|
||||
for i := 0; i < digestMaxLines+3; i++ {
|
||||
entries = append(entries, ExpeditionEntry{
|
||||
Type: "ambient",
|
||||
Summary: fmt.Sprintf("ambient event %d", i),
|
||||
})
|
||||
}
|
||||
got := renderOvernightDigest(entries)
|
||||
|
||||
if n := strings.Count(got, "ambient event"); n != digestMaxLines {
|
||||
t.Errorf("digest carried %d lines, want the cap of %d:\n%s", n, digestMaxLines, got)
|
||||
}
|
||||
// A truncated digest must say so — otherwise it reads as the whole day.
|
||||
if !strings.Contains(got, "and 3 more") {
|
||||
t.Errorf("overflow not reported:\n%s", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRenderOvernightDigest_EmptyWhenNothingNotable(t *testing.T) {
|
||||
entries := []ExpeditionEntry{
|
||||
{Type: "briefing", Summary: "morning briefing"},
|
||||
{Type: "narrative", Summary: "dust everywhere"},
|
||||
}
|
||||
if got := renderOvernightDigest(entries); got != "" {
|
||||
t.Errorf("want no digest block for an unremarkable day, got:\n%s", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAdventureWhoURL_UsesRosterTokenNotHandle(t *testing.T) {
|
||||
// The roster token is salted from a DB-persisted secret.
|
||||
setupZoneRunTestDB(t)
|
||||
uid := id.UserID("@digest-url:example")
|
||||
got := adventureWhoURL(uid)
|
||||
|
||||
want := "/adventure/who/" + eventToken(uid, "roster")
|
||||
if !strings.HasSuffix(got, want) {
|
||||
t.Errorf("who URL = %q, want suffix %q", got, want)
|
||||
}
|
||||
if strings.Contains(got, "digest-url") {
|
||||
t.Errorf("who URL leaked the Matrix handle: %q", got)
|
||||
}
|
||||
}
|
||||
|
||||
func TestAdventureWhoURL_FallsBackToFeedWithoutUser(t *testing.T) {
|
||||
setupZoneRunTestDB(t)
|
||||
if got := adventureWhoURL(""); got != adventureFeedURL() {
|
||||
t.Errorf("empty user should fall back to the feed, got %q", got)
|
||||
}
|
||||
}
|
||||
|
||||
// TestBuildAutoRunDM_NightCampIsSilent — the autopilot's end-of-day digest was
|
||||
// the last recurring second DM of the day. It goes quiet; the camp writes its
|
||||
// own `rest` log entry, so the morning briefing still reports it.
|
||||
func TestBuildAutoRunDM_NightCampIsSilent(t *testing.T) {
|
||||
r := autopilotWalkResult{rooms: 4, reason: stopOK, stream: []string{"…walked…"}}
|
||||
camp := "\n\n⛺ **Autopilot camp** — night"
|
||||
body, ok := buildAutoRunDM("expid", r, camp, autoCampDecision{
|
||||
Kind: CampTypeStandard, Night: true,
|
||||
})
|
||||
if ok || body != "" {
|
||||
t.Errorf("night camp should be silent, got ok=%v body=%q", ok, body)
|
||||
}
|
||||
}
|
||||
|
||||
// The two interactive surfaces the night-camp cut must not touch.
|
||||
func TestBuildAutoRunDM_KeepSetStillSurfaces(t *testing.T) {
|
||||
fork := autopilotWalkResult{rooms: 1, reason: stopFork, finalMsg: "pick a path"}
|
||||
if body, ok := buildAutoRunDM("expid", fork, "", autoCampDecision{}); !ok ||
|
||||
!strings.Contains(body, "pick a path") {
|
||||
t.Errorf("fork must still surface, got ok=%v body=%q", ok, body)
|
||||
}
|
||||
|
||||
hold := autopilotWalkResult{rooms: 2, reason: stopBossSafety}
|
||||
holdCamp := "\n\n⛺ **Rest camp**"
|
||||
if body, ok := buildAutoRunDM("expid", hold, holdCamp, autoCampDecision{
|
||||
Reason: "boss-safety hold — resting before re-engaging",
|
||||
}); !ok || !strings.Contains(body, "Holding before the boss") {
|
||||
t.Errorf("boss-safety hold must still surface, got ok=%v body=%q", ok, body)
|
||||
}
|
||||
}
|
||||
|
||||
// TestDeliverAmbient_SilentButStillLogs — the ambient event still fires and
|
||||
// still records itself; it just stops DMing.
|
||||
func TestDeliverAmbient_SilentButStillLogs(t *testing.T) {
|
||||
setupZoneRunTestDB(t)
|
||||
uid := id.UserID("@digest-ambient:example")
|
||||
defer cleanupExpeditions(uid)
|
||||
|
||||
p := &AdventurePlugin{}
|
||||
sink := installSink(p)
|
||||
|
||||
exp, err := startExpedition(uid, ZoneGoblinWarrens, "",
|
||||
ExpeditionSupplies{Current: 10, Max: 10, DailyBurn: 1, HarshMod: 1})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := p.deliverAmbient(exp, exp.StartDate.Add(time.Hour)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
if dms := sink.dmsTo(uid); len(dms) != 0 {
|
||||
t.Errorf("ambient sent %d DM(s), want 0:\n%s", len(dms), strings.Join(dms, "\n---\n"))
|
||||
}
|
||||
entries, _ := recentExpeditionLog(exp.ID, 10)
|
||||
found := false
|
||||
for _, e := range entries {
|
||||
if e.Type == "ambient" {
|
||||
found = true
|
||||
}
|
||||
}
|
||||
if !found {
|
||||
t.Error("ambient event fired without writing a log entry — the digest would lose it")
|
||||
}
|
||||
}
|
||||
|
||||
// TestDeliverRecap_SilentButStillRunsNightCheck — the recap's mechanical half
|
||||
// (wandering check, threat bump) must survive the message going away.
|
||||
func TestDeliverRecap_SilentButStillRunsNightCheck(t *testing.T) {
|
||||
setupZoneRunTestDB(t)
|
||||
uid := id.UserID("@digest-recap:example")
|
||||
defer cleanupExpeditions(uid)
|
||||
|
||||
p := &AdventurePlugin{}
|
||||
sink := installSink(p)
|
||||
|
||||
exp, err := startExpedition(uid, ZoneGoblinWarrens, "",
|
||||
ExpeditionSupplies{Current: 10, Max: 10, DailyBurn: 1, HarshMod: 1})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
exp.Camp = &CampState{Active: true, Type: CampTypeStandard, EstablishedAt: exp.StartDate}
|
||||
if err := updateCamp(exp.ID, exp.Camp); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
recapAt := exp.StartDate.Add(12 * time.Hour)
|
||||
if err := p.deliverRecap(exp, recapAt); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
if dms := sink.dmsTo(uid); len(dms) != 0 {
|
||||
t.Errorf("recap sent %d DM(s), want 0:\n%s", len(dms), strings.Join(dms, "\n---\n"))
|
||||
}
|
||||
entries, _ := recentExpeditionLog(exp.ID, 10)
|
||||
sawNight, sawRecap := false, false
|
||||
for _, e := range entries {
|
||||
switch e.Type {
|
||||
case "night":
|
||||
sawNight = true
|
||||
case "recap":
|
||||
sawRecap = true
|
||||
}
|
||||
}
|
||||
if !sawNight {
|
||||
t.Error("night wandering check did not run — the recap dropped its mechanics, not just its DM")
|
||||
}
|
||||
if !sawRecap {
|
||||
t.Error("recap log entry missing — the site loses its day boundary")
|
||||
}
|
||||
}
|
||||
|
||||
// TestDeliverBriefing_IsTheOneDailyMessage — the payoff: one DM, carrying the
|
||||
// prior day's silent activity and the reader's own site link.
|
||||
func TestDeliverBriefing_CarriesDigestAndSiteLink(t *testing.T) {
|
||||
setupZoneRunTestDB(t)
|
||||
uid := id.UserID("@digest-briefing:example")
|
||||
defer cleanupExpeditions(uid)
|
||||
|
||||
enablePeteSeam(t)
|
||||
p := &AdventurePlugin{}
|
||||
sink := installSink(p)
|
||||
|
||||
exp, err := startExpedition(uid, ZoneGoblinWarrens, "",
|
||||
ExpeditionSupplies{Current: 10, Max: 10, DailyBurn: 1, HarshMod: 1})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
// Stand in for the day that just went by silently. The day number is
|
||||
// deliberately not CurrentDay: on an event-anchored run the night camp
|
||||
// rolls current_day when it pitches, so entries either side of the
|
||||
// rollover carry different day numbers. The digest windows on time.
|
||||
if err := appendExpeditionLog(exp.ID, exp.CurrentDay+1, "ambient",
|
||||
"ambient: pack_rat — Supplies -0.5", "Something nibbled the stores."); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
if err := p.deliverBriefing(exp, exp.StartDate.Add(20*time.Hour)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
dms := sink.dmsTo(uid)
|
||||
if len(dms) != 1 {
|
||||
t.Fatalf("briefing sent %d DM(s), want exactly 1:\n%s", len(dms), strings.Join(dms, "\n---\n"))
|
||||
}
|
||||
body := dms[0]
|
||||
if !strings.Contains(body, "Since yesterday") {
|
||||
t.Errorf("briefing missing the overnight digest block:\n%s", body)
|
||||
}
|
||||
if !strings.Contains(body, "ambient: pack_rat") {
|
||||
t.Errorf("digest did not carry the silent ambient event:\n%s", body)
|
||||
}
|
||||
if !strings.Contains(body, adventureWhoURL(uid)) {
|
||||
t.Errorf("briefing missing the reader's own site link:\n%s", body)
|
||||
}
|
||||
}
|
||||
@@ -856,9 +856,7 @@ func (p *HangmanPlugin) handleSubmit(ctx MessageContext, phrase string) error {
|
||||
}
|
||||
|
||||
// LLM screening
|
||||
ollamaHost := os.Getenv("OLLAMA_HOST")
|
||||
ollamaModel := os.Getenv("OLLAMA_MODEL")
|
||||
if ollamaHost == "" || ollamaModel == "" {
|
||||
if !llmConfigured() {
|
||||
// No LLM available — add directly
|
||||
if err := p.addPhrase(phrase); err != nil {
|
||||
if err.Error() == "duplicate phrase" {
|
||||
@@ -879,7 +877,7 @@ or
|
||||
|
||||
Phrase: %s`, phrase)
|
||||
|
||||
result, err := callOllama(ollamaHost, ollamaModel, prompt)
|
||||
result, err := callLLM(prompt)
|
||||
if err != nil {
|
||||
slog.Error("hangman: LLM screening failed", "err", err)
|
||||
// Fail open — add it
|
||||
|
||||
@@ -1,24 +1,24 @@
|
||||
package plugin
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"os"
|
||||
"regexp"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"gogobee/internal/db"
|
||||
"gogobee/internal/llm"
|
||||
|
||||
"github.com/chehsunliu/poker"
|
||||
"maunium.net/go/mautrix/id"
|
||||
)
|
||||
|
||||
var holdemTipsClient = &http.Client{Timeout: 60 * time.Second}
|
||||
// holdemTipTimeout preserves the 60s budget the tip rewriter's own http.Client
|
||||
// enforced. Tips are delivered as private messages during a hand, so this sits
|
||||
// between the passive 30s paths and the interactive 120s default.
|
||||
const holdemTipTimeout = 60 * time.Second
|
||||
|
||||
// loadTipsPref loads a user's tip preference from the database.
|
||||
func loadTipsPref(userID id.UserID) bool {
|
||||
@@ -491,11 +491,8 @@ func cardSuitIndex(c poker.Card) int {
|
||||
func generateTip(ctx holdemTipContext) string {
|
||||
base := generateRulesTip(ctx)
|
||||
|
||||
host := os.Getenv("OLLAMA_HOST")
|
||||
model := os.Getenv("OLLAMA_MODEL")
|
||||
|
||||
if host != "" && model != "" {
|
||||
rewritten, err := rewriteTipWithLLM(host, model, ctx, base)
|
||||
if llmConfigured() {
|
||||
rewritten, err := rewriteTipWithLLM(ctx, base)
|
||||
if err != nil {
|
||||
slog.Warn("holdem: LLM tip rewrite failed, using rules tip", "err", err)
|
||||
} else if rewritten != "" {
|
||||
@@ -626,40 +623,19 @@ func buildTipUserPrompt(ctx holdemTipContext) string {
|
||||
// variety. The rules tip is the source of truth — if the rewrite diverges
|
||||
// (empty, action vocabulary changed, etc.) we reject it and the caller falls
|
||||
// back to the original.
|
||||
func rewriteTipWithLLM(host, model string, ctx holdemTipContext, base string) (string, error) {
|
||||
func rewriteTipWithLLM(ctx holdemTipContext, base string) (string, error) {
|
||||
userMsg := buildTipUserPrompt(ctx) + "\nTIP:\n" + base + "\n"
|
||||
req := ollamaChatRequest{
|
||||
Model: model,
|
||||
Messages: []chatMessage{
|
||||
{Role: "system", Content: buildTipSystemPrompt()},
|
||||
{Role: "user", Content: userMsg},
|
||||
},
|
||||
Stream: false,
|
||||
}
|
||||
|
||||
body, err := json.Marshal(req)
|
||||
raw, err := llmGenerate(context.Background(), llm.Request{
|
||||
System: buildTipSystemPrompt(),
|
||||
Prompt: userMsg,
|
||||
Timeout: holdemTipTimeout,
|
||||
})
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("marshal: %w", err)
|
||||
return "", err
|
||||
}
|
||||
|
||||
url := strings.TrimRight(host, "/") + "/api/chat"
|
||||
resp, err := holdemTipsClient.Post(url, "application/json", bytes.NewReader(body))
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("request: %w", err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
respBody, _ := io.ReadAll(resp.Body)
|
||||
return "", fmt.Errorf("status %d: %s", resp.StatusCode, string(respBody))
|
||||
}
|
||||
|
||||
var ollamaResp ollamaChatResponse
|
||||
if err := json.NewDecoder(resp.Body).Decode(&ollamaResp); err != nil {
|
||||
return "", fmt.Errorf("decode: %w", err)
|
||||
}
|
||||
|
||||
tip := extractTipFromResponse(ollamaResp.Message.Content)
|
||||
tip := extractTipFromResponse(raw)
|
||||
if tip == "" {
|
||||
return "", fmt.Errorf("empty response")
|
||||
}
|
||||
|
||||
+14
-61
@@ -1,18 +1,15 @@
|
||||
package plugin
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"os"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"gogobee/internal/db"
|
||||
"gogobee/internal/llm"
|
||||
|
||||
"maunium.net/go/mautrix"
|
||||
"maunium.net/go/mautrix/id"
|
||||
@@ -47,9 +44,7 @@ func (p *HowAmIPlugin) OnMessage(ctx MessageContext) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
ollamaHost := os.Getenv("OLLAMA_HOST")
|
||||
ollamaModel := os.Getenv("OLLAMA_MODEL")
|
||||
if ollamaHost == "" || ollamaModel == "" {
|
||||
if !llmConfigured() {
|
||||
return p.SendReply(ctx.RoomID, ctx.EventID, "LLM is not configured.")
|
||||
}
|
||||
|
||||
@@ -84,9 +79,9 @@ Write the roast now. Do not include any preamble or explanation, just the roast
|
||||
botName, string(target), profile,
|
||||
)
|
||||
|
||||
response, err := callOllama(ollamaHost, ollamaModel, prompt)
|
||||
response, err := callLLM(prompt)
|
||||
if err != nil {
|
||||
slog.Error("howami: ollama call", "err", err)
|
||||
slog.Error("howami: llm call", "err", err)
|
||||
p.SendReply(ctx.RoomID, ctx.EventID, "Couldn't generate the profile. Thanks, Ollama.")
|
||||
return
|
||||
}
|
||||
@@ -181,55 +176,13 @@ func (p *HowAmIPlugin) gatherProfile(userID id.UserID) string {
|
||||
return sb.String()
|
||||
}
|
||||
|
||||
// callOllama sends a prompt to the Ollama generate endpoint and returns the response.
|
||||
func callOllama(host, model, prompt string) (string, error) {
|
||||
apiURL := strings.TrimRight(host, "/") + "/api/generate"
|
||||
|
||||
payload := map[string]interface{}{
|
||||
"model": model,
|
||||
"prompt": prompt,
|
||||
"stream": false,
|
||||
"think": false,
|
||||
"options": map[string]interface{}{
|
||||
"num_ctx": 8192,
|
||||
},
|
||||
}
|
||||
|
||||
data, err := json.Marshal(payload)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("marshal payload: %w", err)
|
||||
}
|
||||
|
||||
client := &http.Client{Timeout: 120 * time.Second}
|
||||
resp, err := client.Post(apiURL, "application/json", bytes.NewReader(data))
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("ollama request: %w", err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
body, err := io.ReadAll(resp.Body)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("read response: %w", err)
|
||||
}
|
||||
|
||||
if resp.StatusCode != 200 {
|
||||
return "", fmt.Errorf("ollama HTTP %d: %s", resp.StatusCode, string(body))
|
||||
}
|
||||
|
||||
var result struct {
|
||||
Response string `json:"response"`
|
||||
}
|
||||
if err := json.Unmarshal(body, &result); err != nil {
|
||||
return "", fmt.Errorf("parse response: %w", err)
|
||||
}
|
||||
|
||||
response := result.Response
|
||||
// Strip <think>...</think> blocks (Qwen 3.5 reasoning)
|
||||
if i := strings.Index(response, "<think>"); i != -1 {
|
||||
if j := strings.Index(response, "</think>"); j != -1 {
|
||||
response = response[:i] + response[j+len("</think>"):]
|
||||
}
|
||||
}
|
||||
|
||||
return strings.TrimSpace(response), nil
|
||||
// callLLM sends a prompt to whichever backend is configured and returns the
|
||||
// completion. Reasoning blocks are stripped by the client. The 8192 context
|
||||
// hint is what this path has always asked Ollama for; vLLM ignores it and uses
|
||||
// the window fixed at server launch.
|
||||
func callLLM(prompt string) (string, error) {
|
||||
return llmGenerate(context.Background(), llm.Request{
|
||||
Prompt: prompt,
|
||||
NumCtx: 8192,
|
||||
})
|
||||
}
|
||||
|
||||
@@ -0,0 +1,45 @@
|
||||
package plugin
|
||||
|
||||
import (
|
||||
"context"
|
||||
"sync"
|
||||
|
||||
"gogobee/internal/llm"
|
||||
)
|
||||
|
||||
// The inference backend is process-wide: one endpoint, one model, selected by
|
||||
// env at startup. Plugins share a single client rather than each rebuilding one
|
||||
// per invocation, and read config through llmConfigured/llmGenerate so that
|
||||
// swapping Ollama for vLLM is a config change rather than a code change.
|
||||
var (
|
||||
llmOnce sync.Once
|
||||
llmShared llm.Client
|
||||
llmCfg llm.Config
|
||||
)
|
||||
|
||||
func llmInit() {
|
||||
llmOnce.Do(func() {
|
||||
llmCfg = llm.ConfigFromEnv()
|
||||
llmShared = llm.New(llmCfg)
|
||||
})
|
||||
}
|
||||
|
||||
// llmConfigured reports whether an endpoint and model are set. Plugins call
|
||||
// this to stay dormant instead of erroring on every invocation — the same role
|
||||
// the old `if ollamaHost == "" || ollamaModel == ""` guards played.
|
||||
func llmConfigured() bool {
|
||||
llmInit()
|
||||
return llmCfg.Configured()
|
||||
}
|
||||
|
||||
// llmClient returns the shared backend client.
|
||||
func llmClient() llm.Client {
|
||||
llmInit()
|
||||
return llmShared
|
||||
}
|
||||
|
||||
// llmGenerate is the one-line path for the common case: a raw prompt in,
|
||||
// visible completion out, reasoning blocks already stripped.
|
||||
func llmGenerate(ctx context.Context, req llm.Request) (string, error) {
|
||||
return llmClient().Generate(ctx, req)
|
||||
}
|
||||
@@ -1,13 +1,11 @@
|
||||
package plugin
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"log/slog"
|
||||
"math/rand"
|
||||
"net/http"
|
||||
"os"
|
||||
"regexp"
|
||||
"strconv"
|
||||
@@ -18,6 +16,7 @@ import (
|
||||
|
||||
"gogobee/internal/db"
|
||||
"gogobee/internal/dreamclient"
|
||||
"gogobee/internal/llm"
|
||||
|
||||
"maunium.net/go/mautrix"
|
||||
"maunium.net/go/mautrix/id"
|
||||
@@ -54,29 +53,31 @@ type queueItem struct {
|
||||
FormattedBody string
|
||||
}
|
||||
|
||||
// LLMPassivePlugin classifies messages using Ollama and reacts accordingly.
|
||||
// classifyTimeout is the per-message budget for passive classification. It
|
||||
// preserves the 30s cap the plugin's own http.Client used to enforce, which is
|
||||
// deliberately tighter than the interactive default: classification runs on
|
||||
// sampled traffic and must never back up the queue.
|
||||
const classifyTimeout = 30 * time.Second
|
||||
|
||||
// LLMPassivePlugin classifies messages using the configured LLM backend and
|
||||
// reacts accordingly.
|
||||
type LLMPassivePlugin struct {
|
||||
Base
|
||||
xp *XPPlugin
|
||||
dict *dreamclient.Client
|
||||
ollamaHost string
|
||||
ollamaModel string
|
||||
sampleRate float64
|
||||
enabled bool
|
||||
xp *XPPlugin
|
||||
dict *dreamclient.Client
|
||||
sampleRate float64
|
||||
enabled bool
|
||||
|
||||
mu sync.Mutex
|
||||
queue []queueItem
|
||||
backoff time.Duration
|
||||
|
||||
httpClient *http.Client
|
||||
stopCh chan struct{}
|
||||
stopCh chan struct{}
|
||||
}
|
||||
|
||||
// NewLLMPassivePlugin creates a new LLM passive classification plugin.
|
||||
func NewLLMPassivePlugin(client *mautrix.Client, xp *XPPlugin, dict *dreamclient.Client) *LLMPassivePlugin {
|
||||
host := os.Getenv("OLLAMA_HOST")
|
||||
model := os.Getenv("OLLAMA_MODEL")
|
||||
enabled := host != "" && model != ""
|
||||
enabled := llmConfigured()
|
||||
|
||||
sampleRate := 0.15
|
||||
if v := os.Getenv("LLM_SAMPLE_RATE"); v != "" {
|
||||
@@ -86,16 +87,13 @@ func NewLLMPassivePlugin(client *mautrix.Client, xp *XPPlugin, dict *dreamclient
|
||||
}
|
||||
|
||||
p := &LLMPassivePlugin{
|
||||
Base: NewBase(client),
|
||||
xp: xp,
|
||||
dict: dict,
|
||||
ollamaHost: host,
|
||||
ollamaModel: model,
|
||||
sampleRate: sampleRate,
|
||||
enabled: enabled,
|
||||
backoff: 5 * time.Second,
|
||||
httpClient: &http.Client{Timeout: 30 * time.Second},
|
||||
stopCh: make(chan struct{}),
|
||||
Base: NewBase(client),
|
||||
xp: xp,
|
||||
dict: dict,
|
||||
sampleRate: sampleRate,
|
||||
enabled: enabled,
|
||||
backoff: 5 * time.Second,
|
||||
stopCh: make(chan struct{}),
|
||||
}
|
||||
|
||||
return p
|
||||
@@ -116,11 +114,11 @@ func (p *LLMPassivePlugin) Commands() []CommandDef {
|
||||
|
||||
func (p *LLMPassivePlugin) Init() error {
|
||||
if p.enabled {
|
||||
slog.Info("llm_passive: enabled", "host", p.ollamaHost, "model", p.ollamaModel, "sample_rate", p.sampleRate)
|
||||
slog.Info("llm_passive: enabled", "backend", llmClient().Backend(),
|
||||
"model", llmClient().Model(), "sample_rate", p.sampleRate)
|
||||
go p.processQueue()
|
||||
} else {
|
||||
slog.Warn("llm_passive: disabled (OLLAMA_HOST or OLLAMA_MODEL not set)",
|
||||
"host", p.ollamaHost, "model", p.ollamaModel)
|
||||
slog.Warn("llm_passive: disabled (LLM endpoint or model not set)")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -329,9 +327,9 @@ func (p *LLMPassivePlugin) classifyAndProcess(item queueItem) error {
|
||||
var todayWOTD string
|
||||
db.Get().QueryRow(`SELECT word FROM wotd_log WHERE date = ?`, today).Scan(&todayWOTD)
|
||||
|
||||
result, err := p.callOllama(item.Body+mentionHint, todayWOTD)
|
||||
result, err := p.classify(item.Body+mentionHint, todayWOTD)
|
||||
if err != nil {
|
||||
return fmt.Errorf("ollama call: %w", err)
|
||||
return fmt.Errorf("llm call: %w", err)
|
||||
}
|
||||
|
||||
// Resolve any display names in LLM targets back to MXIDs
|
||||
@@ -464,21 +462,9 @@ func (p *LLMPassivePlugin) classifyAndProcess(item queueItem) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// ollamaRequest is the request body for the Ollama API.
|
||||
type ollamaRequest struct {
|
||||
Model string `json:"model"`
|
||||
Prompt string `json:"prompt"`
|
||||
Stream bool `json:"stream"`
|
||||
Think bool `json:"think"`
|
||||
}
|
||||
|
||||
// ollamaResponse is the response from the Ollama API.
|
||||
type ollamaResponse struct {
|
||||
Response string `json:"response"`
|
||||
}
|
||||
|
||||
// callOllama sends a classification prompt to Ollama and parses the JSON result.
|
||||
func (p *LLMPassivePlugin) callOllama(messageText, wotd string) (*classificationResult, error) {
|
||||
// classify sends a classification prompt to the configured backend and parses
|
||||
// the JSON result.
|
||||
func (p *LLMPassivePlugin) classify(messageText, wotd string) (*classificationResult, error) {
|
||||
wotdInstruction := `"wotd_used": false`
|
||||
if wotd != "" {
|
||||
wotdInstruction = fmt.Sprintf(`"wotd_used": true | false (whether the message uses the word "%s" correctly and meaningfully — not just mentioning or quoting it)`, wotd)
|
||||
@@ -500,36 +486,18 @@ JSON schema:
|
||||
|
||||
Message: %s`, wotdInstruction, messageText)
|
||||
|
||||
reqBody := ollamaRequest{
|
||||
Model: p.ollamaModel,
|
||||
Prompt: prompt,
|
||||
Stream: false,
|
||||
}
|
||||
|
||||
body, err := json.Marshal(reqBody)
|
||||
slog.Debug("llm_passive: calling backend", "backend", llmClient().Backend(), "model", llmClient().Model())
|
||||
// Classification rides the passive path on every sampled message, so it keeps
|
||||
// the tighter budget it always had rather than the interactive default.
|
||||
raw, err := llmGenerate(context.Background(), llm.Request{
|
||||
Prompt: prompt,
|
||||
Timeout: classifyTimeout,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("marshal request: %w", err)
|
||||
return nil, err
|
||||
}
|
||||
|
||||
url := strings.TrimRight(p.ollamaHost, "/") + "/api/generate"
|
||||
slog.Debug("llm_passive: calling ollama", "url", url, "model", p.ollamaModel)
|
||||
resp, err := p.httpClient.Post(url, "application/json", bytes.NewReader(body))
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("ollama request: %w", err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
respBody, _ := io.ReadAll(resp.Body)
|
||||
return nil, fmt.Errorf("ollama status %d: %s", resp.StatusCode, string(respBody))
|
||||
}
|
||||
|
||||
var ollamaResp ollamaResponse
|
||||
if err := json.NewDecoder(resp.Body).Decode(&ollamaResp); err != nil {
|
||||
return nil, fmt.Errorf("decode ollama response: %w", err)
|
||||
}
|
||||
|
||||
result, err := parseClassification(ollamaResp.Response)
|
||||
result, err := parseClassification(raw)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("parse classification: %w", err)
|
||||
}
|
||||
|
||||
+18
-61
@@ -1,10 +1,9 @@
|
||||
package plugin
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"log/slog"
|
||||
"math"
|
||||
"net/http"
|
||||
@@ -15,6 +14,7 @@ import (
|
||||
"time"
|
||||
|
||||
"gogobee/internal/db"
|
||||
"gogobee/internal/llm"
|
||||
|
||||
"maunium.net/go/mautrix"
|
||||
"maunium.net/go/mautrix/id"
|
||||
@@ -311,9 +311,7 @@ Do not use em dashes. Do not use exclamation marks. Do not offer financial advic
|
||||
If markets are closed or data is stale, note it briefly and move on.`
|
||||
|
||||
func (p *MarketPlugin) generateDailySummary(date string) string {
|
||||
host := os.Getenv("OLLAMA_HOST")
|
||||
model := os.Getenv("OLLAMA_MODEL")
|
||||
if host == "" || model == "" {
|
||||
if !llmConfigured() {
|
||||
return ""
|
||||
}
|
||||
|
||||
@@ -349,7 +347,7 @@ func (p *MarketPlugin) generateDailySummary(date string) string {
|
||||
}
|
||||
prompt.WriteString("\nWrite a 2-3 sentence summary.")
|
||||
|
||||
result, err := p.callOllamaChat(host, model, marketSystemPrompt, prompt.String())
|
||||
result, err := p.chatLLM(marketSystemPrompt, prompt.String())
|
||||
if err != nil {
|
||||
slog.Error("market: ollama summary failed", "err", err)
|
||||
return ""
|
||||
@@ -358,9 +356,7 @@ func (p *MarketPlugin) generateDailySummary(date string) string {
|
||||
}
|
||||
|
||||
func (p *MarketPlugin) generateReportSummary(snapsByDate map[string][]marketSnapshot, dateRange string) string {
|
||||
host := os.Getenv("OLLAMA_HOST")
|
||||
model := os.Getenv("OLLAMA_MODEL")
|
||||
if host == "" || model == "" {
|
||||
if !llmConfigured() {
|
||||
return ""
|
||||
}
|
||||
|
||||
@@ -395,7 +391,7 @@ func (p *MarketPlugin) generateReportSummary(snapsByDate map[string][]marketSnap
|
||||
}
|
||||
prompt.WriteString("\nDescribe the trend in 2-3 sentences. Note any significant moves or divergences between indices.\nBe sardonic but accurate. Reference the VIX trajectory when relevant.")
|
||||
|
||||
result, err := p.callOllamaChat(host, model, marketSystemPrompt, prompt.String())
|
||||
result, err := p.chatLLM(marketSystemPrompt, prompt.String())
|
||||
if err != nil {
|
||||
slog.Warn("market: ollama report summary failed", "err", err)
|
||||
return ""
|
||||
@@ -403,49 +399,14 @@ func (p *MarketPlugin) generateReportSummary(snapsByDate map[string][]marketSnap
|
||||
return result
|
||||
}
|
||||
|
||||
// callOllamaChat calls the Ollama /api/chat endpoint with a system and user message.
|
||||
// Uses the types already defined in holdem_tips.go (same package).
|
||||
func (p *MarketPlugin) callOllamaChat(host, model, systemPrompt, userPrompt string) (string, error) {
|
||||
req := ollamaChatRequest{
|
||||
Model: model,
|
||||
Messages: []chatMessage{
|
||||
{Role: "system", Content: systemPrompt},
|
||||
{Role: "user", Content: userPrompt},
|
||||
},
|
||||
Stream: false,
|
||||
}
|
||||
|
||||
body, err := json.Marshal(req)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("marshal: %w", err)
|
||||
}
|
||||
|
||||
url := strings.TrimRight(host, "/") + "/api/chat"
|
||||
resp, err := p.httpClient.Post(url, "application/json", bytes.NewReader(body))
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("request: %w", err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
respBody, _ := io.ReadAll(resp.Body)
|
||||
return "", fmt.Errorf("status %d: %s", resp.StatusCode, string(respBody))
|
||||
}
|
||||
|
||||
var ollamaResp ollamaChatResponse
|
||||
if err := json.NewDecoder(resp.Body).Decode(&ollamaResp); err != nil {
|
||||
return "", fmt.Errorf("decode: %w", err)
|
||||
}
|
||||
|
||||
text := ollamaResp.Message.Content
|
||||
// Strip <think>...</think> blocks (reasoning models)
|
||||
if i := strings.Index(text, "<think>"); i != -1 {
|
||||
if j := strings.Index(text, "</think>"); j != -1 {
|
||||
text = text[:i] + text[j+len("</think>"):]
|
||||
}
|
||||
}
|
||||
|
||||
return strings.TrimSpace(text), nil
|
||||
// chatLLM sends a system+user pair to the configured backend. Both backends
|
||||
// map this onto their own chat shape, so the market summaries read the same
|
||||
// whichever one is serving.
|
||||
func (p *MarketPlugin) chatLLM(systemPrompt, userPrompt string) (string, error) {
|
||||
return llmGenerate(context.Background(), llm.Request{
|
||||
System: systemPrompt,
|
||||
Prompt: userPrompt,
|
||||
})
|
||||
}
|
||||
|
||||
// ── DB Helpers ───────────────────────────────────────────────────────────────
|
||||
@@ -916,12 +877,10 @@ func (p *MarketPlugin) handleVixReport(ctx MessageContext) error {
|
||||
}
|
||||
|
||||
var summary string
|
||||
host := os.Getenv("OLLAMA_HOST")
|
||||
model := os.Getenv("OLLAMA_MODEL")
|
||||
if host != "" && model != "" {
|
||||
if llmConfigured() {
|
||||
prompt := fmt.Sprintf("VIX (fear index) data over %d days (%s to %s):\n%s\n\nDescribe the fear/greed trajectory in 2-3 sentences. Be sardonic but accurate.",
|
||||
len(entries), entries[0].Date, entries[len(entries)-1].Date, strings.Join(prices, ", "))
|
||||
summary, _ = p.callOllamaChat(host, model, marketSystemPrompt, prompt)
|
||||
summary, _ = p.chatLLM(marketSystemPrompt, prompt)
|
||||
}
|
||||
|
||||
var sb strings.Builder
|
||||
@@ -1020,12 +979,10 @@ func (p *MarketPlugin) handleCompare(ctx MessageContext, args []string) error {
|
||||
for _, e := range entries {
|
||||
prices = append(prices, fmt.Sprintf("%.2f", e.Price))
|
||||
}
|
||||
host := os.Getenv("OLLAMA_HOST")
|
||||
model := os.Getenv("OLLAMA_MODEL")
|
||||
if host != "" && model != "" {
|
||||
if llmConfigured() {
|
||||
prompt := fmt.Sprintf("%s over %d days (%s to %s):\n%s\n\nDescribe the trend in 2-3 sentences. Be sardonic but accurate.",
|
||||
idx.DisplayName, len(entries), entries[0].Date, entries[len(entries)-1].Date, strings.Join(prices, ", "))
|
||||
if summary, err := p.callOllamaChat(host, model, marketSystemPrompt, prompt); err == nil && summary != "" {
|
||||
if summary, err := p.chatLLM(marketSystemPrompt, prompt); err == nil && summary != "" {
|
||||
sb.WriteString(summary)
|
||||
sb.WriteString("\n\n")
|
||||
}
|
||||
|
||||
@@ -1,16 +1,14 @@
|
||||
package plugin
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"os"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"gogobee/internal/llm"
|
||||
"gogobee/internal/peteclient"
|
||||
)
|
||||
|
||||
@@ -39,22 +37,18 @@ const (
|
||||
maxDispatchLede = 800
|
||||
)
|
||||
|
||||
var dispatchHTTP = &http.Client{Timeout: dispatchLLMTimeout}
|
||||
|
||||
// authorDispatch turns a fact into a headline+lede in Pete's voice, or returns
|
||||
// two empty strings if the model is unconfigured, errors, times out, or produces
|
||||
// anything malformed. The fact must already have its FINAL Actors set (post
|
||||
// opt-out anonymisation) — that list is the only set of names the prose may use,
|
||||
// and it is what Pete's guard checks the output against.
|
||||
func authorDispatch(f peteclient.Fact) (headline, lede string) {
|
||||
host := os.Getenv("OLLAMA_HOST")
|
||||
model := os.Getenv("OLLAMA_MODEL")
|
||||
if host == "" || model == "" {
|
||||
if !llmConfigured() {
|
||||
return "", ""
|
||||
}
|
||||
|
||||
prompt := buildDispatchPrompt(f)
|
||||
raw, err := callOllamaDispatch(dispatchHTTP, host, model, prompt)
|
||||
raw, err := callLLMDispatch(dispatchLLMTimeout, prompt)
|
||||
if err != nil {
|
||||
slog.Warn("pete dispatch: LLM authoring failed, Pete will template", "guid", f.GUID, "err", err)
|
||||
return "", ""
|
||||
@@ -128,45 +122,17 @@ The event:
|
||||
%s`, names, facts.String())
|
||||
}
|
||||
|
||||
// callOllamaDispatch posts a single non-streaming generation and returns the raw
|
||||
// completion (think-tags stripped). The client is a parameter because the two
|
||||
// callers have genuinely different patience: a dispatch is authored on a game
|
||||
// chokepoint and must not stall it, while a run summary rides a background
|
||||
// ticker and can afford to wait for a bigger model. See runSummaryHTTP.
|
||||
func callOllamaDispatch(client *http.Client, host, model, prompt string) (string, error) {
|
||||
payload := map[string]interface{}{
|
||||
"model": model,
|
||||
"prompt": prompt,
|
||||
"stream": false,
|
||||
"think": false,
|
||||
"options": map[string]interface{}{
|
||||
"num_ctx": 4096,
|
||||
},
|
||||
}
|
||||
data, err := json.Marshal(payload)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("marshal payload: %w", err)
|
||||
}
|
||||
apiURL := strings.TrimRight(host, "/") + "/api/generate"
|
||||
resp, err := client.Post(apiURL, "application/json", bytes.NewReader(data))
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("ollama request: %w", err)
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
body, err := io.ReadAll(resp.Body)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("read response: %w", err)
|
||||
}
|
||||
if resp.StatusCode != http.StatusOK {
|
||||
return "", fmt.Errorf("ollama HTTP %d: %s", resp.StatusCode, string(body))
|
||||
}
|
||||
var result struct {
|
||||
Response string `json:"response"`
|
||||
}
|
||||
if err := json.Unmarshal(body, &result); err != nil {
|
||||
return "", fmt.Errorf("parse response: %w", err)
|
||||
}
|
||||
return result.Response, nil
|
||||
// callLLMDispatch posts a single non-streaming generation and returns the raw
|
||||
// completion. The timeout is a parameter because the two callers have genuinely
|
||||
// different patience: a dispatch is authored on a game chokepoint and must not
|
||||
// stall it, while a run summary rides a background ticker and can afford to wait
|
||||
// for a bigger model. See runSummaryTimeout.
|
||||
func callLLMDispatch(timeout time.Duration, prompt string) (string, error) {
|
||||
return llmGenerate(context.Background(), llm.Request{
|
||||
Prompt: prompt,
|
||||
NumCtx: 4096,
|
||||
Timeout: timeout,
|
||||
})
|
||||
}
|
||||
|
||||
// parseDispatch pulls {headline, lede} out of the model's completion, tolerating
|
||||
|
||||
@@ -4,8 +4,6 @@ import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"os"
|
||||
"sort"
|
||||
"strings"
|
||||
"sync/atomic"
|
||||
@@ -70,8 +68,6 @@ const maxRunSummary = 1200
|
||||
// for load-then-generate, and a timeout here really does mean the box is down.
|
||||
const runSummaryTimeout = 5 * time.Minute
|
||||
|
||||
var runSummaryHTTP = &http.Client{Timeout: runSummaryTimeout}
|
||||
|
||||
// runSummaryBusy is the whole concurrency story: at most one sweep in flight,
|
||||
// ever. The ticker starts one and moves on, so a cold model loading for minutes
|
||||
// costs the board nothing, and the ticks that fire meanwhile find the flag set
|
||||
@@ -105,7 +101,7 @@ func (p *AdventurePlugin) sweepRunSummaries() {
|
||||
if !peteclient.Enabled() || !newsEmissionOn() {
|
||||
return
|
||||
}
|
||||
if os.Getenv("OLLAMA_HOST") == "" || os.Getenv("OLLAMA_MODEL") == "" {
|
||||
if !llmConfigured() {
|
||||
return // no model, no summary, no wasted queries asking which run needs one
|
||||
}
|
||||
runID := nextRunNeedingSummary()
|
||||
@@ -174,8 +170,7 @@ func authorRunSummary(runID string) (summary, name string) {
|
||||
return "", ""
|
||||
}
|
||||
|
||||
raw, err := callOllamaDispatch(runSummaryHTTP, os.Getenv("OLLAMA_HOST"), os.Getenv("OLLAMA_MODEL"),
|
||||
buildRunSummaryPrompt(name, beats))
|
||||
raw, err := callLLMDispatch(runSummaryTimeout, buildRunSummaryPrompt(name, beats))
|
||||
if err != nil {
|
||||
slog.Warn("run summary: LLM authoring failed", "run", runID, "err", err)
|
||||
return "", name
|
||||
|
||||
+217
-11
@@ -83,9 +83,17 @@ func PluginVersion(p Plugin) string {
|
||||
}
|
||||
|
||||
// dmCache maps user IDs to their DM room IDs to avoid creating duplicate rooms.
|
||||
// It fronts the dm_rooms table; notDMCache is the negative half, marking rooms
|
||||
// LearnDMRoom has already ruled out so group rooms cost one member lookup ever.
|
||||
var (
|
||||
dmCache = make(map[id.UserID]id.RoomID)
|
||||
dmCacheMu sync.Mutex
|
||||
dmCache = make(map[id.UserID]id.RoomID)
|
||||
dmMapped = make(map[id.UserID]bool)
|
||||
notDMCache = make(map[id.RoomID]bool)
|
||||
dmCacheMu sync.Mutex
|
||||
|
||||
// One-shot index of two-person rooms, built by findExistingDMRoom.
|
||||
dmSweepOnce sync.Once
|
||||
dmSweepIndex = make(map[id.UserID]id.RoomID)
|
||||
)
|
||||
|
||||
// Base provides common helpers for plugin implementations.
|
||||
@@ -686,6 +694,180 @@ func (b *Base) SendReact(roomID id.RoomID, eventID id.EventID, emoji string) err
|
||||
return err
|
||||
}
|
||||
|
||||
// rememberDMRoom pins a user's DM room in both the in-process cache and the
|
||||
// database, so a restart doesn't send the bot off creating a duplicate room.
|
||||
func rememberDMRoom(userID id.UserID, roomID id.RoomID) {
|
||||
dmCacheMu.Lock()
|
||||
dmCache[userID] = roomID
|
||||
dmMapped[userID] = true
|
||||
dmCacheMu.Unlock()
|
||||
|
||||
if d := db.Get(); d != nil {
|
||||
_, err := d.Exec(`INSERT INTO dm_rooms (user_id, room_id, updated_at)
|
||||
VALUES (?, ?, unixepoch())
|
||||
ON CONFLICT(user_id) DO UPDATE SET room_id = excluded.room_id, updated_at = excluded.updated_at`,
|
||||
string(userID), string(roomID))
|
||||
if err != nil {
|
||||
slog.Error("persist dm room", "user", userID, "room", roomID, "err", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// forgetDMRoom drops a mapping that no longer resolves to a live shared room.
|
||||
func forgetDMRoom(userID id.UserID) {
|
||||
dmCacheMu.Lock()
|
||||
delete(dmCache, userID)
|
||||
delete(dmMapped, userID)
|
||||
dmCacheMu.Unlock()
|
||||
|
||||
if d := db.Get(); d != nil {
|
||||
if _, err := d.Exec(`DELETE FROM dm_rooms WHERE user_id = ?`, string(userID)); err != nil {
|
||||
slog.Error("forget dm room", "user", userID, "err", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// storedDMRoom reads the persisted DM room for a user, if any.
|
||||
func storedDMRoom(userID id.UserID) (id.RoomID, bool) {
|
||||
d := db.Get()
|
||||
if d == nil {
|
||||
return "", false
|
||||
}
|
||||
var roomID string
|
||||
err := d.QueryRow(`SELECT room_id FROM dm_rooms WHERE user_id = ?`, string(userID)).Scan(&roomID)
|
||||
if err != nil || roomID == "" {
|
||||
return "", false
|
||||
}
|
||||
return id.RoomID(roomID), true
|
||||
}
|
||||
|
||||
// dmRoomUsable reports whether the bot and the user still share the room. The
|
||||
// user counts as present while merely invited — they often never accept, and
|
||||
// treating that as "gone" is what would recreate the room on every send.
|
||||
func (b *Base) dmRoomUsable(roomID id.RoomID, userID id.UserID) bool {
|
||||
ctx := context.Background()
|
||||
|
||||
var self event.MemberEventContent
|
||||
if err := b.Client.StateEvent(ctx, roomID, event.StateMember, string(b.Client.UserID), &self); err != nil {
|
||||
return false
|
||||
}
|
||||
if self.Membership != event.MembershipJoin {
|
||||
return false
|
||||
}
|
||||
|
||||
var other event.MemberEventContent
|
||||
if err := b.Client.StateEvent(ctx, roomID, event.StateMember, string(userID), &other); err != nil {
|
||||
return false
|
||||
}
|
||||
return other.Membership == event.MembershipJoin || other.Membership == event.MembershipInvite
|
||||
}
|
||||
|
||||
// publishDirect appends the room to the bot's m.direct account data, so the
|
||||
// room is labelled as a DM rather than a nameless private room.
|
||||
func (b *Base) publishDirect(userID id.UserID, roomID id.RoomID) {
|
||||
ctx := context.Background()
|
||||
|
||||
dmRooms := map[id.UserID][]id.RoomID{}
|
||||
// A missing m.direct is a 404 — start from empty rather than bailing.
|
||||
_ = b.Client.GetAccountData(ctx, "m.direct", &dmRooms)
|
||||
|
||||
for _, existing := range dmRooms[userID] {
|
||||
if existing == roomID {
|
||||
return
|
||||
}
|
||||
}
|
||||
dmRooms[userID] = append(dmRooms[userID], roomID)
|
||||
|
||||
if err := b.Client.SetAccountData(ctx, "m.direct", dmRooms); err != nil {
|
||||
slog.Warn("publish m.direct", "user", userID, "room", roomID, "err", err)
|
||||
}
|
||||
}
|
||||
|
||||
// RecordDMRoom claims a room as the user's DM room. Used for user-initiated
|
||||
// DM invites, where the invite itself is the intent — it overwrites any older
|
||||
// mapping, since the user just told us which room they want to talk in.
|
||||
func (b *Base) RecordDMRoom(userID id.UserID, roomID id.RoomID) {
|
||||
slog.Info("recorded user-initiated DM room", "user", userID, "room", roomID)
|
||||
rememberDMRoom(userID, roomID)
|
||||
b.publishDirect(userID, roomID)
|
||||
}
|
||||
|
||||
// LearnDMRoom records a room the user messaged the bot in as their DM room,
|
||||
// when it really is a two-person room. This adopts DM rooms that predate the
|
||||
// database mapping instead of leaving them orphaned beside a freshly created one.
|
||||
func (b *Base) LearnDMRoom(userID id.UserID, roomID id.RoomID) {
|
||||
if b == nil || b.Client == nil {
|
||||
return
|
||||
}
|
||||
|
||||
dmCacheMu.Lock()
|
||||
_, known := dmCache[userID]
|
||||
mapped := dmMapped[userID]
|
||||
notDM := notDMCache[roomID]
|
||||
dmCacheMu.Unlock()
|
||||
if known || mapped || notDM {
|
||||
return
|
||||
}
|
||||
if _, ok := storedDMRoom(userID); ok {
|
||||
// Note it as mapped, not as resolved: caching the room here would let
|
||||
// GetDMRoom skip its liveness check on a possibly-stale room.
|
||||
dmCacheMu.Lock()
|
||||
dmMapped[userID] = true
|
||||
dmCacheMu.Unlock()
|
||||
return
|
||||
}
|
||||
|
||||
members, err := b.Client.JoinedMembers(context.Background(), roomID)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
_, botIn := members.Joined[b.Client.UserID]
|
||||
_, userIn := members.Joined[userID]
|
||||
if len(members.Joined) != 2 || !botIn || !userIn {
|
||||
// Group room — remember that so every later message here is free.
|
||||
dmCacheMu.Lock()
|
||||
notDMCache[roomID] = true
|
||||
dmCacheMu.Unlock()
|
||||
return
|
||||
}
|
||||
|
||||
slog.Info("adopted existing DM room", "user", userID, "room", roomID)
|
||||
rememberDMRoom(userID, roomID)
|
||||
}
|
||||
|
||||
// findExistingDMRoom scans the bot's joined rooms for a two-person room shared
|
||||
// with userID. The scan is expensive (one member lookup per room), so it runs
|
||||
// at most once per process and only on the path that would otherwise create a
|
||||
// duplicate room.
|
||||
func (b *Base) findExistingDMRoom(userID id.UserID) (id.RoomID, bool) {
|
||||
dmSweepOnce.Do(func() {
|
||||
ctx := context.Background()
|
||||
joined, err := b.Client.JoinedRooms(ctx)
|
||||
if err != nil {
|
||||
slog.Warn("DM sweep: list joined rooms", "err", err)
|
||||
return
|
||||
}
|
||||
for _, roomID := range joined.JoinedRooms {
|
||||
members, err := b.Client.JoinedMembers(ctx, roomID)
|
||||
if err != nil || len(members.Joined) != 2 {
|
||||
continue
|
||||
}
|
||||
if _, ok := members.Joined[b.Client.UserID]; !ok {
|
||||
continue
|
||||
}
|
||||
for member := range members.Joined {
|
||||
if member != b.Client.UserID {
|
||||
dmSweepIndex[member] = roomID
|
||||
}
|
||||
}
|
||||
}
|
||||
slog.Info("DM sweep complete", "rooms", len(joined.JoinedRooms), "dms", len(dmSweepIndex))
|
||||
})
|
||||
|
||||
roomID, ok := dmSweepIndex[userID]
|
||||
return roomID, ok
|
||||
}
|
||||
|
||||
// GetDMRoom returns the DM room for a user, creating one if needed.
|
||||
func (b *Base) GetDMRoom(userID id.UserID) (id.RoomID, error) {
|
||||
dmCacheMu.Lock()
|
||||
@@ -695,17 +877,42 @@ func (b *Base) GetDMRoom(userID id.UserID) (id.RoomID, error) {
|
||||
}
|
||||
dmCacheMu.Unlock()
|
||||
|
||||
// Check account data for existing DM rooms
|
||||
var dmRooms map[id.UserID][]id.RoomID
|
||||
err := b.Client.GetAccountData(context.Background(), "m.direct", &dmRooms)
|
||||
if err == nil {
|
||||
if rooms, ok := dmRooms[userID]; ok && len(rooms) > 0 {
|
||||
roomID := rooms[len(rooms)-1] // use most recent
|
||||
// Persisted mapping — the authoritative store across restarts.
|
||||
if roomID, ok := storedDMRoom(userID); ok {
|
||||
if b.dmRoomUsable(roomID, userID) {
|
||||
dmCacheMu.Lock()
|
||||
dmCache[userID] = roomID
|
||||
dmCacheMu.Unlock()
|
||||
return roomID, nil
|
||||
}
|
||||
slog.Info("stored DM room no longer usable, recreating", "user", userID, "room", roomID)
|
||||
forgetDMRoom(userID)
|
||||
}
|
||||
|
||||
// Check account data for existing DM rooms
|
||||
var dmRooms map[id.UserID][]id.RoomID
|
||||
err := b.Client.GetAccountData(context.Background(), "m.direct", &dmRooms)
|
||||
if err == nil {
|
||||
for i := len(dmRooms[userID]) - 1; i >= 0; i-- {
|
||||
roomID := dmRooms[userID][i] // most recent first
|
||||
if b.dmRoomUsable(roomID, userID) {
|
||||
rememberDMRoom(userID, roomID)
|
||||
return roomID, nil
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Last resort before creating: sweep the rooms the bot is already in for a
|
||||
// two-person room shared with this user. Users who predate the dm_rooms
|
||||
// table have such a room and nothing pointing at it — without this sweep
|
||||
// their first post-upgrade DM would open yet another duplicate. If several
|
||||
// duplicates exist we cannot tell which is liveliest (no /sync, so no
|
||||
// timestamps), so the newest by room-list order wins.
|
||||
if roomID, ok := b.findExistingDMRoom(userID); ok {
|
||||
slog.Info("recovered pre-existing DM room by member sweep", "user", userID, "room", roomID)
|
||||
rememberDMRoom(userID, roomID)
|
||||
b.publishDirect(userID, roomID)
|
||||
return roomID, nil
|
||||
}
|
||||
|
||||
// No existing DM room — create one
|
||||
@@ -728,9 +935,8 @@ func (b *Base) GetDMRoom(userID id.UserID) (id.RoomID, error) {
|
||||
return "", fmt.Errorf("create DM room: %w", err)
|
||||
}
|
||||
|
||||
dmCacheMu.Lock()
|
||||
dmCache[userID] = resp.RoomID
|
||||
dmCacheMu.Unlock()
|
||||
rememberDMRoom(userID, resp.RoomID)
|
||||
b.publishDirect(userID, resp.RoomID)
|
||||
return resp.RoomID, nil
|
||||
}
|
||||
|
||||
|
||||
@@ -4,7 +4,6 @@ import (
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"math/rand"
|
||||
"os"
|
||||
"strings"
|
||||
|
||||
"gogobee/internal/db"
|
||||
@@ -116,9 +115,7 @@ func (p *TarotPlugin) handleTarot(ctx MessageContext) error {
|
||||
return p.SendReply(ctx.RoomID, ctx.EventID, "You've used up your readings for today. The cards need rest, even if you don't.")
|
||||
}
|
||||
|
||||
ollamaHost := os.Getenv("OLLAMA_HOST")
|
||||
ollamaModel := os.Getenv("OLLAMA_MODEL")
|
||||
if ollamaHost == "" || ollamaModel == "" {
|
||||
if !llmConfigured() {
|
||||
return p.SendReply(ctx.RoomID, ctx.EventID, "Tarot reader is on a union-mandated vacation and will return when morale improves.")
|
||||
}
|
||||
|
||||
@@ -151,7 +148,7 @@ func (p *TarotPlugin) handleTarot(ctx MessageContext) error {
|
||||
}
|
||||
prompt := fmt.Sprintf("%s\n%s\n\n%s\n\nCard drawn: %s%s\nGive the reading.", tarotBasePrompt, tarotSingleSuffix, tarotFewShot, card, extraLines)
|
||||
|
||||
response, err := callOllama(ollamaHost, ollamaModel, prompt)
|
||||
response, err := callLLM(prompt)
|
||||
if err != nil {
|
||||
slog.Error("tarot: ollama call", "err", err)
|
||||
return p.SendReply(ctx.RoomID, ctx.EventID, "Tarot reader is on a union-mandated vacation and will return when morale improves.")
|
||||
@@ -166,9 +163,7 @@ func (p *TarotPlugin) handleSpread(ctx MessageContext) error {
|
||||
return p.SendReply(ctx.RoomID, ctx.EventID, "You've used up your readings for today. The cards need rest, even if you don't.")
|
||||
}
|
||||
|
||||
ollamaHost := os.Getenv("OLLAMA_HOST")
|
||||
ollamaModel := os.Getenv("OLLAMA_MODEL")
|
||||
if ollamaHost == "" || ollamaModel == "" {
|
||||
if !llmConfigured() {
|
||||
return p.SendReply(ctx.RoomID, ctx.EventID, "Tarot reader is on a union-mandated vacation and will return when morale improves.")
|
||||
}
|
||||
|
||||
@@ -202,7 +197,7 @@ func (p *TarotPlugin) handleSpread(ctx MessageContext) error {
|
||||
prompt := fmt.Sprintf("%s\n%s\n\n%s\n\nCards drawn:\n- Past: %s\n- Present: %s\n- Future: %s%s\nGive the reading.",
|
||||
tarotBasePrompt, tarotSpreadSuffix, tarotFewShot, cards[0], cards[1], cards[2], extraLines)
|
||||
|
||||
response, err := callOllama(ollamaHost, ollamaModel, prompt)
|
||||
response, err := callLLM(prompt)
|
||||
if err != nil {
|
||||
slog.Error("tarot: ollama call", "err", err)
|
||||
return p.SendReply(ctx.RoomID, ctx.EventID, "Tarot reader is on a union-mandated vacation and will return when morale improves.")
|
||||
|
||||
@@ -120,9 +120,7 @@ func (p *VibePlugin) resetCooldown(roomID id.RoomID) {
|
||||
}
|
||||
|
||||
func (p *VibePlugin) handleVibe(ctx MessageContext) error {
|
||||
ollamaHost := os.Getenv("OLLAMA_HOST")
|
||||
ollamaModel := os.Getenv("OLLAMA_MODEL")
|
||||
if ollamaHost == "" || ollamaModel == "" {
|
||||
if !llmConfigured() {
|
||||
return p.SendReply(ctx.RoomID, ctx.EventID, "LLM is not configured.")
|
||||
}
|
||||
|
||||
@@ -154,7 +152,7 @@ Describe the room's current vibe:`, botName, transcript)
|
||||
slog.Error("vibe: send thinking", "err", err)
|
||||
}
|
||||
|
||||
response, err := callOllama(ollamaHost, ollamaModel, prompt)
|
||||
response, err := callLLM(prompt)
|
||||
if err != nil {
|
||||
slog.Error("vibe: ollama call", "err", err)
|
||||
p.resetCooldown(ctx.RoomID) // Don't consume cooldown on failure
|
||||
@@ -165,9 +163,7 @@ Describe the room's current vibe:`, botName, transcript)
|
||||
}
|
||||
|
||||
func (p *VibePlugin) handleTLDR(ctx MessageContext) error {
|
||||
ollamaHost := os.Getenv("OLLAMA_HOST")
|
||||
ollamaModel := os.Getenv("OLLAMA_MODEL")
|
||||
if ollamaHost == "" || ollamaModel == "" {
|
||||
if !llmConfigured() {
|
||||
return p.SendReply(ctx.RoomID, ctx.EventID, "LLM is not configured.")
|
||||
}
|
||||
|
||||
@@ -199,7 +195,7 @@ Summary:`, tldrBotName, transcript)
|
||||
slog.Error("vibe: send thinking", "err", err)
|
||||
}
|
||||
|
||||
response, err := callOllama(ollamaHost, ollamaModel, prompt)
|
||||
response, err := callLLM(prompt)
|
||||
if err != nil {
|
||||
slog.Error("vibe: ollama call", "err", err)
|
||||
p.resetCooldown(ctx.RoomID) // Don't consume cooldown on failure
|
||||
|
||||
+17
-78
@@ -1,19 +1,17 @@
|
||||
package plugin
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"os"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"gogobee/internal/db"
|
||||
"gogobee/internal/dreamclient"
|
||||
"gogobee/internal/llm"
|
||||
|
||||
"maunium.net/go/mautrix"
|
||||
"maunium.net/go/mautrix/id"
|
||||
@@ -533,9 +531,7 @@ func (p *WOTDPlugin) trackUsage(ctx MessageContext) {
|
||||
// verifyUsage asks the LLM whether the word was used correctly in context.
|
||||
// Returns false if LLM is not configured or on any error.
|
||||
func (p *WOTDPlugin) verifyUsage(word, message string) bool {
|
||||
host := os.Getenv("OLLAMA_HOST")
|
||||
model := os.Getenv("OLLAMA_MODEL")
|
||||
if host == "" || model == "" {
|
||||
if !llmConfigured() {
|
||||
return false
|
||||
}
|
||||
|
||||
@@ -545,58 +541,30 @@ The Word of the Day is "%s". Was this word used correctly and meaningfully in th
|
||||
|
||||
Respond with ONLY "yes" or "no".`, message, word)
|
||||
|
||||
payload := map[string]interface{}{
|
||||
"model": model,
|
||||
"prompt": prompt,
|
||||
"stream": false,
|
||||
"think": false,
|
||||
}
|
||||
data, err := json.Marshal(payload)
|
||||
if err != nil {
|
||||
return false
|
||||
}
|
||||
|
||||
apiURL := strings.TrimRight(host, "/") + "/api/generate"
|
||||
slog.Debug("wotd: sending LLM verification request", "url", apiURL, "word", word)
|
||||
client := &http.Client{Timeout: 30 * time.Second}
|
||||
resp, err := client.Post(apiURL, "application/json", bytes.NewReader(data))
|
||||
response, err := llmGenerate(context.Background(), llm.Request{
|
||||
Prompt: prompt,
|
||||
Timeout: wotdLLMTimeout,
|
||||
})
|
||||
if err != nil {
|
||||
slog.Error("wotd: LLM verify request failed", "err", err)
|
||||
return false
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
body, err := io.ReadAll(resp.Body)
|
||||
if err != nil || resp.StatusCode != http.StatusOK {
|
||||
return false
|
||||
}
|
||||
|
||||
var result struct {
|
||||
Response string `json:"response"`
|
||||
}
|
||||
if json.Unmarshal(body, &result) != nil {
|
||||
return false
|
||||
}
|
||||
|
||||
response := result.Response
|
||||
// Strip <think>...</think> blocks (Qwen 3.5 reasoning)
|
||||
if i := strings.Index(response, "<think>"); i != -1 {
|
||||
if j := strings.Index(response, "</think>"); j != -1 {
|
||||
response = response[:i] + response[j+len("</think>"):]
|
||||
}
|
||||
}
|
||||
answer := strings.ToLower(strings.TrimSpace(response))
|
||||
accepted := strings.HasPrefix(answer, "yes")
|
||||
slog.Debug("wotd: LLM verification", "word", word, "answer", answer, "accepted", accepted)
|
||||
return accepted
|
||||
}
|
||||
|
||||
// wotdLLMTimeout preserves the 30s cap both WOTD paths enforced with their own
|
||||
// http.Client. These run on message traffic, so they stay well under the
|
||||
// interactive default.
|
||||
const wotdLLMTimeout = 30 * time.Second
|
||||
|
||||
// llmTranslate asks the LLM for a brief English translation of a foreign word.
|
||||
// Returns empty string on failure.
|
||||
func (p *WOTDPlugin) llmTranslate(word, lang string) string {
|
||||
host := os.Getenv("OLLAMA_HOST")
|
||||
model := os.Getenv("OLLAMA_MODEL")
|
||||
if host == "" || model == "" {
|
||||
if !llmConfigured() {
|
||||
return ""
|
||||
}
|
||||
|
||||
@@ -612,44 +580,15 @@ func (p *WOTDPlugin) llmTranslate(word, lang string) string {
|
||||
`Translate the %s word "%s" into English. Reply with ONLY the English translation — one or two words, no explanation, no punctuation.`,
|
||||
langName, word)
|
||||
|
||||
payload := map[string]interface{}{
|
||||
"model": model,
|
||||
"prompt": prompt,
|
||||
"stream": false,
|
||||
"think": false,
|
||||
}
|
||||
data, err := json.Marshal(payload)
|
||||
if err != nil {
|
||||
return ""
|
||||
}
|
||||
|
||||
apiURL := strings.TrimRight(host, "/") + "/api/generate"
|
||||
client := &http.Client{Timeout: 30 * time.Second}
|
||||
resp, err := client.Post(apiURL, "application/json", bytes.NewReader(data))
|
||||
response, err := llmGenerate(context.Background(), llm.Request{
|
||||
Prompt: prompt,
|
||||
Timeout: wotdLLMTimeout,
|
||||
})
|
||||
if err != nil {
|
||||
slog.Error("wotd: LLM translate request failed", "err", err)
|
||||
return ""
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
body, err := io.ReadAll(resp.Body)
|
||||
if err != nil || resp.StatusCode != http.StatusOK {
|
||||
return ""
|
||||
}
|
||||
|
||||
var result struct {
|
||||
Response string `json:"response"`
|
||||
}
|
||||
if json.Unmarshal(body, &result) != nil {
|
||||
return ""
|
||||
}
|
||||
|
||||
response := result.Response
|
||||
if i := strings.Index(response, "<think>"); i != -1 {
|
||||
if j := strings.Index(response, "</think>"); j != -1 {
|
||||
response = response[:i] + response[j+len("</think>"):]
|
||||
}
|
||||
}
|
||||
translation := strings.TrimSpace(response)
|
||||
if translation == "" || len(translation) > 50 {
|
||||
return ""
|
||||
|
||||
@@ -13,6 +13,7 @@ import (
|
||||
"gogobee/internal/bot"
|
||||
"gogobee/internal/db"
|
||||
"gogobee/internal/dreamclient"
|
||||
"gogobee/internal/llm"
|
||||
"gogobee/internal/peteclient"
|
||||
"gogobee/internal/plugin"
|
||||
"gogobee/internal/util"
|
||||
@@ -37,9 +38,11 @@ func main() {
|
||||
logLevel = "info"
|
||||
}
|
||||
util.InitLogger(logLevel)
|
||||
llmCfg := llm.ConfigFromEnv()
|
||||
slog.Info(version.Full(), "level", logLevel,
|
||||
"ollama_host", os.Getenv("OLLAMA_HOST"),
|
||||
"ollama_model", os.Getenv("OLLAMA_MODEL"))
|
||||
"llm_backend", llmCfg.Backend,
|
||||
"llm_endpoint", llmCfg.Endpoint,
|
||||
"llm_model", llmCfg.Model)
|
||||
|
||||
dataDir := os.Getenv("DATA_DIR")
|
||||
if dataDir == "" {
|
||||
@@ -252,6 +255,9 @@ func main() {
|
||||
|
||||
// ---- Set up event handlers ----
|
||||
|
||||
// Minimal Base used only for DM-room bookkeeping from the event handlers.
|
||||
dmLearner := &plugin.Base{Client: client}
|
||||
|
||||
// Auto-join on invite + moderation member tracking
|
||||
sess.OnEventType(event.StateMember, func(ctx context.Context, evt *event.Event) {
|
||||
defer func() {
|
||||
@@ -274,6 +280,11 @@ func main() {
|
||||
slog.Error("failed to join room", "room", evt.RoomID, "err", err)
|
||||
} else {
|
||||
slog.Info("joined room", "room", evt.RoomID)
|
||||
// A user-initiated DM invite: claim it now, before any
|
||||
// outbound DM has a chance to create a rival room.
|
||||
if mem.IsDirect {
|
||||
dmLearner.RecordDMRoom(evt.Sender, evt.RoomID)
|
||||
}
|
||||
}
|
||||
}
|
||||
return
|
||||
@@ -322,6 +333,10 @@ func main() {
|
||||
body = strings.TrimSpace(body[idx+2:])
|
||||
}
|
||||
}
|
||||
// Adopt the room the user is talking in if it's their DM room, so the
|
||||
// bot replies there instead of opening a second one.
|
||||
dmLearner.LearnDMRoom(evt.Sender, evt.RoomID)
|
||||
|
||||
msgCtx := plugin.MessageContext{
|
||||
RoomID: evt.RoomID,
|
||||
EventID: evt.ID,
|
||||
|
||||
Reference in New Issue
Block a user