Normalize every intake path into one event envelope (#283) #78
@@ -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, ""
|
||||
}
|
||||
@@ -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
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
@@ -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)
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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>
|
||||
@@ -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, "<script>") {
|
||||
t.Error("intake title is missing from the page entirely")
|
||||
}
|
||||
}
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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())
|
||||
}
|
||||
}
|
||||
@@ -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 ---
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user