返回 DeepSeek-Reasonix
compact_commit.go
根目录 / internal / agent / compact_commit.go
1 package agent
2
3 import (
4 "context"
5 "errors"
6 "fmt"
7 "time"
8
9 "reasonix/internal/provider"
10 )
11
12 type summaryProjectionCommit struct {
13 canonical, fold, projected []provider.Message
14 result foldSummary
15 transcriptVersion, projectionVersion, generation uint64
16 activeTurn int64
17 trigger, summary, inputHash, outputHash string
18 sourceTokens, projectionTokens int
19 // covered is the canonical length the frozen projection body represents;
20 // messages past it splice live from the transcript.
21 covered int
22 }
23
24 // commitSummaryProjection CAS-installs a checkpoint under compactionMu:
25 // transcript version/hash, projection version, and generation must still match.
26 // The maintenance event is emitted only after the lock is released so a sink
27 // that re-enters ContextMaintenanceSnapshot cannot deadlock.
28 func (a *Agent) commitSummaryProjection(ctx context.Context, commit summaryProjectionCommit) (CompactionState, error) {
29 if ctx == nil {
30 ctx = context.Background()
31 }
32 if err := ctx.Err(); err != nil {
33 return CompactionState{}, err
34 }
35 state := a.summaryProjectionState(commit)
36 a.sess.compactionMu.Lock()
37 // This is the shared commit boundary for ordinary, positional, and fallback
38 // summary projection installs. Cancellation that wins before this point
39 // prevents every variant from publishing a late summary.
40 if err := ctx.Err(); err != nil {
41 a.sess.compactionMu.Unlock()
42 return CompactionState{}, err
43 }
44 current, currentVersion := a.sess.conversation.snapshotMessagesVersion()
45 if currentVersion != commit.transcriptVersion ||
46 len(current) != len(commit.canonical) ||
47 coveredPrefixHash(current, len(current)) != coveredPrefixHash(commit.canonical, len(commit.canonical)) ||
48 a.sess.compactionState.Projection.ProjectionVersion != commit.projectionVersion ||
49 a.sess.compactionState.Generation != commit.generation {
50 a.sess.compactionMu.Unlock()
51 return CompactionState{}, summaryError(errCompressStaleContext)
52 }
53 prev := a.sess.compactionState
54 a.sess.compactionState = state
55 // After installation begins we finish its consistency and durability work;
56 // cancellation may prevent a later batch but cannot tear this batch in half.
57 accepted, err := a.persistInstalledProjectionLocked(context.Background(), state, current)
58 if err != nil {
59 if accepted {
60 a.sess.checkpointState = "pending"
61 a.sess.compactionMu.Unlock()
62 return CompactionState{}, &compactionPersistenceError{fmt.Errorf("persist projection: %w", err)}
63 }
64 a.sess.compactionState = prev
65 a.sess.compactionMu.Unlock()
66 if errors.Is(err, errCompressStaleContext) {
67 return CompactionState{}, err
68 }
69 return CompactionState{}, &compactionPersistenceError{fmt.Errorf("persist projection: %w", err)}
70 }
71 a.sess.checkpointState = "applied"
72 if commit.activeTurn != 0 && commit.trigger != CompactionTriggerManual {
73 a.sess.compaction.lastTurn.Store(commit.activeTurn)
74 }
75 receipt := state.LastReceipt
76 a.sess.compactionMu.Unlock()
77 a.emitContextMaintenance(receipt)
78 return state, nil
79 }
80
81 func (a *Agent) persistInstalledProjectionLocked(ctx context.Context, state CompactionState, canonical []provider.Message) (bool, error) {
82 accepted := false
83 if recorder, ok := a.svc.sessionCheckpointer.(SessionModelContextRecorder); ok {
84 visible := modelVisibleFromProjection(state.Projection, canonical)
85 commit := cloneSessionModelContextCommit(SessionModelContextCommit{
86 OperationID: state.LastReceipt.OperationID,
87 Reason: state.LastReceipt.Action,
88 Messages: visible,
89 })
90 result, err := recorder.RecordSessionModelContext(ctx, commit)
91 accepted = result.Accepted
92 if err != nil {
93 if accepted {
94 a.sess.pendingModelContextCommit = &commit
95 }
96 return accepted, err
97 }
98 if result.Accepted && !result.Durable {
99 a.sess.pendingModelContextCommit = &commit
100 return true, errors.New("model context commit was accepted but is not durable")
101 }
102 }
103 if err := a.persistCompactionStateLocked(); err != nil {
104 if accepted {
105 visible := modelVisibleFromProjection(state.Projection, canonical)
106 commit := cloneSessionModelContextCommit(SessionModelContextCommit{
107 OperationID: state.LastReceipt.OperationID,
108 Reason: state.LastReceipt.Action,
109 Messages: visible,
110 })
111 a.sess.pendingModelContextCommit = &commit
112 }
113 return accepted, err
114 }
115 a.sess.pendingModelContextCommit = nil
116 return accepted, nil
117 }
118
119 func (a *Agent) summaryProjectionState(commit summaryProjectionCommit) CompactionState {
120 projectionVersion := commit.projectionVersion + 1
121 now := time.Now().UTC()
122 summaryHash := summaryContentHash(commit.summary)
123 coveredHash := coveredPrefixHash(commit.canonical, commit.covered)
124 receipt := &ContextMaintenanceReceipt{
125 OperationID: fmt.Sprintf("summary-%d-%s", projectionVersion, commit.outputHash), Status: "applied",
126 Action: "summary", Trigger: commit.trigger, SourceProjection: commit.projectionVersion,
127 ProjectionVersion: projectionVersion, CoveredCount: commit.covered, CoveredPrefixHash: coveredHash,
128 InputHash: commit.inputHash, OutputHash: commit.outputHash, InputTokens: commit.sourceTokens,
129 ResultTokens: commit.projectionTokens, SavedTokens: max(0, commit.sourceTokens-commit.projectionTokens),
130 SummaryHash: summaryHash, CacheBreak: true, CreatedAt: now,
131 }
132 // LastReceipt is authoritative; do not mirror last_trigger/last_mode/token
133 // counters or top-level blocked_* fields (stripped again on save).
134 return CompactionState{
135 SchemaVersion: compactionStateSchemaCurrent, TranscriptVersion: commit.transcriptVersion,
136 Generation: commit.generation + 1, PromptCacheKey: a.currentPromptCacheKey(),
137 Projection: ContextProjection{
138 Messages: commit.projected, TranscriptVersion: commit.transcriptVersion,
139 ProjectionVersion: projectionVersion, CoveredCount: commit.covered, CoveredPrefixHash: coveredHash,
140 PinnedContextHash: pinnedContextCoverageHash(commit.canonical, commit.covered),
141 SummaryHash: summaryHash, SourceTokens: commit.sourceTokens, ProjectionTokens: commit.projectionTokens,
142 ViewInputHash: commit.inputHash, ViewOutputHash: commit.outputHash, CreatedAt: now,
143 },
144 LastReceipt: receipt, UpdatedAt: now,
145 }
146 }
147
147 lines GO