Version, authenticate and fully trace ecosystem calls #84
+90
-17
@@ -61,6 +61,11 @@ func run(args []string) error {
|
||||
mailbox := fs.String("mailbox", "INBOX", "mailbox to read, read-only")
|
||||
interval := fs.Duration("interval", 15*time.Minute, "how often to read the mailbox")
|
||||
lookback := fs.Duration("lookback", 72*time.Hour, "how far back to search on each poll")
|
||||
// -max and -interval are one decision, not two. Every non-bulk message in a
|
||||
// poll is one serialized llama-server call on core's side, and core gates
|
||||
// mail extraction behind voice turns (llm.Gate), so a large batch does not
|
||||
// mute Maven, it just takes a while. Raise -max only alongside whatever
|
||||
// bound core is running.
|
||||
max := fs.Int("max", 25, "most messages to fetch in one poll")
|
||||
timeout := fs.Duration("timeout", 30*time.Second, "IMAP network timeout")
|
||||
statePath := fs.String("state", "", "file remembering which UIDs were read (default: none — every poll re-reads the window)")
|
||||
@@ -125,11 +130,18 @@ func run(args []string) error {
|
||||
log.Printf("mavmaild: bye")
|
||||
return nil
|
||||
case <-t.C:
|
||||
// Core told us mail ingestion is not configured. Nothing will change
|
||||
// without a core restart, and a restart restarts us too, so the
|
||||
// daemon stays up and does nothing at all.
|
||||
//
|
||||
// It does NOT exit. The compose service inherits restart:
|
||||
// unless-stopped, which restarts a clean exit as readily as a crash,
|
||||
// so exiting here produced a loop: log in to IMAP, get refused by
|
||||
// core, exit, restart, log in again. Four IMAP logins an hour
|
||||
// against a mailbox that has nothing to give, and Gmail and Yandex
|
||||
// both rate-limit exactly that.
|
||||
if r.disabled {
|
||||
// Core told us mail ingestion is not configured. Nothing will change
|
||||
// without a core restart, and a restart restarts us too.
|
||||
log.Printf("mavmaild: core does not accept mail — idling")
|
||||
return nil
|
||||
continue
|
||||
}
|
||||
r.pollOnce(ctx, password)
|
||||
}
|
||||
@@ -153,10 +165,17 @@ type reader struct {
|
||||
timeout time.Duration
|
||||
state *seenState
|
||||
|
||||
// dial — connection seam for the tests; nil ⇒ implicit TLS.
|
||||
dial func(addr string, timeout time.Duration) (*email.Conn, error)
|
||||
// fetchMail — the read seam, nil ⇒ the real IMAP read. The tests replace
|
||||
// the whole read rather than the transport: internal/email keeps its dialer
|
||||
// unexported so that no code outside that package can point the reader at a
|
||||
// cleartext socket and hand it the password, and this daemon is code
|
||||
// outside that package.
|
||||
fetchMail func(password string) ([]email.Message, error)
|
||||
|
||||
// disabled — core answered ErrUnknownMethod, i.e. it has no email block.
|
||||
// Written in pollOnce and read in the ticker loop, both on the one
|
||||
// goroutine that run() drives, so it needs no atomic. If a second caller of
|
||||
// pollOnce ever appears, this becomes a race and has to change.
|
||||
disabled bool
|
||||
}
|
||||
|
||||
@@ -179,23 +198,25 @@ func (r *reader) pollOnce(ctx context.Context, password string) {
|
||||
if ctx.Err() != nil {
|
||||
return
|
||||
}
|
||||
if m.Junk {
|
||||
junk++
|
||||
// Marked seen without a model call: the header filter already decided,
|
||||
// and re-classifying it every quarter hour would be pure waste.
|
||||
r.state.mark(m.UID)
|
||||
continue
|
||||
}
|
||||
resp, err := r.core.IngestMail(ctx, ipc.IngestMailReq{
|
||||
req := ipc.IngestMailReq{
|
||||
Mailbox: r.mailbox,
|
||||
UID: m.UID,
|
||||
From: m.From,
|
||||
Subject: m.Subject,
|
||||
Date: m.Date,
|
||||
Body: m.Body,
|
||||
})
|
||||
}
|
||||
if m.Junk {
|
||||
junk++
|
||||
// Core is TOLD, which is what its wire doc says: it counts the bulk
|
||||
// message and answers Skipped without spending the model. The header
|
||||
// filter already decided, so no content is sent with the verdict —
|
||||
// nothing will read it.
|
||||
req = ipc.IngestMailReq{Mailbox: r.mailbox, UID: m.UID, Junk: true}
|
||||
}
|
||||
resp, err := r.core.IngestMail(ctx, req)
|
||||
if errors.Is(err, ipc.ErrUnknownMethod) {
|
||||
log.Printf("mavmaild: core has no email block configured — mail ingestion is off; stopping")
|
||||
log.Printf("mavmaild: core has no email block configured — mail ingestion is off; idling until a restart")
|
||||
r.disabled = true
|
||||
return
|
||||
}
|
||||
@@ -217,6 +238,9 @@ func (r *reader) pollOnce(ctx context.Context, password string) {
|
||||
// fetch reads the mailbox. Messages already in the seen-set are not fetched at
|
||||
// all, so a steady mailbox costs one SEARCH per poll and nothing else.
|
||||
func (r *reader) fetch(password string) ([]email.Message, error) {
|
||||
if r.fetchMail != nil {
|
||||
return r.fetchMail(password)
|
||||
}
|
||||
f := email.FetchSince{
|
||||
Addr: r.addr,
|
||||
User: r.user,
|
||||
@@ -225,8 +249,24 @@ func (r *reader) fetch(password string) ([]email.Message, error) {
|
||||
Since: time.Now().Add(-r.lookback),
|
||||
Max: r.max,
|
||||
Skip: r.state.seen,
|
||||
// Everything below the oldest searchable UID has aged out of the
|
||||
// lookback window and can never be read again. Retiring it is what keeps
|
||||
// one permanently failing message from pinning the high-water mark
|
||||
// forever. See seenState.retire.
|
||||
OnSearch: func(uids []uint32) {
|
||||
if len(uids) == 0 {
|
||||
return
|
||||
}
|
||||
low := uids[0]
|
||||
for _, u := range uids {
|
||||
if u < low {
|
||||
low = u
|
||||
}
|
||||
}
|
||||
r.state.retire(low)
|
||||
},
|
||||
}
|
||||
return f.RunWith(password, r.dial)
|
||||
return f.Run(password)
|
||||
}
|
||||
|
||||
// ---- seen state ------------------------------------------------------------
|
||||
@@ -242,6 +282,14 @@ func (r *reader) fetch(password string) ([]email.Message, error) {
|
||||
// UIDs are per-mailbox and monotonic, so the set is kept as a high-water mark
|
||||
// plus the stragglers above it. If the server ever changes UIDVALIDITY, UIDs
|
||||
// reset and the window is simply re-read once — dedupe absorbs it.
|
||||
//
|
||||
// The high-water mark only advances through a CONTIGUOUS run, so a UID that
|
||||
// never ingests successfully would pin it forever: everything above stays in
|
||||
// the explicit set, and save rewrites all of it every poll. A year of that is
|
||||
// a few hundred thousand entries written every quarter hour, which breaks
|
||||
// nothing loudly and is exactly why it is worth catching. retire is the answer:
|
||||
// a UID that has fallen out of the SEARCH SINCE window can never be fetched
|
||||
// again, so there is nothing left to wait for.
|
||||
type seenState struct {
|
||||
path string
|
||||
high uint32
|
||||
@@ -280,6 +328,31 @@ func (s *seenState) mark(uid uint32) {
|
||||
}
|
||||
}
|
||||
|
||||
// retire records that no UID below floor is reachable any more — they have
|
||||
// aged out of the lookback window, so no poll will ever fetch them. The
|
||||
// high-water mark can jump past the gap they were holding open, and the
|
||||
// stragglers below it leave the explicit set.
|
||||
//
|
||||
// It never moves backwards, so a UIDVALIDITY reset (UIDs restarting low) makes
|
||||
// this a no-op rather than a way to un-see a mailbox.
|
||||
func (s *seenState) retire(floor uint32) {
|
||||
if floor == 0 || floor-1 <= s.high {
|
||||
return
|
||||
}
|
||||
s.high = floor - 1
|
||||
for u := range s.set {
|
||||
if u <= s.high {
|
||||
delete(s.set, u)
|
||||
}
|
||||
}
|
||||
// The run above the new mark may now be contiguous with it.
|
||||
for s.set[s.high+1] {
|
||||
delete(s.set, s.high+1)
|
||||
s.high++
|
||||
}
|
||||
s.dirty = true
|
||||
}
|
||||
|
||||
func (s *seenState) load() error {
|
||||
if s.path == "" {
|
||||
return nil
|
||||
|
||||
+112
-51
@@ -1,13 +1,10 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"context"
|
||||
"fmt"
|
||||
"net"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strconv"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
@@ -16,7 +13,11 @@ import (
|
||||
"github.com/kami/maven/internal/ipc"
|
||||
)
|
||||
|
||||
// ---- a scripted IMAP server, same shape internal/email's tests use ---------
|
||||
// ---- a fake mailbox -------------------------------------------------------
|
||||
//
|
||||
// It fakes the READ, not the protocol: internal/email owns the IMAP tests, and
|
||||
// its dialer is unexported precisely so this package cannot substitute a
|
||||
// transport.
|
||||
|
||||
type fakeIMAP struct {
|
||||
msgs map[uint32]string
|
||||
@@ -24,52 +25,45 @@ type fakeIMAP struct {
|
||||
cmds []string
|
||||
}
|
||||
|
||||
func (f *fakeIMAP) serve(c net.Conn) {
|
||||
defer c.Close()
|
||||
fmt.Fprint(c, "* OK fake ready\r\n")
|
||||
r := bufio.NewReader(c)
|
||||
for {
|
||||
line, err := r.ReadString('\n')
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
parts := strings.SplitN(strings.TrimRight(line, "\r\n"), " ", 2)
|
||||
if len(parts) != 2 {
|
||||
return
|
||||
}
|
||||
tag, cmd := parts[0], parts[1]
|
||||
f.cmds = append(f.cmds, cmd)
|
||||
upper := strings.ToUpper(cmd)
|
||||
switch {
|
||||
case strings.HasPrefix(upper, "LOGIN"), strings.HasPrefix(upper, "EXAMINE"):
|
||||
fmt.Fprintf(c, "%s OK\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))
|
||||
// fetch is the read seam the reader exposes: the daemon cannot reach
|
||||
// internal/email's dialer (it is unexported so nothing outside that package can
|
||||
// point the reader at a cleartext transport), so a test fakes the whole read.
|
||||
// The IMAP protocol itself is covered by internal/email's own tests.
|
||||
func (f *fakeIMAP) fetch(r *reader) func(string) ([]email.Message, error) {
|
||||
return func(string) ([]email.Message, error) {
|
||||
var out []email.Message
|
||||
var low uint32
|
||||
for _, uid := range f.uids {
|
||||
if low == 0 || uid < low {
|
||||
low = uid
|
||||
}
|
||||
fmt.Fprintf(c, "* SEARCH %s\r\n%s OK\r\n", strings.Join(ids, " "), tag)
|
||||
case strings.HasPrefix(upper, "UID FETCH"):
|
||||
uid, _ := strconv.ParseUint(strings.Fields(cmd)[2], 10, 32)
|
||||
if raw, ok := f.msgs[uint32(uid)]; ok {
|
||||
fmt.Fprintf(c, "* 1 FETCH (UID %d BODY[] {%d}\r\n%s)\r\n", uid, len(raw), raw)
|
||||
}
|
||||
fmt.Fprintf(c, "%s OK\r\n", tag)
|
||||
case strings.HasPrefix(upper, "LOGOUT"):
|
||||
fmt.Fprintf(c, "* BYE\r\n%s OK\r\n", tag)
|
||||
return
|
||||
default:
|
||||
fmt.Fprintf(c, "%s BAD\r\n", tag)
|
||||
}
|
||||
if low > 0 {
|
||||
r.state.retire(low)
|
||||
}
|
||||
for i := len(f.uids) - 1; i >= 0; i-- {
|
||||
uid := f.uids[i]
|
||||
if r.state.seen(uid) {
|
||||
continue
|
||||
}
|
||||
raw, ok := f.msgs[uid]
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
f.cmds = append(f.cmds, fmt.Sprintf("UID FETCH %d", uid))
|
||||
msg, err := email.ParseMessage(uid, []byte(raw))
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
out = append(out, msg)
|
||||
if r.max > 0 && len(out) >= r.max {
|
||||
break
|
||||
}
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
}
|
||||
|
||||
func (f *fakeIMAP) dial(_ string, timeout time.Duration) (*email.Conn, error) {
|
||||
cli, srv := net.Pipe()
|
||||
go f.serve(srv)
|
||||
return email.NewConn(cli, timeout)
|
||||
}
|
||||
|
||||
// ---- a fake core -----------------------------------------------------------
|
||||
|
||||
type fakeCore struct {
|
||||
@@ -96,12 +90,13 @@ func mail(subject, body string, extraHeaders ...string) string {
|
||||
|
||||
func newTestReader(t *testing.T, f *fakeIMAP, core *fakeCore, statePath string) *reader {
|
||||
t.Helper()
|
||||
return &reader{
|
||||
r := &reader{
|
||||
core: core, addr: "mail.example:993", user: "kami", mailbox: "INBOX",
|
||||
lookback: 72 * time.Hour, max: 25, timeout: 5 * time.Second,
|
||||
state: newSeenState(statePath),
|
||||
dial: f.dial,
|
||||
}
|
||||
r.fetchMail = f.fetch(r)
|
||||
return r
|
||||
}
|
||||
|
||||
func TestPollHandsMessagesToCore(t *testing.T) {
|
||||
@@ -116,11 +111,26 @@ func TestPollHandsMessagesToCore(t *testing.T) {
|
||||
r := newTestReader(t, f, core, "")
|
||||
r.pollOnce(context.Background(), "secret")
|
||||
|
||||
// The newsletter is filtered before core is asked: only the real mail crosses.
|
||||
if len(core.got) != 1 {
|
||||
t.Fatalf("core saw %d messages, want 1 (the bulk one must not cross): %+v", len(core.got), core.got)
|
||||
// Two calls: the real mail with its text, and the newsletter as a verdict
|
||||
// with no content at all. Core is told about bulk rather than asked, so it
|
||||
// can count it without spending the model.
|
||||
if len(core.got) != 2 {
|
||||
t.Fatalf("core saw %d messages, want 2: %+v", len(core.got), core.got)
|
||||
}
|
||||
var got, bulk ipc.IngestMailReq
|
||||
for _, r := range core.got {
|
||||
if r.Junk {
|
||||
bulk = r
|
||||
} else {
|
||||
got = r
|
||||
}
|
||||
}
|
||||
if bulk.UID != 2 || !bulk.Junk {
|
||||
t.Errorf("bulk req = %+v, want uid 2 flagged junk", bulk)
|
||||
}
|
||||
if bulk.Subject != "" || bulk.Body != "" || bulk.From != "" {
|
||||
t.Errorf("a bulk verdict must carry no mail content: %+v", bulk)
|
||||
}
|
||||
got := core.got[0]
|
||||
if got.UID != 1 || got.Mailbox != "INBOX" || got.Subject != "Счёт" {
|
||||
t.Errorf("ingest req = %+v", got)
|
||||
}
|
||||
@@ -252,3 +262,54 @@ func TestRunRejectsEmptyPasswordFile(t *testing.T) {
|
||||
t.Errorf("an empty password file must be refused before dialling; err = %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
// A UID that never ingests pinned the high-water mark forever, because the mark
|
||||
// only advances through a contiguous run. Once that UID falls out of the
|
||||
// lookback window it can never be fetched again, so there is nothing left to
|
||||
// wait for and everything above it can leave the explicit set.
|
||||
func TestSeenStateRetiresAgedOutUIDs(t *testing.T) {
|
||||
s := newSeenState("")
|
||||
s.mark(1000) // 999 failed and was deliberately not marked
|
||||
s.mark(1001)
|
||||
if s.high != 0 || len(s.set) != 2 {
|
||||
t.Fatalf("high = %d, set = %v; want the mark pinned below the gap", s.high, s.set)
|
||||
}
|
||||
// The next SEARCH window starts at 1000: 999 has aged out.
|
||||
s.retire(1000)
|
||||
if s.high != 1001 {
|
||||
t.Errorf("high = %d, want 1001 once the gap is unreachable", s.high)
|
||||
}
|
||||
if len(s.set) != 0 {
|
||||
t.Errorf("explicit set = %v, want empty", s.set)
|
||||
}
|
||||
if !s.seen(999) || !s.seen(1001) || s.seen(1002) {
|
||||
t.Errorf("seen(999)=%v seen(1001)=%v seen(1002)=%v", s.seen(999), s.seen(1001), s.seen(1002))
|
||||
}
|
||||
}
|
||||
|
||||
// retire never moves the mark backwards: a UIDVALIDITY reset restarts UIDs low,
|
||||
// and that must not un-see a mailbox or re-see one.
|
||||
func TestSeenStateRetireNeverGoesBackwards(t *testing.T) {
|
||||
s := newSeenState("")
|
||||
s.mark(1)
|
||||
s.mark(2)
|
||||
s.retire(1)
|
||||
if s.high != 2 {
|
||||
t.Errorf("high = %d, want 2 unchanged", s.high)
|
||||
}
|
||||
}
|
||||
|
||||
// A poll must not leave the state file growing with UIDs that are already
|
||||
// covered by the high-water mark.
|
||||
func TestPollRetiresThroughTheSearchWindow(t *testing.T) {
|
||||
f := &fakeIMAP{uids: []uint32{100, 101}, msgs: map[uint32]string{100: mail("a", "b"), 101: mail("c", "d")}}
|
||||
core := &fakeCore{}
|
||||
r := newTestReader(t, f, core, "")
|
||||
r.pollOnce(context.Background(), "secret")
|
||||
if r.state.high != 101 {
|
||||
t.Errorf("high = %d, want 101 — everything below the search window is unreachable", r.state.high)
|
||||
}
|
||||
if len(r.state.set) != 0 {
|
||||
t.Errorf("explicit set = %v, want empty", r.state.set)
|
||||
}
|
||||
}
|
||||
|
||||
+9
-2
@@ -131,16 +131,23 @@ services:
|
||||
# "-password-file", "/run/secrets/imap.password",
|
||||
# "-mailbox", "INBOX",
|
||||
# "-interval", "15m",
|
||||
# "-state", "/var/lib/maven/mail-seen.json"]
|
||||
# "-state", "/var/lib/mavmaild/mail-seen.json"]
|
||||
# depends_on: [mavend]
|
||||
# volumes:
|
||||
# - sockets:/run/maven
|
||||
# - dbdata:/var/lib/maven
|
||||
# # Its OWN volume, not dbdata. The whole point of a separate reader is
|
||||
# # that a compromise on either side does not reach the other, and dbdata
|
||||
# # is the encrypted database. The reader needs one JSON file of UIDs and
|
||||
# # gets a volume that holds nothing else, so neither can be restored from
|
||||
# # a backup of the other.
|
||||
# - maildata:/var/lib/mavmaild
|
||||
# - ./deploy/imap.password:/run/secrets/imap.password:ro
|
||||
|
||||
volumes:
|
||||
dbdata:
|
||||
sockets:
|
||||
# maildata — the mail reader's seen-UID file, and nothing else. See mavmaild.
|
||||
maildata:
|
||||
|
||||
networks:
|
||||
default:
|
||||
|
||||
@@ -27,6 +27,13 @@ type FetchSince struct {
|
||||
Since time.Time
|
||||
Max int
|
||||
Skip func(uid uint32) bool
|
||||
// OnSearch, when set, is handed the whole SEARCH result before anything is
|
||||
// fetched, ascending, seen UIDs included. It is how the poller learns which
|
||||
// UIDs are still inside the lookback window: anything below the lowest one
|
||||
// can never be searched for again, and therefore can never be read again.
|
||||
// Without that the poller cannot tell a UID it has not got to yet from one
|
||||
// that has aged out of the window.
|
||||
OnSearch func(uids []uint32)
|
||||
|
||||
// dial — the connection seam, unexported on purpose: see dialer(). Tests
|
||||
// inside this package set it through export_test.go; nothing outside can.
|
||||
@@ -84,6 +91,10 @@ func (f FetchSince) Run(password string) ([]Message, error) {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if f.OnSearch != nil {
|
||||
f.OnSearch(uids)
|
||||
}
|
||||
|
||||
// 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))
|
||||
|
||||
Symlink
+1
@@ -0,0 +1 @@
|
||||
/home/kami/apps/Maven/models/stt
|
||||
Symlink
+1
@@ -0,0 +1 @@
|
||||
/home/kami/apps/Maven/models/tts
|
||||
Reference in New Issue
Block a user