| 103 | } |
| 104 | |
| 105 | func (b *BackgroundManager) Start(ctx context.Context, options map[string]interface{}) error { |
| 106 | err := b.markStart(ctx) |
| 107 | if err != nil { |
| 108 | return err |
| 109 | } |
| 110 | |
| 111 | var processClusterStatus []byte |
| 112 | if b.isClusterAware() { |
| 113 | processClusterStatus, _, err = b.clusterAwareOptions.metadataStore.GetRaw(b.clusterAwareOptions.StatusDocID()) |
| 114 | if err != nil && !base.IsDocNotFoundError(err) { |
| 115 | return pkgerrors.Wrap(err, "Failed to get current process status") |
| 116 | } |
| 117 | } |
| 118 | |
| 119 | b.resetStatus() |
| 120 | b.StartTime = time.Now().UTC() |
| 121 | |
| 122 | err = b.Process.Init(ctx, options, processClusterStatus) |
| 123 | if err != nil { |
| 124 | return err |
| 125 | } |
| 126 | |
| 127 | if b.isClusterAware() { |
| 128 | b.backgroundManagerStatusUpdateWaitGroup.Add(1) |
| 129 | go func(terminator *base.SafeTerminator) { |
| 130 | defer b.backgroundManagerStatusUpdateWaitGroup.Done() |
| 131 | ticker := time.NewTicker(BackgroundManagerStatusUpdateIntervalSecs * time.Second) |
| 132 | for { |
| 133 | select { |
| 134 | case <-ticker.C: |
| 135 | err := b.UpdateStatusClusterAware(ctx) |
| 136 | if err != nil { |
| 137 | base.WarnfCtx(ctx, "Failed to update background manager status: %v", err) |
| 138 | } |
| 139 | case <-terminator.Done(): |
| 140 | ticker.Stop() |
| 141 | return |
| 142 | } |
| 143 | } |
| 144 | }(b.terminator) |
| 145 | } |
| 146 | |
| 147 | go func() { |
| 148 | err := b.Process.Run(ctx, options, b.UpdateStatusClusterAware, b.terminator) |
| 149 | if err != nil { |
| 150 | base.ErrorfCtx(ctx, "Error: %v", err) |
| 151 | b.SetError(err) |
| 152 | } |
| 153 | |
| 154 | b.Terminate() |
| 155 | |
| 156 | b.lock.Lock() |
| 157 | if b.State == BackgroundProcessStateStopping { |
| 158 | b.State = BackgroundProcessStateStopped |
| 159 | } else if b.State != BackgroundProcessStateError { |
| 160 | b.State = BackgroundProcessStateCompleted |
| 161 | } |
| 162 | b.lock.Unlock() |