Normalize every intake path into one event envelope (#283) #78

Closed
claude wants to merge 1 commits from overnight/event-envelope into overnight/coldstart-unlock
19 changed files with 1231 additions and 18 deletions
+235
View File
@@ -0,0 +1,235 @@
// mavend/intake.go — the unified event intake envelope, wired (Vikunja #283).
//
// internal/event defines the envelope and the bounded in-memory journal. This
// file is the one place that FILLS it, and the reason it is one place is worth
// stating, because the alternative was eight patches:
//
// Every intake path in Maven already converges on three writes, and all three
// are ipc.CoreAPI methods —
//
// WriteFact ← POST /api/ambient, mavcaldav, mavpoll's zenmoney + wg reads,
// /api/signal presence probes, the RSS/crawl watermarks
// WriteNote ← the RSS poller, the page crawler, meeting transcripts,
// image descriptions
// CaptureTask ← the voice path, the web form, and the mail reader
//
// — so decorating that ONE interface with a publish covers the lot without a
// caller knowing about events at all. cmd/mavmaild, cmd/mavcaldav, cmd/mavpoll,
// cmd/mavweb and the in-core feed/crawl/capture/vision workers are unchanged:
// they call the same interface they always called, and it now also narrates.
//
// The exception is cmd/mavend/mail.go, which reaches past the interface to
// st.CaptureTask directly. It publishes explicitly; see mailIntake.ingest.
//
// # Production behaviour when nobody is watching
//
// A nil *event.Bus makes Publish a no-op, and newIntakeAPI with a nil bus
// returns the wrapped API unchanged, so there is not even a decorator on the
// call path. The journal is memory-only and is never consulted by the tick
// loop, the router, or delivery — nothing Maven says depends on it. It is a
// read surface (`/events`, `recent_events`) and an observation seam for the
// simulator.
//
// # What is deliberately NOT here
//
// No dispatch. An event is a report that something arrived, never an
// instruction to speak: "a feed item appeared" becoming a notification is the
// nag this repo refuses. Digestion may one day read the journal; it will still
// go through internal/loop's rules and the severity/presence routing table.
package main
import (
"context"
"log"
"strings"
"time"
"github.com/kami/maven/internal/config"
"github.com/kami/maven/internal/event"
"github.com/kami/maven/internal/ipc"
"github.com/kami/maven/internal/store"
)
// newEventBus builds the journal, or returns nil when the operator turned it
// off (a negative config.intake_journal). nil is the "behave exactly as before"
// value all the way down: no decorator, no ring, no /events rows.
func newEventBus(cfg *config.Config) *event.Bus {
if cfg == nil || cfg.IntakeJournal < 0 {
log.Printf("intake journal: off (intake_journal < 0)")
return nil
}
n := cfg.IntakeJournal
if n == 0 {
n = config.DefaultIntakeJournal
}
log.Printf("intake journal: keeping the last %d intake events in memory", n)
return event.NewBus(n)
}
// intakeEventsFn is the daemonAPI.getEvents closure: the bus's ring rendered as
// the wire type. Returns nil for a nil bus, which the daemonAPI reports as an
// empty journal rather than an error.
func intakeEventsFn(bus *event.Bus) func(n int) []ipc.IntakeEvent {
if bus == nil {
return nil
}
return func(n int) []ipc.IntakeEvent {
evs := bus.Recent(n)
out := make([]ipc.IntakeEvent, 0, len(evs))
for _, e := range evs {
out = append(out, ipc.IntakeEvent{
Source: e.Source,
Kind: e.Kind,
EntityIDs: e.EntityIDs,
Title: e.Title,
Body: e.Body,
Priority: e.Priority,
OccurredAt: e.OccurredAt,
})
}
return out
}
}
// intakeAPI decorates a CoreAPI, publishing one envelope per successful
// intake write. Embedding the interface means every other method passes
// through untouched, and a new CoreAPI method is inherited rather than
// silently dropped.
type intakeAPI struct {
ipc.CoreAPI
bus *event.Bus
now func() time.Time
}
// newIntakeAPI wraps api so its intake writes are journalled. A nil bus
// returns api itself — no decorator, no allocation, no behaviour change.
func newIntakeAPI(api ipc.CoreAPI, bus *event.Bus, now func() time.Time) ipc.CoreAPI {
if bus == nil || api == nil {
return api
}
if now == nil {
now = time.Now
}
return &intakeAPI{CoreAPI: api, bus: bus, now: now}
}
// WriteFact journals the fact after it lands. Order matters: an event is a
// report of something that HAPPENED, so a failed write publishes nothing.
func (a *intakeAPI) WriteFact(ctx context.Context, req ipc.WriteFactReq) (int64, error) {
id, err := a.CoreAPI.WriteFact(ctx, req)
if err != nil {
return id, err
}
// OccurredAt is req.Ts, not now: mavpoll's wg read carries the handshake
// instant and the ambient path carries the meeting's start. Flattening
// those to notice-time would make the journal lie about when things
// happened, which is the one thing it is for.
a.bus.Publish(event.Event{
Source: req.Source,
Kind: event.SourceKind(req.Source, event.KindFact),
Title: req.Key,
Body: req.Value,
Priority: factPriority(req),
OccurredAt: req.Ts,
EntityIDs: entityIDs(req.Subject),
}, a.now())
return id, nil
}
// WriteNote journals a note. This is the RSS and crawler path, and also the
// meeting transcript and image description paths, which write their derived
// text as ordinary notes.
func (a *intakeAPI) WriteNote(ctx context.Context, ts time.Time, text string, embedding []float32, source string) (int64, error) {
id, err := a.CoreAPI.WriteNote(ctx, ts, text, embedding, source)
if err != nil {
return id, err
}
title, body := splitFirstLine(text)
a.bus.Publish(event.Event{
Source: source,
Kind: event.SourceKind(source, event.KindNote),
Title: title,
Body: body,
Priority: event.PriorityLow,
OccurredAt: ts,
}, a.now())
return id, nil
}
// CaptureTask journals a captured task, but only when a row was actually
// created. CaptureTask dedupes on normalised text among live rows, so a
// mailbox re-read after a restart must not refill the journal with tasks that
// were already there.
func (a *intakeAPI) CaptureTask(ctx context.Context, req ipc.CaptureTaskReq) (ipc.CaptureTaskResp, error) {
resp, err := a.CoreAPI.CaptureTask(ctx, req)
if err != nil || !resp.Created {
return resp, err
}
a.bus.Publish(publishableTask(store.Task{
CreatedTs: req.Ts,
Text: req.Text,
Source: req.Source,
Evidence: req.Evidence,
Status: req.Status,
Due: req.Due,
}, a.now()), a.now())
return resp, nil
}
// publishableTask is the task→envelope shape, shared with mail.go, which
// captures through the store directly rather than through the interface.
//
// Priority is high for a candidate with a due date and normal otherwise. That
// is the only place this file makes a judgement, and it is a display hint on a
// review page — nothing routes on it.
func publishableTask(t store.Task, now time.Time) event.Event {
occurred := t.CreatedTs
if occurred.IsZero() {
occurred = now
}
prio := event.PriorityNormal
if t.Due != nil {
prio = event.PriorityHigh
}
return event.Event{
Source: t.Source,
Kind: event.KindTask,
Title: t.Text,
Body: t.Evidence,
Priority: prio,
OccurredAt: occurred,
}
}
// factPriority is the attention hint for a fact write. Deliberately crude:
// a low-confidence inference (the ambient notification path writes below 1.0)
// is worth less attention than a read he or a credentialled poller made, and
// nothing else is distinguishable from here.
func factPriority(req ipc.WriteFactReq) string {
if req.Confidence > 0 && req.Confidence < 1.0 {
return event.PriorityLow
}
return event.PriorityNormal
}
// entityIDs turns a fact's free-text Subject into the EntityIDs slot when it
// already looks resolved. Intake runs BEFORE the fact enrichment worker
// resolves a subject against Nexus, so this is almost always empty — the slot
// exists for the paths that do know (the ecosystem acts), not for guessing.
func entityIDs(subject string) []string {
subject = strings.TrimSpace(subject)
if subject == "" || !strings.HasPrefix(subject, "entity:") {
return nil
}
return []string{strings.TrimPrefix(subject, "entity:")}
}
// splitFirstLine renders a note as title + body. Feed and crawl notes are
// written "headline\nsummary\nlink", so the first line is already the title.
func splitFirstLine(text string) (title, body string) {
text = strings.TrimSpace(text)
if i := strings.IndexByte(text, '\n'); i >= 0 {
return strings.TrimSpace(text[:i]), strings.TrimSpace(text[i+1:])
}
return text, ""
}
+178
View File
@@ -0,0 +1,178 @@
package main
import (
"context"
"errors"
"testing"
"time"
"github.com/kami/maven/internal/config"
"github.com/kami/maven/internal/event"
"github.com/kami/maven/internal/ipc"
)
var intakeNow = time.Date(2026, 8, 1, 10, 0, 0, 0, time.UTC)
func intakeClock() time.Time { return intakeNow }
// failingAPI wraps the store adapter, failing the three intake writes on
// demand, so the "a failed write publishes nothing" invariant is testable.
type failingAPI struct {
ipc.CoreAPI
fail bool
}
func (f *failingAPI) WriteFact(ctx context.Context, req ipc.WriteFactReq) (int64, error) {
if f.fail {
return 0, errors.New("injected")
}
return f.CoreAPI.WriteFact(ctx, req)
}
func newIntakeTestAPI(t *testing.T) (ipc.CoreAPI, *event.Bus) {
t.Helper()
st := newTestStore(t)
bus := event.NewBus(32)
return newIntakeAPI(ipc.NewStoreAPI(st), bus, intakeClock), bus
}
func TestIntakeAPIWithoutBusIsTheBareAPI(t *testing.T) {
// The adoption invariant: with the journal off there is not even a
// decorator on the intake path, so production behaves exactly as before.
st := newTestStore(t)
bare := ipc.NewStoreAPI(st)
if got := newIntakeAPI(bare, nil, intakeClock); got != ipc.CoreAPI(bare) {
t.Errorf("newIntakeAPI with a nil bus returned a wrapper, want the bare API")
}
}
func TestNewEventBusOffWhenNegative(t *testing.T) {
if b := newEventBus(&config.Config{IntakeJournal: -1}); b != nil {
t.Error("intake_journal = -1 still built a bus")
}
if b := newEventBus(&config.Config{IntakeJournal: 4}); b == nil {
t.Error("intake_journal = 4 built no bus")
}
}
func TestIntakeJournalsAFactWrite(t *testing.T) {
api, bus := newIntakeTestAPI(t)
ctx := context.Background()
// The ambient path's shape: an env fact below full confidence, timestamped
// at the meeting's start rather than at notice time.
start := intakeNow.Add(2 * time.Hour)
if _, err := api.WriteFact(ctx, ipc.WriteFactReq{
Ts: start, Kind: "env", Key: "calendar_event_20260801_планёрка",
Value: "10:00-11:00 планёрка", Source: "ambient:notif", Confidence: 0.6,
}); err != nil {
t.Fatalf("WriteFact: %v", err)
}
got := bus.Recent(0)
if len(got) != 1 {
t.Fatalf("journal has %d entries, want 1", len(got))
}
e := got[0]
if e.Source != "ambient:notif" || e.Kind != event.KindFact {
t.Errorf("source/kind = %q/%q", e.Source, e.Kind)
}
if e.Title != "calendar_event_20260801_планёрка" {
t.Errorf("title = %q, want the fact key", e.Title)
}
if !e.OccurredAt.Equal(start) {
t.Errorf("occurred_at = %v, want the fact's Ts %v — the journal must not flatten intake to notice time", e.OccurredAt, start)
}
if e.Priority != event.PriorityLow {
t.Errorf("priority = %q, want %q for a sub-1.0 confidence read", e.Priority, event.PriorityLow)
}
}
func TestIntakeDoesNotJournalAFailedWrite(t *testing.T) {
st := newTestStore(t)
bus := event.NewBus(8)
api := newIntakeAPI(&failingAPI{CoreAPI: ipc.NewStoreAPI(st), fail: true}, bus, intakeClock)
if _, err := api.WriteFact(context.Background(), ipc.WriteFactReq{
Ts: intakeNow, Kind: "env", Key: "k", Value: "v", Source: "poll:zenmoney", Confidence: 1,
}); err == nil {
t.Fatal("expected the injected error")
}
if bus.Len() != 0 {
t.Errorf("journal has %d entries after a failed write, want 0 — an event reports something that happened", bus.Len())
}
}
func TestIntakeJournalsANoteAsTitlePlusBody(t *testing.T) {
api, bus := newIntakeTestAPI(t)
// The RSS shape: "headline\nsummary\nlink".
if _, err := api.WriteNote(context.Background(), intakeNow,
"Вышло ядро 6.19\nкраткое содержание\nhttps://example.org/a", nil, "rss:tech"); err != nil {
t.Fatalf("WriteNote: %v", err)
}
got := bus.Recent(1)
if len(got) != 1 {
t.Fatalf("journal has %d entries, want 1", len(got))
}
if got[0].Title != "Вышло ядро 6.19" {
t.Errorf("title = %q, want the headline", got[0].Title)
}
if got[0].Kind != event.KindNote {
t.Errorf("kind = %q, want %q", got[0].Kind, event.KindNote)
}
if got[0].Body == "" {
t.Error("body is empty, want the rest of the note")
}
}
func TestIntakeJournalsOnlyCreatedTasks(t *testing.T) {
api, bus := newIntakeTestAPI(t)
ctx := context.Background()
req := ipc.CaptureTaskReq{Text: "оплатить интернет", Source: "email:inbox", Status: "candidate", Ts: intakeNow}
if _, err := api.CaptureTask(ctx, req); err != nil {
t.Fatalf("CaptureTask: %v", err)
}
// Same text again: CaptureTask dedupes among live rows, and a re-read of a
// mailbox must not refill the journal.
resp, err := api.CaptureTask(ctx, req)
if err != nil {
t.Fatalf("CaptureTask (repeat): %v", err)
}
if resp.Created {
t.Fatal("store did not dedupe; the test cannot check what it means to")
}
if bus.Len() != 1 {
t.Errorf("journal has %d entries, want 1 — a deduped capture must not publish", bus.Len())
}
if got := bus.Recent(1)[0]; got.Kind != event.KindTask || got.Title != "оплатить интернет" {
t.Errorf("entry = %+v, want the captured task", got)
}
}
func TestIntakeEventsFnRendersNewestFirst(t *testing.T) {
api, bus := newIntakeTestAPI(t)
ctx := context.Background()
for _, key := range []string{"a", "b", "c"} {
if _, err := api.WriteFact(ctx, ipc.WriteFactReq{
Ts: intakeNow, Kind: "env", Key: key, Value: "1", Source: "poll:zenmoney", Confidence: 1,
}); err != nil {
t.Fatalf("WriteFact %s: %v", key, err)
}
}
fn := intakeEventsFn(bus)
got := fn(2)
if len(got) != 2 || got[0].Title != "c" || got[1].Title != "b" {
t.Errorf("intakeEventsFn(2) = %+v, want the two newest, newest first", got)
}
if intakeEventsFn(nil) != nil {
t.Error("intakeEventsFn(nil) returned a closure, want nil so daemonAPI reports an empty journal")
}
}
func TestDaemonAPIRecentEventsEmptyWithoutABus(t *testing.T) {
d := &daemonAPI{CoreAPI: ipc.UnimplementedCoreAPI{}}
got, err := d.RecentEvents(context.Background(), 10)
if err != nil {
t.Fatalf("RecentEvents with no journal errored: %v", err)
}
if len(got) != 0 {
t.Errorf("got %d events, want none", len(got))
}
}
+14 -4
View File
@@ -26,6 +26,7 @@ import (
"github.com/kami/maven/internal/config"
"github.com/kami/maven/internal/email"
"github.com/kami/maven/internal/event"
"github.com/kami/maven/internal/ipc"
"github.com/kami/maven/internal/phraser"
"github.com/kami/maven/internal/store"
@@ -42,6 +43,11 @@ type mailIntake struct {
ex *email.Extractor
timeout time.Duration
now func() time.Time
// bus — the unified intake journal (Vikunja #283). This path captures
// through the store directly rather than through ipc.CoreAPI, so the
// decorator in intake.go does not see it and the publish is explicit here.
// nil is a working no-op.
bus *event.Bus
}
// newMailIntake returns nil when mail ingestion must not be available, which is
@@ -52,7 +58,7 @@ type mailIntake struct {
// no keyword fallback: "the subject line became a task" is not extraction,
// it is a mailbox rendered as a to-do list, and it would fill the review
// page faster than he could clear it.
func newMailIntake(st *store.Store, phr phraser.Phraser, cfg *config.Config) *mailIntake {
func newMailIntake(st *store.Store, phr phraser.Phraser, cfg *config.Config, bus *event.Bus) *mailIntake {
if cfg.Email == nil {
return nil
}
@@ -67,7 +73,7 @@ func newMailIntake(st *store.Store, phr phraser.Phraser, cfg *config.Config) *ma
}
ex := email.NewExtractor(llmClientFor(lp, timeout), cfg.Email.MaxTasks, contextBlockFn(cfg, time.Now))
log.Printf("mail intake: enabled (max %d candidates per message, timeout %s)", cfg.Email.MaxTasks, timeout)
return &mailIntake{st: st, ex: ex, timeout: timeout, now: time.Now}
return &mailIntake{st: st, ex: ex, timeout: timeout, now: time.Now, bus: bus}
}
// ingest handles one ipc.MethodIngestMail call.
@@ -128,6 +134,10 @@ func (m *mailIntake) ingest(ctx context.Context, req ipc.IngestMailReq) (ipc.Ing
resp.TaskIDs = append(resp.TaskIDs, id)
if created {
resp.Created++
// Only a row that was actually created. CaptureTask dedupes on
// normalised text among live rows, so a mailbox re-read after a
// restart must not refill the journal with tasks already in it.
m.bus.Publish(publishableTask(t, now), now)
}
}
// Counts only: the log line names the mailbox and the UID, never the subject,
@@ -139,8 +149,8 @@ func (m *mailIntake) ingest(ctx context.Context, req ipc.IngestMailReq) (ipc.Ing
// wireMailIntake installs the IPC hook, or leaves it nil so the method reports
// ErrUnknownMethod. Called on both startup paths (unlocked boot and passkey
// unlock) so mail behaves the same either way.
func wireMailIntake(srv *ipc.Server, st *store.Store, phr phraser.Phraser, cfg *config.Config) {
mi := newMailIntake(st, phr, cfg)
func wireMailIntake(srv *ipc.Server, st *store.Store, phr phraser.Phraser, cfg *config.Config, bus *event.Bus) {
mi := newMailIntake(st, phr, cfg, bus)
if mi == nil {
return
}
+2 -2
View File
@@ -171,12 +171,12 @@ func TestIngestTruncatesEvidence(t *testing.T) {
// exist at all.
func TestNewMailIntakeOffWithoutConfig(t *testing.T) {
st := newTestStore(t)
if mi := newMailIntake(st, nil, &config.Config{}); mi != nil {
if mi := newMailIntake(st, nil, &config.Config{}, nil); mi != nil {
t.Error("no email block must mean no mail intake")
}
// Configured but with a non-LLM phraser: still off — there is no fallback
// extraction, by design.
if mi := newMailIntake(st, nil, &config.Config{Email: &config.EmailConfig{}}); mi != nil {
if mi := newMailIntake(st, nil, &config.Config{Email: &config.EmailConfig{}}, nil); mi != nil {
t.Error("without a llama-server phraser there is nothing to extract with")
}
}
+22 -10
View File
@@ -204,6 +204,16 @@ func run(args []string) error {
crawlWkr *crawlWorker // nil ⇒ no page is watched (the default)
)
// The unified intake journal (Vikunja #283). Built before anything else
// that holds a CoreAPI, because intakeAPI wraps that one interface and
// every intake path in the daemon reaches its sink through it. nil (the
// operator set intake_journal negative) means no decorator at all.
evBus := newEventBus(cfg)
// coreFor is what every in-process holder of a CoreAPI now takes, instead
// of a bare ipc.NewStoreAPI(st). Identical behaviour plus one published
// envelope per successful intake write.
coreFor := func() ipc.CoreAPI { return newIntakeAPI(ipc.NewStoreAPI(st), evBus, time.Now) }
if !locked {
rules = loop.DefaultRules()
gatherer = loop.NewGatherer(st, rules)
@@ -247,7 +257,7 @@ func run(args []string) error {
eco = wireEcosystem(cfg)
// voice
voiceW, err = wireVoice(cfg, ipc.NewStoreAPI(st), phr, st.VectorMemory(), st, eco)
voiceW, err = wireVoice(cfg, coreFor(), phr, st.VectorMemory(), st, eco)
if err != nil {
return fmt.Errorf("wire voice: %w", err)
}
@@ -297,14 +307,15 @@ func run(args []string) error {
tl = newTickLoop(st, gatherer, dispatcher, phr, rules, tickInterval, repeatInterval, autotuneInterval, cfg.Digest, routinesFromConfig(cfg.Routines), config.MorningRoutinesFromConfig(cfg.MorningRoutines), cfg.PatternProposals)
factWorker = newFactEnrichmentWorker(st, eco, time.Duration(cfg.FactEnrichmentInterval))
evalWorker = newMemoryEvalWorker(st, phr, cfg)
feedWkr = newFeedWorker(ipc.NewStoreAPI(st), embedderOf(voiceW), cfg)
crawlWkr = newCrawlWorker(newCrawler(cfg), ipc.NewStoreAPI(st), embedderOf(voiceW), cfg)
feedWkr = newFeedWorker(coreFor(), embedderOf(voiceW), cfg)
crawlWkr = newCrawlWorker(newCrawler(cfg), coreFor(), embedderOf(voiceW), cfg)
coreAPI = &daemonAPI{
CoreAPI: ipc.NewStoreAPI(st),
CoreAPI: coreFor(),
getTrace: tl.trace,
getMorningStatus: func(ctx context.Context) []ipc.MorningRoutineStatus { return tl.morningStatus(ctx, time.Now()) },
getDayPlan: func(ctx context.Context) ipc.DayPlan { return tl.dayPlan(ctx, time.Now()) },
getEvents: intakeEventsFn(evBus),
}
if voiceW != nil && voiceW.handler != nil {
api := coreAPI.(*daemonAPI)
@@ -361,7 +372,7 @@ func run(args []string) error {
// configured and there is a llama-server to extract with, in which case
// ipc.MethodIngestMail reports ErrUnknownMethod.
if !locked {
wireMailIntake(srv, st, phr, cfg)
wireMailIntake(srv, st, phr, cfg, evBus)
wireModelSwap(srv, phr, cfg)
// Vision + the media blob store (Vikunja #252). Both stay dark without a
// media block; MethodDescribeImage answers ErrUnknownMethod then.
@@ -482,7 +493,7 @@ func run(args []string) error {
eco = wireEcosystem(cfg)
voiceW, err = wireVoice(cfg, ipc.NewStoreAPI(st), phr, st.VectorMemory(), st, eco)
voiceW, err = wireVoice(cfg, coreFor(), phr, st.VectorMemory(), st, eco)
if err != nil {
return fmt.Errorf("wire voice: %w", err)
}
@@ -526,22 +537,23 @@ func run(args []string) error {
tl = newTickLoop(st, gatherer, dispatcher, phr, rules, tickInterval, repeatInterval, autotuneInterval, cfg.Digest, routinesFromConfig(cfg.Routines), config.MorningRoutinesFromConfig(cfg.MorningRoutines), cfg.PatternProposals)
factWorker = newFactEnrichmentWorker(st, eco, time.Duration(cfg.FactEnrichmentInterval))
evalWorker = newMemoryEvalWorker(st, phr, cfg)
feedWkr = newFeedWorker(ipc.NewStoreAPI(st), embedderOf(voiceW), cfg)
crawlWkr = newCrawlWorker(newCrawler(cfg), ipc.NewStoreAPI(st), embedderOf(voiceW), cfg)
feedWkr = newFeedWorker(coreFor(), embedderOf(voiceW), cfg)
crawlWkr = newCrawlWorker(newCrawler(cfg), coreFor(), embedderOf(voiceW), cfg)
// Swap the CoreAPI from the locked placeholder to the real store adapter.
newAPI := &daemonAPI{
CoreAPI: ipc.NewStoreAPI(st),
CoreAPI: coreFor(),
getTrace: tl.trace,
getMorningStatus: func(ctx context.Context) []ipc.MorningRoutineStatus { return tl.morningStatus(ctx, time.Now()) },
getDayPlan: func(ctx context.Context) ipc.DayPlan { return tl.dayPlan(ctx, time.Now()) },
getEvents: intakeEventsFn(evBus),
}
if voiceW != nil && voiceW.handler != nil {
newAPI.chatFn = voiceW.handler.handleText
}
srv.SetAPI(newAPI)
srv.Check = (&auth.Gate{Enrollment: auth.NewFloorEnrollment(), Session: passkeySess}).Check
wireMailIntake(srv, st, phr, cfg)
wireMailIntake(srv, st, phr, cfg, evBus)
wireModelSwap(srv, phr, cfg)
keeper := wireVision(ctx, srv, st, embedderOf(voiceW), cfg)
wireCapture(srv, keeper, st, voiceW, phr, cfg)
+12
View File
@@ -942,6 +942,18 @@ type daemonAPI struct {
getDayPlan func(ctx context.Context) ipc.DayPlan
chatFn func(ctx context.Context, text string) string
getMCPServers func() []ipc.MCPServerStatus
getEvents func(n int) []ipc.IntakeEvent
}
// RecentEvents — the unified intake journal (Vikunja #283). Empty, not an
// error, when no bus was wired: "nothing has arrived" and "the journal is off"
// look the same to a reader on purpose, because neither is a fault and the
// page renders both as an empty table.
func (d *daemonAPI) RecentEvents(ctx context.Context, n int) ([]ipc.IntakeEvent, error) {
if d.getEvents == nil {
return nil, nil
}
return d.getEvents(n), nil
}
func (d *daemonAPI) Chat(ctx context.Context, text string) (string, error) {
+24
View File
@@ -0,0 +1,24 @@
{{template "shellTop" "events"}}
<h1>Intake</h1>
<div class=hint>Everything that arrived, newest first — a relayed notification, a mail candidate, a feed
item, a changed page, a spend, a presence probe. One envelope per write; the durable row is still the
fact, note or task itself. Held in memory only, so a restart empties this.</div>
{{if .Err}}<div class=hint>journal unavailable: {{.Err}}</div>{{end}}
{{if and (not .Events) (not .Err)}}
<div class=hint>nothing has arrived yet</div>
{{end}}
{{if .Events}}
<div class=scroll><table class=mono>
<tr><th>when<th>source<th>kind<th>pri<th>what<th>detail</tr>
{{range .Events}}<tr>
<td>{{.OccurredAt.Format "02.01 15:04:05"}}</td>
<td class=gray>{{.Source}}</td>
<td class=gray>{{.Kind}}</td>
<td class=gray>{{.Priority}}</td>
<td>{{.Title}}</td>
<td class=gray>{{.Body}}</td>
</tr>{{end}}
</table></div>
{{end}}
{{template "shellBottom"}}
</html>
+107
View File
@@ -0,0 +1,107 @@
package main
import (
"context"
"errors"
"net/http"
"net/http/httptest"
"strings"
"testing"
"time"
"github.com/kami/maven/internal/ipc"
)
// eventsCore serves a canned intake journal. Embedding
// ipc.UnimplementedCoreAPI means any other call fails loudly.
type eventsCore struct {
ipc.UnimplementedCoreAPI
events []ipc.IntakeEvent
err error
gotN int
}
func (c *eventsCore) RecentEvents(_ context.Context, n int) ([]ipc.IntakeEvent, error) {
c.gotN = n
return c.events, c.err
}
func getEvents(t *testing.T, core ipc.CoreAPI) *httptest.ResponseRecorder {
t.Helper()
w := httptest.NewRecorder()
handleEvents(w, httptest.NewRequest(http.MethodGet, "/events", nil), core)
return w
}
func TestEventsPageRendersTheJournal(t *testing.T) {
core := &eventsCore{events: []ipc.IntakeEvent{
{Source: "rss:tech", Kind: "note", Title: "Вышло ядро 6.19", Priority: "low",
OccurredAt: time.Date(2026, 8, 1, 7, 15, 0, 0, time.UTC)},
{Source: "ambient:notif", Kind: "fact", Title: "calendar_event_20260801_планёрка",
Body: "10:00-11:00 планёрка", Priority: "low",
OccurredAt: time.Date(2026, 8, 1, 10, 0, 0, 0, time.UTC)},
}}
w := getEvents(t, core)
if w.Code != http.StatusOK {
t.Fatalf("status = %d, want 200", w.Code)
}
body := w.Body.String()
for _, want := range []string{"rss:tech", "Вышло ядро 6.19", "ambient:notif", "10:00-11:00 планёрка", "01.08 10:00:00"} {
if !strings.Contains(body, want) {
t.Errorf("page does not mention %q", want)
}
}
if core.gotN != eventsPageLimit {
t.Errorf("asked core for %d events, want %d", core.gotN, eventsPageLimit)
}
}
func TestEventsPageSaysNothingArrived(t *testing.T) {
w := getEvents(t, &eventsCore{})
if w.Code != http.StatusOK {
t.Fatalf("status = %d, want 200", w.Code)
}
if !strings.Contains(w.Body.String(), "nothing has arrived yet") {
t.Error("empty journal did not render the empty-state line")
}
}
func TestEventsPageReportsAReadFailure(t *testing.T) {
// An unreachable journal must say so rather than render an empty table,
// which would imply nothing arrived.
w := getEvents(t, &eventsCore{err: errors.New("core is down")})
if w.Code != http.StatusOK {
t.Fatalf("status = %d, want 200 with the error rendered", w.Code)
}
body := w.Body.String()
if !strings.Contains(body, "journal unavailable") || !strings.Contains(body, "core is down") {
t.Errorf("page did not report the read failure: %s", body)
}
if strings.Contains(body, "nothing has arrived yet") {
t.Error("a failed read rendered as an empty journal")
}
}
func TestEventsPageWithoutCore(t *testing.T) {
w := getEvents(t, nil)
if w.Code != http.StatusServiceUnavailable {
t.Errorf("status = %d, want 503", w.Code)
}
}
func TestEventsPageEscapesIntakeText(t *testing.T) {
// Titles come from outside — a feed headline, a notification. They are shown
// on a page and must never be able to inject markup into it.
core := &eventsCore{events: []ipc.IntakeEvent{{
Source: "rss:x", Kind: "note", Priority: "low",
Title: `<script>alert(1)</script>`,
OccurredAt: time.Date(2026, 8, 1, 7, 0, 0, 0, time.UTC),
}}}
body := getEvents(t, core).Body.String()
if strings.Contains(body, "<script>alert(1)</script>") {
t.Error("intake title was not escaped")
}
if !strings.Contains(body, "&lt;script&gt;") {
t.Error("intake title is missing from the page entirely")
}
}
+42
View File
@@ -71,6 +71,9 @@ var ecosystemHTML string
//go:embed morning.html
var morningHTML string
//go:embed events.html
var eventsHTML string
// ── Ethos Workstation Shell ──
//
// Two template pieces that wrap every page:
@@ -105,6 +108,7 @@ var sidebarSections = []struct {
{Label: "Reminders", URL: "/reminders", Key: "reminders"},
{Label: "Routines", URL: "/routines", Key: "routines"},
{Label: "Morning", URL: "/morning", Key: "morning"},
{Label: "Intake", URL: "/events", Key: "events"},
},
},
{
@@ -318,6 +322,10 @@ var ecosystemTmpl = template.Must(template.New("ecosystem").Funcs(shellFuncs()).
// morning routine (internal/morning). Same shape as trace.html: a plain
// server-rendered page, refreshed on reload — no live-update loop, since
// checklist state changes on the scale of minutes, not seconds.
// eventsTmpl — the unified intake journal (Vikunja #283), read-only. Same
// shape as trace.html and morning.html: server-rendered, refreshed on reload.
var eventsTmpl = template.Must(template.New("events").Funcs(shellFuncs()).Parse(shellTopHTML + eventsHTML + shellBottomHTML))
var morningTmpl = template.Must(template.New("morning").Funcs(shellFuncs()).Parse(shellTopHTML + morningHTML + shellBottomHTML))
func noCache(h http.Handler) http.Handler {
@@ -430,6 +438,9 @@ func main() {
mux.HandleFunc("/morning", func(w http.ResponseWriter, r *http.Request) {
handleMorning(w, r, core)
})
mux.HandleFunc("/events", func(w http.ResponseWriter, r *http.Request) {
handleEvents(w, r, core)
})
ecoURLsCfg := ecoURLs{nexus: *nexusURL, praxis: *praxisURL, hexis: *hexisURL}
mux.HandleFunc("/ecosystem", func(w http.ResponseWriter, r *http.Request) {
handleEcosystem(w, r, ecoURLsCfg)
@@ -1206,6 +1217,37 @@ type morningView struct {
Routines []ipc.MorningRoutineStatus
}
// eventsView — what /events renders. Err is set instead of Events when the
// core could not serve the journal, so the page says why rather than showing an
// empty intake and implying nothing arrived.
type eventsView struct {
Events []ipc.IntakeEvent
Err string
}
// eventsPageLimit — how many envelopes the page shows. The ring holds more; a
// page is for scanning what just happened, not for archaeology.
const eventsPageLimit = 200
func handleEvents(w http.ResponseWriter, r *http.Request, core ipc.CoreAPI) {
if core == nil {
http.Error(w, "intake journal disabled (no -core)", http.StatusServiceUnavailable)
return
}
var view eventsView
evs, err := core.RecentEvents(r.Context(), eventsPageLimit)
if err != nil {
log.Printf("events: %v", err)
view.Err = err.Error()
} else {
view.Events = evs
}
w.Header().Set("Content-Type", "text/html; charset=utf-8")
if err := eventsTmpl.Execute(w, view); err != nil {
log.Printf("events render: %v", err)
}
}
func handleVoice(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "text/html; charset=utf-8")
if err := voiceTmpl.Execute(w, nil); err != nil {
+7 -1
View File
@@ -135,7 +135,13 @@ func Requirement(m ipc.Method) Authority {
ipc.MethodListSpeakers,
// The read side of the model swap: which model is resident, which ones are
// allowlisted. It loads nothing and changes nothing.
ipc.MethodModelStatus:
ipc.MethodModelStatus,
// The unified intake journal (Vikunja #283). AuthRead, and listed
// explicitly rather than inherited so the reasoning is on the record: it
// reports what already arrived — sources, keys, note headlines — which is
// the same material RecentFacts and RecentNotes already return at this
// rung. It writes nothing, and it holds nothing a fact read does not.
ipc.MethodRecentEvents:
return AuthRead
}
// Unknown method ⇒ AuthRead, but ipc.dispatch returns ErrUnknownMethod
+22 -1
View File
@@ -169,6 +169,20 @@ type Config struct {
// live in the reader (cmd/mavmaild), never here.
Email *EmailConfig `json:"email,omitempty"`
// IntakeJournal — how many entries the unified intake journal keeps
// (Vikunja #283): one envelope per thing that arrived, whatever direction it
// came from. Absent ⇒ DefaultIntakeJournal. A NEGATIVE value turns the
// journal off entirely, and then there is no decorator on the intake path at
// all.
//
// Not gated behind an "off unless configured" block like feeds or telegram,
// and the distinction is the one CLAUDE.md draws: that rule exists for
// capabilities that reach OUT — a fetch, a send, a third party. This reaches
// nowhere. It is a bounded in-memory log of writes core already performed,
// it is read only by /events and the simulator, and nothing Maven says
// depends on it.
IntakeJournal int `json:"intake_journal,omitempty"`
// Feeds — RSS/Atom feed reading (Vikunja #258). nil / absent ⇒ no feed is
// ever fetched: reading the outside world is off unless configured, like
// the weather and telegram. See FeedsConfig.
@@ -968,7 +982,11 @@ const (
DefaultRepeatInterval = 5 * time.Minute
DefaultAutotuneInterval = 10 * time.Minute
DefaultRouterThreshold = 0.55
DefaultQueryMinScore = 0.55
// DefaultIntakeJournal — entries kept in the unified intake journal
// (Vikunja #283). A busy day is a few hundred intake writes, so this is
// roughly "today and yesterday" at a few hundred KB of memory.
DefaultIntakeJournal = 512
DefaultQueryMinScore = 0.55
// Read off the margin sweep in internal/memory/recalleval on the e5
// embedder: 0.008 answers 68% of real questions (down from 72%) and cuts
// false recall from 5/5 to 1/5. Every larger delta costs real recall
@@ -1018,6 +1036,9 @@ func Load(path string) (*Config, error) {
}
func (c *Config) applyDefaults() {
if c.IntakeJournal == 0 {
c.IntakeJournal = DefaultIntakeJournal
}
if c.TickInterval == 0 {
c.TickInterval = Duration(DefaultTickInterval)
}
+129
View File
@@ -0,0 +1,129 @@
package event
import (
"sync"
"time"
)
// Bus — the in-memory intake journal: a bounded ring of recent Events plus
// zero or more subscribers.
//
// Two properties are load-bearing, both about not changing production
// behaviour when nobody is watching:
//
// - A nil *Bus is a working no-op. Publish on nil returns immediately, so
// an intake path can call b.Publish(...) unconditionally and a daemon that
// never built a bus behaves exactly as it did before. This is what let
// eight callers adopt the envelope without a config flag each.
// - Publish never blocks on a subscriber and never propagates a panic from
// one. Intake is on the request path of POST /api/ambient and of every
// fact write; a slow or broken observer must not be able to stall or kill
// a write that already succeeded.
//
// The ring is bounded because it is memory that nothing prunes otherwise. Its
// contents are a window, not a record: the durable consequence of an event is
// the fact, note or task the intake path wrote.
type Bus struct {
mu sync.Mutex
ring []Event // len == cap once full; oldest at (next % cap)
next int
n int
subs []func(Event)
}
// DefaultCapacity — how many recent events a bus keeps. A busy day is a few
// hundred intake events (a feed poll is one per new item), so this is roughly
// "today and yesterday" at a few hundred KB.
const DefaultCapacity = 512
// NewBus returns a bus keeping the last capacity events. capacity <= 0 uses
// DefaultCapacity.
func NewBus(capacity int) *Bus {
if capacity <= 0 {
capacity = DefaultCapacity
}
return &Bus{ring: make([]Event, capacity)}
}
// Publish normalizes e, drops it if it is not Valid, appends it to the ring and
// hands it to every subscriber. Safe on a nil receiver and safe from any
// goroutine.
//
// now is passed in rather than read from the clock: the whole point of #284's
// replay is that no time.Now() sits inside a path a scenario drives.
func (b *Bus) Publish(e Event, now time.Time) {
if b == nil {
return
}
e = e.Normalize(now)
if !e.Valid() {
return
}
b.mu.Lock()
b.ring[b.next] = e
b.next = (b.next + 1) % len(b.ring)
if b.n < len(b.ring) {
b.n++
}
subs := make([]func(Event), len(b.subs))
copy(subs, b.subs)
b.mu.Unlock()
for _, fn := range subs {
notify(fn, e)
}
}
// notify calls one subscriber, swallowing a panic. A test double or a page
// renderer must not be able to take down a daemon from the intake path.
func notify(fn func(Event), e Event) {
defer func() { _ = recover() }()
fn(e)
}
// Subscribe registers fn to be called for every subsequent event, in publish
// order. There is no unsubscribe: subscribers are wired at startup and live as
// long as the daemon. Safe on a nil receiver (the subscription is dropped,
// which is the honest outcome when there is no bus to subscribe to).
func (b *Bus) Subscribe(fn func(Event)) {
if b == nil || fn == nil {
return
}
b.mu.Lock()
defer b.mu.Unlock()
b.subs = append(b.subs, fn)
}
// Recent returns up to limit events, newest first. limit <= 0 returns
// everything held. Safe on a nil receiver (returns nil).
func (b *Bus) Recent(limit int) []Event {
if b == nil {
return nil
}
b.mu.Lock()
defer b.mu.Unlock()
if b.n == 0 {
return nil
}
if limit <= 0 || limit > b.n {
limit = b.n
}
out := make([]Event, 0, limit)
// next points one past the newest; walk backwards.
for i := 0; i < limit; i++ {
idx := (b.next - 1 - i + len(b.ring)*2) % len(b.ring)
out = append(out, b.ring[idx])
}
return out
}
// Len reports how many events the ring currently holds. Safe on nil.
func (b *Bus) Len() int {
if b == nil {
return 0
}
b.mu.Lock()
defer b.mu.Unlock()
return b.n
}
+198
View File
@@ -0,0 +1,198 @@
// Package event is the unified intake envelope (Vikunja #283,
// 20-07-2026-BACKLOG.md item 1).
//
// # The problem it solves
//
// Things arrive at Maven from a lot of directions: a relayed Android
// notification (POST /api/ambient), a mail the reader extracted candidates
// from (ingest_mail), an RSS item, a changed page the crawler noticed, a
// zenmoney spend, a CalDAV event, a wg handshake that means he is home, a
// photo he sent, a meeting she was asked to record. Each of those grew its own
// shape, its own storage decision and its own log line. Nothing could answer
// "what came in today, from where" without reading eight packages.
//
// An Event is that answer: one flat, source-agnostic description of "something
// arrived". It is deliberately NOT a new storage layer and NOT a new transport.
// Every intake path keeps writing exactly what it wrote before — a fact, a
// note, a candidate task — and additionally describes what it did as an Event.
// The envelope is a VIEW over intake, not a replacement for it, which is why
// adopting it did not require touching eight callers.
//
// # What it is not
//
// - Not a command. An Event is a report of something that happened; nothing
// in Maven executes one. Digestion may read them; it may not be driven by
// an event alone, because "a thing arrived" is not "a thing must be said".
// - Not durable. The bus is a bounded in-memory ring. An event's durable
// consequence is the fact/note/task the intake path already wrote; the
// envelope is the recent-history window on top. A restart losing the ring
// loses nothing that mattered.
// - Not a secret store. Body carries what the intake path was already willing
// to log or store. Nothing puts a mail body, an IMAP password or a
// voiceprint in here, and callers must keep it that way.
package event
import (
"encoding/json"
"strings"
"time"
)
// Event — one thing that arrived, normalized.
//
// The field set is the one recorded in the backlog, and it is intentionally
// small: anything source-specific goes in Payload, so adding a source never
// widens the struct and never breaks a reader.
type Event struct {
// Source — provenance, in the facts vocabulary already used across the
// repo: "ambient:notif", "caldav:personal", "poll:zenmoney", "rss:<feed>",
// "crawl:<watch>", "email:<mailbox>", "infer:wg", "tap:voice". Same string
// the fact or note was written under, so an event and its row can be
// matched up by eye.
Source string `json:"source"`
// Kind — what sort of thing arrived, from the closed set below. This is the
// field digestion switches on; Source is for provenance and display.
Kind string `json:"kind"`
// EntityIDs — Nexus entity ids this event is about, when the intake path
// knew any. Usually empty: most intake happens before enrichment resolves a
// subject to an entity.
EntityIDs []string `json:"entity_ids,omitempty"`
// Title — one short line, safe to show on a page. For a fact it is the key,
// for a note the first line, for a task the task text.
Title string `json:"title"`
// Body — optional detail, already truncated by the caller.
Body string `json:"body,omitempty"`
// Priority — one of PriorityLow / PriorityNormal / PriorityHigh. It is a
// hint about attention, not a delivery instruction: nothing here decides
// whether Maven speaks. That stays with internal/loop and internal/delivery,
// where the severity/presence routing table lives.
Priority string `json:"priority"`
// OccurredAt — when the thing happened, NOT when Maven noticed it. A wg
// handshake carries the handshake instant; an RSS item carries its publish
// time. Intake paths already make this distinction when writing facts, and
// the envelope must not flatten it.
OccurredAt time.Time `json:"occurred_at"`
// Payload — source-specific extra, opaque here. Optional.
Payload json.RawMessage `json:"payload,omitempty"`
}
// Kinds. Closed set: a reader may switch on these exhaustively. A new intake
// path picks the closest existing kind before it adds one — the point of the
// envelope is that digestion has a small stable input.
const (
// KindFact — something was written to the facts table: a calendar read, a
// zenmoney window, a presence probe, a crawler watermark.
KindFact = "fact"
// KindNote — something was written to the notes table: an RSS item, a
// changed page, a meeting transcript, an image description.
KindNote = "note"
// KindTask — a candidate task was captured: the mail reader, the web form,
// the voice path.
KindTask = "task"
// KindMessage — an inbound message on a reach channel. Nothing produces
// this yet (telegram is send-only today); the kind exists so the bridge,
// when it lands, is a constructor and not a schema change.
KindMessage = "message"
// KindHealth — a service or probe reported its own state.
KindHealth = "health"
)
// Priorities.
const (
PriorityLow = "low"
PriorityNormal = "normal"
PriorityHigh = "high"
)
// TitleMaxRunes / BodyMaxRunes bound what an envelope carries. The ring is
// in memory and served to a web page; a 40 KB crawled article has no business
// in either. Cut on a rune boundary — most of this text is Russian and half a
// cyrillic letter is a broken line.
const (
TitleMaxRunes = 120
BodyMaxRunes = 400
)
// Normalize returns e with its fields put in range: whitespace collapsed out
// of Title, Title and Body truncated, an unknown or empty Priority forced to
// PriorityNormal, and a zero OccurredAt filled from now.
//
// It takes now as a parameter rather than reading the clock, so the whole
// package stays pure and the simulator (Vikunja #284) can replay intake against
// a scripted clock.
func (e Event) Normalize(now time.Time) Event {
e.Title = truncateRunes(strings.Join(strings.Fields(e.Title), " "), TitleMaxRunes)
e.Body = truncateRunes(strings.TrimSpace(e.Body), BodyMaxRunes)
if !validPriority(e.Priority) {
e.Priority = PriorityNormal
}
if e.Kind == "" {
e.Kind = KindFact
}
if e.OccurredAt.IsZero() {
e.OccurredAt = now
}
return e
}
// Valid reports whether e carries the minimum a reader can rely on: a source,
// a known kind, a title and a time. The bus drops anything that fails — an
// envelope with no provenance is worse than no envelope, because it looks like
// evidence.
func (e Event) Valid() bool {
return e.Source != "" && validKind(e.Kind) && e.Title != "" && !e.OccurredAt.IsZero()
}
func validKind(k string) bool {
switch k {
case KindFact, KindNote, KindTask, KindMessage, KindHealth:
return true
}
return false
}
func validPriority(p string) bool {
switch p {
case PriorityLow, PriorityNormal, PriorityHigh:
return true
}
return false
}
// truncateRunes cuts s to n runes, marking the cut.
func truncateRunes(s string, n int) string {
r := []rune(s)
if len(r) <= n {
return s
}
return string(r[:n]) + "…"
}
// SourceKind guesses the Kind for a source string when the caller has not said
// otherwise. It exists so the one intake decorator in cmd/mavend does not need
// a switch per writer: the source prefix already tells you what arrived.
//
// Unknown prefixes get fallback, which is what the caller was going to write
// anyway (a WriteFact call knows it is a fact).
func SourceKind(source, fallback string) string {
switch {
case strings.HasPrefix(source, "rss:"), strings.HasPrefix(source, "crawl:"):
return KindNote
case strings.HasPrefix(source, "email:"):
return KindTask
case strings.HasPrefix(source, "probe:"), strings.HasPrefix(source, "health:"):
return KindHealth
}
return fallback
}
+185
View File
@@ -0,0 +1,185 @@
package event
import (
"strings"
"sync"
"testing"
"time"
)
var testNow = time.Date(2026, 8, 1, 9, 30, 0, 0, time.UTC)
func TestNormalizeFillsDefaults(t *testing.T) {
got := Event{Source: "poll:zenmoney", Title: " spent today "}.Normalize(testNow)
if got.Title != "spent today" {
t.Errorf("title = %q, want collapsed whitespace", got.Title)
}
if got.Priority != PriorityNormal {
t.Errorf("priority = %q, want %q", got.Priority, PriorityNormal)
}
if got.Kind != KindFact {
t.Errorf("kind = %q, want %q", got.Kind, KindFact)
}
if !got.OccurredAt.Equal(testNow) {
t.Errorf("occurred_at = %v, want %v", got.OccurredAt, testNow)
}
}
func TestNormalizeKeepsRealOccurredAt(t *testing.T) {
// A wg handshake carries the handshake instant, not "now". Flattening that
// would make every intake look like it happened at notice time.
real := testNow.Add(-3 * time.Hour)
got := Event{Source: "infer:wg", Title: "wg_handshake", OccurredAt: real}.Normalize(testNow)
if !got.OccurredAt.Equal(real) {
t.Errorf("occurred_at = %v, want the supplied %v", got.OccurredAt, real)
}
}
func TestNormalizeTruncatesOnRuneBoundary(t *testing.T) {
long := strings.Repeat("я", TitleMaxRunes+50)
got := Event{Source: "rss:x", Title: long}.Normalize(testNow)
r := []rune(got.Title)
if len(r) != TitleMaxRunes+1 { // +1 for the ellipsis marker
t.Fatalf("title runes = %d, want %d", len(r), TitleMaxRunes+1)
}
if r[len(r)-1] != '…' {
t.Errorf("truncated title does not mark the cut: %q", string(r[len(r)-3:]))
}
for _, c := range r[:TitleMaxRunes] {
if c != 'я' {
t.Fatalf("truncation broke a rune: got %q", c)
}
}
}
func TestNormalizeRejectsUnknownPriority(t *testing.T) {
got := Event{Source: "s", Title: "t", Priority: "URGENT!!"}.Normalize(testNow)
if got.Priority != PriorityNormal {
t.Errorf("priority = %q, want %q", got.Priority, PriorityNormal)
}
}
func TestValid(t *testing.T) {
base := Event{Source: "rss:tech", Kind: KindNote, Title: "заголовок", OccurredAt: testNow}
if !base.Valid() {
t.Fatal("well-formed event reported invalid")
}
for name, mut := range map[string]func(Event) Event{
"no source": func(e Event) Event { e.Source = ""; return e },
"no title": func(e Event) Event { e.Title = ""; return e },
"no time": func(e Event) Event { e.OccurredAt = time.Time{}; return e },
"bad kind": func(e Event) Event { e.Kind = "whatever"; return e },
} {
if mut(base).Valid() {
t.Errorf("%s: reported valid", name)
}
}
}
func TestSourceKind(t *testing.T) {
cases := map[string]string{
"rss:tech": KindNote,
"crawl:kernel": KindNote,
"email:inbox": KindTask,
"probe:netdata": KindHealth,
"ambient:notif": KindFact,
"tap:voice": KindFact,
}
for src, want := range cases {
if got := SourceKind(src, KindFact); got != want {
t.Errorf("SourceKind(%q) = %q, want %q", src, got, want)
}
}
}
func TestBusNilIsANoOp(t *testing.T) {
// The whole adoption story depends on this: an intake path calls Publish
// unconditionally, and a daemon with no bus behaves as it did before.
var b *Bus
b.Publish(Event{Source: "s", Kind: KindFact, Title: "t"}, testNow)
b.Subscribe(func(Event) { t.Error("nil bus delivered to a subscriber") })
if got := b.Recent(10); got != nil {
t.Errorf("Recent on nil bus = %v, want nil", got)
}
if got := b.Len(); got != 0 {
t.Errorf("Len on nil bus = %d, want 0", got)
}
}
func TestBusRecentIsNewestFirst(t *testing.T) {
b := NewBus(8)
for _, title := range []string{"one", "two", "three"} {
b.Publish(Event{Source: "rss:t", Kind: KindNote, Title: title}, testNow)
}
got := b.Recent(0)
if len(got) != 3 {
t.Fatalf("len = %d, want 3", len(got))
}
want := []string{"three", "two", "one"}
for i, w := range want {
if got[i].Title != w {
t.Errorf("Recent()[%d] = %q, want %q", i, got[i].Title, w)
}
}
if lim := b.Recent(2); len(lim) != 2 || lim[0].Title != "three" {
t.Errorf("Recent(2) = %v, want the two newest", lim)
}
}
func TestBusRingEvicts(t *testing.T) {
b := NewBus(3)
for _, title := range []string{"a", "b", "c", "d", "e"} {
b.Publish(Event{Source: "s", Kind: KindFact, Title: title}, testNow)
}
if b.Len() != 3 {
t.Fatalf("Len = %d, want the capacity 3", b.Len())
}
got := b.Recent(0)
want := []string{"e", "d", "c"}
for i, w := range want {
if got[i].Title != w {
t.Errorf("Recent()[%d] = %q, want %q", i, got[i].Title, w)
}
}
}
func TestBusDropsInvalid(t *testing.T) {
b := NewBus(4)
b.Publish(Event{Kind: KindFact, Title: "no source"}, testNow)
b.Publish(Event{Source: "s", Kind: KindFact}, testNow)
if b.Len() != 0 {
t.Errorf("Len = %d, want 0 — an envelope with no provenance must not be kept", b.Len())
}
}
func TestBusSubscriberPanicDoesNotBreakIntake(t *testing.T) {
b := NewBus(4)
var seen int
b.Subscribe(func(Event) { panic("observer is broken") })
b.Subscribe(func(Event) { seen++ })
b.Publish(Event{Source: "s", Kind: KindFact, Title: "t"}, testNow)
if seen != 1 {
t.Errorf("healthy subscriber called %d times, want 1", seen)
}
if b.Len() != 1 {
t.Errorf("event not recorded despite a panicking subscriber")
}
}
func TestBusConcurrentPublish(t *testing.T) {
b := NewBus(256)
var wg sync.WaitGroup
for i := 0; i < 16; i++ {
wg.Add(1)
go func() {
defer wg.Done()
for j := 0; j < 10; j++ {
b.Publish(Event{Source: "s", Kind: KindFact, Title: "t"}, testNow)
}
}()
}
wg.Wait()
if b.Len() != 160 {
t.Errorf("Len = %d, want 160", b.Len())
}
}
+25
View File
@@ -660,6 +660,31 @@ type CoreAPI interface {
// (router → dialogue → action → replier) and returns the reply text.
// No audio or stt/tts — for text channels (mavweb, telegram).
Chat(ctx context.Context, text string) (string, error)
// RecentEvents returns the daemon's unified intake journal, newest first
// (Vikunja #283) — one envelope per thing that arrived, whatever direction
// it came from: a relayed notification, a mail candidate, a feed item, a
// changed page, a spend, a presence probe.
//
// Read-only and daemon-cached, the same shape as TickTrace and DayPlan:
// the store adapter returns an error, because the journal is a bounded
// in-memory ring and not a table. Its contents are a window over intake,
// never the durable record — that is still the fact, note or task the
// intake path wrote.
RecentEvents(ctx context.Context, n int) ([]IntakeEvent, error)
}
// IntakeEvent — one entry of the unified intake journal on the wire. Mirrors
// event.Event field for field; the ipc package does not import internal/event
// so the wire shape stays independent of the in-process type.
type IntakeEvent struct {
Source string `json:"source"`
Kind string `json:"kind"`
EntityIDs []string `json:"entity_ids,omitempty"`
Title string `json:"title"`
Body string `json:"body,omitempty"`
Priority string `json:"priority"`
OccurredAt time.Time `json:"occurred_at"`
}
// --- Rule trace / explanation DTOs ---
+9
View File
@@ -74,6 +74,7 @@ var readOnlyMethods = map[Method]bool{
MethodMorningStatus: true,
MethodMCPServers: true,
MethodDayPlan: true,
MethodRecentEvents: true,
}
// Dial connects to a core socket at path and returns a Client. The module
@@ -593,6 +594,14 @@ func (c *Client) TickTrace(ctx context.Context) (TickTrace, error) {
return t, nil
}
func (c *Client) RecentEvents(ctx context.Context, n int) ([]IntakeEvent, error) {
var e []IntakeEvent
if err := c.call(ctx, MethodRecentEvents, nReq{N: n}, &e); err != nil {
return nil, err
}
return e, nil
}
func (c *Client) MCPServers(ctx context.Context) ([]MCPServerStatus, error) {
var s []MCPServerStatus
if err := c.call(ctx, MethodMCPServers, nil, &s); err != nil {
+16
View File
@@ -211,6 +211,12 @@ func (a *storeAPI) MorningStatus(ctx context.Context) ([]MorningRoutineStatus, e
return nil, errors.New("store: morning status not available via direct store API")
}
// RecentEvents — same shape as TickTrace: the intake journal is a bounded ring
// in the daemon's memory, not a table, so a bare store cannot serve it.
func (a *storeAPI) RecentEvents(ctx context.Context, n int) ([]IntakeEvent, error) {
return nil, errors.New("store: intake events not available via direct store API")
}
func (a *storeAPI) MCPServers(ctx context.Context) ([]MCPServerStatus, error) {
return nil, nil // no manager behind a bare store: nothing configured
}
@@ -884,6 +890,16 @@ var methodTable = map[Method]handlerFunc{
MethodMorningStatus: withoutParams(func(ctx context.Context, api CoreAPI) ([]MorningRoutineStatus, error) {
return api.MorningStatus(ctx)
}),
MethodRecentEvents: withParams(func(ctx context.Context, api CoreAPI, p nReq) ([]IntakeEvent, error) {
out, err := api.RecentEvents(ctx, p.N)
if err != nil {
return nil, err
}
if out == nil {
out = []IntakeEvent{}
}
return out, nil
}),
MethodMCPServers: withoutParams(func(ctx context.Context, api CoreAPI) ([]MCPServerStatus, error) {
out, err := api.MCPServers(ctx)
if err != nil {
+3
View File
@@ -122,6 +122,9 @@ func (UnimplementedCoreAPI) TickTrace(ctx context.Context) (TickTrace, error) {
func (UnimplementedCoreAPI) MorningStatus(ctx context.Context) ([]MorningRoutineStatus, error) {
return nil, ErrNotImplemented
}
func (UnimplementedCoreAPI) RecentEvents(ctx context.Context, n int) ([]IntakeEvent, error) {
return nil, ErrNotImplemented
}
func (UnimplementedCoreAPI) MCPServers(ctx context.Context) ([]MCPServerStatus, error) {
return nil, ErrNotImplemented
}
+1
View File
@@ -62,6 +62,7 @@ const (
MethodEnrollSpeaker Method = "enroll_speaker"
MethodListSpeakers Method = "list_speakers"
MethodForgetSpeaker Method = "forget_speaker"
MethodRecentEvents Method = "recent_events"
)
// Request — one frame from module to core. Params is the JSON-encoded argument