返回 DeepSeek-Reasonix
maintenance_commit.go
根目录 / internal / agent / maintenance_commit.go
1 package agent
2
3 import (
4 "context"
5 "errors"
6 "fmt"
7 "time"
8
9 "reasonix/internal/provider"
10 )
11
12 // maintenanceInstall is one free projection rewrite (no summarizer call): the
13 // visible view it started from and the projected view replacing it.
14 type maintenanceInstall struct {
15 trigger, action string
16 state CompactionState
17 canonical []provider.Message
18 transcriptVersion uint64
19 visible, projected []provider.Message
20 affected int
21 }
22
23 // installMaintenanceProjection CAS-installs a free projection under
24 // compactionMu. The caller owns compactionRunMu for the whole maintenance run;
25 // canonical storage, including RawContent, is never modified.
26 func (a *Agent) installMaintenanceProjection(ctx context.Context, in maintenanceInstall) (bool, error) {
27 if ctx == nil {
28 ctx = context.Background()
29 }
30 if err := ctx.Err(); err != nil {
31 return false, err
32 }
33 projected := projectionMessagesPreservingPinnedContext(in.projected)
34 projected, _, err := rebasePinnedContextProjection(projected, in.canonical, len(in.canonical))
35 if err != nil {
36 return false, err
37 }
38 if err := ctx.Err(); err != nil {
39 return false, err
40 }
41 sourceTokens := a.estimatedVisibleRequestTokens(in.visible)
42 resultTokens := a.estimatedVisibleRequestTokens(projected)
43 inputHash := a.contextMaintenanceInputHash(modelInputMessages(in.visible))
44 outputHash := providerVisibleFingerprint(modelInputMessages(projected))
45 projectionVersion := in.state.Projection.ProjectionVersion + 1
46 now := time.Now().UTC()
47 coveredHash := coveredPrefixHash(in.canonical, len(in.canonical))
48 receipt := &ContextMaintenanceReceipt{
49 OperationID: fmt.Sprintf("%s-%d-%s", in.action, projectionVersion, outputHash), Status: "applied", Action: in.action,
50 Trigger: in.trigger, SourceProjection: in.state.Projection.ProjectionVersion, ProjectionVersion: projectionVersion,
51 CoveredCount: len(in.canonical), CoveredPrefixHash: coveredHash, InputHash: inputHash, OutputHash: outputHash,
52 InputTokens: sourceTokens, ResultTokens: resultTokens, SavedTokens: max(0, sourceTokens-resultTokens),
53 AffectedToolResults: in.affected, CacheBreak: true, CreatedAt: now,
54 }
55 next := in.state
56 next.SchemaVersion = compactionStateSchemaCurrent
57 next.TranscriptVersion = in.transcriptVersion
58 next.Generation++
59 next.PromptCacheKey = a.currentPromptCacheKey()
60 next.Projection = ContextProjection{
61 Messages: projected, TranscriptVersion: in.transcriptVersion, ProjectionVersion: projectionVersion,
62 CoveredCount: len(in.canonical), CoveredPrefixHash: coveredHash, SourceTokens: sourceTokens,
63 PinnedContextHash: pinnedContextCoverageHash(in.canonical, len(in.canonical)),
64 ProjectionTokens: resultTokens, ViewInputHash: inputHash, ViewOutputHash: outputHash, CreatedAt: now,
65 }
66 next.LastReceipt = receipt
67 next.UpdatedAt = now
68
69 a.sess.compactionMu.Lock()
70 // Cancellation and the projection compare-and-swap share this lock boundary.
71 // Once the state is installed, persistence finishes atomically with respect
72 // to this batch; cancellation can only prevent a later batch.
73 if err := ctx.Err(); err != nil {
74 a.sess.compactionMu.Unlock()
75 return false, err
76 }
77 current, currentVersion := a.sess.conversation.snapshotMessagesVersion()
78 if currentVersion != in.transcriptVersion || len(current) != len(in.canonical) ||
79 coveredPrefixHash(current, len(current)) != coveredHash ||
80 a.sess.compactionState.Projection.ProjectionVersion != in.state.Projection.ProjectionVersion ||
81 a.sess.compactionState.Generation != in.state.Generation {
82 a.sess.compactionMu.Unlock()
83 return false, errCompressStaleContext
84 }
85 previous := a.sess.compactionState
86 a.sess.compactionState = next
87 accepted, err := a.persistInstalledProjectionLocked(context.Background(), next, current)
88 if err != nil {
89 if accepted {
90 a.sess.checkpointState = "pending"
91 a.sess.compactionMu.Unlock()
92 return false, &compactionPersistenceError{fmt.Errorf("persist %s projection: %w", in.action, err)}
93 }
94 a.sess.compactionState = previous
95 a.sess.compactionMu.Unlock()
96 if errors.Is(err, errCompressStaleContext) {
97 return false, err
98 }
99 return false, &compactionPersistenceError{fmt.Errorf("persist %s projection: %w", in.action, err)}
100 }
101 a.sess.checkpointState = "applied"
102 a.sess.compactionMu.Unlock()
103 a.emitContextMaintenance(receipt)
104 return true, nil
105 }
106
106 lines GO