Compare commits

...

2 Commits

Author SHA1 Message Date
kami b4646155b4 Read a mailbox read-only, in a client small enough to audit (#246)
internal/email is the reading half of the email reader: a ~200-line IMAP
client (LOGIN, EXAMINE, UID SEARCH SINCE, UID FETCH BODY.PEEK, LOGOUT), a
MIME-to-plaintext converter, and a header-only junk filter.

Two protocol choices are the design, not shortcuts. EXAMINE instead of
SELECT means the session is read-only at the protocol level, so no command
in it can flip a flag or expunge anything by mistake. BODY.PEEK instead of
BODY means reading a message does not mark it \Seen — Maven reads his mail
and leaves no trace of having done so, and the unread state in his own
client stays his.

Hand-rolled rather than go-imap because this is the one path that holds his
mailbox credential and reads his private mail: five commands with no
dependencies is auditable in a sitting. No IDLE and no cleartext/STARTTLS
either — an option to send his password over a plain socket is an option to
get it wrong once.

Junk is decided by headers alone, before any model is involved:
List-Unsubscribe/List-Id, Precedence: bulk, Auto-Submitted, the spam
headers, and Gmail's own category labels. Sender lists and subject keywords
are deliberately absent — they age badly and they would put his contacts in
a config file. A junk verdict only means "do not spend the model on this";
nothing is deleted and no server flag is touched.

Nothing here logs a body, a subject or an address, the junk reason names a
header rather than content, and an undecodable charset degrades to
headers-only instead of feeding the model mojibake. Verified against
recorded .eml fixtures and an in-process fake IMAP server.
2026-08-01 02:59:24 +04:00
kami da647e87d0 Read spending from zenmoney in the poller, answer it from facts (#125)
The trust boundary is zenmoney, not maven — they already hold his bank
sessions. So the poller reads /v8/diff/ and writes totals as
facts(kind=env, source=poll:zenmoney); core reads those back when he asks
and never sees the token.

internal/zenmoney sums transactions per currency over a window, skipping
tombstoned rows and transfers between his own accounts, and refuses to
encode a summary built from zero transactions. That refusal is the whole
design: a failed or empty read writes nothing and leaves the last good
total alone, because a zero recited as fact is worse than silence. No
currency conversion either — a figure he can check against his bank beats
one he cannot.

Off unless configured, and the token is read from a FILE rather than a
flag so it never lands in `ps`, in docker-compose.yml, or in shell
history. Nothing about the money is search input, no tick rule reads the
keys, and the log lines name keys, never figures.

The live-credential half is BLOCKED: there is no zenmoney account or token
here, so everything is verified against a recorded diff fixture.
2026-08-01 02:50:27 +04:00
26 changed files with 2328 additions and 3 deletions
+2
View File
@@ -34,6 +34,8 @@ deps
deploy/db_key.env deploy/db_key.env
# Deploy secret (telegram bot token + chat id) — never commit # Deploy secret (telegram bot token + chat id) — never commit
deploy/telegram.env deploy/telegram.env
# zenmoney API token, read by mavpoll (never in argv, never committed)
deploy/zenmoney.token
# Temp files # Temp files
/tmp/ /tmp/
+78
View File
@@ -0,0 +1,78 @@
package main
import (
"context"
"log"
"github.com/kami/maven/internal/ipc"
"github.com/kami/maven/internal/router"
"github.com/kami/maven/internal/zenmoney"
)
// Money questions (Vikunja #125).
//
// This is the whole read side: mavpoll holds the zenmoney token and writes
// facts(kind=env, source=poll:zenmoney); core reads them back when he asks.
// Core never sees the token, never calls zenmoney, and has no rule on these
// keys — a total is never a reason for Maven to speak first. Maven is not a
// nag, least of all about his money.
//
// Nothing here can reach the external search capability: the figures are read
// from the store and rendered locally, and his financial data is never search
// input.
// queryMoney — "сколько я потратил сегодня?", "покажи мои траты".
//
// Answers only from the latest fact the poller wrote. Three honest outcomes and
// no fourth: the figure, "the fact is old and here is its date", or "money
// tracking is not connected". It never computes, estimates or rounds a total of
// its own — an invented number about his money is the worst thing this could do.
func (h *reactiveHandler) queryMoney(ctx context.Context, t *queryTurn) (string, bool) {
window, ok := router.ParseMoneyQuery(t.dec.Utterance)
if !ok {
return "", false
}
key, phrase := zenmoney.KeySpentMonth, "в этом месяце"
if window == router.MoneyToday {
key, phrase = zenmoney.KeySpentToday, "сегодня"
}
fact, err := h.api.LatestFactBySource(ctx, key, zenmoney.Source)
if err != nil {
// No fact at all is the normal state when the capability is off. Claim
// the turn anyway: falling through to recall would answer a question
// about money with whatever note happens to be nearest.
if !isNoFactErr(err) {
log.Printf("voice: money fact: %v", err)
}
return "я не отслеживаю траты — не подключено.", true
}
val, err := zenmoney.ParseFactValue(fact.Value)
if err != nil {
log.Printf("voice: money fact: decode: %v", err)
return "не получилось прочитать траты.", true
}
reply := val.FormatRU(phrase)
if reply == "" {
return "по тратам пока нечего сказать.", true
}
// A stale fact is reported as stale rather than spoken as today's number.
if h.now().Sub(fact.Ts) > zenmoney.StaleAfter {
return "данные от " + fact.Ts.Local().Format("02.01") + ": " + reply, true
}
return reply, true
}
// isNoFactErr — ErrNoFact survives the wire wrapped, so unwrap for it.
func isNoFactErr(err error) bool {
for e := err; e != nil; {
if e == ipc.ErrNoFact {
return true
}
u, ok := e.(interface{ Unwrap() error })
if !ok {
return false
}
e = u.Unwrap()
}
return false
}
+139
View File
@@ -0,0 +1,139 @@
package main
import (
"context"
"strings"
"testing"
"time"
"github.com/kami/maven/internal/ipc"
"github.com/kami/maven/internal/router"
"github.com/kami/maven/internal/zenmoney"
)
// moneyAPI answers only LatestFactBySource; everything else is unimplemented,
// which is the assertion that answering a money question costs no model call
// and reaches no network.
type moneyAPI struct {
ipc.UnimplementedCoreAPI
fact ipc.Fact
err error
gotKey string
gotSrc string
callCnt int
}
func (a *moneyAPI) LatestFactBySource(_ context.Context, key, source string) (ipc.Fact, error) {
a.gotKey, a.gotSrc = key, source
a.callCnt++
return a.fact, a.err
}
func moneyNow() time.Time { return time.Date(2026, 8, 15, 20, 0, 0, 0, time.UTC) }
func moneyFact(ts time.Time, val string) ipc.Fact {
return ipc.Fact{Kind: "env", Key: zenmoney.KeySpentMonth, Value: val, Source: zenmoney.Source, Ts: ts}
}
func TestQueryMoneyAnswersFromTheFact(t *testing.T) {
api := &moneyAPI{fact: moneyFact(moneyNow(), `{"spent":[{"currency":"RUB","amount":1749.5}],"count":3}`)}
h := &reactiveHandler{api: api, now: moneyNow}
reply, ok := h.queryMoney(context.Background(), &queryTurn{
dec: router.Decision{Utterance: "сколько я потратил в этом месяце?"},
})
if !ok {
t.Fatal("the money source must claim a money question")
}
if api.gotKey != zenmoney.KeySpentMonth || api.gotSrc != zenmoney.Source {
t.Errorf("read %q/%q, want the month key from the poller's source", api.gotKey, api.gotSrc)
}
if !strings.Contains(reply, "1749.5") {
t.Errorf("reply = %q, want the exact figure", reply)
}
if !strings.Contains(reply, "в этом месяце") {
t.Errorf("reply = %q, want the window named", reply)
}
}
func TestQueryMoneyPicksTodaysKey(t *testing.T) {
api := &moneyAPI{fact: moneyFact(moneyNow(), `{"spent":[{"currency":"RUB","amount":250}],"count":1}`)}
h := &reactiveHandler{api: api, now: moneyNow}
if _, ok := h.queryMoney(context.Background(), &queryTurn{
dec: router.Decision{Utterance: "сколько я потратил сегодня?"},
}); !ok {
t.Fatal("expected the source to claim it")
}
if api.gotKey != zenmoney.KeySpentToday {
t.Errorf("key = %q, want today's", api.gotKey)
}
}
// The capability is off unless configured, and then there is no fact. She says
// so instead of letting the recall pass answer a money question from a note.
func TestQueryMoneySaysNotConnected(t *testing.T) {
h := &reactiveHandler{api: &moneyAPI{err: ipc.ErrNoFact}, now: moneyNow}
reply, ok := h.queryMoney(context.Background(), &queryTurn{
dec: router.Decision{Utterance: "сколько я потратил?"},
})
if !ok {
t.Fatal("expected the source to claim it")
}
if !strings.Contains(reply, "не подключено") {
t.Errorf("reply = %q, want an honest 'not connected'", reply)
}
// No number of any kind in that answer.
for _, d := range []string{"0", "1", "2", "3", "4", "5", "6", "7", "8", "9"} {
if strings.Contains(reply, d) {
t.Errorf("reply %q contains a digit — nothing was read, so there is no figure", reply)
}
}
}
// A fact older than the staleness bound is dated rather than spoken as if it
// were current: the poller can be down, and last week's total presented as
// today's is a lie by omission.
func TestQueryMoneyDatesAStaleFact(t *testing.T) {
old := moneyNow().Add(-72 * time.Hour)
api := &moneyAPI{fact: moneyFact(old, `{"spent":[{"currency":"RUB","amount":100}],"count":1}`)}
h := &reactiveHandler{api: api, now: moneyNow}
reply, _ := h.queryMoney(context.Background(), &queryTurn{
dec: router.Decision{Utterance: "сколько я потратил?"},
})
if !strings.Contains(reply, "данные от") {
t.Errorf("reply = %q, want the stale fact dated", reply)
}
}
func TestQueryMoneyPassesOtherQuestions(t *testing.T) {
api := &moneyAPI{}
h := &reactiveHandler{api: api, now: moneyNow}
for _, u := range []string{"какая погода?", "я потратил весь день на это", "какие у меня задачи?"} {
if _, ok := h.queryMoney(context.Background(), &queryTurn{dec: router.Decision{Utterance: u}}); ok {
t.Errorf("the money source claimed %q", u)
}
}
if api.callCnt != 0 {
t.Error("a non-money question must not read the money facts")
}
}
// Money must be answered before the recall sources, or a question about
// spending gets answered by the nearest note.
func TestQuerySourcesOrderMoneyBeforeRecall(t *testing.T) {
moneyAt, notesAt := -1, -1
for i, src := range querySources {
switch src.name {
case "money":
moneyAt = i
case "notes":
notesAt = i
}
}
if moneyAt < 0 || notesAt < 0 {
t.Fatalf("sources missing: money=%d notes=%d", moneyAt, notesAt)
}
if moneyAt > notesAt {
t.Errorf("money source at %d, after notes at %d", moneyAt, notesAt)
}
}
+5
View File
@@ -62,6 +62,11 @@ var querySources = []querySource{
// matcher requires a task noun or an explicit "что … сделать", so a // matcher requires a task noun or an explicit "что … сделать", so a
// date-bearing question still reaches the calendar. // date-bearing question still reaches the calendar.
{"tasks", (*reactiveHandler).queryTasks}, {"tasks", (*reactiveHandler).queryTasks},
// Before the recall sources too: "сколько я потратил?" is a question about
// the money facts the poller wrote, and the notes pass would otherwise
// answer it from whatever he once said about spending. Its matcher needs a
// money noun plus an actual ask, so "я потратил весь день" is untouched.
{"money", (*reactiveHandler).queryMoney},
{"calendar", (*reactiveHandler).queryCalendar}, {"calendar", (*reactiveHandler).queryCalendar},
{"weather", (*reactiveHandler).queryWeather}, {"weather", (*reactiveHandler).queryWeather},
{"embed", (*reactiveHandler).queryEmbed}, {"embed", (*reactiveHandler).queryEmbed},
+124 -3
View File
@@ -9,6 +9,12 @@
// Two sources, each its own provenance (the loop's rules trust source): // Two sources, each its own provenance (the loop's rules trust source):
// - netdata → poll:netdata resource alarms (disk/mem/cert/temp) // - netdata → poll:netdata resource alarms (disk/mem/cert/temp)
// - kuma → poll:uptimekuma service up/down (the source of truth for it) // - kuma → poll:uptimekuma service up/down (the source of truth for it)
// - zenmoney → poll:zenmoney spending/income totals (Vikunja #125)
//
// The zenmoney source is why the token lives HERE and not in core: the poller
// already owns every other third-party credential, it holds no store key, and
// core never needs to know an account exists to answer a question about a fact
// the poller wrote. It is off unless -zenmoney-token-file is given.
// //
// Netdata needs no auth over the wg-fronted net. Kuma's /metrics needs an API // Netdata needs no auth over the wg-fronted net. Kuma's /metrics needs an API
// key (basic-auth); without -kuma the whole kuma path is skipped (netdata-only // key (basic-auth); without -kuma the whole kuma path is skipped (netdata-only
@@ -37,6 +43,7 @@ import (
"time" "time"
"github.com/kami/maven/internal/ipc" "github.com/kami/maven/internal/ipc"
"github.com/kami/maven/internal/zenmoney"
) )
func main() { func main() {
@@ -52,6 +59,9 @@ func run(args []string) error {
netdataURL := fs.String("netdata", "http://127.0.0.1:19999", "netdata base URL ('' to disable)") netdataURL := fs.String("netdata", "http://127.0.0.1:19999", "netdata base URL ('' to disable)")
kumaURL := fs.String("kuma", "", "uptime-kuma metrics URL, e.g. http://127.0.0.1:3001/metrics ('' to disable)") kumaURL := fs.String("kuma", "", "uptime-kuma metrics URL, e.g. http://127.0.0.1:3001/metrics ('' to disable)")
kumaKey := fs.String("kuma-key", "", "uptime-kuma API key (basic-auth username)") kumaKey := fs.String("kuma-key", "", "uptime-kuma API key (basic-auth username)")
zenTokenFile := fs.String("zenmoney-token-file", "", "file holding the zenmoney API token ('' disables money tracking)")
zenURL := fs.String("zenmoney-url", zenmoney.DefaultBaseURL, "zenmoney API base URL (tests/self-hosted proxies)")
zenInterval := fs.Duration("zenmoney-interval", time.Hour, "how often to read zenmoney (money does not move every minute)")
wgIface := fs.String("wg", "", "wireguard interface for the presence signal, e.g. wg0 or 'all' ('' to disable)") wgIface := fs.String("wg", "", "wireguard interface for the presence signal, e.g. wg0 or 'all' ('' to disable)")
wgCmd := fs.String("wg-cmd", "wg", "wg binary (use e.g. 'sudo wg' if the poller lacks CAP_NET_ADMIN)") wgCmd := fs.String("wg-cmd", "wg", "wg binary (use e.g. 'sudo wg' if the poller lacks CAP_NET_ADMIN)")
interval := fs.Duration("interval", 60*time.Second, "poll cadence") interval := fs.Duration("interval", 60*time.Second, "poll cadence")
@@ -62,8 +72,24 @@ func run(args []string) error {
if *socket == "" { if *socket == "" {
return fmt.Errorf("-socket is required") return fmt.Errorf("-socket is required")
} }
if *netdataURL == "" && *kumaURL == "" && *wgIface == "" { if *netdataURL == "" && *kumaURL == "" && *wgIface == "" && *zenTokenFile == "" {
return fmt.Errorf("nothing to poll: set -netdata, -kuma and/or -wg") return fmt.Errorf("nothing to poll: set -netdata, -kuma, -wg and/or -zenmoney-token-file")
}
// The token is read from a file, never taken as a flag value: an argv token
// is visible in `ps` to every user on the box and lands in the compose file
// and the shell history. Read once at start — a rotated token means a
// restart, which is cheaper than re-reading his credential every hour.
var zen *zenmoney.Client
if *zenTokenFile != "" {
raw, err := os.ReadFile(*zenTokenFile)
if err != nil {
return fmt.Errorf("read zenmoney token: %w", err)
}
zen, err = zenmoney.New(strings.TrimSpace(string(raw)), *zenURL, *timeout*3)
if err != nil {
return err
}
} }
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
@@ -83,9 +109,13 @@ func run(args []string) error {
kumaKey: *kumaKey, kumaKey: *kumaKey,
wgIface: *wgIface, wgIface: *wgIface,
wgCmd: *wgCmd, wgCmd: *wgCmd,
zen: zen,
zenEvery: *zenInterval,
} }
log.Printf("mavpoll: polling every %s (netdata=%q kuma=%q wg=%q)", *interval, *netdataURL, *kumaURL, *wgIface) // The token is never logged, not even its length.
log.Printf("mavpoll: polling every %s (netdata=%q kuma=%q wg=%q zenmoney=%v every %s)",
*interval, *netdataURL, *kumaURL, *wgIface, zen != nil, *zenInterval)
p.pollOnce(ctx) // fire immediately; don't idle a full interval on start p.pollOnce(ctx) // fire immediately; don't idle a full interval on start
t := time.NewTicker(*interval) t := time.NewTicker(*interval)
defer t.Stop() defer t.Stop()
@@ -108,6 +138,12 @@ type poller struct {
kumaKey string kumaKey string
wgIface string wgIface string
wgCmd string wgCmd string
// zen is nil unless a token file was configured — money tracking is a
// capability, off by default like weather and telegram.
zen *zenmoney.Client
zenEvery time.Duration
zenLast time.Time
} }
// pollOnce — one sweep of both sources. A failure in one source logs and does // pollOnce — one sweep of both sources. A failure in one source logs and does
@@ -129,6 +165,66 @@ func (p *poller) pollOnce(ctx context.Context) {
log.Printf("mavpoll: wg: %v", err) log.Printf("mavpoll: wg: %v", err)
} }
} }
// Money on its own, much slower cadence: a bank feed that updates hourly
// polled every minute is 60 pointless reads of his financial history.
if p.zen != nil && now.Sub(p.zenLast) >= p.zenEvery {
p.zenLast = now
if err := p.pollZenmoney(ctx, now); err != nil {
log.Printf("mavpoll: zenmoney: %v", err)
}
}
}
// ---- zenmoney: spending/income totals → money facts ------------------------
// pollZenmoney reads today's and this month's totals and writes them as
// facts(kind=env, source=poll:zenmoney) (Vikunja #125).
//
// Two properties this function exists to hold:
//
// - An empty or failed read writes NOTHING. zenmoney.Summary.Value() refuses
// to encode a summary built from zero transactions, so a poller that cannot
// reach the API leaves the last good fact in place rather than overwriting
// it with a zero Maven would then recite as fact.
// - Nothing about the money leaves the box except the diff request itself, to
// the service that already holds his bank sessions. The totals are written
// to the store and read back only when he asks; they are never search input
// and no tick rule fires on them.
//
// Both windows are read from one diff call each. Two calls an hour against an
// API whose whole job is this is not worth caching.
// moneyWindow — one fact key and the period it covers.
type moneyWindow struct {
key string
from, to time.Time
}
func (p *poller) pollZenmoney(ctx context.Context, now time.Time) error {
dFrom, dTo := zenmoney.DayWindow(now)
mFrom, mTo := zenmoney.MonthWindow(now)
windows := []moneyWindow{
{zenmoney.KeySpentToday, dFrom, dTo},
{zenmoney.KeySpentMonth, mFrom, mTo},
}
var firstErr error
for _, w := range windows {
sum, err := p.zen.Since(ctx, w.from, w.to)
if err != nil {
if firstErr == nil {
firstErr = err
}
continue
}
val, ok := sum.Value()
if !ok {
// Nothing read. Silence, not a zero.
continue
}
if err := p.writeIfChangedRaw(ctx, w.key, zenmoney.Source, val, now); err != nil && firstErr == nil {
firstErr = err
}
}
return firstErr
} }
// ---- wireguard: latest handshake → presence signal ------------------------- // ---- wireguard: latest handshake → presence signal -------------------------
@@ -301,6 +397,31 @@ func (p *poller) writeIfChanged(ctx context.Context, key, source, val string, no
return nil return nil
} }
// writeIfChangedRaw is writeIfChanged for values that are already JSON (the
// money facts store an object, not a string). Kept separate rather than
// generalising writeIfChanged, because the string-valued env facts encoding
// their own value is the convention the rules rely on.
//
// The log line names the key and the source, never the figures: mavpoll's log
// is not the place his spending ends up.
func (p *poller) writeIfChangedRaw(ctx context.Context, key, source, jsonVal string, now time.Time) error {
prev, err := p.core.LatestFactBySource(ctx, key, source)
switch {
case err == nil && prev.Value == jsonVal:
return nil
case err != nil && err != ipc.ErrNoFact && !isNoFact(err):
return fmt.Errorf("read %s: %w", key, err)
}
if _, err := p.core.WriteFact(ctx, ipc.WriteFactReq{
Ts: now, Kind: "env", Key: key, Value: jsonVal,
Source: source, Confidence: 1.0,
}); err != nil {
return fmt.Errorf("write %s: %w", key, err)
}
log.Printf("mavpoll: %s updated (%s)", key, source)
return nil
}
// isNoFact — ErrNoFact rehydrated over the wire is wrapped (fmt.Errorf %w), so // isNoFact — ErrNoFact rehydrated over the wire is wrapped (fmt.Errorf %w), so
// errors.Is is the right check; keep a helper so the switch above reads clean. // errors.Is is the right check; keep a helper so the switch above reads clean.
func isNoFact(err error) bool { func isNoFact(err error) bool {
+126
View File
@@ -1,8 +1,17 @@
package main package main
import ( import (
"context"
"encoding/json" "encoding/json"
"net/http"
"net/http/httptest"
"os"
"strings"
"testing" "testing"
"time"
"github.com/kami/maven/internal/ipc"
"github.com/kami/maven/internal/zenmoney"
) )
func TestMaxSeverity(t *testing.T) { func TestMaxSeverity(t *testing.T) {
@@ -62,3 +71,120 @@ func TestParseMaxHandshake(t *testing.T) {
} }
} }
} }
// ---- zenmoney (Vikunja #125) ----------------------------------------------
// factCore records the facts the poller wrote and answers "no fact yet".
type factCore struct {
ipc.UnimplementedCoreAPI
written []ipc.WriteFactReq
prev map[string]string
}
func (c *factCore) LatestFactBySource(_ context.Context, key, source string) (ipc.Fact, error) {
if v, ok := c.prev[key+"|"+source]; ok {
return ipc.Fact{Key: key, Source: source, Value: v}, nil
}
return ipc.Fact{}, ipc.ErrNoFact
}
func (c *factCore) WriteFact(_ context.Context, req ipc.WriteFactReq) (int64, error) {
c.written = append(c.written, req)
return int64(len(c.written)), nil
}
func zenFixtureServer(t *testing.T, body []byte, status int) *httptest.Server {
t.Helper()
return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if status != http.StatusOK {
w.WriteHeader(status)
return
}
w.Write(body)
}))
}
func TestPollZenmoneyWritesMoneyFacts(t *testing.T) {
body, err := os.ReadFile("../../internal/zenmoney/testdata/diff.json")
if err != nil {
t.Fatal(err)
}
srv := zenFixtureServer(t, body, http.StatusOK)
defer srv.Close()
zen, err := zenmoney.New("tok", srv.URL, time.Second)
if err != nil {
t.Fatal(err)
}
core := &factCore{}
p := &poller{core: core, zen: zen}
now := time.Date(2026, 8, 1, 21, 0, 0, 0, time.UTC)
if err := p.pollZenmoney(context.Background(), now); err != nil {
t.Fatal(err)
}
if len(core.written) != 2 {
t.Fatalf("wrote %d facts, want today + month", len(core.written))
}
for _, f := range core.written {
if f.Kind != "env" || f.Source != zenmoney.Source {
t.Errorf("fact = %+v, want kind=env source=%s", f, zenmoney.Source)
}
if _, err := zenmoney.ParseFactValue(f.Value); err != nil {
t.Errorf("fact value %q does not decode: %v", f.Value, err)
}
}
}
// A read that returns nothing for the window writes NOTHING. Silence, not a
// zero: an invented 0 would be recited back to him as fact.
func TestPollZenmoneyWritesNothingWhenEmpty(t *testing.T) {
srv := zenFixtureServer(t, []byte(`{"serverTimestamp":1,"instrument":[],"transaction":[]}`), http.StatusOK)
defer srv.Close()
zen, _ := zenmoney.New("tok", srv.URL, time.Second)
core := &factCore{}
p := &poller{core: core, zen: zen}
if err := p.pollZenmoney(context.Background(), time.Now()); err != nil {
t.Fatal(err)
}
if len(core.written) != 0 {
t.Errorf("wrote %+v, want no fact at all", core.written)
}
}
// An API failure must not overwrite the last good total either.
func TestPollZenmoneyFailureWritesNothing(t *testing.T) {
srv := zenFixtureServer(t, nil, http.StatusUnauthorized)
defer srv.Close()
zen, _ := zenmoney.New("bad", srv.URL, time.Second)
core := &factCore{}
p := &poller{core: core, zen: zen}
if err := p.pollZenmoney(context.Background(), time.Now()); err == nil {
t.Error("want the 401 reported")
}
if len(core.written) != 0 {
t.Errorf("wrote %+v on a failed read", core.written)
}
}
// Unchanged totals do not churn the facts table.
func TestWriteIfChangedRawSkipsUnchanged(t *testing.T) {
core := &factCore{prev: map[string]string{
zenmoney.KeySpentToday + "|" + zenmoney.Source: `{"count":1}`,
}}
p := &poller{core: core}
if err := p.writeIfChangedRaw(context.Background(), zenmoney.KeySpentToday, zenmoney.Source, `{"count":1}`, time.Now()); err != nil {
t.Fatal(err)
}
if len(core.written) != 0 {
t.Errorf("wrote %+v for an unchanged value", core.written)
}
}
// Money tracking is off unless configured: no token file, no zenmoney client,
// and the poller still refuses to start with nothing at all to poll.
func TestRunRequiresSomethingToPoll(t *testing.T) {
err := run([]string{"-socket", "/tmp/nope.sock", "-netdata", "", "-kuma", "", "-wg", ""})
if err == nil || !strings.Contains(err.Error(), "nothing to poll") {
t.Errorf("err = %v, want a 'nothing to poll' refusal", err)
}
}
+8
View File
@@ -101,9 +101,17 @@ services:
"-netdata", "http://127.0.0.1:19999", "-netdata", "http://127.0.0.1:19999",
"-kuma", "http://127.0.0.1:3001/metrics", "-kuma", "http://127.0.0.1:3001/metrics",
"-kuma-key", "uk5_mavpoll-key"] "-kuma-key", "uk5_mavpoll-key"]
# Money tracking (Vikunja #125) is OFF: it needs a zenmoney token,
# which mavpoll reads from a FILE so it never appears in `ps`, in
# this file, or in shell history. To enable, mount the token and
# append: "-zenmoney-token-file", "/run/secrets/zenmoney.token"
# (optionally "-zenmoney-interval", "1h"). Core never sees the
# token — the poller writes facts(kind=env, source=poll:zenmoney)
# and mavend only reads those back when he asks.
depends_on: [mavend] depends_on: [mavend]
volumes: volumes:
- sockets:/run/maven - sockets:/run/maven
# - ./deploy/zenmoney.token:/run/secrets/zenmoney.token:ro
volumes: volumes:
dbdata: dbdata:
+96
View File
@@ -0,0 +1,96 @@
package email
import (
"fmt"
"time"
)
// FetchSince is the whole read path in one call: connect, log in, examine the
// mailbox read-only, list what arrived since a date, fetch and parse the ones
// the caller has not seen, log out.
//
// It is a function rather than a long-lived object because a mail poller should
// not hold an authenticated session (and therefore his credential in a live TLS
// state) between polls. Connect, read, drop.
//
// skip decides which UIDs are already known — the poller's seen-set. max bounds
// one poll: a mailbox that received 400 messages overnight must not turn into
// 400 LLM calls, and the newest max are the ones a task could still be hiding
// in. Junk messages are returned too, flagged, so the caller can mark them seen
// without a second protocol round.
type FetchSince struct {
Addr string // host or host:993
User string
Mailbox string // e.g. "INBOX"
Timeout time.Duration
Since time.Time
Max int
Skip func(uid uint32) bool
// dial is the connection seam. nil means Dial (implicit TLS); the tests set
// it to an in-process fake. Unexported so no configuration path can point
// the reader at a non-TLS transport.
dial func(addr string, timeout time.Duration) (*Conn, error)
}
// Run performs one read. password is passed here, not stored in the struct, so
// the configuration of a mailbox and the secret for it are never the same value
// sitting in the same place.
func (f FetchSince) Run(password string) ([]Message, error) {
if f.Addr == "" || f.User == "" || f.Mailbox == "" {
return nil, fmt.Errorf("email: mailbox not configured (addr/user/mailbox)")
}
dial := f.dial
if dial == nil {
dial = Dial
}
c, err := dial(f.Addr, f.Timeout)
if err != nil {
return nil, err
}
defer c.Close()
if err := c.Login(f.User, password); err != nil {
return nil, err
}
defer c.Logout()
if err := c.Select(f.Mailbox); err != nil {
return nil, err
}
uids, err := c.SearchSince(f.Since)
if err != nil {
return nil, err
}
// Newest UIDs first — IMAP hands them back ascending, and when Max clips the
// list the recent mail is what matters.
wanted := make([]uint32, 0, len(uids))
for i := len(uids) - 1; i >= 0; i-- {
if f.Skip != nil && f.Skip(uids[i]) {
continue
}
wanted = append(wanted, uids[i])
if f.Max > 0 && len(wanted) >= f.Max {
break
}
}
out := make([]Message, 0, len(wanted))
for _, uid := range wanted {
raw, err := c.Fetch(uid)
if err != nil {
// One unreadable message does not abandon the poll; the rest of the
// mailbox is still worth reading. The error names the UID, not the
// message.
return out, fmt.Errorf("email: fetch uid %d: %w", uid, err)
}
if len(raw) == 0 {
continue // vanished between SEARCH and FETCH
}
msg, err := ParseMessage(uid, raw)
if err != nil {
continue // unparsable headers — nothing to review, skip silently
}
out = append(out, msg)
}
return out, nil
}
+50
View File
@@ -0,0 +1,50 @@
package email
import (
"net"
"strings"
"testing"
"time"
)
func TestFetchSinceRun(t *testing.T) {
mk := func(subject string) string {
return "Subject: " + subject + "\r\nContent-Type: text/plain; charset=utf-8\r\n\r\nbody\r\n"
}
f := &fakeIMAP{
uids: []uint32{1, 2, 3},
msgs: map[uint32]string{1: mk("one"), 2: mk("two"), 3: mk("three")},
}
fs := FetchSince{
Addr: "mail.example:993", User: "kami", Mailbox: "INBOX",
Timeout: 5 * time.Second,
Since: time.Date(2026, 7, 30, 0, 0, 0, 0, time.UTC),
Max: 2,
Skip: func(uid uint32) bool { return uid == 3 },
dial: func(addr string, timeout time.Duration) (*Conn, error) {
cli, srv := net.Pipe()
go f.serve(t, srv)
return NewConn(cli, timeout)
},
}
msgs, err := fs.Run("secret")
if err != nil {
t.Fatalf("run: %v", err)
}
// Newest first, the already-seen UID skipped, Max respected.
if len(msgs) != 2 {
t.Fatalf("got %d messages, want 2: %+v", len(msgs), msgs)
}
if msgs[0].Subject != "two" || msgs[1].Subject != "one" {
t.Errorf("subjects = %q,%q, want two,one (newest first)", msgs[0].Subject, msgs[1].Subject)
}
if strings.Contains(strings.Join(f.cmds, " "), "UID FETCH 3") {
t.Error("a skipped UID must not be fetched again")
}
}
func TestFetchSinceRequiresConfig(t *testing.T) {
if _, err := (FetchSince{}).Run("secret"); err == nil {
t.Fatal("an unconfigured mailbox must not be read")
}
}
+280
View File
@@ -0,0 +1,280 @@
package email
import (
"bufio"
"crypto/tls"
"fmt"
"io"
"net"
"regexp"
"strconv"
"strings"
"time"
)
// A minimal IMAP4rev1 client — LOGIN, SELECT, UID SEARCH, UID FETCH with
// BODY.PEEK, LOGOUT, and nothing else.
//
// Why hand-rolled instead of go-imap: the whole surface Maven needs is five
// commands, and this is the one code path that holds his mailbox credential and
// reads his private mail. A ~200-line client with no dependencies is auditable
// in one sitting; a general-purpose IMAP library is a much larger amount of
// code doing much more than we asked, in the most sensitive place in the tree.
// If IDLE, CONDSTORE or server-side threading ever become worth having, that
// trade should be re-made deliberately.
//
// BODY.PEEK[] rather than BODY[] is load-bearing: Maven reads his mail and must
// leave no trace of having done so. Reading a message here does not mark it
// \Seen, so the unread state in his own mail client stays his.
// DefaultIMAPPort — implicit-TLS IMAP. There is no cleartext and no STARTTLS
// path in this client: an option to send his password over a plain socket is an
// option to get it wrong once.
const DefaultIMAPPort = "993"
// Conn — one authenticated IMAP connection. Not safe for concurrent use; the
// poller drives one connection at a time.
type Conn struct {
rwc io.ReadWriteCloser
r *bufio.Reader
tag int
timeout time.Duration
}
// Dial opens an implicit-TLS connection and reads the server greeting.
func Dial(addr string, timeout time.Duration) (*Conn, error) {
host, _, err := net.SplitHostPort(addr)
if err != nil {
host, addr = addr, net.JoinHostPort(addr, DefaultIMAPPort)
}
d := &net.Dialer{Timeout: timeout}
// ServerName is set from the host we asked for: certificate verification is
// the only thing standing between his password and a MITM on the way out.
c, err := tls.DialWithDialer(d, "tcp", addr, &tls.Config{ServerName: host, MinVersion: tls.VersionTLS12})
if err != nil {
return nil, fmt.Errorf("email: dial %s: %w", addr, err)
}
return NewConn(c, timeout)
}
// NewConn wraps an already-open stream (the tests speak IMAP over a pipe) and
// consumes the greeting.
func NewConn(rwc io.ReadWriteCloser, timeout time.Duration) (*Conn, error) {
c := &Conn{rwc: rwc, r: bufio.NewReaderSize(rwc, 64<<10), timeout: timeout}
line, err := c.readLine()
if err != nil {
return nil, fmt.Errorf("email: greeting: %w", err)
}
if !strings.HasPrefix(line, "* OK") && !strings.HasPrefix(line, "* PREAUTH") {
c.rwc.Close()
return nil, fmt.Errorf("email: server refused connection: %s", line)
}
return c, nil
}
func (c *Conn) Close() error { return c.rwc.Close() }
// Login authenticates with LOGIN. The password is passed as an argument and
// never stored on the Conn: nothing in this package keeps a credential alive
// past the command that uses it, so no struct dump or panic trace can carry it.
func (c *Conn) Login(user, pass string) error {
// The command line itself is never logged (see exec) — a LOGIN line IS the
// credential.
if _, err := c.exec(fmt.Sprintf("LOGIN %s %s", quote(user), quote(pass))); err != nil {
return fmt.Errorf("email: login: %w", err)
}
return nil
}
// Select opens a mailbox read-only. EXAMINE, not SELECT: read-only at the
// protocol level means no command in this session can change a flag, expunge a
// message, or move anything, even by mistake.
func (c *Conn) Select(mailbox string) error {
if _, err := c.exec(fmt.Sprintf("EXAMINE %s", quote(mailbox))); err != nil {
return fmt.Errorf("email: examine %s: %w", mailbox, err)
}
return nil
}
// SearchSince returns the UIDs of messages received on or after since. An
// unlimited search is not offered: the first poll against a years-old mailbox
// would otherwise fetch everything and hand a decade of mail to the model.
//
// The IMAP SINCE key has date granularity (and compares the server's internal
// date), so the result can include messages slightly older than since. The
// caller dedupes by UID anyway, so a wider window costs one extra fetch.
func (c *Conn) SearchSince(since time.Time) ([]uint32, error) {
cmd := fmt.Sprintf("UID SEARCH SINCE %s", since.Format("2-Jan-2006"))
lines, err := c.exec(cmd)
if err != nil {
return nil, fmt.Errorf("email: search: %w", err)
}
var uids []uint32
for _, l := range lines {
rest, ok := untagged(l, "SEARCH")
if !ok {
continue
}
for _, f := range strings.Fields(rest) {
n, err := strconv.ParseUint(f, 10, 32)
if err == nil {
uids = append(uids, uint32(n))
}
}
}
return uids, nil
}
var literalSize = regexp.MustCompile(`\{(\d+)\}$`)
// Fetch returns the raw RFC 5322 bytes of one message, by UID.
//
// Returns (nil, nil) when the UID no longer exists — a message he deleted
// between SEARCH and FETCH is normal, not an error.
func (c *Conn) Fetch(uid uint32) ([]byte, error) {
tag := c.nextTag()
if err := c.send(fmt.Sprintf("%s UID FETCH %d (BODY.PEEK[])", tag, uid)); err != nil {
return nil, err
}
var raw []byte
for {
line, err := c.readLine()
if err != nil {
return nil, fmt.Errorf("email: fetch %d: %w", uid, err)
}
if done, err := c.tagged(tag, line); done {
if err != nil {
return nil, fmt.Errorf("email: fetch %d: %w", uid, err)
}
return raw, nil
}
m := literalSize.FindStringSubmatch(strings.TrimSpace(line))
if m == nil {
continue
}
n, err := strconv.Atoi(m[1])
if err != nil {
continue
}
buf := make([]byte, n)
if _, err := io.ReadFull(c.r, buf); err != nil {
return nil, fmt.Errorf("email: fetch %d: literal: %w", uid, err)
}
if raw == nil {
raw = buf
}
}
}
// Logout ends the session politely. A failure is not worth reporting — the
// connection is being closed either way.
func (c *Conn) Logout() {
_, _ = c.exec("LOGOUT")
}
// ---- protocol plumbing -----------------------------------------------------
func (c *Conn) nextTag() string {
c.tag++
return fmt.Sprintf("a%03d", c.tag)
}
// exec sends one command and returns the untagged response lines.
//
// Neither the command nor the response is ever logged here. LOGIN goes through
// this function, and a debug line "sent: a001 LOGIN ..." is how a credential
// ends up in a log file forever.
func (c *Conn) exec(cmd string) ([]string, error) {
tag := c.nextTag()
if err := c.send(tag + " " + cmd); err != nil {
return nil, err
}
var lines []string
for {
line, err := c.readLine()
if err != nil {
return nil, err
}
if done, err := c.tagged(tag, line); done {
return lines, err
}
lines = append(lines, line)
// A response line may carry a literal (e.g. a header FETCH). Nothing we
// send asks for one outside Fetch, but skip it if it appears so the
// stream stays aligned.
if m := literalSize.FindStringSubmatch(strings.TrimSpace(line)); m != nil {
if n, err := strconv.Atoi(m[1]); err == nil {
if _, err := io.CopyN(io.Discard, c.r, int64(n)); err != nil {
return nil, err
}
}
}
}
}
// tagged reports whether line completes the command with this tag, and turns a
// NO/BAD completion into an error. The error text is the server's, which never
// echoes a password.
func (c *Conn) tagged(tag, line string) (bool, error) {
if !strings.HasPrefix(line, tag+" ") {
return false, nil
}
rest := strings.TrimSpace(line[len(tag):])
switch {
case strings.HasPrefix(rest, "OK"):
return true, nil
case strings.HasPrefix(rest, "NO"), strings.HasPrefix(rest, "BAD"):
return true, fmt.Errorf("server said: %s", rest)
default:
return true, fmt.Errorf("unexpected completion: %s", rest)
}
}
func (c *Conn) send(line string) error {
c.setDeadline()
if _, err := io.WriteString(c.rwc, line+"\r\n"); err != nil {
return fmt.Errorf("email: write: %w", err)
}
return nil
}
func (c *Conn) readLine() (string, error) {
c.setDeadline()
line, err := c.r.ReadString('\n')
if err != nil {
return "", err
}
return strings.TrimRight(line, "\r\n"), nil
}
// setDeadline applies the per-connection timeout when the transport supports
// one. A hung IMAP server must not park the poller forever.
func (c *Conn) setDeadline() {
if c.timeout <= 0 {
return
}
if d, ok := c.rwc.(interface{ SetDeadline(time.Time) error }); ok {
_ = d.SetDeadline(time.Now().Add(c.timeout))
}
}
// untagged splits "* SEARCH 1 2 3" into its payload when the key matches.
func untagged(line, key string) (string, bool) {
if !strings.HasPrefix(line, "* ") {
return "", false
}
rest := strings.TrimSpace(line[2:])
if !strings.HasPrefix(rest, key) {
return "", false
}
return strings.TrimSpace(rest[len(key):]), true
}
// quote renders an IMAP quoted string. Passwords routinely contain characters
// that would otherwise end the argument early, and CR/LF are stripped rather
// than escaped because there is no legal way to send them — a credential file
// with a stray newline must not become a second command.
func quote(s string) string {
s = strings.NewReplacer("\r", "", "\n", "").Replace(s)
return `"` + strings.NewReplacer(`\`, `\\`, `"`, `\"`).Replace(s) + `"`
}
+162
View File
@@ -0,0 +1,162 @@
package email
import (
"bufio"
"fmt"
"net"
"strconv"
"strings"
"testing"
"time"
)
// fakeIMAP is a scripted server: enough of IMAP to exercise the client, and
// nothing more. It records the commands it received so a test can assert on the
// protocol (BODY.PEEK rather than BODY, EXAMINE rather than SELECT).
type fakeIMAP struct {
msgs map[uint32]string
uids []uint32
cmds []string
failOn string // substring of a command to answer NO
}
func (f *fakeIMAP) serve(t *testing.T, c net.Conn) {
t.Helper()
defer c.Close()
fmt.Fprint(c, "* OK fake IMAP ready\r\n")
r := bufio.NewReader(c)
for {
line, err := r.ReadString('\n')
if err != nil {
return
}
line = strings.TrimRight(line, "\r\n")
parts := strings.SplitN(line, " ", 2)
if len(parts) != 2 {
return
}
tag, cmd := parts[0], parts[1]
f.cmds = append(f.cmds, cmd)
if f.failOn != "" && strings.Contains(cmd, f.failOn) {
fmt.Fprintf(c, "%s NO computer says no\r\n", tag)
continue
}
upper := strings.ToUpper(cmd)
switch {
case strings.HasPrefix(upper, "LOGIN"), strings.HasPrefix(upper, "EXAMINE"):
fmt.Fprintf(c, "%s OK done\r\n", tag)
case strings.HasPrefix(upper, "UID SEARCH"):
var ids []string
for _, u := range f.uids {
ids = append(ids, strconv.FormatUint(uint64(u), 10))
}
fmt.Fprintf(c, "* SEARCH %s\r\n", strings.Join(ids, " "))
fmt.Fprintf(c, "%s OK search done\r\n", tag)
case strings.HasPrefix(upper, "UID FETCH"):
uid64, _ := strconv.ParseUint(strings.Fields(cmd)[2], 10, 32)
raw, ok := f.msgs[uint32(uid64)]
if ok {
fmt.Fprintf(c, "* 1 FETCH (UID %d BODY[] {%d}\r\n", uid64, len(raw))
fmt.Fprint(c, raw)
fmt.Fprint(c, ")\r\n")
}
fmt.Fprintf(c, "%s OK fetch done\r\n", tag)
case strings.HasPrefix(upper, "LOGOUT"):
fmt.Fprint(c, "* BYE\r\n")
fmt.Fprintf(c, "%s OK bye\r\n", tag)
return
default:
fmt.Fprintf(c, "%s BAD unknown\r\n", tag)
}
}
}
// dialFake wires a client Conn to an in-process server over net.Pipe.
func dialFake(t *testing.T, f *fakeIMAP) *Conn {
t.Helper()
cli, srv := net.Pipe()
go f.serve(t, srv)
c, err := NewConn(cli, 5*time.Second)
if err != nil {
t.Fatalf("greeting: %v", err)
}
t.Cleanup(func() { c.Close() })
return c
}
func TestIMAPRoundTrip(t *testing.T) {
body := "Subject: hello\r\nContent-Type: text/plain; charset=utf-8\r\n\r\nCall the bank.\r\n"
f := &fakeIMAP{uids: []uint32{4, 9}, msgs: map[uint32]string{4: body, 9: body}}
c := dialFake(t, f)
if err := c.Login("kami", `pa"ss\word`); err != nil {
t.Fatalf("login: %v", err)
}
if err := c.Select("INBOX"); err != nil {
t.Fatalf("select: %v", err)
}
uids, err := c.SearchSince(time.Date(2026, 8, 1, 0, 0, 0, 0, time.UTC))
if err != nil {
t.Fatalf("search: %v", err)
}
if len(uids) != 2 || uids[0] != 4 || uids[1] != 9 {
t.Fatalf("uids = %v, want [4 9]", uids)
}
raw, err := c.Fetch(9)
if err != nil {
t.Fatalf("fetch: %v", err)
}
if string(raw) != body {
t.Errorf("fetched %q, want the literal verbatim", raw)
}
c.Logout()
joined := strings.Join(f.cmds, "\n")
// Read-only at the protocol level, and peeking — Maven must leave no trace
// of having read his mail.
if !strings.Contains(joined, "EXAMINE") || strings.Contains(joined, "SELECT ") {
t.Errorf("want EXAMINE (read-only), got:\n%s", joined)
}
if !strings.Contains(joined, "BODY.PEEK[]") {
t.Errorf("want BODY.PEEK, got:\n%s", joined)
}
// The password must have been quoted and escaped, not truncated at the quote.
if !strings.Contains(joined, `"pa\"ss\\word"`) {
t.Errorf("password not quoted correctly:\n%s", joined)
}
// SINCE must carry the IMAP date form.
if !strings.Contains(joined, "SINCE 1-Aug-2026") {
t.Errorf("want a SINCE date, got:\n%s", joined)
}
}
func TestIMAPServerNoIsAnError(t *testing.T) {
f := &fakeIMAP{failOn: "LOGIN"}
c := dialFake(t, f)
err := c.Login("kami", "wrong")
if err == nil {
t.Fatal("a NO completion must be an error")
}
// The error is the server's text; it must not echo the credential.
if strings.Contains(err.Error(), "wrong") {
t.Errorf("error leaks the password: %v", err)
}
}
func TestIMAPFetchMissingUID(t *testing.T) {
f := &fakeIMAP{uids: []uint32{1}, msgs: map[uint32]string{}}
c := dialFake(t, f)
raw, err := c.Fetch(1)
if err != nil {
t.Fatalf("fetch: %v", err)
}
if raw != nil {
t.Errorf("a vanished UID should give nil, got %q", raw)
}
}
func TestQuoteStripsNewlines(t *testing.T) {
if got := quote("pass\r\nA1 LOGOUT"); strings.ContainsAny(got, "\r\n") {
t.Errorf("quote kept a line break: %q", got)
}
}
+80
View File
@@ -0,0 +1,80 @@
package email
import (
"net/mail"
"strings"
)
// The junk filter — the cheapest and most important half of reading mail.
//
// A mailbox is mostly machine-generated: newsletters, receipts nobody acts on,
// social notifications, marketing. Sending all of it to a 1.7B and asking "is
// there a task here" produces confident nonsense at a rate proportional to the
// volume, so junk is decided by HEADERS, before any model sees the message.
//
// The rules are all bulk-mail markers that senders set on themselves, never
// guesses about content:
//
// - List-Unsubscribe / List-Id — by definition a mailing list. If he can
// unsubscribe from it, it is not asking him to do anything.
// - Precedence: bulk|junk|list — the sender declaring itself bulk.
// - Auto-Submitted other than "no" (RFC 3834) — generated by a machine.
// - X-Spam-Flag: YES, X-Spam-Status: Yes — the spam filter upstream already
// decided; we do not second-guess it in the other direction.
// - X-GM-LABELS / X-Gmail-Labels containing a Gmail category — Gmail's own
// Promotions/Social/Forums/Spam classification, when the server sends it.
//
// Deliberately NOT here: sender allow/deny lists and subject keyword matching.
// Both are configuration that ages badly and both would be a place for his
// contacts to end up in a config file. If a real correspondent's mail is being
// dropped, the fix is a rule about a header, not a list of names.
//
// A junk verdict never deletes anything and never touches a flag on the server.
// It means "do not spend the model on this", nothing more.
// junkHeaders — headers whose mere presence marks bulk mail.
var junkPresence = []string{"List-Unsubscribe", "List-Id", "List-Post"}
// gmailCategories — Gmail's category labels, lowercased as they appear in
// X-GM-LABELS. "important" and "inbox" are labels too, and are NOT categories.
// Matching is by these exact tokens (substring is fine — they are namespaced
// and cannot appear in a hand-made label by accident), so a user label named
// "Social Club" is not mistaken for Gmail's Social category.
var gmailCategories = []string{
"category_promotions", "category_social", "category_forums", "category_updates",
`\spam`, `\junk`,
}
// classifyJunk returns whether the message is bulk/automated and why. The
// reason is a short header name, safe to log — it names the marker, never the
// sender or the subject.
func classifyJunk(h mail.Header) (bool, string) {
for _, name := range junkPresence {
if strings.TrimSpace(h.Get(name)) != "" {
return true, strings.ToLower(name)
}
}
switch strings.ToLower(strings.TrimSpace(h.Get("Precedence"))) {
case "bulk", "junk", "list":
return true, "precedence"
}
if v := strings.ToLower(strings.TrimSpace(h.Get("Auto-Submitted"))); v != "" && v != "no" {
return true, "auto-submitted"
}
if strings.EqualFold(strings.TrimSpace(h.Get("X-Spam-Flag")), "yes") {
return true, "x-spam-flag"
}
if v := strings.ToLower(strings.TrimSpace(h.Get("X-Spam-Status"))); strings.HasPrefix(v, "yes") {
return true, "x-spam-status"
}
labels := strings.ToLower(h.Get("X-GM-LABELS") + " " + h.Get("X-Gmail-Labels"))
for _, c := range gmailCategories {
if c == "" {
continue
}
if strings.Contains(labels, c) {
return true, "gmail-category"
}
}
return false, ""
}
+59
View File
@@ -0,0 +1,59 @@
package email
import (
"net/mail"
"strings"
"testing"
)
func headers(t *testing.T, raw string) mail.Header {
t.Helper()
m, err := mail.ReadMessage(strings.NewReader(strings.ReplaceAll(raw, "\n", "\r\n") + "\r\n\r\nbody\r\n"))
if err != nil {
t.Fatalf("read headers: %v", err)
}
return m.Header
}
func TestClassifyJunk(t *testing.T) {
cases := []struct {
name string
raw string
junk bool
reason string
}{
{"personal", "From: a@b.c\nSubject: привет", false, ""},
{"list-unsubscribe", "From: a@b.c\nList-Unsubscribe: <mailto:u@b.c>", true, "list-unsubscribe"},
{"list-id", "From: a@b.c\nList-Id: <golang-nuts.example>", true, "list-id"},
{"precedence bulk", "From: a@b.c\nPrecedence: bulk", true, "precedence"},
{"auto-submitted", "From: a@b.c\nAuto-Submitted: auto-generated", true, "auto-submitted"},
{"auto-submitted no", "From: a@b.c\nAuto-Submitted: no", false, ""},
{"spam flag", "From: a@b.c\nX-Spam-Flag: YES", true, "x-spam-flag"},
{"spam status", "From: a@b.c\nX-Spam-Status: Yes, score=9.1", true, "x-spam-status"},
{"spam status no", "From: a@b.c\nX-Spam-Status: No, score=0.1", false, ""},
{"gmail promo", "From: a@b.c\nX-Gmail-Labels: Inbox,CATEGORY_PROMOTIONS", true, "gmail-category"},
{"user label", "From: a@b.c\nX-Gmail-Labels: Social Club,Important", false, ""},
}
for _, c := range cases {
t.Run(c.name, func(t *testing.T) {
junk, reason := classifyJunk(headers(t, c.raw))
if junk != c.junk || reason != c.reason {
t.Errorf("classifyJunk = (%v, %q), want (%v, %q)", junk, reason, c.junk, c.reason)
}
})
}
}
func TestNewsletterFixtureIsJunk(t *testing.T) {
msg, err := ParseMessage(9, fixture(t, "newsletter.eml"))
if err != nil {
t.Fatalf("parse: %v", err)
}
if !msg.Junk {
t.Fatal("a newsletter with List-Unsubscribe + Precedence: bulk must be junk")
}
// The reason is what gets logged, so it must never carry mail content.
if strings.Contains(msg.JunkReason, "@") || strings.Contains(msg.JunkReason, "Скидки") {
t.Errorf("junk reason leaks content: %q", msg.JunkReason)
}
}
+258
View File
@@ -0,0 +1,258 @@
// Package email is the reading half of the email reader (Vikunja #246,
// docs/plans/01-email-reader.md): a small IMAP client, a MIME-to-plaintext
// converter, and the junk filter that decides a message is not worth reading at
// all. Extraction lives in extract.go and writes nothing itself.
//
// Two constraints shape everything here, both from CLAUDE.md:
//
// - Mail is personal. Nothing in this package logs a body, a subject, or an
// address; callers get the text and decide. Mail text is never search input
// — no function here reaches the network except the IMAP connection itself.
// - Off unless configured. There is no default host, no default account, and
// no fallback that would make a mailbox get read because a field was empty.
//
// The IMAP subset is deliberately tiny (LOGIN, SELECT, UID SEARCH, UID FETCH
// with BODY.PEEK, LOGOUT). No IDLE: a poll every few minutes is what a task
// candidate needs, and IDLE would mean holding a connection and a credential
// open forever for latency nobody is waiting on.
package email
import (
"encoding/base64"
"fmt"
"io"
"mime"
"mime/multipart"
"mime/quotedprintable"
"net/mail"
"regexp"
"strings"
)
// MaxBodyBytes — how much of one message body is kept. A task hides in the
// first screenful; the rest is signature, quoted history and legal boilerplate,
// and it would only spend the resident model's 4096-token context.
const MaxBodyBytes = 4000
// Message — one mail, reduced to the fields extraction and review need.
//
// Raw is deliberately absent: once a message is parsed the original bytes are
// dropped, so no caller can accidentally log or forward the whole mail.
type Message struct {
UID uint32
From string
Subject string
Date string // as sent, unparsed — display only
Body string // plaintext, decoded, HTML-stripped, truncated
// Junk is set by the junk filter (see junk.go). A junk message is carried
// rather than dropped so the poller can count it and still mark it seen.
Junk bool
JunkReason string
}
// ParseMessage turns one RFC 5322 message into a Message.
//
// It never fails on a body it cannot understand: an unparsable or
// unsupported-charset body yields an empty Body and the headers still come
// through, because a subject line alone is often the whole task ("Счёт за
// интернет"). Only a message whose headers cannot be read at all is an error.
func ParseMessage(uid uint32, raw []byte) (Message, error) {
m, err := mail.ReadMessage(strings.NewReader(string(raw)))
if err != nil {
return Message{}, fmt.Errorf("email: parse message: %w", err)
}
msg := Message{
UID: uid,
From: decodeHeader(m.Header.Get("From")),
Subject: decodeHeader(m.Header.Get("Subject")),
Date: m.Header.Get("Date"),
}
msg.Junk, msg.JunkReason = classifyJunk(m.Header)
body, err := plaintextBody(m.Header.Get("Content-Type"), m.Header.Get("Content-Transfer-Encoding"), m.Body)
if err == nil {
msg.Body = truncate(collapse(body), MaxBodyBytes)
}
return msg, nil
}
// plaintextBody walks the MIME tree and returns the best plaintext it can.
//
// Preference order inside a multipart: text/plain first, text/html stripped
// only when there is no plain part. multipart/mixed attachments are skipped
// wholesale — an attachment is a file, not a sentence, and reading one would
// mean parsing arbitrary formats from the network.
func plaintextBody(contentType, encoding string, body io.Reader) (string, error) {
mediaType, params, err := mime.ParseMediaType(contentType)
if contentType == "" || err != nil {
// No Content-Type at all is legal and means text/plain; a broken one is
// treated the same rather than dropping the message.
mediaType, params = "text/plain", nil
}
switch {
case strings.HasPrefix(mediaType, "multipart/"):
boundary := params["boundary"]
if boundary == "" {
return "", fmt.Errorf("email: multipart without boundary")
}
return multipartText(multipart.NewReader(body, boundary))
case mediaType == "text/html":
raw, err := decodeBody(body, encoding, params["charset"])
if err != nil {
return "", err
}
return stripHTML(raw), nil
case mediaType == "text/plain":
return decodeBody(body, encoding, params["charset"])
default:
// A single-part non-text message (a bare PDF, say). No body, headers only.
return "", nil
}
}
// multipartText reads one multipart level, recursing into nested multiparts.
// Returns the plain part if any part yielded one, else the stripped HTML.
func multipartText(mr *multipart.Reader) (string, error) {
var plain, html string
for {
part, err := mr.NextPart()
if err == io.EOF {
break
}
if err != nil {
// A truncated multipart still gives up whatever came before it.
break
}
if part.FileName() != "" {
part.Close()
continue // attachment
}
ct := part.Header.Get("Content-Type")
mediaType, _, _ := mime.ParseMediaType(ct)
text, err := plaintextBody(ct, part.Header.Get("Content-Transfer-Encoding"), part)
part.Close()
if err != nil || strings.TrimSpace(text) == "" {
continue
}
if mediaType == "text/html" && !strings.HasPrefix(mediaType, "multipart/") {
if html == "" {
html = text
}
continue
}
if plain == "" {
plain = text
}
}
if strings.TrimSpace(plain) != "" {
return plain, nil
}
return html, nil
}
// decodeBody applies the transfer encoding, then the charset.
//
// Charset support is UTF-8 (and ASCII, its subset) only, on purpose: x/text's
// encoding tables are not vendored here, and guessing at windows-1251 bytes
// would feed the model mojibake it would happily extract a task from. An
// unsupported charset returns an error, which ParseMessage turns into an empty
// body — subject-only, which is honest.
func decodeBody(r io.Reader, encoding, charset string) (string, error) {
switch strings.ToLower(strings.TrimSpace(encoding)) {
case "quoted-printable":
r = quotedprintable.NewReader(r)
case "base64":
r = newBase64Reader(r)
}
b, err := io.ReadAll(io.LimitReader(r, 1<<20))
if err != nil && len(b) == 0 {
return "", fmt.Errorf("email: read body: %w", err)
}
switch cs := strings.ToLower(strings.TrimSpace(charset)); cs {
case "", "utf-8", "utf8", "us-ascii", "ascii":
return string(b), nil
default:
return "", fmt.Errorf("email: unsupported charset %q", cs)
}
}
// decodeHeader decodes RFC 2047 encoded words ("=?utf-8?B?...?="), which is how
// every Russian subject line arrives. Undecodable headers come back as-is
// rather than empty: a mangled subject is still a hint, and it is only ever
// shown to him as evidence.
func decodeHeader(v string) string {
dec := new(mime.WordDecoder)
out, err := dec.DecodeHeader(v)
if err != nil {
return collapse(v)
}
return collapse(out)
}
var (
scriptStyle = regexp.MustCompile(`(?is)<(script|style)\b[^>]*>.*?</\s*(script|style)\s*>`)
htmlBreak = regexp.MustCompile(`(?i)<\s*(br\s*/?|/p|/div|/tr|/li|/h[1-6])\s*>`)
htmlTag = regexp.MustCompile(`(?s)<[^>]*>`)
htmlComment = regexp.MustCompile(`(?s)<!--.*?-->`)
)
// stripHTML reduces an HTML part to text. A regex stripper, not a parser:
// x/net/html is not vendored, and the consumer is a model reading prose — a
// stray angle bracket costs nothing, whereas a new dependency for the privacy-
// sensitive path costs review.
func stripHTML(s string) string {
s = scriptStyle.ReplaceAllString(s, " ")
s = htmlComment.ReplaceAllString(s, " ")
s = htmlBreak.ReplaceAllString(s, "\n")
s = htmlTag.ReplaceAllString(s, " ")
return unescapeEntities(s)
}
var entities = strings.NewReplacer(
"&nbsp;", " ", "&amp;", "&", "&lt;", "<", "&gt;", ">",
"&quot;", `"`, "&#39;", "'", "&apos;", "'", "&mdash;", "—", "&ndash;", "",
)
func unescapeEntities(s string) string { return entities.Replace(s) }
// collapse squeezes runs of whitespace, keeping single newlines. Mail bodies
// arrive with hard-wrapped lines and blocks of blank space; the model does not
// need them and they are pure context budget.
func collapse(s string) string {
lines := strings.Split(strings.ReplaceAll(s, "\r\n", "\n"), "\n")
var out []string
blank := 0
for _, l := range lines {
l = strings.TrimSpace(strings.Join(strings.Fields(l), " "))
if l == "" {
blank++
if blank > 1 {
continue
}
out = append(out, "")
continue
}
blank = 0
out = append(out, l)
}
return strings.TrimSpace(strings.Join(out, "\n"))
}
// truncate cuts to n bytes on a rune boundary.
func truncate(s string, n int) string {
if len(s) <= n {
return s
}
cut := s[:n]
for len(cut) > 0 && !isRuneStart(cut[len(cut)-1]) {
cut = cut[:len(cut)-1]
}
return strings.TrimSpace(cut) + "…"
}
func isRuneStart(b byte) bool { return b&0xC0 != 0x80 }
// newBase64Reader — base64.NewDecoder already skips the CRLFs mail bodies wrap
// with, so this is only a named seam for decodeBody to read cleanly.
func newBase64Reader(r io.Reader) io.Reader {
return base64.NewDecoder(base64.StdEncoding, r)
}
+111
View File
@@ -0,0 +1,111 @@
package email
import (
"os"
"path/filepath"
"strings"
"testing"
)
func fixture(t *testing.T, name string) []byte {
t.Helper()
b, err := os.ReadFile(filepath.Join("testdata", name))
if err != nil {
t.Fatalf("read fixture %s: %v", name, err)
}
return b
}
func TestParsePlainRussian(t *testing.T) {
msg, err := ParseMessage(7, fixture(t, "plain_ru.eml"))
if err != nil {
t.Fatalf("parse: %v", err)
}
if msg.UID != 7 {
t.Errorf("uid = %d, want 7", msg.UID)
}
if want := "Нужно закрыть задачу"; msg.Subject != want {
t.Errorf("subject = %q, want %q", msg.Subject, want)
}
if !strings.Contains(msg.From, "Антон") {
t.Errorf("from = %q, want the decoded display name", msg.From)
}
if !strings.Contains(msg.Body, "Надо отправить акт до пятницы.") {
t.Errorf("body = %q, want the quoted-printable text decoded", msg.Body)
}
if msg.Junk {
t.Errorf("a personal mail must not be junk (%s)", msg.JunkReason)
}
}
func TestParseHTMLOnlyIsStripped(t *testing.T) {
msg, err := ParseMessage(1, fixture(t, "html_only.eml"))
if err != nil {
t.Fatalf("parse: %v", err)
}
if strings.Contains(msg.Body, "<") || strings.Contains(msg.Body, "color:red") || strings.Contains(msg.Body, "x()") {
t.Errorf("body still has markup/script/style: %q", msg.Body)
}
for _, want := range []string{"Счёт за интернет: 700", "Оплатить до 5 августа."} {
if !strings.Contains(msg.Body, want) {
t.Errorf("body = %q, want it to contain %q", msg.Body, want)
}
}
// &nbsp; must have become a real space, not vanished into the number.
if strings.Contains(msg.Body, "&nbsp;") {
t.Errorf("entity left unescaped: %q", msg.Body)
}
}
func TestParsePrefersPlainAndSkipsAttachments(t *testing.T) {
msg, err := ParseMessage(2, fixture(t, "mixed_attachment.eml"))
if err != nil {
t.Fatalf("parse: %v", err)
}
if got := strings.TrimSpace(msg.Body); got != "Sign the contract before Monday." {
t.Errorf("body = %q, want the text/plain alternative only", got)
}
if strings.Contains(msg.Body, "PDF") {
t.Errorf("attachment bytes leaked into the body: %q", msg.Body)
}
}
// An unsupported charset must degrade to headers-only rather than to mojibake
// the model would then extract a task from.
func TestParseUnsupportedCharsetKeepsHeaders(t *testing.T) {
msg, err := ParseMessage(3, fixture(t, "cp1251.eml"))
if err != nil {
t.Fatalf("parse: %v", err)
}
if msg.Subject != "Legacy" {
t.Errorf("subject = %q, want Legacy", msg.Subject)
}
if msg.Body != "" {
t.Errorf("body = %q, want empty for an undecodable charset", msg.Body)
}
}
func TestParseTruncatesLongBody(t *testing.T) {
var b strings.Builder
b.WriteString("Subject: long\r\nContent-Type: text/plain; charset=utf-8\r\n\r\n")
for i := 0; i < 2000; i++ {
b.WriteString("длинная строка ")
}
msg, err := ParseMessage(4, []byte(b.String()))
if err != nil {
t.Fatalf("parse: %v", err)
}
if len(msg.Body) > MaxBodyBytes+8 {
t.Errorf("body kept %d bytes, want ≤ %d", len(msg.Body), MaxBodyBytes)
}
if !strings.HasSuffix(msg.Body, "…") {
t.Errorf("truncated body should be marked: %q", msg.Body[len(msg.Body)-20:])
}
}
func TestCollapseSqueezesBlankLines(t *testing.T) {
got := collapse(" a b \r\n\r\n\r\n\r\n c \r\n")
if got != "a b\n\nc" {
t.Errorf("collapse = %q, want %q", got, "a b\n\nc")
}
}
+7
View File
@@ -0,0 +1,7 @@
From: legacy@example.org
To: kami@example.org
Subject: Legacy
Date: Fri, 01 Aug 2026 05:00:00 +0400
Content-Type: text/plain; charset="windows-1251"
Ï
+13
View File
@@ -0,0 +1,13 @@
From: billing@isp.example
To: kami@example.org
Subject: =?utf-8?B?0KHRh9GR0YIg0LfQsCDQuNC90YLQtdGA0L3QtdGC?=
Date: Fri, 01 Aug 2026 08:00:00 +0400
MIME-Version: 1.0
Content-Type: multipart/alternative; boundary="B1"
--B1
Content-Type: text/html; charset="utf-8"
Content-Transfer-Encoding: base64
PGh0bWw+PGhlYWQ+PHN0eWxlPnB7Y29sb3I6cmVkfTwvc3R5bGU+PC9oZWFkPjxib2R5PjxwPtCh0YfRkdGCINC30LAg0LjQvdGC0LXRgNC90LXRgjogNzAwJm5ic3A74oK9PC9wPjxwPtCe0L/Qu9Cw0YLQuNGC0Ywg0LTQviA1INCw0LLQs9GD0YHRgtCwLjwvcD48c2NyaXB0PngoKTwvc2NyaXB0PjwvYm9keT48L2h0bWw+
--B1--
+26
View File
@@ -0,0 +1,26 @@
From: hr@work.example
To: kami@example.org
Subject: Contract
Date: Fri, 01 Aug 2026 07:00:00 +0400
MIME-Version: 1.0
Content-Type: multipart/mixed; boundary="M1"
--M1
Content-Type: multipart/alternative; boundary="A1"
--A1
Content-Type: text/plain; charset="utf-8"
Sign the contract before Monday.
--A1
Content-Type: text/html; charset="utf-8"
<p>Sign the contract before Monday.</p>
--A1--
--M1
Content-Type: application/pdf; name="contract.pdf"
Content-Disposition: attachment; filename="contract.pdf"
Content-Transfer-Encoding: base64
JVBERi0xLjQgbm90IHJlYWxseSBhIHBkZg==
--M1--
+9
View File
@@ -0,0 +1,9 @@
From: news@shop.example
To: kami@example.org
Subject: =?utf-8?B?0KHQutC40LTQutC4INGC0L7Qu9GM0LrQviDRgdC10LPQvtC00L3Rjw==?=
Date: Fri, 01 Aug 2026 06:00:00 +0400
List-Unsubscribe: <mailto:unsub@shop.example>
Precedence: bulk
Content-Type: text/plain; charset="utf-8"
Sale!
+14
View File
@@ -0,0 +1,14 @@
From: =?utf-8?B?0JDQvdGC0L7QvQ==?= <anton@example.org>
To: kami@example.org
Subject: =?utf-8?B?0J3Rg9C20L3QviDQt9Cw0LrRgNGL0YLRjCDQt9Cw0LTQsNGH0YM=?=
Date: Fri, 01 Aug 2026 09:12:00 +0400
Content-Type: text/plain; charset="utf-8"
Content-Transfer-Encoding: quoted-printable
Message-ID: <plain-ru@example.org>
=D0=9F=D1=80=D0=B8=D0=B2=D0=B5=D1=82! =D0=9D=D0=B0=D0=B4=D0=BE =D0=BE=D1=82=
=D0=BF=D1=80=D0=B0=D0=B2=D0=B8=D1=82=D1=8C =D0=B0=D0=BA=D1=82 =D0=B4=D0=BE =
=D0=BF=D1=8F=D1=82=D0=BD=D0=B8=D1=86=D1=8B.
--
Anton
+71
View File
@@ -0,0 +1,71 @@
package router
import "strings"
// Money questions, matched deterministically (Vikunja #125).
//
// No new intent, for the same reason as tasks: the intent enum is a contract
// with the relabelling prompt. "сколько я потратил?" is a query; which figure
// it asks for is a lookup, not something to ask a 1.7B — and a model asked to
// invent a spending total will happily do it.
// MoneyWindow — which period a money question asks about.
type MoneyWindow int
const (
MoneyNone MoneyWindow = iota
MoneyToday
MoneyMonth
)
// moneyNouns — the words that make a question be about his money.
var moneyNouns = []string{
"потратил", "потратила", "тратил", "траты", "трат", "расходы", "расходов",
"заработал", "потрачено", "денег", "spend", "spent", "expenses",
}
// ParseMoneyQuery reports whether an utterance asks about spending or income,
// and over which window. Defaults to the month: "сколько я потратил?" without a
// period is the month-to-date question, which is the one worth answering.
//
// Narrow on purpose. A money noun alone is not enough — "я потратил весь день
// на это" is him talking about his day, so an amount word or an explicit
// question word has to be there too.
func ParseMoneyQuery(text string) (MoneyWindow, bool) {
toks := planTokens(text)
if len(toks) == 0 {
return MoneyNone, false
}
hasNoun := false
for _, t := range toks {
for _, n := range moneyNouns {
if t == n {
hasNoun = true
}
}
}
if !hasNoun {
return MoneyNone, false
}
// "весь день", "время", "силы" — spending that is not money.
for _, t := range toks {
switch t {
case "день", "дня", "время", "времени", "силы", "сил", "нервы":
return MoneyNone, false
}
}
asking := hasTok(toks, "сколько") || hasTok(toks, "какие") || hasTok(toks, "покажи") ||
hasTok(toks, "how") || hasTok(toks, "much") || hasTok(toks, "my") ||
hasTok(toks, "мои") || hasTok(toks, "траты") || hasTok(toks, "расходы")
if !asking {
return MoneyNone, false
}
lower := strings.ToLower(text)
switch {
case hasTok(toks, "сегодня") || strings.Contains(lower, "today"):
return MoneyToday, true
case hasTok(toks, "месяц") || hasTok(toks, "месяце") || strings.Contains(lower, "month"):
return MoneyMonth, true
}
return MoneyMonth, true
}
+31
View File
@@ -0,0 +1,31 @@
package router
import "testing"
func TestParseMoneyQuery(t *testing.T) {
cases := []struct {
in string
window MoneyWindow
ok bool
}{
{"сколько я потратил сегодня?", MoneyToday, true},
{"сколько я потратил в этом месяце?", MoneyMonth, true},
{"сколько я потратил?", MoneyMonth, true}, // month-to-date by default
{"покажи мои траты", MoneyMonth, true},
{"какие у меня расходы за месяц", MoneyMonth, true},
{"how much did I spend today", MoneyToday, true},
{"сколько я заработал в этом месяце", MoneyMonth, true},
// Not about money.
{"я потратил весь день на это", MoneyNone, false},
{"потратил много сил", MoneyNone, false},
{"какая погода?", MoneyNone, false},
{"я купил молоко", MoneyNone, false},
{"", MoneyNone, false},
}
for _, c := range cases {
w, ok := ParseMoneyQuery(c.in)
if ok != c.ok || w != c.window {
t.Errorf("ParseMoneyQuery(%q) = (%v, %v), want (%v, %v)", c.in, w, ok, c.window, c.ok)
}
}
}
+267
View File
@@ -0,0 +1,267 @@
// Package zenmoney reads spending and income from ZenMoney's /v8/diff/ API
// (Vikunja #125).
//
// Trust boundary: ZenMoney, not Maven. They already hold his bank sessions —
// this package only reads back what they have, over a token that lives in the
// poller module and is never handed to core. Nothing here writes to ZenMoney;
// diff is called read-only (an empty change set in, a change set out).
//
// Two rules the code exists to enforce:
//
// - NEVER invent a number. Every figure in a Summary is a sum of amounts the
// API returned. A request that fails, or returns nothing, produces no
// summary and therefore no fact — silence, not a zero. A confidently wrong
// "ты потратил 0" is worse than no answer.
// - His money is never search input. This package holds no notes, no
// utterances and no persona text, and it has no path to the external search
// capability. The only thing that leaves the box here is the diff request
// itself, to the service that already has the data.
package zenmoney
import (
"bytes"
"context"
"encoding/json"
"fmt"
"io"
"net/http"
"sort"
"strings"
"time"
)
// DefaultBaseURL — ZenMoney's API root. Overridable so the tests can point at
// an httptest server replaying a recorded response.
const DefaultBaseURL = "https://api.zenmoney.ru"
// Client is a ZenMoney diff reader. The token is held here, in the poller's
// address space; core never receives it and never learns it exists.
type Client struct {
BaseURL string
Token string
HTTP *http.Client
}
// New returns a client with a bounded HTTP timeout. An empty token is a
// programming error the caller must catch — the capability is off unless
// configured, so a client is only ever built when a token was supplied.
func New(token, baseURL string, timeout time.Duration) (*Client, error) {
if strings.TrimSpace(token) == "" {
return nil, fmt.Errorf("zenmoney: empty token")
}
if baseURL == "" {
baseURL = DefaultBaseURL
}
if timeout <= 0 {
timeout = 20 * time.Second
}
return &Client{
BaseURL: strings.TrimRight(baseURL, "/"),
Token: token,
HTTP: &http.Client{Timeout: timeout},
}, nil
}
// diffRequest — the smallest body /v8/diff/ accepts. serverTimestamp is the
// incremental cursor: the server returns objects changed at or after it.
type diffRequest struct {
CurrentClientTimestamp int64 `json:"currentClientTimestamp"`
ServerTimestamp int64 `json:"serverTimestamp"`
}
// diffResponse — only the fields spending needs. ZenMoney returns a dozen more
// object types (tags, merchants, budgets, reminders); decoding them would mean
// holding more of his financial life in memory than the question needs.
type diffResponse struct {
ServerTimestamp int64 `json:"serverTimestamp"`
Instrument []instrument `json:"instrument"`
Transaction []transaction `json:"transaction"`
}
type instrument struct {
ID int64 `json:"id"`
ShortTitle string `json:"shortTitle"`
}
type transaction struct {
ID string `json:"id"`
Date string `json:"date"` // "2026-07-15"
Deleted bool `json:"deleted"`
Income float64 `json:"income"`
Outcome float64 `json:"outcome"`
IncomeInstrument int64 `json:"incomeInstrument"`
OutcomeInstrmnt int64 `json:"outcomeInstrument"`
IncomeAccount string `json:"incomeAccount"`
OutcomeAccount string `json:"outcomeAccount"`
}
// Money — an amount in one currency. Kept as the currency's own short title
// ("RUB", "EUR") rather than converted: ZenMoney's rates are a snapshot, and
// converting would turn a figure he can check against his bank into one he
// cannot.
type Money struct {
Currency string `json:"currency"`
Amount float64 `json:"amount"`
}
// Summary — what was spent and earned over a window, per currency, plus how
// many transactions it was computed from. Count is the honesty check: a
// summary built from zero transactions is not "you spent nothing", it is "there
// was nothing to read", and callers treat it as no answer.
type Summary struct {
From, To time.Time
Spent []Money `json:"spent"`
Earned []Money `json:"earned"`
Count int `json:"count"`
// ServerTimestamp — the cursor the API returned, for the caller to log or
// carry. Not used as an incremental cursor for summaries; see Since.
ServerTimestamp int64 `json:"-"`
}
// Since returns the summary of transactions dated in [from, to).
//
// The diff cursor is set to `from` so the server only sends objects changed
// since then, which for a "this month" window is everything filed this month.
// The caveat, deliberately accepted: a transaction he EDITED this month but
// dated last month arrives too, and is then excluded by date — so editing old
// records cannot inflate this month's total. The reverse case (a transaction
// dated this month, filed and last changed before `from`) cannot exist.
func (c *Client) Since(ctx context.Context, from, to time.Time) (Summary, error) {
resp, err := c.diff(ctx, from.Unix())
if err != nil {
return Summary{}, err
}
return summarize(resp, from, to), nil
}
func (c *Client) diff(ctx context.Context, serverTimestamp int64) (diffResponse, error) {
body, err := json.Marshal(diffRequest{
CurrentClientTimestamp: time.Now().Unix(),
ServerTimestamp: serverTimestamp,
})
if err != nil {
return diffResponse{}, err
}
req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.BaseURL+"/v8/diff/", bytes.NewReader(body))
if err != nil {
return diffResponse{}, err
}
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", "Bearer "+c.Token)
hc := c.HTTP
if hc == nil {
hc = &http.Client{Timeout: 20 * time.Second}
}
res, err := hc.Do(req)
if err != nil {
return diffResponse{}, err
}
defer res.Body.Close()
raw, err := io.ReadAll(io.LimitReader(res.Body, 32<<20))
if err != nil {
return diffResponse{}, err
}
if res.StatusCode != http.StatusOK {
// The status only. The body of a failed diff can echo account data, and
// this string reaches the log.
return diffResponse{}, fmt.Errorf("zenmoney diff: %s", res.Status)
}
var out diffResponse
if err := json.Unmarshal(raw, &out); err != nil {
return diffResponse{}, fmt.Errorf("zenmoney diff: decode: %w", err)
}
return out, nil
}
// summarize sums the transactions dated inside the window.
//
// Excluded, in order: deleted rows (ZenMoney tombstones rather than removes),
// transfers and currency exchanges (income and outcome both non-zero — moving
// his own money between his own accounts is not spending), and anything dated
// outside the window.
func summarize(resp diffResponse, from, to time.Time) Summary {
cur := map[int64]string{}
for _, in := range resp.Instrument {
cur[in.ID] = in.ShortTitle
}
spent := map[string]float64{}
earned := map[string]float64{}
count := 0
for _, t := range resp.Transaction {
if t.Deleted {
continue
}
d, err := time.ParseInLocation("2006-01-02", t.Date, from.Location())
if err != nil {
continue // an undated row is not a number we can place
}
if d.Before(from) || !d.Before(to) {
continue
}
if t.Income > 0 && t.Outcome > 0 {
continue // transfer / exchange
}
switch {
case t.Outcome > 0:
spent[currency(cur, t.OutcomeInstrmnt)] += t.Outcome
count++
case t.Income > 0:
earned[currency(cur, t.IncomeInstrument)] += t.Income
count++
}
}
return Summary{
From: from, To: to,
Spent: sortMoney(spent), Earned: sortMoney(earned),
Count: count, ServerTimestamp: resp.ServerTimestamp,
}
}
// currency names the instrument, or says it does not know. An unknown id keeps
// the amount rather than dropping it: a sum without a currency label is still
// his money, and silently discarding it would understate the total.
func currency(names map[int64]string, id int64) string {
if s := names[id]; s != "" {
return s
}
return "?"
}
// sortMoney gives the amounts a stable order (largest first) so the rendered
// string and the written fact do not churn between polls.
func sortMoney(m map[string]float64) []Money {
out := make([]Money, 0, len(m))
for c, a := range m {
out = append(out, Money{Currency: c, Amount: a})
}
sort.Slice(out, func(i, j int) bool {
if out[i].Amount != out[j].Amount {
return out[i].Amount > out[j].Amount
}
return out[i].Currency < out[j].Currency
})
return out
}
// Empty reports whether the summary rests on no transactions at all. Callers
// must treat an empty summary as "nothing to say", never as a zero: the
// difference between "he spent nothing" and "the read returned nothing" is the
// difference between an answer and an invented one.
func (s Summary) Empty() bool { return s.Count == 0 }
// MonthWindow — the first instant of now's month, and now's own day-end
// exclusive bound, in now's location. The window a "сколько я потратил в этом
// месяце?" question means.
func MonthWindow(now time.Time) (from, to time.Time) {
loc := now.Location()
from = time.Date(now.Year(), now.Month(), 1, 0, 0, 0, 0, loc)
to = time.Date(now.Year(), now.Month(), now.Day(), 0, 0, 0, 0, loc).AddDate(0, 0, 1)
return from, to
}
// DayWindow — today, in now's location.
func DayWindow(now time.Time) (from, to time.Time) {
loc := now.Location()
from = time.Date(now.Year(), now.Month(), now.Day(), 0, 0, 0, 0, loc)
return from, from.AddDate(0, 0, 1)
}
+178
View File
@@ -0,0 +1,178 @@
package zenmoney
import (
"context"
"encoding/json"
"net/http"
"net/http/httptest"
"os"
"strings"
"testing"
"time"
)
// fixtureServer replays testdata/diff.json and records the request, so the
// tests can assert the wire contract (Bearer token, POST, /v8/diff/) without a
// ZenMoney account.
func fixtureServer(t *testing.T, got *diffRequest, auth *string) *httptest.Server {
t.Helper()
body, err := os.ReadFile("testdata/diff.json")
if err != nil {
t.Fatal(err)
}
return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
t.Errorf("method = %s, want POST", r.Method)
}
if r.URL.Path != "/v8/diff/" {
t.Errorf("path = %s, want /v8/diff/", r.URL.Path)
}
if auth != nil {
*auth = r.Header.Get("Authorization")
}
if got != nil {
if err := json.NewDecoder(r.Body).Decode(got); err != nil {
t.Errorf("decode request: %v", err)
}
}
w.Header().Set("Content-Type", "application/json")
w.Write(body)
}))
}
func aug(day int) time.Time { return time.Date(2026, 8, day, 0, 0, 0, 0, time.UTC) }
func TestSinceSumsSpendingPerCurrency(t *testing.T) {
var req diffRequest
var auth string
srv := fixtureServer(t, &req, &auth)
defer srv.Close()
c, err := New("tok", srv.URL, time.Second)
if err != nil {
t.Fatal(err)
}
s, err := c.Since(context.Background(), aug(1), aug(6))
if err != nil {
t.Fatal(err)
}
if auth != "Bearer tok" {
t.Errorf("Authorization = %q", auth)
}
if req.ServerTimestamp != aug(1).Unix() {
t.Errorf("serverTimestamp = %d, want the window start", req.ServerTimestamp)
}
// 1500 + 249.5 RUB spent, 12 EUR spent, 3000 RUB in. The transfer (t4), the
// deleted row (t6) and July's salary (t3) are all excluded.
want := map[string]float64{"RUB": 1749.5, "EUR": 12}
if len(s.Spent) != 2 {
t.Fatalf("spent = %+v, want two currencies", s.Spent)
}
for _, m := range s.Spent {
if want[m.Currency] != m.Amount {
t.Errorf("spent %s = %v, want %v", m.Currency, m.Amount, want[m.Currency])
}
}
if len(s.Earned) != 1 || s.Earned[0].Amount != 3000 || s.Earned[0].Currency != "RUB" {
t.Errorf("earned = %+v, want 3000 RUB (July's salary is outside the window)", s.Earned)
}
if s.Count != 4 {
t.Errorf("count = %d, want 4 counted transactions", s.Count)
}
// Largest first, so the fact value does not churn between polls.
if s.Spent[0].Currency != "RUB" {
t.Errorf("spent order = %+v, want the largest amount first", s.Spent)
}
}
// A window with nothing in it is NOT a zero. No transactions means no answer,
// and the caller must be able to tell the difference.
func TestSinceEmptyWindowIsNotAZero(t *testing.T) {
srv := fixtureServer(t, nil, nil)
defer srv.Close()
c, _ := New("tok", srv.URL, time.Second)
s, err := c.Since(context.Background(), time.Date(2026, 9, 1, 0, 0, 0, 0, time.UTC), time.Date(2026, 9, 30, 0, 0, 0, 0, time.UTC))
if err != nil {
t.Fatal(err)
}
if !s.Empty() {
t.Fatalf("summary = %+v, want empty", s)
}
if _, ok := s.Value(); ok {
t.Error("an empty summary must not produce a fact value")
}
}
func TestSinceReportsHTTPFailureWithoutTheBody(t *testing.T) {
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
http.Error(w, `{"account":"acc-card","secret":"leaky"}`, http.StatusUnauthorized)
}))
defer srv.Close()
c, _ := New("tok", srv.URL, time.Second)
_, err := c.Since(context.Background(), aug(1), aug(6))
if err == nil {
t.Fatal("want an error on 401")
}
if strings.Contains(err.Error(), "acc-card") || strings.Contains(err.Error(), "leaky") {
t.Errorf("error %q echoes the response body — it reaches the log", err)
}
}
func TestNewRequiresAToken(t *testing.T) {
if _, err := New(" ", "", 0); err == nil {
t.Error("want an error for an empty token — the capability is off unless configured")
}
}
func TestMonthAndDayWindows(t *testing.T) {
now := time.Date(2026, 8, 15, 21, 30, 0, 0, time.UTC)
from, to := MonthWindow(now)
if from != time.Date(2026, 8, 1, 0, 0, 0, 0, time.UTC) || to != time.Date(2026, 8, 16, 0, 0, 0, 0, time.UTC) {
t.Errorf("month window = %v..%v", from, to)
}
from, to = DayWindow(now)
if from != time.Date(2026, 8, 15, 0, 0, 0, 0, time.UTC) || to != time.Date(2026, 8, 16, 0, 0, 0, 0, time.UTC) {
t.Errorf("day window = %v..%v", from, to)
}
}
func TestFactValueRoundTripAndFormat(t *testing.T) {
s := Summary{Spent: []Money{{"RUB", 1749.5}}, Earned: []Money{{"RUB", 3000}}, Count: 3}
raw, ok := s.Value()
if !ok {
t.Fatal("want a fact value")
}
v, err := ParseFactValue(raw)
if err != nil {
t.Fatal(err)
}
got := v.FormatRU("в этом месяце")
if !strings.Contains(got, "1749.5 RUB") || !strings.Contains(got, "3000 RUB") {
t.Errorf("reply = %q, want the exact figures", got)
}
// Persona: informal, feminine, no commentary on his spending.
for _, bad := range []string{"вы", "ваш", "милый", "дорогой", "рад ", "слишком", "много"} {
if strings.Contains(got, bad) {
t.Errorf("reply %q contains %q", got, bad)
}
}
if strings.Contains(got, "он ") {
t.Errorf("reply %q talks about him in the third person", got)
}
}
// An empty fact value renders to nothing, so a caller cannot accidentally
// speak a zero.
func TestFormatRUEmptyRendersNothing(t *testing.T) {
if got := (FactValue{}).FormatRU("сегодня"); got != "" {
t.Errorf("reply = %q, want empty", got)
}
}
func TestFormatAmountKeepsTheTruth(t *testing.T) {
for in, want := range map[float64]string{1500: "1500", 249.5: "249.5", 0.99: "0.99", 1749.55: "1749.55"} {
if got := formatAmount(in); got != want {
t.Errorf("formatAmount(%v) = %q, want %q", in, got, want)
}
}
}
+99
View File
@@ -0,0 +1,99 @@
package zenmoney
import (
"encoding/json"
"fmt"
"strings"
"time"
)
// Fact keys the poller writes, all under source "poll:zenmoney". Two windows,
// because they are the two questions he actually asks; a per-category
// breakdown would mean storing what he bought, and the store is not a ledger.
const (
KeySpentToday = "money_today"
KeySpentMonth = "money_month"
)
// Source — the provenance every money fact carries. The loop's rules trust
// source, and nothing in Maven has a rule on these keys: they are read when he
// asks, never a reason to speak. Maven is not a nag, least of all about money.
const Source = "poll:zenmoney"
// FactValue — the JSON stored in a money fact. A wire shape of its own rather
// than the Summary struct so From/To (which carry a timezone and a clock) stay
// out of the store; the key already says which window it is.
type FactValue struct {
Spent []Money `json:"spent"`
Earned []Money `json:"earned"`
Count int `json:"count"`
}
// Value encodes the summary for the facts table. Returns ok=false for an empty
// summary: no transactions read means no fact written, so that a failed or
// empty poll can never be recited back to him as a zero.
func (s Summary) Value() (string, bool) {
if s.Empty() {
return "", false
}
b, err := json.Marshal(FactValue{Spent: s.Spent, Earned: s.Earned, Count: s.Count})
if err != nil {
return "", false
}
return string(b), true
}
// ParseFactValue decodes a stored money fact.
func ParseFactValue(raw string) (FactValue, error) {
var v FactValue
if err := json.Unmarshal([]byte(raw), &v); err != nil {
return FactValue{}, err
}
return v, nil
}
// FormatRU renders a money fact the way Maven says it — feminine, informal,
// and only about numbers that came from ZenMoney. window is the Russian phrase
// for the period ("сегодня", "в этом месяце").
//
// No commentary. She reports the figure and stops: an opinion about his
// spending is exactly the nagging Maven is not for.
func (v FactValue) FormatRU(window string) string {
if v.Count == 0 {
return ""
}
var parts []string
if len(v.Spent) > 0 {
parts = append(parts, "потратил "+joinMoney(v.Spent))
}
if len(v.Earned) > 0 {
parts = append(parts, "получил "+joinMoney(v.Earned))
}
if len(parts) == 0 {
return ""
}
return window + " ты " + strings.Join(parts, ", ") + "."
}
func joinMoney(ms []Money) string {
parts := make([]string, 0, len(ms))
for _, m := range ms {
parts = append(parts, fmt.Sprintf("%s %s", formatAmount(m.Amount), m.Currency))
}
return strings.Join(parts, " и ")
}
// formatAmount — whole units when the amount is whole, two decimals otherwise.
// Never rounded to something prettier than the truth.
func formatAmount(a float64) string {
if a == float64(int64(a)) {
return fmt.Sprintf("%d", int64(a))
}
return strings.TrimRight(strings.TrimRight(fmt.Sprintf("%.2f", a), "0"), ".")
}
// StaleAfter — how old a money fact may be and still be worth reciting. The
// poller is off unless configured and can be down; answering with last week's
// total as if it were today's would be a lie by omission, so a stale fact is
// reported as stale.
const StaleAfter = 26 * time.Hour
+35
View File
@@ -0,0 +1,35 @@
{
"serverTimestamp": 1785312000,
"instrument": [
{"id": 2, "title": "Российский рубль", "shortTitle": "RUB", "symbol": "₽", "rate": 1},
{"id": 3, "title": "Евро", "shortTitle": "EUR", "symbol": "€", "rate": 100}
],
"account": [
{"id": "acc-card", "title": "карта", "instrument": 2},
{"id": "acc-cash", "title": "наличные", "instrument": 2},
{"id": "acc-eur", "title": "евро", "instrument": 3}
],
"transaction": [
{"id": "t1", "date": "2026-08-01", "changed": 1785300000, "income": 0, "outcome": 1500,
"incomeInstrument": 2, "outcomeInstrument": 2, "incomeAccount": "acc-card", "outcomeAccount": "acc-card",
"payee": "пятёрочка", "deleted": false},
{"id": "t2", "date": "2026-08-01", "changed": 1785300001, "income": 0, "outcome": 249.5,
"incomeInstrument": 2, "outcomeInstrument": 2, "incomeAccount": "acc-card", "outcomeAccount": "acc-card",
"payee": "метро", "deleted": false},
{"id": "t3", "date": "2026-07-20", "changed": 1785300002, "income": 120000, "outcome": 0,
"incomeInstrument": 2, "outcomeInstrument": 2, "incomeAccount": "acc-card", "outcomeAccount": "acc-card",
"payee": "зарплата", "deleted": false},
{"id": "t4", "date": "2026-08-02", "changed": 1785300003, "income": 5000, "outcome": 5000,
"incomeInstrument": 2, "outcomeInstrument": 2, "incomeAccount": "acc-cash", "outcomeAccount": "acc-card",
"payee": "", "comment": "снял наличные", "deleted": false},
{"id": "t5", "date": "2026-08-03", "changed": 1785300004, "income": 0, "outcome": 12,
"incomeInstrument": 3, "outcomeInstrument": 3, "incomeAccount": "acc-eur", "outcomeAccount": "acc-eur",
"payee": "hosting", "deleted": false},
{"id": "t6", "date": "2026-08-04", "changed": 1785300005, "income": 0, "outcome": 999,
"incomeInstrument": 2, "outcomeInstrument": 2, "incomeAccount": "acc-card", "outcomeAccount": "acc-card",
"payee": "удалённая", "deleted": true},
{"id": "t7", "date": "2026-08-05", "changed": 1785300006, "income": 3000, "outcome": 0,
"incomeInstrument": 2, "outcomeInstrument": 2, "incomeAccount": "acc-card", "outcomeAccount": "acc-card",
"payee": "возврат", "deleted": false}
]
}