diff --git a/cmd/mavend/intake.go b/cmd/mavend/intake.go new file mode 100644 index 0000000..0e7542d --- /dev/null +++ b/cmd/mavend/intake.go @@ -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, "" +} diff --git a/cmd/mavend/intake_test.go b/cmd/mavend/intake_test.go new file mode 100644 index 0000000..1105628 --- /dev/null +++ b/cmd/mavend/intake_test.go @@ -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)) + } +} diff --git a/cmd/mavend/mail.go b/cmd/mavend/mail.go index ebde48a..5167e71 100644 --- a/cmd/mavend/mail.go +++ b/cmd/mavend/mail.go @@ -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 } diff --git a/cmd/mavend/mail_test.go b/cmd/mavend/mail_test.go index d06bb70..05a5112 100644 --- a/cmd/mavend/mail_test.go +++ b/cmd/mavend/mail_test.go @@ -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") } } diff --git a/cmd/mavend/main.go b/cmd/mavend/main.go index 0702ac0..4913d91 100644 --- a/cmd/mavend/main.go +++ b/cmd/mavend/main.go @@ -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) diff --git a/cmd/mavend/tick.go b/cmd/mavend/tick.go index 746cbe7..b3b7c7d 100644 --- a/cmd/mavend/tick.go +++ b/cmd/mavend/tick.go @@ -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) { diff --git a/cmd/mavweb/events.html b/cmd/mavweb/events.html new file mode 100644 index 0000000..df655e9 --- /dev/null +++ b/cmd/mavweb/events.html @@ -0,0 +1,24 @@ +{{template "shellTop" "events"}} +

Intake

+
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.
+{{if .Err}}
journal unavailable: {{.Err}}
{{end}} +{{if and (not .Events) (not .Err)}} +
nothing has arrived yet
+{{end}} +{{if .Events}} +
+ +{{range .Events}} + + + + + + +{{end}} +
whensourcekindpriwhatdetail
{{.OccurredAt.Format "02.01 15:04:05"}}{{.Source}}{{.Kind}}{{.Priority}}{{.Title}}{{.Body}}
+{{end}} +{{template "shellBottom"}} + diff --git a/cmd/mavweb/events_test.go b/cmd/mavweb/events_test.go new file mode 100644 index 0000000..2d86102 --- /dev/null +++ b/cmd/mavweb/events_test.go @@ -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: ``, + OccurredAt: time.Date(2026, 8, 1, 7, 0, 0, 0, time.UTC), + }}} + body := getEvents(t, core).Body.String() + if strings.Contains(body, "") { + t.Error("intake title was not escaped") + } + if !strings.Contains(body, "<script>") { + t.Error("intake title is missing from the page entirely") + } +} diff --git a/cmd/mavweb/main.go b/cmd/mavweb/main.go index 29be4de..17e49af 100644 --- a/cmd/mavweb/main.go +++ b/cmd/mavweb/main.go @@ -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 { diff --git a/internal/auth/policy.go b/internal/auth/policy.go index cc480e6..f940f80 100644 --- a/internal/auth/policy.go +++ b/internal/auth/policy.go @@ -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 diff --git a/internal/config/config.go b/internal/config/config.go index 7a0b6cf..8865711 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -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) } diff --git a/internal/event/bus.go b/internal/event/bus.go new file mode 100644 index 0000000..5b29a16 --- /dev/null +++ b/internal/event/bus.go @@ -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 +} diff --git a/internal/event/event.go b/internal/event/event.go new file mode 100644 index 0000000..add6420 --- /dev/null +++ b/internal/event/event.go @@ -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:", + // "crawl:", "email:", "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 +} diff --git a/internal/event/event_test.go b/internal/event/event_test.go new file mode 100644 index 0000000..0eaec48 --- /dev/null +++ b/internal/event/event_test.go @@ -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()) + } +} diff --git a/internal/ipc/api.go b/internal/ipc/api.go index 7c43b5f..30cc4ab 100644 --- a/internal/ipc/api.go +++ b/internal/ipc/api.go @@ -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 --- diff --git a/internal/ipc/client.go b/internal/ipc/client.go index ad48f3b..fae2f4d 100644 --- a/internal/ipc/client.go +++ b/internal/ipc/client.go @@ -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 { diff --git a/internal/ipc/server.go b/internal/ipc/server.go index dce43bb..cee5355 100644 --- a/internal/ipc/server.go +++ b/internal/ipc/server.go @@ -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 { diff --git a/internal/ipc/unimplemented.go b/internal/ipc/unimplemented.go index 8c289cb..20e1b9f 100644 --- a/internal/ipc/unimplemented.go +++ b/internal/ipc/unimplemented.go @@ -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 } diff --git a/internal/ipc/wire.go b/internal/ipc/wire.go index a3d1e53..fc0d7f5 100644 --- a/internal/ipc/wire.go +++ b/internal/ipc/wire.go @@ -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