feat(session): add CompactOldTurns, SessionGraph, and session lifecycle (TASKS-3 Phase 2)
- CompactOldTurns: turn-granularity compaction replacing full-rewrite path - SessionGraph + TurnWriter: thin wrapper for future staged migration - Startup Prune (7d TTL) + periodic prune (6h) in flushLoop - AgentLoop gcLoop: 30min GC of idle sessionLocks entries - Wire CompactOldTurns into summarizeSession with fallback Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
parent
934027affb
commit
df4e65069b
8 changed files with 536 additions and 27 deletions
10
CLAUDE.md
10
CLAUDE.md
|
|
@ -26,16 +26,18 @@ Lint: `golangci-lint run`
|
||||||
- **Sandbox/Spawn**: `pkg/tools/sandbox.go`, `pkg/tools/spawn.go` 実装済み。
|
- **Sandbox/Spawn**: `pkg/tools/sandbox.go`, `pkg/tools/spawn.go` 実装済み。
|
||||||
- **AgentReporter**: `orch.AgentReporter` / `orch.Noop` / `orch.Broadcaster` で統一。main/heartbeat/subagent 全セッションが同一 Broadcaster に発火。Mini App は `agentLoop.GetOrchBroadcaster()` → `handler.SetOrchBroadcaster()` で受信。
|
- **AgentReporter**: `orch.AgentReporter` / `orch.Noop` / `orch.Broadcaster` で統一。main/heartbeat/subagent 全セッションが同一 Broadcaster に発火。Mini App は `agentLoop.GetOrchBroadcaster()` → `handler.SetOrchBroadcaster()` で受信。
|
||||||
|
|
||||||
## Session DAG (Phase 0–1 実装済み)
|
## Session DAG (Phase 0–2 実装済み)
|
||||||
|
|
||||||
- **SQLite SessionStore**: `pkg/session/sqlite.go` — `modernc.org/sqlite` (CGO不要)、WAL モード、`sessions` + `turns` テーブル
|
- **SQLite SessionStore**: `pkg/session/sqlite.go` — `modernc.org/sqlite` (CGO不要)、WAL モード、`sessions` + `turns` テーブル
|
||||||
- **SessionStore interface**: `pkg/session/store.go` — Create/Get/List/Append/Turns/Compact/Fork/Prune 等15メソッド
|
- **SessionStore interface**: `pkg/session/store.go` — Create/Get/List/Append/Turns/Compact/Fork/Prune 等15メソッド
|
||||||
- **LegacyAdapter**: `pkg/session/legacy_adapter.go` — SessionStore をラップし SessionManager と同一 API を提供。`Store()` / `AdvanceStored()` で直接 DAG 操作も可能
|
- **LegacyAdapter**: `pkg/session/legacy_adapter.go` — SessionStore をラップし SessionManager と同一 API を提供。`Store()` / `AdvanceStored()` で直接 DAG 操作も可能
|
||||||
|
- **CompactOldTurns**: `LegacyAdapter.CompactOldTurns(key, keepLast, summary)` — flush → turn 単位で cut point 算出 → SQLite Compact → キャッシュ更新。`summarizeSession()` から呼び出し (fallback 付き)
|
||||||
|
- **SessionGraph**: `pkg/session/graph.go` — `SessionGraph` + `BeginTurn()` / `TurnWriter`。`LegacyAdapter.Graph()` で取得。将来の段階移行準備
|
||||||
- **JSON → SQLite migration**: `pkg/session/migrate.go` — 起動時に `sessions/*.json` を検出 → SQLite import → `.json.migrated` にリネーム
|
- **JSON → SQLite migration**: `pkg/session/migrate.go` — 起動時に `sessions/*.json` を検出 → SQLite import → `.json.migrated` にリネーム
|
||||||
- **配線**: `pkg/agent/instance.go` の `Sessions` 型が `*LegacyAdapter` に変更。`sessions.db` を workspace 直下に生成
|
- **配線**: `pkg/agent/instance.go` の `Sessions` 型が `*LegacyAdapter` に変更。`sessions.db` を workspace 直下に生成
|
||||||
- **SessionRecorder**: `pkg/tools/session_recorder.go` (interface) + `pkg/agent/session_recorder.go` (impl) — SubagentManager から Fork/Turn/Completion/Report を記録。循環依存回避のためインターフェースは pkg/tools 側
|
- **SessionRecorder**: `pkg/tools/session_recorder.go` (interface) + `pkg/agent/session_recorder.go` (impl) — SubagentManager から Fork/Turn/Completion/Report を記録
|
||||||
- **Fork/Report フロー**: `SubagentManager.Spawn()` で Fork 記録、`runTask()` で Turn + Completion 記録、`processSystemMessage()` で TurnReport を直接 store に書き込み + `AdvanceStored` で二重書き込み防止
|
- **Prune**: 起動時 `store.Prune(7d)` + `flushLoop` 内 6h 定期 prune。`AgentLoop.gcLoop` で 30分毎に idle `sessionLocks` を GC
|
||||||
- **Phase 2以降**: SessionGraph 直接呼び出し、Compaction、Mini App 可視化 → `todo/TASKS-3.md` 参照
|
- **Phase 3以降**: UI/CLI — Mini App セッショングラフ可視化、`/session` コマンド → `todo/TASKS-3.md` 参照
|
||||||
|
|
||||||
## コードの匂い — チェックリスト
|
## コードの匂い — チェックリスト
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -143,6 +143,12 @@ func NewAgentInstance(
|
||||||
log.Printf("session migration: %d sessions migrated to SQLite", n)
|
log.Printf("session migration: %d sessions migrated to SQLite", n)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if n, perr := store.Prune(session.DefaultPruneTTL); perr != nil {
|
||||||
|
log.Printf("session prune error: %v", perr)
|
||||||
|
} else if n > 0 {
|
||||||
|
log.Printf("session prune: %d old sessions removed", n)
|
||||||
|
}
|
||||||
|
|
||||||
sessionsManager := session.NewLegacyAdapter(store)
|
sessionsManager := session.NewLegacyAdapter(store)
|
||||||
|
|
||||||
contextBuilder := NewContextBuilder(workspace)
|
contextBuilder := NewContextBuilder(workspace)
|
||||||
|
|
|
||||||
|
|
@ -106,6 +106,7 @@ type AgentLoop struct {
|
||||||
onHeartbeatThreadUpdate func(int)
|
onHeartbeatThreadUpdate func(int)
|
||||||
orchBroadcaster *orch.Broadcaster // nil when --orchestration not set
|
orchBroadcaster *orch.Broadcaster // nil when --orchestration not set
|
||||||
orchReporter orch.AgentReporter // always non-nil (Noop when disabled)
|
orchReporter orch.AgentReporter // always non-nil (Noop when disabled)
|
||||||
|
done chan struct{} // closed by Close() to stop background goroutines
|
||||||
}
|
}
|
||||||
|
|
||||||
// processOptions configures how a message is processed
|
// processOptions configures how a message is processed
|
||||||
|
|
@ -178,11 +179,14 @@ func NewAgentLoop(
|
||||||
sessions: NewSessionTracker(),
|
sessions: NewSessionTracker(),
|
||||||
orchBroadcaster: orchBroadcaster,
|
orchBroadcaster: orchBroadcaster,
|
||||||
orchReporter: orchReporter,
|
orchReporter: orchReporter,
|
||||||
|
done: make(chan struct{}),
|
||||||
}
|
}
|
||||||
|
|
||||||
// Register shared tools to all agents (needs al for reporter injection).
|
// Register shared tools to all agents (needs al for reporter injection).
|
||||||
registerSharedTools(cfg, msgBus, registry, provider, al)
|
registerSharedTools(cfg, msgBus, registry, provider, al)
|
||||||
|
|
||||||
|
go al.gcLoop()
|
||||||
|
|
||||||
return al
|
return al
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -513,6 +517,12 @@ func (al *AgentLoop) Stop() {
|
||||||
// Close releases resources held by the loop (e.g. flushes write-behind stats
|
// Close releases resources held by the loop (e.g. flushes write-behind stats
|
||||||
// and dirty session data). Should be called during graceful shutdown.
|
// and dirty session data). Should be called during graceful shutdown.
|
||||||
func (al *AgentLoop) Close() {
|
func (al *AgentLoop) Close() {
|
||||||
|
select {
|
||||||
|
case <-al.done:
|
||||||
|
// already closed
|
||||||
|
default:
|
||||||
|
close(al.done)
|
||||||
|
}
|
||||||
if al.stats != nil {
|
if al.stats != nil {
|
||||||
al.stats.Close()
|
al.stats.Close()
|
||||||
}
|
}
|
||||||
|
|
@ -523,6 +533,35 @@ func (al *AgentLoop) Close() {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// gcLoop periodically cleans up stale sessionLock entries.
|
||||||
|
func (al *AgentLoop) gcLoop() {
|
||||||
|
ticker := time.NewTicker(30 * time.Minute)
|
||||||
|
defer ticker.Stop()
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-ticker.C:
|
||||||
|
al.gcSessionLocks()
|
||||||
|
case <-al.done:
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// gcSessionLocks removes unlocked (idle) sessionSemaphore entries from the map.
|
||||||
|
func (al *AgentLoop) gcSessionLocks() {
|
||||||
|
al.sessionLocks.Range(func(key, val any) bool {
|
||||||
|
sem := val.(*sessionSemaphore)
|
||||||
|
select {
|
||||||
|
case <-sem.ch:
|
||||||
|
// Was unlocked — safe to remove
|
||||||
|
al.sessionLocks.Delete(key)
|
||||||
|
default:
|
||||||
|
// Currently locked — in use, keep
|
||||||
|
}
|
||||||
|
return true
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
func (al *AgentLoop) RegisterTool(tool tools.Tool) {
|
func (al *AgentLoop) RegisterTool(tool tools.Tool) {
|
||||||
for _, agentID := range al.registry.ListAgentIDs() {
|
for _, agentID := range al.registry.ListAgentIDs() {
|
||||||
if agent, ok := al.registry.GetAgent(agentID); ok {
|
if agent, ok := al.registry.GetAgent(agentID); ok {
|
||||||
|
|
@ -3441,11 +3480,15 @@ func (al *AgentLoop) summarizeSession(agent *AgentInstance, sessionKey string) {
|
||||||
}
|
}
|
||||||
|
|
||||||
if finalSummary != "" {
|
if finalSummary != "" {
|
||||||
|
if err := agent.Sessions.CompactOldTurns(sessionKey, 4, finalSummary); err != nil {
|
||||||
|
logger.ErrorCF("agent", "CompactOldTurns failed, falling back",
|
||||||
|
map[string]any{"error": err.Error()})
|
||||||
agent.Sessions.SetSummary(sessionKey, finalSummary)
|
agent.Sessions.SetSummary(sessionKey, finalSummary)
|
||||||
agent.Sessions.TruncateHistory(sessionKey, 4)
|
agent.Sessions.TruncateHistory(sessionKey, 4)
|
||||||
agent.Sessions.Save(sessionKey)
|
agent.Sessions.Save(sessionKey)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// summarizeBatch summarizes a batch of messages.
|
// summarizeBatch summarizes a batch of messages.
|
||||||
func (al *AgentLoop) summarizeBatch(
|
func (al *AgentLoop) summarizeBatch(
|
||||||
|
|
|
||||||
103
pkg/session/graph.go
Normal file
103
pkg/session/graph.go
Normal file
|
|
@ -0,0 +1,103 @@
|
||||||
|
package session
|
||||||
|
|
||||||
|
import (
|
||||||
|
"errors"
|
||||||
|
"sync"
|
||||||
|
|
||||||
|
"github.com/sipeed/picoclaw/pkg/providers"
|
||||||
|
)
|
||||||
|
|
||||||
|
// SessionGraph is a thin wrapper around SessionStore that provides
|
||||||
|
// structured turn-writing via BeginTurn/TurnWriter.
|
||||||
|
// It does NOT replace LegacyAdapter — existing call sites remain unchanged.
|
||||||
|
// Future phases will migrate callers to use SessionGraph directly.
|
||||||
|
type SessionGraph struct {
|
||||||
|
store SessionStore
|
||||||
|
}
|
||||||
|
|
||||||
|
// NewSessionGraph creates a SessionGraph backed by the given store.
|
||||||
|
func NewSessionGraph(store SessionStore) *SessionGraph {
|
||||||
|
return &SessionGraph{store: store}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Messages returns all messages for the session by reading turns from the store.
|
||||||
|
func (g *SessionGraph) Messages(sessionKey string) ([]providers.Message, error) {
|
||||||
|
turns, err := g.store.Turns(sessionKey, 0)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
var msgs []providers.Message
|
||||||
|
for _, t := range turns {
|
||||||
|
msgs = append(msgs, t.Messages...)
|
||||||
|
}
|
||||||
|
if msgs == nil {
|
||||||
|
msgs = []providers.Message{}
|
||||||
|
}
|
||||||
|
return msgs, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// BeginTurn starts a new turn that can be built up incrementally
|
||||||
|
// and committed atomically.
|
||||||
|
func (g *SessionGraph) BeginTurn(sessionKey string, kind TurnKind) *TurnWriter {
|
||||||
|
return &TurnWriter{
|
||||||
|
store: g.store,
|
||||||
|
sessionKey: sessionKey,
|
||||||
|
turn: Turn{
|
||||||
|
SessionKey: sessionKey,
|
||||||
|
Kind: kind,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// TurnWriter accumulates messages for a single turn and commits them atomically.
|
||||||
|
type TurnWriter struct {
|
||||||
|
mu sync.Mutex
|
||||||
|
store SessionStore
|
||||||
|
sessionKey string
|
||||||
|
turn Turn
|
||||||
|
committed bool
|
||||||
|
discarded bool
|
||||||
|
}
|
||||||
|
|
||||||
|
// Add appends a message to the pending turn.
|
||||||
|
func (tw *TurnWriter) Add(msg providers.Message) {
|
||||||
|
tw.mu.Lock()
|
||||||
|
defer tw.mu.Unlock()
|
||||||
|
tw.turn.Messages = append(tw.turn.Messages, msg)
|
||||||
|
}
|
||||||
|
|
||||||
|
// SetOrigin sets the origin session key for this turn (e.g. subagent source).
|
||||||
|
func (tw *TurnWriter) SetOrigin(sessionKey string) {
|
||||||
|
tw.mu.Lock()
|
||||||
|
defer tw.mu.Unlock()
|
||||||
|
tw.turn.OriginKey = sessionKey
|
||||||
|
}
|
||||||
|
|
||||||
|
// SetAuthor sets the author field for this turn.
|
||||||
|
func (tw *TurnWriter) SetAuthor(author string) {
|
||||||
|
tw.mu.Lock()
|
||||||
|
defer tw.mu.Unlock()
|
||||||
|
tw.turn.Author = author
|
||||||
|
}
|
||||||
|
|
||||||
|
// Commit writes the accumulated turn to the store.
|
||||||
|
// Returns an error if already committed or discarded.
|
||||||
|
func (tw *TurnWriter) Commit() error {
|
||||||
|
tw.mu.Lock()
|
||||||
|
defer tw.mu.Unlock()
|
||||||
|
if tw.committed {
|
||||||
|
return errors.New("turn already committed")
|
||||||
|
}
|
||||||
|
if tw.discarded {
|
||||||
|
return errors.New("turn already discarded")
|
||||||
|
}
|
||||||
|
tw.committed = true
|
||||||
|
return tw.store.Append(tw.sessionKey, &tw.turn)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Discard marks the turn as abandoned — nothing is written.
|
||||||
|
func (tw *TurnWriter) Discard() {
|
||||||
|
tw.mu.Lock()
|
||||||
|
defer tw.mu.Unlock()
|
||||||
|
tw.discarded = true
|
||||||
|
}
|
||||||
145
pkg/session/graph_test.go
Normal file
145
pkg/session/graph_test.go
Normal file
|
|
@ -0,0 +1,145 @@
|
||||||
|
package session
|
||||||
|
|
||||||
|
import (
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
"github.com/sipeed/picoclaw/pkg/providers"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestSessionGraph_Messages(t *testing.T) {
|
||||||
|
store := newTestStore(t)
|
||||||
|
if err := store.Create("g1", nil); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if err := store.Append("g1", &Turn{
|
||||||
|
Kind: TurnNormal,
|
||||||
|
Messages: []providers.Message{{Role: "user", Content: "hello"}},
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if err := store.Append("g1", &Turn{
|
||||||
|
Kind: TurnNormal,
|
||||||
|
Messages: []providers.Message{
|
||||||
|
{Role: "assistant", Content: "hi"},
|
||||||
|
{Role: "user", Content: "how are you"},
|
||||||
|
},
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
|
||||||
|
g := NewSessionGraph(store)
|
||||||
|
msgs, err := g.Messages("g1")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if len(msgs) != 3 {
|
||||||
|
t.Fatalf("expected 3 messages, got %d", len(msgs))
|
||||||
|
}
|
||||||
|
if msgs[0].Content != "hello" || msgs[1].Content != "hi" || msgs[2].Content != "how are you" {
|
||||||
|
t.Errorf("unexpected messages: %+v", msgs)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestSessionGraph_Messages_Empty(t *testing.T) {
|
||||||
|
store := newTestStore(t)
|
||||||
|
if err := store.Create("empty", nil); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
g := NewSessionGraph(store)
|
||||||
|
msgs, err := g.Messages("empty")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if msgs == nil || len(msgs) != 0 {
|
||||||
|
t.Errorf("expected empty slice, got %v", msgs)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestTurnWriter_Commit(t *testing.T) {
|
||||||
|
store := newTestStore(t)
|
||||||
|
if err := store.Create("tw1", nil); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
|
||||||
|
g := NewSessionGraph(store)
|
||||||
|
tw := g.BeginTurn("tw1", TurnNormal)
|
||||||
|
tw.Add(providers.Message{Role: "user", Content: "msg1"})
|
||||||
|
tw.Add(providers.Message{Role: "assistant", Content: "msg2"})
|
||||||
|
tw.SetOrigin("parent-key")
|
||||||
|
tw.SetAuthor("agent-1")
|
||||||
|
|
||||||
|
if err := tw.Commit(); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
|
||||||
|
turns, err := store.Turns("tw1", 0)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if len(turns) != 1 {
|
||||||
|
t.Fatalf("expected 1 turn, got %d", len(turns))
|
||||||
|
}
|
||||||
|
if len(turns[0].Messages) != 2 {
|
||||||
|
t.Fatalf("expected 2 messages, got %d", len(turns[0].Messages))
|
||||||
|
}
|
||||||
|
if turns[0].OriginKey != "parent-key" {
|
||||||
|
t.Errorf("expected origin 'parent-key', got %q", turns[0].OriginKey)
|
||||||
|
}
|
||||||
|
if turns[0].Author != "agent-1" {
|
||||||
|
t.Errorf("expected author 'agent-1', got %q", turns[0].Author)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestTurnWriter_Discard(t *testing.T) {
|
||||||
|
store := newTestStore(t)
|
||||||
|
if err := store.Create("tw2", nil); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
|
||||||
|
g := NewSessionGraph(store)
|
||||||
|
tw := g.BeginTurn("tw2", TurnNormal)
|
||||||
|
tw.Add(providers.Message{Role: "user", Content: "should not persist"})
|
||||||
|
tw.Discard()
|
||||||
|
|
||||||
|
turns, err := store.Turns("tw2", 0)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if len(turns) != 0 {
|
||||||
|
t.Errorf("expected 0 turns after discard, got %d", len(turns))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestTurnWriter_DoubleCommit(t *testing.T) {
|
||||||
|
store := newTestStore(t)
|
||||||
|
if err := store.Create("tw3", nil); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
|
||||||
|
g := NewSessionGraph(store)
|
||||||
|
tw := g.BeginTurn("tw3", TurnNormal)
|
||||||
|
tw.Add(providers.Message{Role: "user", Content: "once"})
|
||||||
|
|
||||||
|
if err := tw.Commit(); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if err := tw.Commit(); err == nil {
|
||||||
|
t.Error("expected error on double commit")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestTurnWriter_CommitAfterDiscard(t *testing.T) {
|
||||||
|
store := newTestStore(t)
|
||||||
|
if err := store.Create("tw4", nil); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
|
||||||
|
g := NewSessionGraph(store)
|
||||||
|
tw := g.BeginTurn("tw4", TurnNormal)
|
||||||
|
tw.Add(providers.Message{Role: "user", Content: "x"})
|
||||||
|
tw.Discard()
|
||||||
|
|
||||||
|
if err := tw.Commit(); err == nil {
|
||||||
|
t.Error("expected error on commit after discard")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -459,11 +459,89 @@ func (la *LegacyAdapter) Save(key string) error {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// DefaultPruneTTL is the default time-to-live for session pruning.
|
||||||
|
const DefaultPruneTTL = 7 * 24 * time.Hour
|
||||||
|
|
||||||
|
// CompactOldTurns flushes pending writes, then compacts SQLite turns
|
||||||
|
// keeping only the last keepLast messages. Sets session summary to the given value.
|
||||||
|
func (la *LegacyAdapter) CompactOldTurns(key string, keepLast int, summary string) error {
|
||||||
|
// 1. Flush pending messages to SQLite
|
||||||
|
if err := la.Save(key); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
// 2. Query all turns
|
||||||
|
turns, err := la.store.Turns(key, 0)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
// 3. Count total messages, find cut point
|
||||||
|
totalMsgs := 0
|
||||||
|
for _, t := range turns {
|
||||||
|
totalMsgs += len(t.Messages)
|
||||||
|
}
|
||||||
|
if keepLast >= totalMsgs {
|
||||||
|
// Nothing to compact, just update summary
|
||||||
|
if err := la.store.SetSummary(key, summary); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
la.mu.Lock()
|
||||||
|
if c, ok := la.cache[key]; ok {
|
||||||
|
c.summary = summary
|
||||||
|
}
|
||||||
|
la.mu.Unlock()
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
dropCount := totalMsgs - keepLast
|
||||||
|
accumulated := 0
|
||||||
|
cutSeq := 0
|
||||||
|
for _, t := range turns {
|
||||||
|
accumulated += len(t.Messages)
|
||||||
|
if accumulated <= dropCount {
|
||||||
|
cutSeq = t.Seq
|
||||||
|
} else {
|
||||||
|
break
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if cutSeq == 0 {
|
||||||
|
if err := la.store.SetSummary(key, summary); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
la.mu.Lock()
|
||||||
|
if c, ok := la.cache[key]; ok {
|
||||||
|
c.summary = summary
|
||||||
|
}
|
||||||
|
la.mu.Unlock()
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
// 4. Compact in SQLite
|
||||||
|
if err := la.store.Compact(key, cutSeq, summary); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
// 5. Update in-memory cache
|
||||||
|
la.mu.Lock()
|
||||||
|
defer la.mu.Unlock()
|
||||||
|
if c, ok := la.cache[key]; ok {
|
||||||
|
if keepLast < len(c.messages) {
|
||||||
|
c.messages = c.messages[len(c.messages)-keepLast:]
|
||||||
|
}
|
||||||
|
c.stored = len(c.messages)
|
||||||
|
c.replaced = false
|
||||||
|
c.dirty = false
|
||||||
|
c.summary = summary
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
// Store returns the underlying SessionStore for direct DAG operations.
|
// Store returns the underlying SessionStore for direct DAG operations.
|
||||||
func (la *LegacyAdapter) Store() SessionStore {
|
func (la *LegacyAdapter) Store() SessionStore {
|
||||||
return la.store
|
return la.store
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Graph returns a SessionGraph backed by the underlying store.
|
||||||
|
func (la *LegacyAdapter) Graph() *SessionGraph {
|
||||||
|
return NewSessionGraph(la.store)
|
||||||
|
}
|
||||||
|
|
||||||
// AdvanceStored increments the stored counter for a session by delta,
|
// AdvanceStored increments the stored counter for a session by delta,
|
||||||
// preventing the flush loop from re-persisting messages already written
|
// preventing the flush loop from re-persisting messages already written
|
||||||
// directly to the store (e.g. TurnReport).
|
// directly to the store (e.g. TurnReport).
|
||||||
|
|
@ -494,18 +572,17 @@ func (la *LegacyAdapter) Close() {
|
||||||
}
|
}
|
||||||
|
|
||||||
func (la *LegacyAdapter) flushLoop() {
|
func (la *LegacyAdapter) flushLoop() {
|
||||||
ticker := time.NewTicker(5 * time.Minute)
|
flushTicker := time.NewTicker(5 * time.Minute)
|
||||||
|
pruneTicker := time.NewTicker(6 * time.Hour)
|
||||||
defer ticker.Stop()
|
defer flushTicker.Stop()
|
||||||
|
defer pruneTicker.Stop()
|
||||||
for {
|
for {
|
||||||
select {
|
select {
|
||||||
case <-ticker.C:
|
case <-flushTicker.C:
|
||||||
|
|
||||||
la.FlushDirty()
|
la.FlushDirty()
|
||||||
|
case <-pruneTicker.C:
|
||||||
|
_, _ = la.store.Prune(DefaultPruneTTL)
|
||||||
case <-la.done:
|
case <-la.done:
|
||||||
|
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -482,3 +482,134 @@ func TestBackend_IncrementalSave(t *testing.T) {
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestCompactOldTurns(t *testing.T) {
|
||||||
|
dbPath := filepath.Join(t.TempDir(), "test.db")
|
||||||
|
store, err := OpenSQLiteStore(dbPath)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
la := NewLegacyAdapter(store)
|
||||||
|
defer la.Close()
|
||||||
|
|
||||||
|
la.GetOrCreate("k1")
|
||||||
|
// Turn 1: 2 messages
|
||||||
|
la.AddMessage("k1", "user", "a")
|
||||||
|
la.AddMessage("k1", "assistant", "b")
|
||||||
|
la.Save("k1")
|
||||||
|
// Turn 2: 3 messages
|
||||||
|
la.AddMessage("k1", "user", "c")
|
||||||
|
la.AddMessage("k1", "assistant", "d")
|
||||||
|
la.AddMessage("k1", "user", "e")
|
||||||
|
la.Save("k1")
|
||||||
|
// Turn 3: 2 messages
|
||||||
|
la.AddMessage("k1", "user", "f")
|
||||||
|
la.AddMessage("k1", "assistant", "g")
|
||||||
|
la.Save("k1")
|
||||||
|
|
||||||
|
// Total: 7 messages across 3 turns. keepLast=2 → drop 5 → compact turns 1+2 (5 msgs)
|
||||||
|
if err := la.CompactOldTurns("k1", 2, "test summary"); err != nil {
|
||||||
|
t.Fatalf("CompactOldTurns: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
h := la.GetHistory("k1")
|
||||||
|
if len(h) != 2 {
|
||||||
|
t.Fatalf("expected 2 messages in cache, got %d", len(h))
|
||||||
|
}
|
||||||
|
if h[0].Content != "f" || h[1].Content != "g" {
|
||||||
|
t.Errorf("unexpected messages: %+v", h)
|
||||||
|
}
|
||||||
|
if s := la.GetSummary("k1"); s != "test summary" {
|
||||||
|
t.Errorf("expected summary 'test summary', got %q", s)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Verify in SQLite: only turn 3 remains
|
||||||
|
turns, _ := store.Turns("k1", 0)
|
||||||
|
if len(turns) != 1 {
|
||||||
|
t.Fatalf("expected 1 turn in SQLite, got %d", len(turns))
|
||||||
|
}
|
||||||
|
if len(turns[0].Messages) != 2 {
|
||||||
|
t.Errorf("expected 2 messages in remaining turn, got %d", len(turns[0].Messages))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestCompactOldTurns_NothingToCompact(t *testing.T) {
|
||||||
|
dbPath := filepath.Join(t.TempDir(), "test.db")
|
||||||
|
store, err := OpenSQLiteStore(dbPath)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
la := NewLegacyAdapter(store)
|
||||||
|
defer la.Close()
|
||||||
|
|
||||||
|
la.GetOrCreate("k1")
|
||||||
|
la.AddMessage("k1", "user", "a")
|
||||||
|
la.AddMessage("k1", "assistant", "b")
|
||||||
|
la.Save("k1")
|
||||||
|
|
||||||
|
// keepLast=10 >= total 2 → nothing compacted, summary still updated
|
||||||
|
if err := la.CompactOldTurns("k1", 10, "new summary"); err != nil {
|
||||||
|
t.Fatalf("CompactOldTurns: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
h := la.GetHistory("k1")
|
||||||
|
if len(h) != 2 {
|
||||||
|
t.Fatalf("expected 2 messages, got %d", len(h))
|
||||||
|
}
|
||||||
|
if s := la.GetSummary("k1"); s != "new summary" {
|
||||||
|
t.Errorf("expected 'new summary', got %q", s)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestCompactOldTurns_SingleTurn(t *testing.T) {
|
||||||
|
dbPath := filepath.Join(t.TempDir(), "test.db")
|
||||||
|
store, err := OpenSQLiteStore(dbPath)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
la := NewLegacyAdapter(store)
|
||||||
|
defer la.Close()
|
||||||
|
|
||||||
|
la.GetOrCreate("k1")
|
||||||
|
la.AddMessage("k1", "user", "a")
|
||||||
|
la.AddMessage("k1", "assistant", "b")
|
||||||
|
la.AddMessage("k1", "user", "c")
|
||||||
|
la.Save("k1")
|
||||||
|
|
||||||
|
// Single turn with 3 messages, keepLast=2 → dropCount=1, but first turn has 3 msgs
|
||||||
|
// accumulated(3) > dropCount(1) on first turn → cutSeq=0 → no compaction
|
||||||
|
if err := la.CompactOldTurns("k1", 2, "sum"); err != nil {
|
||||||
|
t.Fatalf("CompactOldTurns: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
h := la.GetHistory("k1")
|
||||||
|
if len(h) != 3 {
|
||||||
|
t.Fatalf("expected 3 messages (no compaction), got %d", len(h))
|
||||||
|
}
|
||||||
|
if s := la.GetSummary("k1"); s != "sum" {
|
||||||
|
t.Errorf("expected 'sum', got %q", s)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestCompactOldTurns_Graph(t *testing.T) {
|
||||||
|
dbPath := filepath.Join(t.TempDir(), "test.db")
|
||||||
|
store, err := OpenSQLiteStore(dbPath)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
la := NewLegacyAdapter(store)
|
||||||
|
defer la.Close()
|
||||||
|
|
||||||
|
la.GetOrCreate("k1")
|
||||||
|
la.AddMessage("k1", "user", "hello")
|
||||||
|
la.Save("k1")
|
||||||
|
|
||||||
|
g := la.Graph()
|
||||||
|
msgs, err := g.Messages("k1")
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if len(msgs) != 1 || msgs[0].Content != "hello" {
|
||||||
|
t.Errorf("unexpected graph messages: %+v", msgs)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
|
||||||
|
|
@ -188,21 +188,23 @@ func (tw *TurnWriter) Discard()
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
### Phase 2: AgentLoop 直接移行
|
### Phase 2: Compaction + Session Lifecycle + SessionGraph Prep ✅
|
||||||
|
|
||||||
8. **SessionGraph 直接呼び出し**
|
8. ✅ **SessionGraph 薄ラッパー** (段階化: 既存 call site 変更なし)
|
||||||
- `runAgentLoop()` を `SessionGraph.BeginTurn()` / `TurnWriter` 経由に変更
|
- `pkg/session/graph.go`: `SessionGraph` + `BeginTurn()` / `TurnWriter`
|
||||||
- `processSystemMessage()` を Report ターン生成に変更
|
- `LegacyAdapter.Graph()` アクセサ追加
|
||||||
- LegacyAdapter 廃止
|
- 将来の段階移行の準備のみ — LegacyAdapter は廃止しない
|
||||||
|
|
||||||
9. **Compaction**
|
9. ✅ **CompactOldTurns**
|
||||||
- 古いターンの messages を空にして summary で置換
|
- `LegacyAdapter.CompactOldTurns(key, keepLast, summary)` — SQLite 直接 compact
|
||||||
- context window 管理と連動
|
- `summarizeSession()` から呼び出し (fallback 付き)
|
||||||
|
- full-rewrite (`replaced=true → Compact+Append`) を回避
|
||||||
|
|
||||||
10. **セッションライフサイクル**
|
10. ✅ **セッションライフサイクル**
|
||||||
- `Delete(key)` + TTL エビクション (Prune)
|
- 起動時 `store.Prune(DefaultPruneTTL)` — 7日超過セッション削除
|
||||||
- `sessionLocks sync.Map` の GC
|
- `flushLoop` に定期 prune ticker (6h 間隔)
|
||||||
- 起動時の遅延ロード (SQLite なので自然に実現)
|
- `AgentLoop.gcLoop()` — 30分ごとに `sessionLocks` の idle エントリ GC
|
||||||
|
- `AgentLoop.done` チャネル + `Close()` で GC ループ停止
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue