返回 DeepSeek-Reasonix
historical_import.go
根目录 / desktop / historical_import.go
1 package main
2
3 import (
4 "context"
5 "errors"
6 "io"
7 "log/slog"
8 "os"
9 "path/filepath"
10 "sort"
11 "strings"
12 "sync"
13 "time"
14
15 "reasonix/desktop/internal/workspacestate"
16 "reasonix/internal/agent"
17 "reasonix/internal/historywork"
18 "reasonix/internal/identitylock"
19 "reasonix/internal/session"
20 "reasonix/internal/store"
21 )
22
23 var errHistoricalSourceBusy = errors.New("historical session is in use; close the other instance and retry")
24
25 type HistoricalSessionView struct {
26 ID string `json:"id"`
27 Title string `json:"title"`
28 Format string `json:"format"`
29 Status string `json:"status"`
30 ErrorCode string `json:"errorCode,omitempty"`
31 ErrorDetail string `json:"errorDetail,omitempty"`
32 Session *session.SessionRef `json:"session,omitempty"`
33 Source *SessionSourceRef `json:"source,omitempty"`
34 }
35
36 type HistoricalImportStatus struct {
37 Items []HistoricalSessionView `json:"items"`
38 Running bool `json:"running"`
39 Paused bool `json:"paused"`
40 Remaining int `json:"remaining"`
41 Completed int `json:"completed"`
42 Blocked int `json:"blocked"`
43 Failed int `json:"failed"`
44 }
45
46 type historicalSource struct {
47 path, format, scope, root, head, version string
48 }
49 type historicalImportCall struct {
50 operationID string
51 sourceKey string
52 ctx context.Context
53 cancel context.CancelFunc
54 done chan struct{}
55 result SessionRestoreResult
56 err error
57 status string
58 errorCode string
59 errorDetail string
60 revision uint64
61 interactive, batch bool
62 }
63 type historicalImportCoordinator struct {
64 mu sync.Mutex
65 discoveryMu sync.Mutex
66 discoveryPending bool
67 legacyReconcile historicalLegacyReconcileState
68 catalogEnabled bool
69 catalogAt time.Time
70 catalogRevision uint64
71 catalog []historicalCatalogEntry
72 sources map[string]historicalSource
73 views map[string]HistoricalSessionView
74 calls map[string]*historicalImportCall
75 operations map[string]*historicalImportCall
76 updates map[string]*historicalSourceUpdateCall
77 updateWorker chan struct{}
78 revision uint64
79 ctx context.Context
80 cancel context.CancelFunc
81 queue []string
82 current string
83 queueLoaded bool
84 queueRevision uint64
85 queueRelease func()
86 presentations map[string]historicalSourcePresentation
87 unavailableSources map[string]bool
88 running, paused, stopped bool
89 wake chan struct{}
90 workers sync.WaitGroup
91 }
92
93 func (a *App) GetHistoricalImportStatus() HistoricalImportStatus {
94 c := &a.historicalImports
95 c.mu.Lock()
96 defer c.mu.Unlock()
97 return c.status()
98 }
99
100 func (a *App) historicalPreparationStatus(sourceKey string) string {
101 c := &a.historicalImports
102 c.mu.Lock()
103 defer c.mu.Unlock()
104 if view, ok := c.views[sourceKey]; ok {
105 return view.Status
106 }
107 return "available"
108 }
109
110 // Listing reads directory entries and registry metadata only.
111 func (a *App) ListHistoricalSessions() (HistoricalImportStatus, error) {
112 return a.listHistoricalSessions(a.bootContext())
113 }
114
115 func scanHistoricalRoot(ctx context.Context, source desktopMigrationSource, format string, add func(string, string, string, string, string), coordinators ...*historywork.Coordinator) error {
116 f, err := os.Open(source.root)
117 if os.IsNotExist(err) {
118 return nil
119 }
120 if err != nil {
121 return err
122 }
123 defer f.Close()
124 for {
125 if err := ctx.Err(); err != nil {
126 return err
127 }
128 release := func(int64) {}
129 if len(coordinators) > 0 && coordinators[0] != nil {
130 release, err = coordinators[0].BackgroundSlice(ctx, false)
131 if err != nil {
132 return err
133 }
134 }
135 bytes, count := int64(0), 0
136 started := time.Now()
137 var readErr error
138 for count < historywork.BatchEntries && bytes+historywork.ReadChunk <= historywork.BatchBytes && time.Since(started) < historywork.SliceDuration {
139 if readErr = ctx.Err(); readErr != nil {
140 break
141 }
142 var entries []os.DirEntry
143 entries, readErr = f.ReadDir(1)
144 if readErr != nil {
145 break
146 }
147 count++
148 entry := entries[0]
149 if strings.HasPrefix(entry.Name(), ".") || entry.Type()&os.ModeSymlink != 0 {
150 continue
151 }
152 if format == "canonical" && !entry.IsDir() {
153 continue
154 }
155 if format == "legacy" && (entry.IsDir() || !store.IsSessionTranscriptName(entry.Name())) {
156 continue
157 }
158 path := filepath.Join(source.root, entry.Name())
159 if format == "canonical" && !hasHistoricalSessionArtifacts(path) {
160 continue
161 }
162 if format == "legacy" {
163 // The bounded head sidecar is the only payload this discovery reads.
164 // Charge its maximum size so errors and concurrent changes cannot
165 // exceed the shared rate allowance.
166 bytes += historywork.ReadChunk
167 if addIndexedHistoricalHeads(path, source, add) {
168 continue
169 }
170 }
171 add(path, format, source.scope, source.workspaceRoot, "")
172 }
173 release(bytes)
174 if errors.Is(readErr, io.EOF) {
175 return nil
176 }
177 if readErr != nil {
178 return readErr
179 }
180 if len(coordinators) == 0 || coordinators[0] == nil {
181 if err := historywork.Pause(ctx, bytes); err != nil {
182 return err
183 }
184 }
185 }
186 }
187
188 func addIndexedHistoricalHeads(path string, source desktopMigrationSource, add func(string, string, string, string, string)) bool {
189 if info, err := os.Stat(store.SessionEventIndex(path)); err != nil || info.Size() > historywork.ReadChunk {
190 return false
191 }
192 index, err := agent.ReadSessionHeadIndex(path)
193 if err != nil || index == nil || !index.Current(path) {
194 return false
195 }
196 selected := ""
197 for _, head := range index.Heads {
198 if !head.Retired && head.Selected {
199 selected = head.ID
200 }
201 }
202 add(path, "legacy", source.scope, source.workspaceRoot, "")
203 for _, head := range index.Heads {
204 if !head.Retired && head.ID != "" && head.ID != selected {
205 add(path, "legacy", source.scope, source.workspaceRoot, head.ID)
206 }
207 }
208 return true
209 }
210
211 func addHistoricalRegistrySources(state workspacestate.State, add func(string, string, string, string, string)) {
212 workspaceSource := func(path, format, workspaceID, head string) {
213 w := state.Workspaces[workspaceID]
214 scope := "project"
215 if workspaceID == "global" {
216 scope = "global"
217 }
218 add(path, format, scope, w.Root, head)
219 }
220 for _, mapping := range state.SourceMappings {
221 workspaceSource(mapping.Path, mapping.Format, mapping.WorkspaceID, mapping.HeadID)
222 }
223 for _, op := range state.PendingOperations {
224 if op.Mapping != nil && (op.Kind == "import" || op.Kind == "restore") {
225 workspaceSource(op.Mapping.Path, op.Mapping.Format, op.WorkspaceID, op.Mapping.HeadID)
226 }
227 }
228 }
229
230 func historicalImportView(state workspacestate.State, id string, source historicalSource, view HistoricalSessionView) HistoricalSessionView {
231 if view.ID == "" {
232 view = HistoricalSessionView{ID: id, Title: filepath.Base(source.path), Format: source.format, Status: "available"}
233 }
234 view.Source = &SessionSourceRef{HostID: localDesktopHostID, SourceKey: desktopSourceKey(source.path, source.head), Path: source.path, HeadID: source.head}
235 mapping, ok, err := historicalMappingForSource(state, id)
236 if err != nil {
237 view.Status, view.ErrorCode, view.ErrorDetail, view.Session = "failed", "target_changed", "", nil
238 return view
239 }
240 if !ok {
241 return view
242 }
243 view.Status = "imported"
244 ref := session.SessionRef{HostID: localDesktopHostID, SessionID: mapping.SessionID}
245 view.Session = &ref
246 if lifecycle := state.SessionStates[mapping.SessionID].Lifecycle; lifecycle == workspacestate.Deleted || lifecycle == workspacestate.Archived {
247 view.Status = strings.ToLower(lifecycle)
248 view.Session = nil
249 }
250 return view
251 }
252
253 func (c *historicalImportCoordinator) initialize(ctx context.Context) {
254 if c.sources != nil {
255 if !c.stopped && !c.running && len(c.calls) == 0 && c.ctx.Err() != nil {
256 c.ctx, c.cancel = context.WithCancel(ctx)
257 }
258 return
259 }
260 c.sources = map[string]historicalSource{}
261 c.views = map[string]HistoricalSessionView{}
262 c.calls = map[string]*historicalImportCall{}
263 c.operations = map[string]*historicalImportCall{}
264 c.updates = map[string]*historicalSourceUpdateCall{}
265 c.updateWorker = make(chan struct{}, 1)
266 c.presentations = map[string]historicalSourcePresentation{}
267 c.unavailableSources = map[string]bool{}
268 c.ctx, c.cancel = context.WithCancel(ctx)
269 c.wake = make(chan struct{}, 1)
270 }
271 func (c *historicalImportCoordinator) status() HistoricalImportStatus {
272 out := HistoricalImportStatus{Items: []HistoricalSessionView{}, Running: c.running, Paused: c.paused, Remaining: len(c.queue)}
273 if c.current != "" {
274 out.Remaining++
275 }
276 for _, view := range c.views {
277 out.Items = append(out.Items, view)
278 switch view.Status {
279 case "imported":
280 out.Completed++
281 case "blocked":
282 out.Blocked++
283 case "failed":
284 out.Failed++
285 }
286 }
287 sort.Slice(out.Items, func(i, j int) bool { return out.Items[i].ID < out.Items[j].ID })
288 return out
289 }
290
291 func (a *App) ImportHistoricalSession(id string) (SessionRestoreResult, error) {
292 _, listErr := a.ListHistoricalSessions()
293 if listErr != nil {
294 // A damaged or inaccessible historical root must not make healthy
295 // sources unusable. The requested source is checked below; callers can
296 // still inspect the list's per-source status for the affected root.
297 c := &a.historicalImports
298 c.mu.Lock()
299 _, known := c.sources[id]
300 c.mu.Unlock()
301 if !known {
302 return SessionRestoreResult{}, listErr
303 }
304 }
305 call, err := a.prepareHistoricalSession(id, true, false)
306 if err != nil {
307 return SessionRestoreResult{}, err
308 }
309 return waitHistoricalImport(call)
310 }
311
312 // Duplicate requests join one import. Interactive navigation and the bulk queue
313 // hold independent demands so cancelling one cannot abort the other.
314 func (a *App) prepareHistoricalSession(id string, interactive, batch bool) (*historicalImportCall, error) {
315 c := &a.historicalImports
316 c.mu.Lock()
317 if c.stopped || a.shuttingDown.Load() {
318 c.mu.Unlock()
319 return nil, context.Canceled
320 }
321 source, ok := c.sources[id]
322 if !ok {
323 c.mu.Unlock()
324 return nil, errors.New("historical session is unavailable; refresh the list")
325 }
326 if call := c.calls[id]; call != nil {
327 call.interactive = call.interactive || interactive
328 call.batch = call.batch || batch
329 c.mu.Unlock()
330 return call, nil
331 }
332 for operationID, previous := range c.operations {
333 if previous.sourceKey == id {
334 if previous.status == "ready" {
335 c.mu.Unlock()
336 return previous, nil
337 }
338 delete(c.operations, operationID)
339 }
340 }
341 ctx, cancel := context.WithCancel(c.ctx)
342 c.revision++
343 operationID := "prepare-" + strings.TrimPrefix(newTabID(), "tab_")
344 call := &historicalImportCall{operationID: operationID, sourceKey: id, ctx: ctx, cancel: cancel,
345 done: make(chan struct{}), status: "queued", revision: c.revision, interactive: interactive, batch: batch}
346 c.calls[id] = call
347 c.operations[call.operationID] = call
348 view := c.views[id]
349 view.Status, view.ErrorCode, view.ErrorDetail = "queued", "", ""
350 c.views[id] = view
351 c.workers.Add(1)
352 c.mu.Unlock()
353 go a.runHistoricalPreparation(call, id, source)
354 return call, nil
355 }
356
357 func waitHistoricalImport(call *historicalImportCall) (SessionRestoreResult, error) {
358 select {
359 case <-call.done:
360 return call.result, call.err
361 case <-call.ctx.Done():
362 <-call.done
363 return call.result, call.err
364 }
365 }
366
367 func (a *App) runHistoricalPreparation(call *historicalImportCall, id string, source historicalSource) {
368 c := &a.historicalImports
369 defer c.workers.Done()
370 c.mu.Lock()
371 c.revision++
372 call.status, call.revision = "preparing", c.revision
373 view := c.views[id]
374 view.Status, view.ErrorCode, view.ErrorDetail = "importing", "", ""
375 c.views[id] = view
376 c.mu.Unlock()
377 result, err := a.importHistoricalSource(call.ctx, id, source)
378 if err != nil && !errors.Is(err, context.Canceled) && !historicalSourceBusyError(err) {
379 // Preparation can fail before the archive/open mutation is reached.
380 // Keep its cause in the local host log, correlated by opaque source ID.
381 slog.Warn("desktop: historical session preparation failed", "source_key", id,
382 "operation", call.operationID, "format", source.format, "err", err)
383 }
384 var presentationErr error
385 if err == nil {
386 presentationErr = a.applyHistoricalSourcePresentation(desktopSourceKey(source.path, source.head), result.Session)
387 }
388 c.mu.Lock()
389 if c.calls[id] != call {
390 call.result, call.err = result, err
391 close(call.done)
392 c.mu.Unlock()
393 return
394 }
395 call.result, call.err = result, err
396 view = c.views[id]
397 if err == nil {
398 view.Status, call.status = "imported", "ready"
399 ref := result.Session
400 view.Session = &ref
401 if presentationErr != nil {
402 view.ErrorCode = "presentation_pending"
403 }
404 } else {
405 code := historicalImportFailureCode(err)
406 view.Status, call.status = "failed", "failed"
407 switch code {
408 case "source_busy":
409 view.Status, call.status = "blocked", "blocked"
410 err = errHistoricalSourceBusy
411 case "cancelled":
412 view.Status, call.status = "available", "cancelled"
413 }
414 view.ErrorCode, call.errorCode = code, code
415 view.ErrorDetail, call.errorDetail = historicalImportFailureDetail(err), historicalImportFailureDetail(err)
416 call.err = err
417 }
418 c.revision++
419 call.revision = c.revision
420 c.views[id] = view
421 delete(c.calls, id)
422 close(call.done)
423 c.mu.Unlock()
424 a.emitProjectTreeChanged()
425 }
426
427 func (a *App) saveHistoricalSourcePresentation(sourceKey string, update func(*historicalSourcePresentation)) error {
428 c := &a.historicalImports
429 c.mu.Lock()
430 defer c.mu.Unlock()
431 c.initialize(a.bootContext())
432 if !c.queueLoaded {
433 c.loadQueueLocked()
434 }
435 return updateHistoricalSidecar(func(saved *historicalImportQueueSidecar) error {
436 presentation := saved.Presentations[sourceKey]
437 update(&presentation)
438 saved.Presentations[sourceKey] = presentation
439 c.presentations = saved.Presentations
440 return nil
441 })
442 }
443
444 func (a *App) applyHistoricalSourcePresentation(sourceKey string, ref session.SessionRef) error {
445 saved, err := readHistoricalSidecar()
446 if err != nil {
447 return err
448 }
449 presentation, ok := saved.Presentations[sourceKey]
450 if !ok {
451 return nil
452 }
453 var joined error
454 if presentation.Title != "" {
455 joined = errors.Join(joined, a.desktopSessionService("").SetTitle(a.bootContext(), ref, presentation.Title))
456 }
457 if presentation.Pinned != nil {
458 joined = errors.Join(joined, a.workspaceRegistry().UpdatePresentation(a.bootContext(), []string{ref.SessionID}, nil, presentation.Pinned))
459 }
460 return joined
461 }
462
463 func historicalSourceBusyError(err error) bool {
464 return errors.Is(err, identitylock.ErrHeld) ||
465 errors.Is(err, errHistoricalSourceBusy) ||
466 errors.Is(err, agent.ErrSessionLeaseHeld) ||
467 errors.Is(err, session.ErrWriterOwned)
468 }
469
470 func (a *App) importHistoricalSource(ctx context.Context, id string, source historicalSource) (SessionRestoreResult, error) {
471 state, err := a.workspaceRegistry().Load(ctx)
472 if err != nil {
473 return SessionRestoreResult{}, err
474 }
475 if result, handled, err := a.resumeReadyHistoricalImport(ctx, state, id, source); handled {
476 return result, err
477 }
478 release, err := acquireHistoricalSource(ctx, id, source)
479 if err != nil {
480 return SessionRestoreResult{}, err
481 }
482 defer release()
483 // Another process may have committed between the initial read and our claim.
484 state, err = a.workspaceRegistry().Load(ctx)
485 if err != nil {
486 return SessionRestoreResult{}, err
487 }
488 if result, handled, err := a.resumeReadyHistoricalImport(ctx, state, id, source); handled {
489 return result, err
490 }
491 if source.version != "" {
492 current, fingerprintErr := desktopSourceFingerprint(source.path)
493 if fingerprintErr != nil {
494 return SessionRestoreResult{}, fingerprintErr
495 }
496 if current != source.version {
497 return SessionRestoreResult{}, newSessionOperationError("target_changed", "The historical source changed. Check for updates again.")
498 }
499 }
500 workspace, err := a.ensureDesktopWorkspace(ctx, source.scope, source.root)
501 if err != nil {
502 return SessionRestoreResult{}, err
503 }
504 if result, handled, err := a.resumeConflictingHistoricalVersion(ctx, state, source, workspace); handled {
505 return result, err
506 }
507 migration := desktopMigrationSource{scope: source.scope, workspaceRoot: source.root, headID: source.head, versionFingerprint: source.version, registeredSourceKey: id}
508 if resume := pendingHistoricalOperation(state, id); resume != nil {
509 migration.operationID = resume.ID
510 migration.registeredSourceKey = resume.Mapping.SourceKey
511 }
512 err = a.convertHistoricalSource(ctx, source, migration, workspace)
513 if err != nil {
514 return SessionRestoreResult{}, err
515 }
516 state, err = a.workspaceRegistry().Load(ctx)
517 if err != nil {
518 return SessionRestoreResult{}, err
519 }
520 mapping, ok, err := state.ResolveSource(id)
521 if err != nil {
522 return SessionRestoreResult{}, err
523 }
524 if !ok {
525 return SessionRestoreResult{}, errors.New("historical import has not committed")
526 }
527 if lifecycle := state.SessionStates[mapping.SessionID].Lifecycle; lifecycle != workspacestate.Active {
528 return SessionRestoreResult{}, historicalRetiredError(lifecycle)
529 }
530 return SessionRestoreResult{Session: session.SessionRef{HostID: localDesktopHostID, SessionID: mapping.SessionID}, WorkspaceID: mapping.WorkspaceID, Generation: state.Generation}, nil
531 }
532
533 func historicalOperationRank(op workspacestate.Operation) int {
534 switch op.Phase {
535 case "content_ready":
536 return 0
537 case "prepared":
538 return 1
539 default:
540 return 2
541 }
542 }
543
544 func (a *App) convertHistoricalSource(ctx context.Context, source historicalSource, migration desktopMigrationSource, workspace string) (err error) {
545 if source.format == "canonical" {
546 migration.root = filepath.Dir(source.path)
547 old, openErr := session.NewService("migration-source", session.NewFilesystemPersistence(migration.root))
548 if openErr != nil {
549 return openErr
550 }
551 defer func() { err = errors.Join(err, old.Shutdown(context.Background())) }()
552 return a.migrateCanonicalSession(ctx, old, migration, workspace, filepath.Base(source.path))
553 }
554 if source.format == "legacy" || source.format == "legacy-trash" {
555 return a.migrateLegacySession(ctx, source.path, migration, workspace)
556 }
557 return errors.New("historical format is unsupported")
558 }
559
560 // StartHistoricalImport snapshots the requested set; later discoveries are not
561 // silently added. Empty means all currently available/failed/busy sources.
562 func (a *App) StartHistoricalImport(ids []string) (HistoricalImportStatus, error) {
563 _, listErr := a.ListHistoricalSessions()
564 c := &a.historicalImports
565 c.mu.Lock()
566 defer c.mu.Unlock()
567 if listErr != nil && len(ids) > 0 {
568 for _, id := range ids {
569 if _, ok := c.sources[id]; !ok {
570 return c.status(), listErr
571 }
572 }
573 }
574 if c.running || c.stopped || a.shuttingDown.Load() {
575 return c.status(), errors.New("historical import is already running or stopping")
576 }
577 if err := c.claimQueueLocked(); err != nil {
578 return c.status(), err
579 }
580 started := false
581 defer func() {
582 if !started {
583 c.releaseQueueLocked()
584 }
585 }()
586 if len(ids) == 0 {
587 for id, v := range c.views {
588 if v.Status != "imported" && v.Status != "deleted" && v.Status != "archived" {
589 ids = append(ids, id)
590 }
591 }
592 sort.Strings(ids)
593 }
594 seen := map[string]bool{}
595 queue := []string{}
596 for _, id := range ids {
597 if _, ok := c.sources[id]; !ok {
598 return c.status(), errors.New("historical source is unavailable")
599 }
600 if !seen[id] {
601 queue = append(queue, id)
602 seen[id] = true
603 }
604 }
605 c.queue, c.running, c.paused = queue, true, false
606 c.current = ""
607 if err := c.saveQueueLocked(); err != nil {
608 c.queue, c.running = nil, false
609 return c.status(), err
610 }
611 c.workers.Add(1)
612 started = true
613 go a.runHistoricalImportQueue()
614 return c.status(), nil
615 }
616
617 func (a *App) runHistoricalImportQueue() {
618 c := &a.historicalImports
619 defer c.workers.Done()
620 defer func() {
621 c.mu.Lock()
622 c.running = false
623 c.releaseQueueLocked()
624 c.mu.Unlock()
625 }()
626 for {
627 c.mu.Lock()
628 if len(c.queue) == 0 || c.stopped || c.ctx.Err() != nil {
629 c.mu.Unlock()
630 return
631 }
632 if c.paused {
633 wake, ctx := c.wake, c.ctx
634 c.mu.Unlock()
635 select {
636 case <-wake:
637 case <-ctx.Done():
638 }
639 continue
640 }
641 id := c.queue[0]
642 c.queue = c.queue[1:]
643 c.current = id
644 if err := c.saveQueueLocked(); err != nil {
645 c.queue = append([]string{id}, c.queue...)
646 c.current, c.paused = "", true
647 c.mu.Unlock()
648 return
649 }
650 c.mu.Unlock()
651 call, err := a.prepareHistoricalSession(id, false, true)
652 if err == nil {
653 _, _ = waitHistoricalImport(call)
654 }
655 c.mu.Lock()
656 // Shutdown preserves the last durable selection, including the current
657 // item. A committed item is idempotently resolved on manual continuation.
658 if c.stopped || c.ctx.Err() != nil {
659 c.paused = true
660 c.mu.Unlock()
661 return
662 }
663 if c.current == id {
664 c.current = ""
665 }
666 if err := c.saveQueueLocked(); err != nil {
667 c.current, c.paused = id, true
668 c.mu.Unlock()
669 return
670 }
671 c.mu.Unlock()
672 }
673 }
674
675 // Pause finishes the current item. Cancel also interrupts its source work;
676 // durable prepared/content_ready records remain available to the next request.
677 func (a *App) ControlHistoricalImport(action string) (HistoricalImportStatus, error) {
678 c := &a.historicalImports
679 c.mu.Lock()
680 c.initialize(a.bootContext())
681 if c.stopped || a.shuttingDown.Load() {
682 c.mu.Unlock()
683 return HistoricalImportStatus{Items: []HistoricalSessionView{}}, context.Canceled
684 }
685 if err := c.claimQueueLocked(); err != nil {
686 status := c.status()
687 c.mu.Unlock()
688 return status, err
689 }
690 startWorker := false
691 defer func() {
692 c.mu.Lock()
693 if !c.running {
694 c.releaseQueueLocked()
695 }
696 c.mu.Unlock()
697 }()
698 switch action {
699 case "pause":
700 c.paused = true
701 case "resume":
702 c.paused = false
703 if !c.running && (len(c.queue) > 0 || c.current != "") {
704 if c.current != "" {
705 c.queue = append([]string{c.current}, c.queue...)
706 c.current = ""
707 }
708 c.running, startWorker = true, true
709 }
710 select {
711 case c.wake <- struct{}{}:
712 default:
713 }
714 case "cancel":
715 c.queue = nil
716 c.current = ""
717 c.paused = false
718 for _, call := range c.calls {
719 call.batch = false
720 if !call.interactive {
721 call.cancel()
722 }
723 }
724 default:
725 status := c.status()
726 c.mu.Unlock()
727 return status, errors.New("unknown historical import action")
728 }
729 if err := c.saveQueueLocked(); err != nil {
730 if startWorker {
731 c.running = false
732 }
733 status := c.status()
734 c.mu.Unlock()
735 return status, err
736 }
737 status := c.status()
738 if startWorker {
739 c.workers.Add(1)
740 }
741 c.mu.Unlock()
742 if startWorker {
743 go a.runHistoricalImportQueue()
744 }
745 return status, nil
746 }
747
748 func (a *App) stopHistoricalImports() {
749 c := &a.historicalImports
750 c.mu.Lock()
751 c.initialize(a.bootContext())
752 c.stopped = true
753 c.cancel()
754 for _, call := range c.calls {
755 call.cancel()
756 }
757 c.mu.Unlock()
758 // Cancellation and draining happen before the runtime shutdown barrier.
759 c.workers.Wait()
760 }
761
762 // Legacy recovery RPCs share cancellation/draining with the on-demand queue.
763 func (a *App) beginHistoricalRecovery() (context.Context, func(), error) {
764 c := &a.historicalImports
765 c.mu.Lock()
766 defer c.mu.Unlock()
767 c.initialize(a.bootContext())
768 if c.stopped || a.shuttingDown.Load() {
769 return nil, nil, context.Canceled
770 }
771 c.workers.Add(1)
772 return c.ctx, c.workers.Done, nil
773 }
774
774 lines GO