Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| b4646155b4 | |||
| da647e87d0 |
@@ -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/
|
||||||
|
|||||||
@@ -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
|
||||||
|
}
|
||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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
@@ -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 {
|
||||||
|
|||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -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:
|
||||||
|
|||||||
@@ -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
|
||||||
|
}
|
||||||
@@ -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")
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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) + `"`
|
||||||
|
}
|
||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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, ""
|
||||||
|
}
|
||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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(
|
||||||
|
" ", " ", "&", "&", "<", "<", ">", ">",
|
||||||
|
""", `"`, "'", "'", "'", "'", "—", "—", "–", "–",
|
||||||
|
)
|
||||||
|
|
||||||
|
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)
|
||||||
|
}
|
||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
// must have become a real space, not vanished into the number.
|
||||||
|
if strings.Contains(msg.Body, " ") {
|
||||||
|
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")
|
||||||
|
}
|
||||||
|
}
|
||||||
Vendored
+7
@@ -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
@@ -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
@@ -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
@@ -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!
|
||||||
Vendored
+14
@@ -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
|
||||||
@@ -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
|
||||||
|
}
|
||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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)
|
||||||
|
}
|
||||||
@@ -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)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -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
|
||||||
Vendored
+35
@@ -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}
|
||||||
|
]
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user