返回 DeepSeek-Reasonix
historical_archive.go
根目录 / desktop / historical_archive.go
1 package main
2
3 import (
4 "context"
5 "errors"
6 "fmt"
7 "log/slog"
8 "path/filepath"
9 "strings"
10
11 "reasonix/desktop/internal/workspacestate"
12 "reasonix/internal/agent"
13 "reasonix/internal/session"
14 )
15
16 // Explicit source archives stage content as archived from the beginning. They
17 // must not use PrepareSession, which publishes an active/openable conversation.
18 func (a *App) archiveHistoricalSource(selector SessionSelector) (SessionMutationResult, error) {
19 if selector.Source != nil && selector.Source.HostID != "" && selector.Source.HostID != localDesktopHostID {
20 return SessionMutationResult{}, newSessionOperationError("unsupported", "This source belongs to another host.")
21 }
22 ctx, done, err := a.beginHistoricalRecovery()
23 if err != nil {
24 return SessionMutationResult{}, err
25 }
26 defer done()
27 id, source, err := a.historicalSourceForSelector(selector)
28 if err != nil {
29 return SessionMutationResult{}, err
30 }
31 operationID := "archive-source-" + strings.TrimPrefix(newTabID(), "tab_")
32 result, err := a.archiveHistoricalSourceLocked(ctx, id, source, operationID, "archive historical source")
33 if err != nil {
34 if !errors.Is(err, errTopicArchiveBusy) {
35 slog.Warn("desktop: historical archive failed", "source_key", id, "operation", operationID, "err", err)
36 }
37 return SessionMutationResult{}, sessionOperationErrorForTarget(err, id, operationID)
38 }
39 archivedPaths := []string{source.path}
40 siblings := legacyHeadVersions(source)
41 recovered, listErr := recoveredLegacySiblings(source)
42 pending := 0
43 if listErr != nil {
44 slog.Warn("desktop: recovery lineage unreadable", "source_key", id, "cause", "lineage_unreadable", "err", listErr)
45 pending++
46 }
47 isRecovered := map[string]bool{}
48 for _, path := range recovered {
49 isRecovered[path] = true
50 siblings = append(siblings, SessionSourceRef{HostID: localDesktopHostID, Path: path})
51 }
52 state, err := a.workspaceRegistry().Load(ctx)
53 if err != nil {
54 slog.Warn("desktop: sibling lifecycle unavailable", "source_key", id, "err", err)
55 }
56 for _, sibling := range siblings {
57 counted := isRecovered[sibling.Path] && sibling.HeadID == ""
58 siblingID, siblingSource, err := a.historicalSourceForSelector(SessionSelector{Source: &sibling})
59 if err != nil {
60 slog.Warn("desktop: sibling archive skipped", "source_key", id, "err", err)
61 if counted {
62 pending++
63 }
64 continue
65 }
66 if historicalSourceMapped(state, siblingID) {
67 continue
68 }
69 siblingResult, err := a.archiveHistoricalSourceLocked(ctx, siblingID, siblingSource,
70 "archive-source-"+strings.TrimPrefix(newTabID(), "tab_"), "archive historical sibling")
71 if err != nil {
72 slog.Warn("desktop: sibling archive failed", "source_key", siblingID, "cause", "sibling_archive_failed", "err", err)
73 if counted {
74 pending++
75 }
76 continue
77 }
78 result.IdentityAliases = append(result.IdentityAliases, siblingResult.IdentityAliases...)
79 archivedPaths = append(archivedPaths, siblingSource.path)
80 }
81 c := &a.historicalImports
82 c.mu.Lock()
83 for i := range c.catalog {
84 for _, path := range archivedPaths {
85 if c.catalog[i].node.Source != nil && sameDesktopPath(c.catalog[i].node.Source.Path, path) {
86 c.catalog[i].sourceChanged = false
87 }
88 }
89 }
90 c.mu.Unlock()
91 a.emitProjectTreeChanged()
92 if pending != 0 {
93 result.Outcome = sessionOutcomeArchivedPartial
94 result.PendingSiblings = pending
95 }
96 return result, nil
97 }
98
99 // Each source is archived under its own runtime-mutation hold so a long
100 // lineage never keeps other tabs from sending between snapshots.
101 func (a *App) archiveHistoricalSourceLocked(ctx context.Context, id string, source historicalSource, operationID, operation string) (SessionMutationResult, error) {
102 release, ok := a.tryLockRuntimeMutationBounded(operation)
103 if !ok {
104 return SessionMutationResult{}, errTopicArchiveBusy
105 }
106 defer release()
107 return a.archiveHistoricalSourceWithOperation(ctx, id, source, operationID)
108 }
109
110 // A sibling with a durable mapping already has its own lifecycle. Archiving a
111 // different version must not retire an active session or copy an old version
112 // again after another head changes the shared legacy transcript.
113 func historicalSourceMapped(state workspacestate.State, id string) bool {
114 _, mapped, err := state.ResolveSource(id)
115 return err == nil && mapped
116 }
117
118 // headStartsConversation reports a head the user split off under a name. A
119 // rewind, an unnamed fork ("fork from this turn") and a second writer's
120 // concurrent head all continue their parent head's conversation.
121 func headStartsConversation(head agent.SessionHead) bool {
122 return head.ParentHead == "" || head.Kind == agent.HeadKindFork && strings.TrimSpace(head.Name) != ""
123 }
124
125 // legacyHeadVersions lists the other live heads of a legacy transcript that are
126 // versions of the selected head's conversation. The sidebar lists every version
127 // as its own row with the same title; named forks stay.
128 func legacyHeadVersions(source historicalSource) []SessionSourceRef {
129 if source.format != "legacy" || source.head == "" {
130 return nil
131 }
132 heads, err := agent.ListSessionHeads(source.path)
133 if err != nil {
134 return nil
135 }
136 byID := make(map[string]agent.SessionHead, len(heads))
137 for _, head := range heads {
138 byID[head.ID] = head
139 }
140 origin := func(id string) string {
141 for range heads {
142 head, ok := byID[id]
143 if !ok || headStartsConversation(head) {
144 break
145 }
146 id = head.ParentHead
147 }
148 return id
149 }
150 if _, ok := byID[source.head]; !ok {
151 return nil
152 }
153 want := origin(source.head)
154 var versions []SessionSourceRef
155 for _, head := range heads {
156 if head.ID == source.head || head.Retired || head.Kind == agent.HeadKindConcurrent || origin(head.ID) != want {
157 continue
158 }
159 versions = append(versions, SessionSourceRef{HostID: localDesktopHostID, Path: source.path, HeadID: head.ID})
160 }
161 return versions
162 }
163
164 func (a *App) archiveHistoricalSourceWithOperation(ctx context.Context, id string, source historicalSource, operationID string) (SessionMutationResult, error) {
165 release, err := acquireHistoricalSource(ctx, id, source)
166 if err != nil {
167 return SessionMutationResult{}, err
168 }
169 defer release()
170 fingerprint, err := desktopSourceFingerprint(source.path)
171 if err != nil {
172 return SessionMutationResult{}, err
173 }
174 mapping, dependencies, err := a.stageHistoricalArchive(ctx, id, source, fingerprint)
175 if err != nil {
176 return SessionMutationResult{}, err
177 }
178 state, err := a.workspaceRegistry().Load(ctx)
179 if err != nil {
180 return SessionMutationResult{}, err
181 }
182 ref := session.SessionRef{HostID: localDesktopHostID, SessionID: mapping.SessionID}
183 if len(dependencies) != 0 {
184 operationID = "archive-" + dependencies[0]
185 }
186 lifecycle := state.SessionStates[ref.SessionID]
187 outcome := "archived"
188 if lifecycle.Lifecycle == workspacestate.Deleted {
189 outcome = "already_removed"
190 if err := retireRemovedSourceTopic(state, source); err != nil {
191 return SessionMutationResult{}, fmt.Errorf("retire removed source topic: %w", err)
192 }
193 } else if lifecycle.Lifecycle == workspacestate.Archived {
194 if _, err := a.desktopSessionService("").Query().Snapshot(ctx, ref); err != nil {
195 return SessionMutationResult{}, err
196 }
197 } else if lifecycle.Lifecycle != workspacestate.Archived {
198 verify := func(ctx context.Context, current workspacestate.State) error {
199 fp, err := desktopSourceFingerprint(source.path)
200 if err != nil || fp != fingerprint {
201 return errors.Join(err, workspacestate.ErrMutationConflict)
202 }
203 for _, dependency := range dependencies {
204 if err := a.validateHistoricalArchiveContent(ctx, source, current.PendingOperations[dependency]); err != nil {
205 return err
206 }
207 }
208 return nil
209 }
210 if err := a.archiveSessionRefsWithOperationConditional([]session.SessionRef{ref}, operationID, verify, dependencies...); err != nil {
211 return SessionMutationResult{}, fmt.Errorf("commit historical archive: %w", err)
212 }
213 state, err = a.workspaceRegistry().Load(ctx)
214 if err != nil {
215 return SessionMutationResult{}, err
216 }
217 lifecycle = state.SessionStates[ref.SessionID]
218 }
219 if outcome != "already_removed" && source.format == "canonical" && ref.SessionID != filepath.Base(source.path) {
220 outcome = "archived_copy"
221 }
222 projectionPending := false
223 if outcome != "already_removed" {
224 if err := a.applyHistoricalSourcePresentation(desktopSourceKey(source.path, source.head), ref); err != nil {
225 // Archive is durable, just as import can succeed with presentation
226 // pending. Retrying the source action reapplies its saved overlay.
227 projectionPending = true
228 slog.Warn("desktop: archived source presentation pending", "source_key", id, "operation", operationID, "err", err)
229 }
230 }
231 target := SessionTarget{SessionRef: ref, Scope: source.scope, WorkspaceRoot: source.root}
232 aliases := a.sessionTargetIdentityAliases(target)
233 aliases = append(aliases, "source\x00local\x00"+id)
234 return SessionMutationResult{TargetKey: target.key(), OperationID: operationID, Committed: true,
235 Outcome: outcome, LifecycleGeneration: lifecycle.Generation, IdentityAliases: aliases, ProjectionPending: projectionPending}, nil
236 }
237
237 lines GO