Skip to content

Commit d1ba73e

Browse files
authored
Add observability for memory + agent-loop lifecycle signals (#18)
1 parent 6d2cecc commit d1ba73e

17 files changed

Lines changed: 1055 additions & 22 deletions

File tree

AGENTS.md

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -43,11 +43,11 @@ cmd/odek/
4343
*_test.go 200+ unit + E2E tests covering all tools
4444
internal/
4545
llm/ OpenAI-compatible HTTP client with reasoning_content support
46-
loop/ ReAct engine: observe → think → parallel-act → repeat
46+
loop/ ReAct engine: observe → think → parallel-act → repeat. signal.go — SignalEvent observability (context_trimmed, tool_recovery).
4747
tool/ Thread-safe tool registry, clarify.go, send_message.go
4848
danger/ Command/URL classification + bypass-resistant tokenizer (substitution, $IFS, wrappers, basenames). TTYApprover with friction mode.
4949
auth/ Interactive approval system
50-
memory/ MemoryManager (facts, buffer, episodes, merge, scan). EpisodeProvenance — tainted episodes never auto-replayed.
50+
memory/ MemoryManager (facts, buffer, episodes, merge, scan). EpisodeProvenance — tainted episodes never auto-replayed. notifier.go — MemoryEvent lifecycle observability (fact/episode events fan out to terminal/WebUI/Telegram/programmatic handler).
5151
session/ Session store (CRUD, trim, cleanup, compact JSON). AuditStore + divergence heuristic.
5252
skills/ Skill system (types, loader, triggers, self-improve, curator, import, cache). SkillProvenance gate — NeedsReview skills pinned to Lazy.
5353
config/ Config file loading, env vars, secrets.env, priority merge

README.md

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -45,7 +45,7 @@ Parallel OS-process sub-agents via `delegate_tasks`. True isolation — each sub
4545
Skill-matched `SKILL.md` files load on-demand. Auto-learns from patterns every session — detects multi-step procedures, error recoveries, repeated actions, and user corrections. **LLM-enhanced**: each detected pattern is enriched with an LLM-generated name, description, trigger keywords, and structured body with overview, steps, pitfalls, and verification sections. Use `--no-learn` to disable. Import skills from any URI with automatic LLM risk assessment. [docs/CLI.md#skills](docs/CLI.md#skills)
4646

4747
### 💾 Persistent Memory
48-
Three tiers: **facts** (agent-managed durable entries), **session buffer** (auto-appended turn summaries), **episodes** (LLM-extracted knowledge from past sessions). Merge-on-write via go-vector RandomProjections — cosine >0.7 auto-merges, <0.3 auto-adds. Saves ~80% LLM calls. [docs/MEMORY.md](docs/MEMORY.md)
48+
Three tiers: **facts** (agent-managed durable entries), **session buffer** (auto-appended turn summaries), **episodes** (LLM-extracted knowledge from past sessions). Merge-on-write via go-vector RandomProjections — cosine >0.7 auto-merges, <0.3 auto-adds. Saves ~80% LLM calls. Every lifecycle moment (fact add/merge/consolidate, episode store/dedup/evict/promote) emits an observable event surfaced in the terminal (verbose), Web UI, Telegram, or a programmatic `MemoryEventHandler`. [docs/MEMORY.md](docs/MEMORY.md)
4949

5050
### 🔧 Multi-Turn Sessions
5151
Save, resume, list, trim, and clean up conversations. Sessions persist as JSON in `~/.odek/sessions/`. Continue any session with `odek continue`. [docs/SESSIONS.md](docs/SESSIONS.md)

cmd/odek/main.go

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -827,6 +827,11 @@ func run(args []string) error {
827827
rend.WithSkillVerbose(resolved.Skills.Verbose)
828828
}
829829

830+
// Surface memory lifecycle + agent-signal notifications in verbose mode so
831+
// fact/episode activity and silent recoveries (context trim, tool recovery)
832+
// are observable without flooding the default terminal output.
833+
rend.WithMemoryVerbose(resolved.InteractionMode == "verbose")
834+
830835
// Resolve skills config pointer (only when learn mode is enabled)
831836
var skillsCfg *skills.SkillsConfig
832837
if resolved.Skills.Learn {

cmd/odek/serve.go

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -384,6 +384,27 @@ func newServeAgent(resolved config.ResolvedConfig, system string, sendFn func(v
384384
"heuristic": event.Heuristic,
385385
})
386386
},
387+
MemoryEventHandler: func(event memory.MemoryEvent) {
388+
sendFn(map[string]any{
389+
"type": "memory_event",
390+
"event": event.Type,
391+
"target": event.Target,
392+
"session_id": event.SessionID,
393+
"content": event.Content,
394+
"count": event.Count,
395+
"new_count": event.NewCount,
396+
"untrusted": event.Untrusted,
397+
})
398+
},
399+
AgentSignalHandler: func(event loop.SignalEvent) {
400+
sendFn(map[string]any{
401+
"type": "agent_signal",
402+
"event": event.Type,
403+
"detail": event.Detail,
404+
"tool": event.Tool,
405+
"count": event.Count,
406+
})
407+
},
387408
// Stream thinking/reasoning content to the WebUI.
388409
// Only fire for pre-tool iterations (reasoning before tool calls);
389410
// post-tool callbacks have no new reasoning to display.

cmd/odek/telegram.go

Lines changed: 46 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@ import (
2121
"github.com/BackendStack21/odek/internal/config"
2222
"github.com/BackendStack21/odek/internal/llm"
2323
"github.com/BackendStack21/odek/internal/loop"
24+
"github.com/BackendStack21/odek/internal/memory"
2425
"github.com/BackendStack21/odek/internal/render"
2526
"github.com/BackendStack21/odek/internal/schedule"
2627
"github.com/BackendStack21/odek/internal/session"
@@ -1526,6 +1527,51 @@ func handleChatMessage(
15261527
&telegram.SendOpts{ReplyMarkup: replyMarkup, ParseMode: "Markdown", ReplyToMessageID: messageID})
15271528
}
15281529
},
1530+
MemoryEventHandler: func(event memory.MemoryEvent) {
1531+
// Only surface memory activity in verbose mode, and only the
1532+
// operationally meaningful events (skip high-frequency/internal ones
1533+
// like episode_stored and episode_deduped to avoid chat noise).
1534+
if skillsCfg == nil || !skillsCfg.Verbose {
1535+
return
1536+
}
1537+
var msg string
1538+
switch event.Type {
1539+
case "fact_added":
1540+
msg = "🧠 Memory fact added (" + event.Target + ")"
1541+
case "fact_merged":
1542+
msg = "🧠 Memory fact merged (" + event.Target + ")"
1543+
case "fact_replaced":
1544+
msg = "🧠 Memory fact updated (" + event.Target + ")"
1545+
case "fact_removed":
1546+
msg = "🧠 Memory fact removed (" + event.Target + ")"
1547+
case "fact_consolidated":
1548+
msg = fmt.Sprintf("🧠 Memory consolidated (%s: %d → %d)", event.Target, event.Count, event.NewCount)
1549+
case "episode_evicted":
1550+
msg = fmt.Sprintf("💾 %d episode(s) evicted", event.Count)
1551+
case "episode_pending_review":
1552+
msg = "🔒 Episode pending review (untrusted): " + event.SessionID
1553+
case "episode_promoted":
1554+
msg = "💾 Episode promoted: " + event.SessionID
1555+
default:
1556+
return
1557+
}
1558+
sendAsync(bot, chatID, msg, &telegram.SendOpts{ReplyToMessageID: messageID})
1559+
},
1560+
AgentSignalHandler: func(event loop.SignalEvent) {
1561+
if skillsCfg == nil || !skillsCfg.Verbose {
1562+
return
1563+
}
1564+
var msg string
1565+
switch event.Type {
1566+
case "context_trimmed":
1567+
msg = fmt.Sprintf("✂️ Context trimmed (%s): %d group(s) dropped", event.Detail, event.Count)
1568+
case "tool_recovery":
1569+
msg = "🔁 Tool recovery: " + event.Tool
1570+
default:
1571+
return
1572+
}
1573+
sendAsync(bot, chatID, msg, &telegram.SendOpts{ReplyToMessageID: messageID})
1574+
},
15291575
Approver: approver,
15301576
DangerousConfig: &resolved.Dangerous,
15311577
}

cmd/odek/ui/app.js

Lines changed: 58 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -316,6 +316,14 @@ function connect() {
316316
case 'skill_event':
317317
handleSkillEvent(event);
318318
break;
319+
320+
case 'memory_event':
321+
handleMemoryEvent(event);
322+
break;
323+
324+
case 'agent_signal':
325+
handleAgentSignal(event);
326+
break;
319327
}
320328
};
321329
}
@@ -1079,6 +1087,56 @@ function handleSkillEvent(event) {
10791087
}
10801088
}
10811089

1090+
// ── Memory Events ──
1091+
function handleMemoryEvent(event) {
1092+
switch (event.event) {
1093+
case 'fact_added':
1094+
showToast('🧠 Memory fact added (' + (event.target || '') + ')');
1095+
break;
1096+
case 'fact_merged':
1097+
showToast('🧠 Memory fact merged (' + (event.target || '') + ')');
1098+
break;
1099+
case 'fact_replaced':
1100+
showToast('🧠 Memory fact updated (' + (event.target || '') + ')');
1101+
break;
1102+
case 'fact_removed':
1103+
showToast('🧠 Memory fact removed (' + (event.target || '') + ')');
1104+
break;
1105+
case 'fact_consolidated':
1106+
showToast('🧠 Memory consolidated (' + (event.target || '') + ': ' +
1107+
(event.count || 0) + ' → ' + (event.new_count || 0) + ')');
1108+
break;
1109+
case 'episode_stored':
1110+
// Silent by default — fires after every qualifying session.
1111+
break;
1112+
case 'episode_promoted':
1113+
showToast('💾 ✓ Episode promoted: ' + (event.session_id || ''));
1114+
break;
1115+
case 'episode_evicted':
1116+
showToast('💾 ✗ ' + (event.count || 0) + ' episode(s) evicted');
1117+
break;
1118+
case 'episode_pending_review':
1119+
showToast('🔒 Episode pending review (untrusted): ' + (event.session_id || ''));
1120+
break;
1121+
case 'episode_deduped':
1122+
// Silent — internal dedup detail.
1123+
break;
1124+
}
1125+
}
1126+
1127+
// ── Agent Signals ──
1128+
function handleAgentSignal(event) {
1129+
switch (event.event) {
1130+
case 'context_trimmed':
1131+
showToast('✂️ Context trimmed (' + (event.detail || '') + '): ' +
1132+
(event.count || 0) + ' group(s) dropped');
1133+
break;
1134+
case 'tool_recovery':
1135+
showToast('🔁 Tool recovery: ' + (event.tool || ''));
1136+
break;
1137+
}
1138+
}
1139+
10821140
// ── New Session ──
10831141
// Saved on first load so we can restore the empty state after clearing.
10841142
let savedEmptyStateNode = null;

docs/MEMORY.md

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -166,6 +166,39 @@ All memory content is scanned on write for:
166166

167167
Rejected content returns an error to the agent.
168168

169+
## Observability (lifecycle events)
170+
171+
Every memory lifecycle moment emits a `memory.MemoryEvent` so operators can see
172+
activity that was previously silent. Events fan out (via `MultiMemoryNotifier`)
173+
to whichever surfaces are wired:
174+
175+
- **Terminal** — shown in verbose interaction mode (`--interaction verbose`),
176+
e.g. `🧠 memory[user] added: ...`, `🧠 consolidated memory[env] (5 → 2 entries)`.
177+
- **Web UI** — streamed over the WebSocket as `memory_event` and surfaced as toasts.
178+
- **Telegram** — posted in the chat when the bot runs verbose.
179+
- **Programmatic** — set `Config.MemoryEventHandler` to receive every event.
180+
181+
| Event | Fired when | Key fields |
182+
|---|---|---|
183+
| `fact_added` | a new durable fact is appended (not on a silent dedup) | `Target`, `Content` |
184+
| `fact_merged` | merge-on-write folds a fact into a near-duplicate | `Target`, `Content`, `Similarity` |
185+
| `fact_replaced` | an existing fact is replaced | `Target`, `Content` |
186+
| `fact_removed` | a fact is removed | `Target`, `Content` |
187+
| `fact_consolidated` | LLM consolidation merges entries | `Target`, `Count``NewCount` |
188+
| `episode_stored` | a session episode is extracted + persisted | `SessionID`, `Count` (turns), `Untrusted` |
189+
| `episode_deduped` | a new episode replaces a near-duplicate | `SessionID`, `Similarity` |
190+
| `episode_evicted` | episodes pruned by TTL / count cap | `Sessions`, `Count` |
191+
| `episode_promoted` | a tainted episode is user-approved | `SessionID` |
192+
| `episode_pending_review` | an untrusted episode is stored but excluded from recall | `SessionID` |
193+
194+
Notifiers must be non-blocking — fact writes fire mid-loop and episode events
195+
fire from the post-session background goroutines.
196+
197+
The agent loop also emits `loop.SignalEvent`s for previously-silent self-healing
198+
(`context_trimmed` when message groups are dropped to fit the context window,
199+
`tool_recovery` when a repeatedly-failing tool triggers a corrective hint),
200+
surfaced the same way via `Config.AgentSignalHandler`.
201+
169202
## Architecture
170203

171204
### Episode Index Caching

internal/loop/loop.go

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -70,6 +70,7 @@ type Engine struct {
7070
episodeCtx EpisodeContextFunc // optional: per-turn episode search
7171

7272
toolEventHandler ToolEventHandler // optional: fires during tool execution
73+
signalHandler SignalHandler // optional: fires on internal loop signals
7374

7475
// interactionMode controls how progress is surfaced to the user.
7576
// "engaging" (default), "verbose", "enhance", or "off" (silent).
@@ -367,6 +368,12 @@ func (e *Engine) trimContext(messages []llm.Message, toolDefs []llm.ToolDef) []l
367368
newMsgs = append(newMsgs, trimMsg)
368369
newMsgs = append(newMsgs, messages[1:]...)
369370
messages = newMsgs
371+
372+
e.emitSignal(SignalEvent{
373+
Type: "context_trimmed",
374+
Detail: "proactive",
375+
Count: droppedGroups,
376+
})
370377
}
371378

372379
return messages
@@ -672,6 +679,11 @@ func (e *Engine) runLoop(ctx context.Context, messages []llm.Message) (string, [
672679
if isContextLengthError(err) {
673680
trimmed := trimToSurvival(messages)
674681
if len(trimmed) < len(messages) {
682+
e.emitSignal(SignalEvent{
683+
Type: "context_trimmed",
684+
Detail: "survival",
685+
Count: len(messages) - len(trimmed),
686+
})
675687
messages = trimmed
676688
// Reset memory index — trimToSurvival drops it.
677689
e.memMsgIdx = -1
@@ -1029,6 +1041,11 @@ func (e *Engine) runLoop(ctx context.Context, messages []llm.Message) (string, [
10291041
toolName)
10301042
}
10311043
corrections = append(corrections, correction)
1044+
e.emitSignal(SignalEvent{
1045+
Type: "tool_recovery",
1046+
Tool: toolName,
1047+
Detail: correction,
1048+
})
10321049
// Reset counter after injecting suggestion
10331050
e.maxConsecutiveToolErrors[toolName] = 0
10341051
}

internal/loop/signal.go

Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,45 @@
1+
package loop
2+
3+
import "time"
4+
5+
// SignalEvent represents an internal agent-loop signal that was previously
6+
// invisible to the operator — moments where the engine silently intervened to
7+
// keep the session alive or productive. Surfacing these closes observability
8+
// gaps around context management and tool-failure recovery.
9+
//
10+
// Not every field is set for every Type; the zero value means "not applicable".
11+
type SignalEvent struct {
12+
// Type is the signal kind. One of:
13+
// "context_trimmed" — prior message groups were dropped to fit the token
14+
// budget (Count = groups dropped, Detail = "proactive"
15+
// for the pre-call budget trim or "survival" for the
16+
// post-error nuclear trim)
17+
// "tool_recovery" — a tool failed repeatedly and the engine injected a
18+
// corrective hint so the model changes approach
19+
// (Tool = failing tool, Detail = the correction)
20+
Type string
21+
Detail string // human-readable detail (mode, correction text, etc.)
22+
Tool string // tool name for tool_recovery
23+
Count int // groups dropped (context_trimmed)
24+
Timestamp time.Time // when the signal fired (UTC)
25+
}
26+
27+
// SignalHandler receives agent-loop signal events. Implementations must be
28+
// non-blocking — signals fire inside the hot loop.
29+
type SignalHandler func(event SignalEvent)
30+
31+
// emitSignal fires the engine's signal handler if one is configured, stamping
32+
// the timestamp when the caller left it zero. Safe to call unconditionally.
33+
func (e *Engine) emitSignal(ev SignalEvent) {
34+
if e.signalHandler == nil {
35+
return
36+
}
37+
if ev.Timestamp.IsZero() {
38+
ev.Timestamp = time.Now().UTC()
39+
}
40+
e.signalHandler(ev)
41+
}
42+
43+
// SetSignalHandler sets the optional agent-loop signal callback. Passing nil
44+
// disables signal emission.
45+
func (e *Engine) SetSignalHandler(cb SignalHandler) { e.signalHandler = cb }

internal/loop/signal_test.go

Lines changed: 42 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,42 @@
1+
package loop
2+
3+
import "testing"
4+
5+
func TestEmitSignal_NilHandlerIsSafe(t *testing.T) {
6+
e := &Engine{}
7+
// No handler set — must be a no-op, not a panic.
8+
e.emitSignal(SignalEvent{Type: "context_trimmed"})
9+
}
10+
11+
func TestSetSignalHandler_ReceivesEventsAndStampsTime(t *testing.T) {
12+
e := &Engine{}
13+
var got []SignalEvent
14+
e.SetSignalHandler(func(ev SignalEvent) { got = append(got, ev) })
15+
16+
e.emitSignal(SignalEvent{Type: "context_trimmed", Detail: "survival", Count: 3})
17+
e.emitSignal(SignalEvent{Type: "tool_recovery", Tool: "shell", Detail: "try a different approach"})
18+
19+
if len(got) != 2 {
20+
t.Fatalf("expected 2 signals, got %d", len(got))
21+
}
22+
if got[0].Type != "context_trimmed" || got[0].Count != 3 || got[0].Detail != "survival" {
23+
t.Errorf("unexpected first signal: %+v", got[0])
24+
}
25+
if got[0].Timestamp.IsZero() {
26+
t.Error("expected timestamp to be stamped on emit")
27+
}
28+
if got[1].Tool != "shell" {
29+
t.Errorf("expected tool=shell, got %q", got[1].Tool)
30+
}
31+
}
32+
33+
func TestSetSignalHandler_NilDisables(t *testing.T) {
34+
e := &Engine{}
35+
called := false
36+
e.SetSignalHandler(func(SignalEvent) { called = true })
37+
e.SetSignalHandler(nil)
38+
e.emitSignal(SignalEvent{Type: "tool_recovery"})
39+
if called {
40+
t.Error("handler should not fire after being set to nil")
41+
}
42+
}

0 commit comments

Comments
 (0)