• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

kubevirt / kubevirt / 40152778-4fd9-4c20-936f-30cdcf398926

10 Jul 2026 10:14PM UTC coverage: 72.342% (+0.02%) from 72.32%
40152778-4fd9-4c20-936f-30cdcf398926

push

prow

web-flow
Merge pull request #18217 from awels/unquarantine_decentralize_lm_cancel_tests

Unquarantine decentralized live migration cancel tests

47 of 60 new or added lines in 2 files covered. (78.33%)

207 existing lines in 2 files now uncovered.

83915 of 115997 relevant lines covered (72.34%)

587.19 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

72.95
/pkg/synchronization-controller/synchronization-controller.go
1
/*
2
 * This file is part of the KubeVirt project
3
 *
4
 * Licensed under the Apache License, Version 2.0 (the "License");
5
 * you may not use this file except in compliance with the License.
6
 * You may obtain a copy of the License at
7
 *
8
 *     http://www.apache.org/licenses/LICENSE-2.0
9
 *
10
 * Unless required by applicable law or agreed to in writing, software
11
 * distributed under the License is distributed on an "AS IS" BASIS,
12
 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.  * See the License for the specific language governing permissions and
13
 * limitations under the License.
14
 *
15
 * Copyright The KubeVirt Authors.
16
 *
17
 */
18

19
package synchronization
20

21
import (
22
        "crypto/tls"
23
        "crypto/x509"
24
        "encoding/json"
25
        "errors"
26
        "fmt"
27
        "net"
28
        "strconv"
29
        "sync"
30
        "time"
31

32
        k8sv1 "k8s.io/api/core/v1"
33
        apiequality "k8s.io/apimachinery/pkg/api/equality"
34
        metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
35
        "k8s.io/apimachinery/pkg/types"
36
        "k8s.io/apimachinery/pkg/util/wait"
37
        "k8s.io/client-go/tools/cache"
38
        "k8s.io/client-go/util/workqueue"
39

40
        virtv1 "kubevirt.io/api/core/v1"
41

42
        "kubevirt.io/kubevirt/pkg/apimachinery/patch"
43

44
        "kubevirt.io/kubevirt/pkg/controller"
45

46
        context "golang.org/x/net/context"
47
        "google.golang.org/grpc"
48
        "google.golang.org/grpc/credentials"
49
        "kubevirt.io/client-go/kubecli"
50
        "kubevirt.io/client-go/log"
51

52
        syncv1 "kubevirt.io/kubevirt/pkg/synchronizer-com/synchronization/v1"
53
)
54

55
const (
56
        defaultTimeout = 30
57

58
        MyPodIP = "MY_POD_IP"
59

60
        noSourceStatusErrorMsg                        = "must pass source status"
61
        noTargetStatusErrorMsg                        = "must pass target status"
62
        sourceUnableToLocateVMIMigrationIDErrorMsg    = "source: unable to locate VMI for migrationID %s"
63
        targetUnableToLocateVMIMigrationIDErrorMsg    = "target: unable to locate VMI for migrationID %s"
64
        sourceUnableToLocateVMIMigrationIDErrorMsgVMI = "source: unable to locate VMI for migrationID %s, vmi: %s"
65
        targetUnableToLocateVMIMigrationIDErrorMsgVMI = "target: unable to locate VMI for migrationID %s, vmi: %s"
66

67
        waitingForSyncErrorMessage = "waiting for incoming synchronization, unable to proceed"
68

69
        successMessage = "success"
70

71
        maxCloseRetries = 10
72

73
        SynchronizationFinalizer = "synchronization.kubevirt.io/migrationFinalizer"
74
)
75

76
type SynchronizationController struct {
77
        client kubecli.KubevirtClient
78

79
        vmiInformer       cache.SharedIndexInformer
80
        migrationInformer cache.SharedIndexInformer
81

82
        listener        net.Listener
83
        bindAddress     string
84
        bindPort        int
85
        ip              string
86
        clientTLSConfig *tls.Config
87
        serverTLSConfig *tls.Config
88
        timeout         int
89

90
        queue     workqueue.TypedRateLimitingInterface[string]
91
        hasSynced func() bool
92

93
        syncOutboundConnectionMap  *sync.Map
94
        syncReceivingConnectionMap *sync.Map
95
        failedCloseConnections     *sync.Map
96
        grpcServer                 *grpc.Server
97
}
98

99
func NewSynchronizationController(
100
        client kubecli.KubevirtClient,
101
        vmiInformer cache.SharedIndexInformer,
102
        migrationInformer cache.SharedIndexInformer,
103
        clientTLSConfig,
104
        serverTLSConfig *tls.Config,
105
        bindAddress string,
106
        bindPort int,
107
        ip string,
108
) (*SynchronizationController, error) {
104✔
109
        syncController := &SynchronizationController{
104✔
110
                vmiInformer:       vmiInformer,
104✔
111
                migrationInformer: migrationInformer,
104✔
112
                clientTLSConfig:   clientTLSConfig,
104✔
113
                serverTLSConfig:   serverTLSConfig,
104✔
114
                timeout:           defaultTimeout,
104✔
115
                bindAddress:       bindAddress,
104✔
116
                bindPort:          bindPort,
104✔
117
                client:            client,
104✔
118
                ip:                ip,
104✔
119
        }
104✔
120

104✔
121
        queue := workqueue.NewTypedRateLimitingQueueWithConfig[string](
104✔
122
                workqueue.DefaultTypedControllerRateLimiter[string](),
104✔
123
                workqueue.TypedRateLimitingQueueConfig[string]{Name: "sync-vmi-status"},
104✔
124
        )
104✔
125
        syncController.queue = queue
104✔
126

104✔
127
        syncController.hasSynced = func() bool {
104✔
128
                return vmiInformer.HasSynced() && migrationInformer.HasSynced()
×
129
        }
×
130

131
        syncController.syncOutboundConnectionMap = &sync.Map{}
104✔
132
        syncController.syncReceivingConnectionMap = &sync.Map{}
104✔
133
        syncController.failedCloseConnections = &sync.Map{}
104✔
134

104✔
135
        _, err := vmiInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
104✔
136
                AddFunc:    syncController.addVmiFunc,
104✔
137
                DeleteFunc: syncController.deleteVmiFunc,
104✔
138
                UpdateFunc: syncController.updateVmiFunc,
104✔
139
        })
104✔
140
        if err != nil {
104✔
141
                return nil, err
×
142
        }
×
143

144
        if err := syncController.migrationInformer.AddIndexers(map[string]cache.IndexFunc{
104✔
145
                "byUID":               indexByMigrationUID,
104✔
146
                "byActiveVMIName":     indexByActiveVmiName,
104✔
147
                "byTargetMigrationID": indexByTargetMigrationID,
104✔
148
                "bySourceMigrationID": indexBySourceMigrationID,
104✔
149
        }); err != nil {
104✔
150
                return nil, err
×
151
        }
×
152

153
        if _, err := syncController.migrationInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
104✔
154
                AddFunc:    syncController.addMigrationFunc,
104✔
155
                DeleteFunc: syncController.deleteMigrationFunc,
104✔
156
                UpdateFunc: syncController.updateMigrationFunc,
104✔
157
        }); err != nil {
104✔
158
                return nil, err
×
159
        }
×
160

161
        syncController.grpcServer = grpc.NewServer(grpc.Creds(credentials.NewTLS(serverTLSConfig)))
104✔
162
        syncv1.RegisterSynchronizeServer(syncController.grpcServer, syncController)
104✔
163

104✔
164
        return syncController, nil
104✔
165
}
166

167
func (s *SynchronizationController) addVmiFunc(addObj interface{}) {
×
168
        s.enqueueVirtualMachineInstance(addObj)
×
169
}
×
170

171
func (s *SynchronizationController) deleteVmiFunc(addObj interface{}) {
×
172
        s.enqueueVirtualMachineInstance(addObj)
×
173
}
×
174

175
func (s *SynchronizationController) updateVmiFunc(_, curr interface{}) {
×
176
        s.enqueueVirtualMachineInstance(curr)
×
177
}
×
178

179
func (s *SynchronizationController) enqueueVirtualMachineInstance(obj interface{}) {
×
180
        vmi, ok := obj.(*virtv1.VirtualMachineInstance)
×
181
        if ok {
×
182
                key, err := controller.KeyFunc(vmi)
×
183
                if err != nil {
×
184
                        log.Log.Object(vmi).Reason(err).Error("failed to extract key from virtualmachine.")
×
185
                        return
×
186
                }
×
187
                s.queue.Add(key)
×
188
        }
189
}
190

191
func (s *SynchronizationController) addMigrationFunc(addObj interface{}) {
1✔
192
        s.enqueueVirtualMachineInstanceFromMigration(addObj)
1✔
193
}
1✔
194

195
func (s *SynchronizationController) deleteMigrationFunc(delObj interface{}) {
2✔
196
        // Clean up any synchronization connections in the map.
2✔
197
        s.enqueueVirtualMachineInstanceFromMigration(delObj)
2✔
198
        // Close any connections associated with this migration.
2✔
199
        migration, ok := delObj.(*virtv1.VirtualMachineInstanceMigration)
2✔
200
        if ok {
4✔
201
                if !migration.IsDecentralized() {
2✔
202
                        return
×
203
                }
×
204
                if migration.Spec.Receive != nil {
3✔
205
                        log.Log.V(4).Object(migration).Infof("closing receiving connection for migrationID %s", migration.Spec.Receive.MigrationID)
1✔
206
                        if err := s.closeConnectionForMigrationID(s.syncReceivingConnectionMap, migration.Spec.Receive.MigrationID); err != nil {
1✔
207
                                log.Log.Reason(err).Infof("unable to close connection for migrationID %s, possibly leaked connection", migration.Spec.Receive.MigrationID)
×
208
                        }
×
209
                } else if migration.Spec.SendTo != nil {
2✔
210
                        log.Log.V(4).Object(migration).Infof("closing outbound connection for migrationID %s", migration.Spec.SendTo.MigrationID)
1✔
211
                        if err := s.closeConnectionForMigrationID(s.syncOutboundConnectionMap, migration.Spec.SendTo.MigrationID); err != nil {
1✔
212
                                log.Log.Reason(err).Infof("unable to close connection for migrationID %s, possibly leaked connection", migration.Spec.SendTo.MigrationID)
×
213
                        }
×
214
                }
215
        }
216
}
217

218
func (s *SynchronizationController) closeConnectionForMigrationID(syncMap *sync.Map, migrationID string) error {
5✔
219
        obj, loaded := syncMap.LoadAndDelete(migrationID)
5✔
220
        if loaded {
9✔
221
                log.Log.V(4).Infof("closing connection associated with migrationID %s", migrationID)
4✔
222
                outboundConnection, ok := obj.(*SynchronizationConnection)
4✔
223
                if ok {
7✔
224
                        if err := outboundConnection.Close(); err != nil {
3✔
225
                                log.Log.Warningf("unable to close connection for migrationID %s, %v", migrationID, err)
×
226
                                s.failedCloseConnections.Store(outboundConnection, 0)
×
227
                                return err
×
228
                        }
×
229
                } else {
1✔
230
                        log.Log.Warningf("unable to close connection for migrationID %s, type is %v", migrationID, obj)
1✔
231
                        return fmt.Errorf("unknown type %v", obj)
1✔
232
                }
1✔
233
        }
234
        return nil
4✔
235
}
236

237
func (s *SynchronizationController) updateMigrationFunc(_, curr interface{}) {
1✔
238
        s.enqueueVirtualMachineInstanceFromMigration(curr)
1✔
239
}
1✔
240

241
func (s *SynchronizationController) enqueueVirtualMachineInstanceFromMigration(obj interface{}) {
4✔
242
        migration, ok := obj.(*virtv1.VirtualMachineInstanceMigration)
4✔
243
        if ok {
8✔
244
                key := controller.NamespacedKey(migration.Namespace, migration.Spec.VMIName)
4✔
245
                s.queue.Add(key)
4✔
246
        }
4✔
247
}
248

249
func (s *SynchronizationController) Run(threadiness int, stopCh <-chan struct{}) error {
×
250
        defer controller.HandlePanic()
×
251
        defer s.queue.ShutDown()
×
252
        defer s.closeConnections()
×
253

×
254
        log.Log.Info("starting vmi status synchronization controller.")
×
255

×
256
        // Wait for cache sync before we start the pod controller
×
257
        cache.WaitForCacheSync(stopCh, s.hasSynced)
×
258

×
259
        // Start the actual work
×
260
        for i := 0; i < threadiness; i++ {
×
261
                go wait.Until(s.runWorker, time.Second, stopCh)
×
262
        }
×
263
        go wait.Until(s.runConnectionCleanup, 5*time.Second, stopCh)
×
264

×
265
        conn, err := s.createTcpListener()
×
266
        if err != nil {
×
267
                log.Log.Criticalf("received error %v, exiting", err)
×
268
                return err
×
269
        } else {
×
270
                go func() {
×
271
                        s.grpcServer.Serve(conn)
×
272
                }()
×
273
        }
274
        if err := s.rebuildConnectionsAndUpdateSyncAddress(); err != nil {
×
275
                return err
×
276
        }
×
277

278
        log.Log.Info("waiting on stop signal")
×
279
        <-stopCh
×
280
        log.Log.Info("normally stopping vmi status synchronization controller.")
×
281
        return nil
×
282
}
283

284
func (s *SynchronizationController) closeConnections() {
31✔
285
        log.Log.V(1).Info("closing listener and grpcserver")
31✔
286
        if s.listener != nil {
62✔
287
                s.listener.Close()
31✔
288
        }
31✔
289
        if s.grpcServer != nil {
62✔
290
                s.grpcServer.GracefulStop()
31✔
291
        }
31✔
292
        log.Log.V(1).Infof("closing outbound connections")
31✔
293
        s.syncOutboundConnectionMap.Range(closeMapConnections)
31✔
294
        log.Log.V(1).Infof("closing inbound connections")
31✔
295
        s.syncReceivingConnectionMap.Range(closeMapConnections)
31✔
296
}
297

298
func closeMapConnections(k, obj interface{}) bool {
12✔
299
        outboundConnection, ok := obj.(*SynchronizationConnection)
12✔
300
        if ok && outboundConnection != nil {
24✔
301
                log.Log.V(1).Infof("closing connection for migration ID: %s", outboundConnection.migrationID)
12✔
302
                if err := outboundConnection.Close(); err != nil {
12✔
303
                        log.Log.Warningf("unable to close connection for VMI %s during shutdown, %v", k, err)
×
304
                }
×
305
        } else {
×
306
                log.Log.Warningf("unable to close connection for VMI %s during shutdown", k)
×
307
        }
×
308
        return true
12✔
309
}
310

311
func (s *SynchronizationController) runWorker() {
×
312
        for s.Execute() {
×
313
        }
×
314
}
315

316
func (s *SynchronizationController) Execute() bool {
1✔
317
        key, quit := s.queue.Get()
1✔
318
        if quit {
1✔
319
                return false
×
320
        }
×
321

322
        defer s.queue.Done(key)
1✔
323
        err := s.execute(key)
1✔
324

1✔
325
        if err != nil {
1✔
326
                log.Log.Reason(err).Infof("reenqueuing VirtualMachineInstance %v", key)
×
327
                s.queue.AddRateLimited(key)
×
328
        } else {
1✔
329
                log.Log.V(4).Infof("processed VirtualMachineInstance %v", key)
1✔
330
                s.queue.Forget(key)
1✔
331
        }
1✔
332
        return true
1✔
333
}
334

335
func (s *SynchronizationController) execute(key string) error {
10✔
336
        // Fetch the latest VMI state from cache
10✔
337
        obj, exists, _ := s.vmiInformer.GetStore().GetByKey(key)
10✔
338
        if !exists {
10✔
339
                return nil
×
340
        }
×
341
        vmi := obj.(*virtv1.VirtualMachineInstance)
10✔
342

10✔
343
        migration, err := s.getMigrationForVMI(vmi)
10✔
344
        if err != nil {
10✔
345
                return err
×
346
        }
×
347
        if migration != nil && migration.IsDecentralized() {
18✔
348
                if err := s.handleMigrationFinalizer(migration); err != nil {
8✔
349
                        return err
×
350
                }
×
351
                if migration.IsDecentralizedSource() {
12✔
352
                        if migration.DeletionTimestamp != nil {
5✔
353
                                log.Log.V(2).Object(migration).Infof("migration is being deleted, syncing source state before canceling target migration")
1✔
354
                                if err := s.handleSourceState(vmi.DeepCopy(), migration); err != nil {
1✔
NEW
355
                                        return s.updateDecentralizedFailureOnSource(vmi, migration, err)
×
NEW
356
                                }
×
357
                                if err := s.cancelTargetRemoteMigration(vmi, migration); err != nil {
1✔
358
                                        return err
×
359
                                }
×
360
                                return nil
1✔
361
                        }
362
                        err := s.handleSourceState(vmi.DeepCopy(), migration)
3✔
363
                        return s.updateDecentralizedFailureOnSource(vmi, migration, err)
3✔
364
                }
365
                if migration.IsDecentralizedTarget() {
8✔
366
                        if migration.DeletionTimestamp != nil {
5✔
367
                                log.Log.V(2).Object(migration).Infof("migration is being deleted, informing the source that the migration is canceled")
1✔
368
                                // migration is being deleted, inform the target that the migration is canceled
1✔
369
                                if err := s.cancelSourceRemoteMigration(vmi, migration); err != nil {
1✔
370
                                        return err
×
371
                                }
×
372
                                return nil
1✔
373
                        }
374
                        err := s.handleTargetState(vmi.DeepCopy(), migration)
3✔
375
                        return s.updateDecentralizedFailureOnTarget(vmi, migration, err)
3✔
376
                }
377
                return nil
×
378
        } else {
2✔
379
                // After a target-initiated cancel the source migration object may be removed
2✔
380
                // before the source VMI records abort completion. Push the final source state
2✔
381
                // once more so the target VMI receives EndTimestamp and can finish abort cleanup.
2✔
382
                if needsOrphanSourceAbortSync(vmi) {
3✔
383
                        log.Log.Object(vmi).Info("migration object gone after source abort, pushing final source state to target")
1✔
384
                        if err := s.handleSourceState(vmi.DeepCopy(), nil); err != nil {
1✔
NEW
385
                                return err
×
NEW
386
                        }
×
387
                }
388
                // No migration found don't do anything
389
                // We should only clear the condition if we are not waiting for synchronization.
390
                if err := s.clearDecentralizedLiveMigrationFailure(vmi); err != nil {
2✔
391
                        return err
×
392
                }
×
393
                log.Log.Object(vmi).V(4).Info("no active decentralized migration found for VMI")
2✔
394
                return nil
2✔
395
        }
396
}
397

398
func needsOrphanSourceAbortSync(vmi *virtv1.VirtualMachineInstance) bool {
6✔
399
        ms := vmi.Status.MigrationState
6✔
400
        if ms == nil || ms.SourceState == nil || ms.TargetState == nil {
7✔
401
                return false
1✔
402
        }
1✔
403
        if ms.MigrationUID != ms.SourceState.MigrationUID {
6✔
404
                return false
1✔
405
        }
1✔
406
        return ms.AbortRequested && ms.EndTimestamp != nil
4✔
407
}
408

409
func (s *SynchronizationController) updateDecentralizedFailureOnSource(vmi *virtv1.VirtualMachineInstance, migration *virtv1.VirtualMachineInstanceMigration, opErr error) error {
6✔
410
        if opErr != nil {
8✔
411
                if err := s.setDecentralizedLiveMigrationFailure(vmi, getErrorMessageForDecentralizedLiveMigrationFailure(opErr)); err != nil {
2✔
412
                        return err
×
413
                }
×
414
                return opErr
2✔
415
        }
416
        return s.clearDecentralizedLiveMigrationFailure(vmi)
4✔
417
}
418

419
func (s *SynchronizationController) updateDecentralizedFailureOnTarget(vmi *virtv1.VirtualMachineInstance, migration *virtv1.VirtualMachineInstanceMigration, opErr error) error {
8✔
420
        // failure from handleTargetState
8✔
421
        if opErr != nil {
11✔
422
                return s.setDecentralizedLiveMigrationFailure(vmi, getErrorMessageForDecentralizedLiveMigrationFailure(opErr))
3✔
423
        }
3✔
424

425
        // success: special case waiting for sync
426
        if migration.Status.Phase == virtv1.MigrationWaitingForSync {
6✔
427
                return s.setDecentralizedLiveMigrationFailure(vmi, waitingForSyncErrorMessage)
1✔
428
        }
1✔
429

430
        // success and not waiting for sync: clear condition
431
        return s.clearDecentralizedLiveMigrationFailure(vmi)
4✔
432
}
433

434
func getErrorMessageForDecentralizedLiveMigrationFailure(err error) string {
7✔
435
        log.Log.V(1).Infof("error message: %v", err)
7✔
436
        // Once we upgrade to golang 1.26, we no longer need to check for x509.HostnameError, since https://github.com/golang/go/issues/76445 will be fixed.
7✔
437
        if errors.As(err, &x509.HostnameError{}) {
10✔
438
                return "x509 hostname error"
3✔
439
        }
3✔
440
        message := ""
4✔
441
        if err != nil {
8✔
442
                message = err.Error()
4✔
443
        }
4✔
444
        return message
4✔
445
}
446

447
func (s *SynchronizationController) setDecentralizedLiveMigrationFailure(vmi *virtv1.VirtualMachineInstance, errorMessage string) error {
9✔
448
        orgVmi := vmi.DeepCopy()
9✔
449
        condition := &virtv1.VirtualMachineInstanceCondition{
9✔
450
                Type:               virtv1.VirtualMachineInstanceDecentralizedLiveMigrationFailure,
9✔
451
                Status:             k8sv1.ConditionTrue,
9✔
452
                Reason:             virtv1.VirtualMachineInstanceReasonDecentralizedNotMigratable,
9✔
453
                Message:            errorMessage,
9✔
454
                LastTransitionTime: metav1.Now(),
9✔
455
        }
9✔
456
        controller.NewVirtualMachineInstanceConditionManager().UpdateCondition(vmi, condition)
9✔
457
        if err2 := s.patchVMIConditions(context.Background(), orgVmi, vmi); err2 != nil {
9✔
458
                log.Log.Reason(err2).Infof("unable to patch VMI conditions after decentralized live migration failure")
×
459
                return err2
×
460
        }
×
461
        return nil
9✔
462
}
463

464
func (s *SynchronizationController) clearDecentralizedLiveMigrationFailure(vmi *virtv1.VirtualMachineInstance) error {
11✔
465
        orgVmi := vmi.DeepCopy()
11✔
466
        controller.NewVirtualMachineInstanceConditionManager().RemoveCondition(vmi, virtv1.VirtualMachineInstanceDecentralizedLiveMigrationFailure)
11✔
467
        if err := s.patchVMIConditions(context.Background(), orgVmi, vmi); err != nil {
11✔
468
                log.Log.Reason(err).Infof("unable to patch VMI conditions after clearing decentralized live migration failure")
×
469
                return err
×
470
        }
×
471
        return nil
11✔
472
}
473

474
func (s *SynchronizationController) handleMigrationFinalizer(migration *virtv1.VirtualMachineInstanceMigration) error {
8✔
475
        originalMigration := migration.DeepCopy()
8✔
476
        if !migration.IsFinal() && migration.DeletionTimestamp == nil {
14✔
477
                controller.AddFinalizer(migration, SynchronizationFinalizer)
6✔
478
        } else {
8✔
479
                controller.RemoveFinalizer(migration, SynchronizationFinalizer)
2✔
480
        }
2✔
481
        if !apiequality.Semantic.DeepEqual(originalMigration.ObjectMeta, migration.ObjectMeta) {
12✔
482
                log.Log.V(4).Object(migration).Infof("adding or removing finalizer to migration, %v", migration.Finalizers)
4✔
483
                patchSet := patch.New()
4✔
484
                patchSet.AddOption(
4✔
485
                        patch.WithReplace("/metadata/finalizers", migration.Finalizers),
4✔
486
                )
4✔
487
                patchBytes, err := patchSet.GeneratePayload()
4✔
488
                if err != nil {
4✔
489
                        return err
×
490
                }
×
491
                if _, err := s.client.VirtualMachineInstanceMigration(migration.Namespace).Patch(context.Background(), migration.Name, types.JSONPatchType, patchBytes, metav1.PatchOptions{}); err != nil {
4✔
492
                        return err
×
493
                }
×
494
        }
495
        return nil
8✔
496
}
497

498
func (s *SynchronizationController) cancelSourceRemoteMigration(vmi *virtv1.VirtualMachineInstance, migration *virtv1.VirtualMachineInstanceMigration) error {
1✔
499
        if vmi == nil || vmi.Status.MigrationState == nil || vmi.Status.MigrationState.SourceState == nil || vmi.Status.MigrationState.SourceState.MigrationUID == "" {
2✔
500
                return nil
1✔
501
        }
1✔
502
        if migration.UID == vmi.Status.MigrationState.SourceState.MigrationUID {
×
503
                return fmt.Errorf("source migration UID %s is the same as the VMI's migration UID %s", migration.UID, vmi.Status.MigrationState.SourceState.MigrationUID)
×
504
        }
×
505
        log.Log.V(4).Object(migration).Infof("cancelling source remote migration for VMI %s/%s", vmi.Namespace, vmi.Name)
×
506
        return s.cancelRemoteMigration(vmi.Status.MigrationState.SourceState.MigrationUID, migration.Spec.Receive.MigrationID, s.syncOutboundConnectionMap)
×
507
}
508

509
func (s *SynchronizationController) cancelTargetRemoteMigration(vmi *virtv1.VirtualMachineInstance, migration *virtv1.VirtualMachineInstanceMigration) error {
1✔
510
        if vmi == nil || vmi.Status.MigrationState == nil || vmi.Status.MigrationState.TargetState == nil || vmi.Status.MigrationState.TargetState.MigrationUID == "" {
1✔
511
                return nil
×
512
        }
×
513
        if migration.UID == vmi.Status.MigrationState.TargetState.MigrationUID {
1✔
514
                return fmt.Errorf("target migration UID %s is the same as the VMI's migration UID %s", migration.UID, vmi.Status.MigrationState.TargetState.MigrationUID)
×
515
        }
×
516
        log.Log.V(4).Object(migration).Infof("cancelling target remote migration for VMI %s/%s", vmi.Namespace, vmi.Name)
1✔
517
        return s.cancelRemoteMigration(vmi.Status.MigrationState.TargetState.MigrationUID, migration.Spec.SendTo.MigrationID, s.syncReceivingConnectionMap)
1✔
518
}
519

520
func (s *SynchronizationController) cancelRemoteMigration(migrationUID types.UID, migrationID string, connectionMap *sync.Map) error {
1✔
521
        obj, ok := connectionMap.Load(migrationID)
1✔
522
        if !ok {
2✔
523
                // No connection found, don't do anything
1✔
524
                return nil
1✔
525
        }
1✔
526
        conn, ok := obj.(*SynchronizationConnection)
×
527
        if !ok {
×
528
                return fmt.Errorf("found unknown object in outbound connection cache %#v", conn)
×
529
        }
×
530
        client := syncv1.NewSynchronizeClient(conn.grpcClientConnection)
×
531
        ctx, cancel := context.WithTimeout(context.Background(), time.Duration(s.timeout)*time.Second)
×
532
        defer cancel()
×
533
        _, err := client.CancelMigration(ctx, &syncv1.MigrationCancelRequest{
×
534
                MigrationUID: string(migrationUID),
×
535
        })
×
536
        if err != nil {
×
537
                return err
×
538
        }
×
539
        return nil
×
540
}
541

542
func (s *SynchronizationController) getMigrationIDFromUID(migrationUID types.UID) (string, error) {
17✔
543
        objs, err := s.migrationInformer.GetIndexer().ByIndex("byUID", string(migrationUID))
17✔
544
        if err != nil {
17✔
545
                return "", err
×
546
        }
×
547
        if len(objs) > 1 {
17✔
548
                return "", fmt.Errorf("found more than one migration with same UID")
×
549
        }
×
550
        if len(objs) == 0 {
19✔
551
                return "", nil
2✔
552
        }
2✔
553
        migration, ok := objs[0].(*virtv1.VirtualMachineInstanceMigration)
15✔
554
        if !ok {
15✔
555
                return "", fmt.Errorf("found unknown object in migration cache")
×
556
        }
×
557
        var migrationID string
15✔
558
        if migration.Spec.Receive != nil {
23✔
559
                migrationID = migration.Spec.Receive.MigrationID
8✔
560
        }
8✔
561
        if migration.Spec.SendTo != nil {
22✔
562
                migrationID = migration.Spec.SendTo.MigrationID
7✔
563
        }
7✔
564
        return migrationID, nil
15✔
565
}
566

567
func (s *SynchronizationController) getOutboundSourceConnection(vmi *virtv1.VirtualMachineInstance, migrationState *virtv1.VirtualMachineInstanceMigrationState) (*SynchronizationConnection, error) {
9✔
568
        if migrationState.TargetState == nil || migrationState.TargetState.SyncAddress == nil || *migrationState.TargetState.SyncAddress == "" {
10✔
569
                return nil, nil
1✔
570
        }
1✔
571
        migrationID, err := s.resolveSourceOutboundMigrationID(migrationState)
8✔
572
        if err != nil {
8✔
NEW
573
                return nil, err
×
NEW
574
        }
×
575
        if migrationID == "" {
8✔
NEW
576
                return nil, nil
×
NEW
577
        }
×
578
        return s.getOutboundConnectionByMigrationID(vmi, migrationID, *migrationState.TargetState.SyncAddress, s.syncOutboundConnectionMap)
8✔
579
}
580

581
func (s *SynchronizationController) resolveSourceOutboundMigrationID(migrationState *virtv1.VirtualMachineInstanceMigrationState) (string, error) {
9✔
582
        if migrationState == nil || migrationState.SourceState == nil {
9✔
NEW
583
                return "", nil
×
NEW
584
        }
×
585
        migrationID, err := s.getMigrationIDFromUID(migrationState.SourceState.MigrationUID)
9✔
586
        if err != nil || migrationID != "" {
16✔
587
                return migrationID, err
7✔
588
        }
7✔
589
        // The source migration object may already be deleted after a target-initiated cancel.
590
        if migrationState.TargetState != nil && migrationState.TargetState.MigrationUID != "" {
4✔
591
                return s.getMigrationIDFromUID(migrationState.TargetState.MigrationUID)
2✔
592
        }
2✔
NEW
593
        return "", nil
×
594
}
595

596
func (s *SynchronizationController) getOutboundTargetConnection(vmi *virtv1.VirtualMachineInstance, migrationState *virtv1.VirtualMachineInstanceMigrationState) (*SynchronizationConnection, error) {
7✔
597
        if migrationState.SourceState == nil || migrationState.SourceState.SyncAddress == nil || *migrationState.SourceState.SyncAddress == "" {
8✔
598
                return nil, nil
1✔
599
        }
1✔
600
        return s.getOutboundConnection(vmi, migrationState.TargetState.MigrationUID, *migrationState.SourceState.SyncAddress, s.syncReceivingConnectionMap)
6✔
601
}
602

603
func (s *SynchronizationController) getOutboundConnection(vmi *virtv1.VirtualMachineInstance, migrationUID types.UID, syncAddress string, connectionMap *sync.Map) (*SynchronizationConnection, error) {
6✔
604
        if migrationUID == "" {
6✔
605
                return nil, nil
×
606
        }
×
607
        migrationID, err := s.getMigrationIDFromUID(migrationUID)
6✔
608
        if err != nil {
6✔
609
                return nil, err
×
610
        }
×
611
        return s.getOutboundConnectionByMigrationID(vmi, migrationID, syncAddress, connectionMap)
6✔
612
}
613

614
func (s *SynchronizationController) getOutboundConnectionByMigrationID(vmi *virtv1.VirtualMachineInstance, migrationID string, syncAddress string, connectionMap *sync.Map) (*SynchronizationConnection, error) {
14✔
615
        if migrationID == "" {
14✔
NEW
616
                return nil, nil
×
NEW
617
        }
×
618
        log.Log.Object(vmi).V(4).Infof("found migration ID %s", migrationID)
14✔
619
        obj, ok := connectionMap.Load(migrationID)
14✔
620
        if !ok {
26✔
621
                grpcClientConnection, err := s.createOutboundConnection(syncAddress)
12✔
622
                if err != nil {
12✔
623
                        return nil, err
×
624
                }
×
625
                conn := &SynchronizationConnection{
12✔
626
                        migrationID:          migrationID,
12✔
627
                        grpcClientConnection: grpcClientConnection,
12✔
628
                }
12✔
629
                connectionMap.Store(migrationID, conn)
12✔
630
                return conn, nil
12✔
631
        }
632
        outboundSyncConnection, ok := obj.(*SynchronizationConnection)
2✔
633
        if !ok {
2✔
634
                return nil, fmt.Errorf("found unknown object in outbound connection cache %#v", outboundSyncConnection)
×
635
        }
×
636
        return outboundSyncConnection, nil
2✔
637
}
638

639
func (s *SynchronizationController) handleSourceState(vmi *virtv1.VirtualMachineInstance, migration *virtv1.VirtualMachineInstanceMigration) error {
10✔
640
        var outboundConnection *SynchronizationConnection
10✔
641
        var err error
10✔
642
        if vmi.Status.MigrationState == nil {
11✔
643
                // No migration state, don't do anything
1✔
644
                return nil
1✔
645
        }
1✔
646
        if vmi.Status.MigrationState.SourceState == nil || vmi.Status.MigrationState.TargetState == nil {
11✔
647
                // No migration state, don't do anything
2✔
648
                return nil
2✔
649
        }
2✔
650

651
        sourceState := vmi.Status.MigrationState.SourceState
7✔
652
        if sourceState.SyncAddress == nil || *sourceState.SyncAddress == "" {
10✔
653
                syncAddress, err := s.getLocalSynchronizationAddress()
3✔
654
                if err != nil {
3✔
655
                        return err
×
656
                }
×
657
                sourceState.SyncAddress = &syncAddress
3✔
658
        }
659
        targetState := vmi.Status.MigrationState.TargetState
7✔
660
        if targetState.SyncAddress != nil && sourceState.MigrationUID != "" {
14✔
661
                if outboundConnection, err = s.getOutboundSourceConnection(vmi, vmi.Status.MigrationState); err != nil {
7✔
662
                        return err
×
663
                }
×
664
        }
665
        if outboundConnection == nil {
7✔
666
                log.Log.Object(vmi).V(4).Info("no synchronization connection found for source, doing nothing")
×
667
                return nil
×
668
        }
×
669
        vmiStatusJson, err := json.Marshal(vmi.Status)
7✔
670
        if err != nil {
7✔
671
                return err
×
672
        }
×
673
        client := syncv1.NewSynchronizeClient(outboundConnection.grpcClientConnection)
7✔
674
        ctx, cancel := context.WithTimeout(context.Background(), time.Duration(s.timeout)*time.Second)
7✔
675
        defer cancel()
7✔
676

7✔
677
        if _, err := client.SyncSourceMigrationStatus(ctx, &syncv1.VMIStatusRequest{
7✔
678
                MigrationID: outboundConnection.migrationID,
7✔
679
                VmiStatus: &syncv1.VMIStatus{
7✔
680
                        VmiStatusJson: vmiStatusJson,
7✔
681
                },
7✔
682
        }); err != nil {
8✔
683
                return err
1✔
684
        }
1✔
685
        if migration != nil && migration.IsFinal() {
6✔
686
                if migration.Spec.SendTo != nil {
×
687
                        log.Log.Object(migration).Infof("completed migration for VMI %s/%s, closing outbound connections", migration.Namespace, migration.Spec.VMIName)
×
688
                        s.closeConnectionForMigrationID(s.syncOutboundConnectionMap, migration.Spec.SendTo.MigrationID)
×
689
                }
×
690
        }
691

692
        return nil
6✔
693
}
694

695
func (s *SynchronizationController) handleTargetState(vmi *virtv1.VirtualMachineInstance, migration *virtv1.VirtualMachineInstanceMigration) error {
8✔
696
        if vmi.Status.MigrationState == nil {
9✔
697
                // No migration state, don't do anything
1✔
698
                return nil
1✔
699
        }
1✔
700
        if vmi.Status.MigrationState.TargetState == nil || vmi.Status.MigrationState.SourceState == nil {
9✔
701
                // No migration state, don't do anything
2✔
702
                return nil
2✔
703
        }
2✔
704

705
        var outboundConnection *SynchronizationConnection
5✔
706
        var err error
5✔
707
        sourceState := vmi.Status.MigrationState.SourceState
5✔
708
        targetState := vmi.Status.MigrationState.TargetState
5✔
709
        if targetState.SyncAddress == nil || *targetState.SyncAddress == "" {
8✔
710
                syncAddress, err := s.getLocalSynchronizationAddress()
3✔
711
                if err != nil {
3✔
712
                        return err
×
713
                }
×
714
                targetState.SyncAddress = &syncAddress
3✔
715
        }
716

717
        if sourceState.SyncAddress != nil && targetState.MigrationUID != "" {
10✔
718
                if outboundConnection, err = s.getOutboundTargetConnection(vmi, vmi.Status.MigrationState); err != nil {
5✔
719
                        return err
×
720
                }
×
721
        }
722
        if outboundConnection == nil {
5✔
723
                log.Log.Object(vmi).V(4).Info("no synchronization connection found for target, doing nothing")
×
724
                return nil
×
725
        }
×
726

727
        vmiStatusJson, err := json.Marshal(vmi.Status)
5✔
728
        if err != nil {
5✔
729
                return err
×
730
        }
×
731
        client := syncv1.NewSynchronizeClient(outboundConnection.grpcClientConnection)
5✔
732
        ctx, cancel := context.WithTimeout(context.Background(), time.Duration(s.timeout)*time.Second)
5✔
733
        defer cancel()
5✔
734

5✔
735
        _, err = client.SyncTargetMigrationStatus(ctx, &syncv1.VMIStatusRequest{
5✔
736
                MigrationID: outboundConnection.migrationID,
5✔
737
                VmiStatus: &syncv1.VMIStatus{
5✔
738
                        VmiStatusJson: vmiStatusJson,
5✔
739
                },
5✔
740
        })
5✔
741
        if err != nil {
6✔
742
                return err
1✔
743
        }
1✔
744
        if migration.IsFinal() {
4✔
745
                if migration.Spec.Receive != nil {
×
746
                        log.Log.Object(migration).Infof("completed migration for VMI %s/%s, closing receiving connections", migration.Namespace, migration.Spec.VMIName)
×
747
                        s.closeConnectionForMigrationID(s.syncReceivingConnectionMap, migration.Spec.Receive.MigrationID)
×
748
                }
×
749
        }
750

751
        return nil
4✔
752
}
753

754
func (s *SynchronizationController) getMigrationForVMI(vmi *virtv1.VirtualMachineInstance) (*virtv1.VirtualMachineInstanceMigration, error) {
10✔
755
        objects, err := s.migrationInformer.GetIndexer().ByIndex("byActiveVMIName", vmi.Name)
10✔
756
        if err != nil {
10✔
757
                return nil, err
×
758
        }
×
759
        if len(objects) > 0 {
18✔
760
                count := 0
8✔
761
                var res *virtv1.VirtualMachineInstanceMigration
8✔
762
                for _, migrationObj := range objects {
16✔
763
                        migration, ok := migrationObj.(*virtv1.VirtualMachineInstanceMigration)
8✔
764
                        if !ok {
8✔
765
                                return nil, fmt.Errorf("not a virtual machine instance migration")
×
766
                        }
×
767
                        if migration.Namespace == vmi.Namespace {
16✔
768
                                if migration.IsDecentralizedSource() {
12✔
769
                                        if vmi.Status.MigrationState != nil && vmi.Status.MigrationState.SourceState != nil && migration.UID == vmi.Status.MigrationState.SourceState.MigrationUID {
8✔
770
                                                count++
4✔
771
                                                res = migration
4✔
772
                                        }
4✔
773
                                } else if migration.IsDecentralizedTarget() {
8✔
774
                                        if vmi.Status.MigrationState != nil && vmi.Status.MigrationState.TargetState != nil && migration.UID == vmi.Status.MigrationState.TargetState.MigrationUID {
8✔
775
                                                count++
4✔
776
                                                res = migration
4✔
777
                                        }
4✔
778
                                }
779
                        }
780
                }
781
                if count > 1 {
8✔
782
                        return nil, fmt.Errorf("found more than one migration pointing to same VMI")
×
783
                } else if count == 0 {
8✔
784
                        return nil, nil
×
785
                }
×
786
                return res, nil
8✔
787
        }
788
        return nil, nil
2✔
789
}
790

791
func (s *SynchronizationController) rebuildConnectionsAndUpdateSyncAddress() error {
7✔
792
        // Go and find all active migration resources, if they are decentralized rebuild either
7✔
793
        // the incoming or outbound connections, and call sync to update the remote with the new
7✔
794
        // address.
7✔
795
        objs := s.migrationInformer.GetStore().List()
7✔
796
        log.Log.V(4).Infof("rebuilding any connections, and updating remote VMIs, found %d migrations", len(objs))
7✔
797
        for _, obj := range objs {
13✔
798
                migration, ok := obj.(*virtv1.VirtualMachineInstanceMigration)
6✔
799
                if !ok {
6✔
800
                        return fmt.Errorf("unknown object in migration store %v", obj)
×
801
                }
×
802
                if isOnGoingMigration(migration) {
11✔
803
                        vmi, err := s.getVMIFromMigration(migration)
5✔
804
                        if err != nil {
5✔
805
                                return err
×
806
                        }
×
807
                        if vmi == nil {
6✔
808
                                // No VMI found, can't update it, so skip it.
1✔
809
                                continue
1✔
810
                        }
811
                        // ongoing migration.
812
                        if migration.Spec.Receive != nil {
6✔
813
                                // We are the target
2✔
814
                                log.Log.Object(migration).Object(vmi).Info("found ongoing target migration for vmi, rebuilding connection")
2✔
815
                                if err := s.rebuildTargetConnection(migration, vmi); err != nil {
2✔
816
                                        return err
×
817
                                }
×
818
                        } else if migration.Spec.SendTo != nil {
4✔
819
                                // We are the source
2✔
820
                                log.Log.Object(migration).Object(vmi).Info("found ongoing source migration for vmi, rebuilding connection")
2✔
821
                                if err := s.rebuildSourceConnection(migration, vmi); err != nil {
2✔
822
                                        return err
×
823
                                }
×
824
                        }
825
                }
826
        }
827
        return nil
7✔
828
}
829

830
func isOnGoingMigration(migration *virtv1.VirtualMachineInstanceMigration) bool {
6✔
831
        return migration.IsDecentralized() && migration.Status.Phase != virtv1.MigrationFailed && migration.Status.Phase != virtv1.MigrationSucceeded
6✔
832
}
6✔
833

834
func (s *SynchronizationController) rebuildTargetConnection(migration *virtv1.VirtualMachineInstanceMigration, vmi *virtv1.VirtualMachineInstance) error {
2✔
835
        conn, err := s.getOutboundTargetConnection(vmi, vmi.Status.MigrationState)
2✔
836
        if err != nil {
2✔
837
                return err
×
838
        }
×
839
        if conn == nil {
3✔
840
                return nil
1✔
841
        }
1✔
842
        s.syncReceivingConnectionMap.Store(migration.Spec.Receive.MigrationID, conn)
1✔
843
        if vmi.Status.MigrationState != nil && vmi.Status.MigrationState.TargetState != nil {
2✔
844
                url, err := s.getLocalSynchronizationAddress()
1✔
845
                if err != nil {
1✔
846
                        return err
×
847
                }
×
848
                origVMI := vmi.DeepCopy()
1✔
849
                vmi.Status.MigrationState.TargetState.SyncAddress = &url
1✔
850
                // patching will cause reconcile loop to connect to remote to update
1✔
851
                if err := s.patchVMI(context.Background(), origVMI, vmi); err != nil {
1✔
852
                        return err
×
853
                }
×
854
        }
855
        return nil
1✔
856
}
857

858
func (s *SynchronizationController) rebuildSourceConnection(migration *virtv1.VirtualMachineInstanceMigration, vmi *virtv1.VirtualMachineInstance) error {
2✔
859
        conn, err := s.getOutboundSourceConnection(vmi, vmi.Status.MigrationState)
2✔
860
        if err != nil {
2✔
861
                return err
×
862
        }
×
863
        if conn == nil {
3✔
864
                return nil
1✔
865
        }
1✔
866
        s.syncOutboundConnectionMap.Store(migration.Spec.SendTo.MigrationID, conn)
1✔
867
        if vmi.Status.MigrationState != nil && vmi.Status.MigrationState.SourceState != nil {
2✔
868
                url, err := s.getLocalSynchronizationAddress()
1✔
869
                if err != nil {
1✔
870
                        return err
×
871
                }
×
872
                origVMI := vmi.DeepCopy()
1✔
873
                vmi.Status.MigrationState.SourceState.SyncAddress = &url
1✔
874
                // patching will cause reconcile loop to connect to remote to update
1✔
875
                if err := s.patchVMI(context.Background(), origVMI, vmi); err != nil {
1✔
876
                        return err
×
877
                }
×
878
        }
879
        return nil
1✔
880
}
881

882
func (s *SynchronizationController) getVMIFromMigration(migration *virtv1.VirtualMachineInstanceMigration) (*virtv1.VirtualMachineInstance, error) {
7✔
883
        key := controller.NamespacedKey(migration.Namespace, migration.Spec.VMIName)
7✔
884
        obj, exists, err := s.vmiInformer.GetStore().GetByKey(key)
7✔
885
        if err != nil {
7✔
886
                return nil, err
×
887
        }
×
888
        if !exists {
9✔
889
                return nil, nil
2✔
890
        }
2✔
891
        return obj.(*virtv1.VirtualMachineInstance).DeepCopy(), nil
5✔
892
}
893

894
func (s *SynchronizationController) getLocalSynchronizationAddress() (string, error) {
45✔
895
        if s.ip != "" {
49✔
896
                return net.JoinHostPort(s.ip, strconv.Itoa(s.bindPort)), nil
4✔
897
        }
4✔
898
        // TODO figure out how to get my URL with or without submariner (url changes based on export)
899
        return s.listener.Addr().String(), nil
41✔
900
}
901

902
func (s *SynchronizationController) createOutboundConnection(connectionURL string) (*grpc.ClientConn, error) {
12✔
903
        logger := log.Log.With("outbound", connectionURL)
12✔
904
        logger.Info("creating new synchronization grpc connection")
12✔
905

12✔
906
        client, err := grpc.NewClient(connectionURL, grpc.WithTransportCredentials(credentials.NewTLS(s.clientTLSConfig)))
12✔
907
        return client, err
12✔
908
}
12✔
909

910
func (s *SynchronizationController) createTcpListener() (net.Listener, error) {
31✔
911
        if s.listener != nil {
31✔
912
                return s.listener, nil
×
913
        }
×
914
        var ln net.Listener
31✔
915
        var err error
31✔
916
        addr := net.JoinHostPort(s.bindAddress, strconv.Itoa(s.bindPort))
31✔
917
        ln, err = net.Listen("tcp", addr)
31✔
918
        if err != nil {
31✔
919
                log.Log.Reason(err).Error("failed to create tcp listener")
×
920
                return nil, err
×
921
        }
×
922
        s.listener = ln
31✔
923
        return ln, nil
31✔
924
}
925

926
func (s *SynchronizationController) findTargetMigrationFromMigrationID(migrationID string) (*virtv1.VirtualMachineInstanceMigration, error) {
28✔
927
        return s.findMigrationFromMigrationIDByIndex("byTargetMigrationID", migrationID)
28✔
928
}
28✔
929

930
func (s *SynchronizationController) findSourceMigrationFromMigrationID(migrationID string) (*virtv1.VirtualMachineInstanceMigration, error) {
24✔
931
        return s.findMigrationFromMigrationIDByIndex("bySourceMigrationID", migrationID)
24✔
932
}
24✔
933

934
func (s *SynchronizationController) findMigrationFromMigrationIDByIndex(indexName, migrationID string) (*virtv1.VirtualMachineInstanceMigration, error) {
53✔
935
        objs, err := s.migrationInformer.GetIndexer().ByIndex(indexName, migrationID)
53✔
936
        if err != nil {
53✔
937
                return nil, err
×
938
        }
×
939
        if len(objs) > 1 {
53✔
940
                log.Log.Warningf("found multiple migrations for migrationID %s, picking first one", migrationID)
×
941
        }
×
942
        for _, obj := range objs {
103✔
943
                migration, _ := obj.(*virtv1.VirtualMachineInstanceMigration)
50✔
944
                return migration, nil
50✔
945
        }
50✔
946
        return nil, nil
3✔
947
}
948

949
func (s *SynchronizationController) SyncSourceMigrationStatus(ctx context.Context, request *syncv1.VMIStatusRequest) (*syncv1.VMIStatusResponse, error) {
18✔
950
        if request.VmiStatus == nil || len(request.VmiStatus.VmiStatusJson) == 0 {
19✔
951
                return &syncv1.VMIStatusResponse{
1✔
952
                        Message: noSourceStatusErrorMsg,
1✔
953
                }, fmt.Errorf(noSourceStatusErrorMsg)
1✔
954
        }
1✔
955
        migration, err := s.findTargetMigrationFromMigrationID(request.MigrationID)
17✔
956
        if migration == nil {
18✔
957
                return &syncv1.VMIStatusResponse{
1✔
958
                        Message: fmt.Sprintf(sourceUnableToLocateVMIMigrationIDErrorMsg, request.MigrationID),
1✔
959
                }, fmt.Errorf(sourceUnableToLocateVMIMigrationIDErrorMsg, request.MigrationID)
1✔
960
        }
1✔
961
        key := controller.NamespacedKey(migration.Namespace, migration.Spec.VMIName)
16✔
962
        log.Log.Object(migration).V(5).Infof("looking up VMI %s", key)
16✔
963
        obj, exists, err := s.vmiInformer.GetStore().GetByKey(key)
16✔
964
        if err != nil || !exists {
18✔
965
                if err == nil {
4✔
966
                        err = fmt.Errorf(sourceUnableToLocateVMIMigrationIDErrorMsgVMI, request.MigrationID, key)
2✔
967
                }
2✔
968
                return &syncv1.VMIStatusResponse{
2✔
969
                        Message: fmt.Sprintf(sourceUnableToLocateVMIMigrationIDErrorMsgVMI, request.MigrationID, key),
2✔
970
                }, err
2✔
971
        }
972
        vmi := obj.(*virtv1.VirtualMachineInstance)
14✔
973
        remoteStatus := &virtv1.VirtualMachineInstanceStatus{}
14✔
974
        if err := json.Unmarshal(request.VmiStatus.VmiStatusJson, remoteStatus); err != nil {
15✔
975
                return &syncv1.VMIStatusResponse{
1✔
976
                        Message: fmt.Sprintf("unable to unmarshal vmistatus for migrationID %s", request.MigrationID),
1✔
977
                }, err
1✔
978
        }
1✔
979
        if remoteStatus.MigrationState == nil {
14✔
980
                return &syncv1.VMIStatusResponse{
1✔
981
                        Message: noSourceStatusErrorMsg,
1✔
982
                }, fmt.Errorf(noSourceStatusErrorMsg)
1✔
983
        }
1✔
984
        newVMI := vmi.DeepCopy()
12✔
985
        if newVMI.Status.MigrationState == nil {
17✔
986
                newVMI.Status.MigrationState = &virtv1.VirtualMachineInstanceMigrationState{}
5✔
987
        }
5✔
988

989
        // Only update SourceState if this migration is still active and matches the VMI's current migration.
990
        // This prevents stale updates from a completed decentralized migration from interfering with a new compute migration.
991
        if migration.IsFinal() {
13✔
992
                log.Log.Object(migration).Infof("Migration is final, ignoring source state update for VMI %s/%s", vmi.Namespace, vmi.Name)
1✔
993
                return &syncv1.VMIStatusResponse{
1✔
994
                        Message: successMessage,
1✔
995
                }, nil
1✔
996
        }
1✔
997

998
        // Check if the VMI's current migration matches this migration
999
        if newVMI.Status.MigrationState.MigrationUID != "" && newVMI.Status.MigrationState.MigrationUID != migration.UID {
12✔
1000
                log.Log.Object(migration).Warningf("VMI %s/%s has different migration UID %s, ignoring source state update for migration %s",
1✔
1001
                        vmi.Namespace, vmi.Name, newVMI.Status.MigrationState.MigrationUID, migration.UID)
1✔
1002
                return &syncv1.VMIStatusResponse{
1✔
1003
                        Message: successMessage,
1✔
1004
                }, nil
1✔
1005
        }
1✔
1006

1007
        log.Log.Object(newVMI).V(5).Infof("vmi migration source state: %#v", newVMI.Status.MigrationState.SourceState)
10✔
1008
        log.Log.Object(newVMI).V(5).Infof("remote migration source state: %#v", remoteStatus.MigrationState.SourceState)
10✔
1009
        newVMI.Status.MigrationState.SourceState = remoteStatus.MigrationState.SourceState.DeepCopy()
10✔
1010
        copyLegacySourceFields(newVMI, remoteStatus.MigrationState)
10✔
1011
        if len(remoteStatus.MigratedVolumes) > 0 {
10✔
1012
                log.Log.Object(newVMI).V(5).Infof("SyncSourceMigrationStatus: Copying migrated volumes to target state, %#v", newVMI.Status.MigratedVolumes)
×
1013
                newVMI.Status.MigratedVolumes = getMergedSourceMigratedVolumes(newVMI.Status.MigratedVolumes, remoteStatus.MigratedVolumes)
×
1014
        } else {
10✔
1015
                log.Log.Object(newVMI).V(5).Info("SyncSourceMigrationStatus: No source migrated volumes found")
10✔
1016
        }
10✔
1017
        newVMI.Status.MigrationMethod = remoteStatus.MigrationMethod
10✔
1018
        newVMI.Status.MigrationTransport = remoteStatus.MigrationTransport
10✔
1019
        if !apiequality.Semantic.DeepEqual(vmi.Status, newVMI.Status) {
20✔
1020
                if err := s.patchVMI(ctx, vmi, newVMI); err != nil {
11✔
1021
                        return &syncv1.VMIStatusResponse{
1✔
1022
                                Message: fmt.Sprintf("unable to synchronize VMI for migrationID %s", request.MigrationID),
1✔
1023
                        }, err
1✔
1024
                }
1✔
1025
                log.Log.Object(newVMI).With("MigrationID", request.MigrationID).V(5).Info("successfully patched VMI with source state")
9✔
1026
        }
1027
        log.Log.Object(newVMI).V(5).Info("returning success to grpc caller, source")
9✔
1028
        return &syncv1.VMIStatusResponse{
9✔
1029
                Message: successMessage,
9✔
1030
        }, nil
9✔
1031
}
1032

1033
func getMergedTargetMigratedVolumes(vmiMigratedVolumes []virtv1.StorageMigratedVolumeInfo, remoteMigratedVolumes []virtv1.StorageMigratedVolumeInfo) []virtv1.StorageMigratedVolumeInfo {
16✔
1034
        remoteVolumeMap := make(map[string]virtv1.StorageMigratedVolumeInfo)
16✔
1035
        for _, volume := range remoteMigratedVolumes {
25✔
1036
                remoteVolumeMap[volume.VolumeName] = volume
9✔
1037
        }
9✔
1038
        mergedVolumes := make([]virtv1.StorageMigratedVolumeInfo, 0)
16✔
1039
        for _, volume := range vmiMigratedVolumes {
26✔
1040
                if remoteVolume, ok := remoteVolumeMap[volume.VolumeName]; ok {
18✔
1041
                        mergedVolume := virtv1.StorageMigratedVolumeInfo{
8✔
1042
                                VolumeName: volume.VolumeName,
8✔
1043
                        }
8✔
1044
                        if remoteVolume.DestinationPVCInfo != nil {
15✔
1045
                                mergedVolume.DestinationPVCInfo = remoteVolume.DestinationPVCInfo.DeepCopy()
7✔
1046
                        }
7✔
1047
                        if volume.SourcePVCInfo != nil {
15✔
1048
                                mergedVolume.SourcePVCInfo = volume.SourcePVCInfo.DeepCopy()
7✔
1049
                        }
7✔
1050
                        mergedVolumes = append(mergedVolumes, mergedVolume)
8✔
1051
                } else {
2✔
1052
                        mergedVolumes = append(mergedVolumes, volume)
2✔
1053
                }
2✔
1054
        }
1055
        return mergedVolumes
16✔
1056
}
1057

1058
func getMergedSourceMigratedVolumes(vmiMigratedVolumes []virtv1.StorageMigratedVolumeInfo, remoteMigratedVolumes []virtv1.StorageMigratedVolumeInfo) []virtv1.StorageMigratedVolumeInfo {
11✔
1059
        remoteVolumeMap := make(map[string]virtv1.StorageMigratedVolumeInfo)
11✔
1060
        for _, volume := range remoteMigratedVolumes {
24✔
1061
                remoteVolumeMap[volume.VolumeName] = volume
13✔
1062
        }
13✔
1063
        mergedVolumes := make([]virtv1.StorageMigratedVolumeInfo, 0)
11✔
1064
        for _, vmiVolume := range vmiMigratedVolumes {
22✔
1065
                if remoteVolume, ok := remoteVolumeMap[vmiVolume.VolumeName]; ok {
21✔
1066
                        log.Log.V(5).Infof("Merging volume %s", vmiVolume.VolumeName)
10✔
1067
                        // Found a match merge the current target volume with the incoming source volume
10✔
1068
                        mergedVolume := virtv1.StorageMigratedVolumeInfo{
10✔
1069
                                VolumeName: vmiVolume.VolumeName,
10✔
1070
                        }
10✔
1071
                        if vmiVolume.SourcePVCInfo != nil {
17✔
1072
                                mergedVolume.SourcePVCInfo = vmiVolume.SourcePVCInfo.DeepCopy()
7✔
1073
                        } else {
10✔
1074
                                mergedVolume.SourcePVCInfo = remoteVolume.SourcePVCInfo.DeepCopy()
3✔
1075
                        }
3✔
1076
                        if vmiVolume.DestinationPVCInfo != nil {
19✔
1077
                                mergedVolume.DestinationPVCInfo = vmiVolume.DestinationPVCInfo.DeepCopy()
9✔
1078
                        }
9✔
1079
                        mergedVolumes = append(mergedVolumes, mergedVolume)
10✔
1080
                        delete(remoteVolumeMap, vmiVolume.VolumeName)
10✔
1081
                }
1082
        }
1083
        for _, volume := range remoteVolumeMap {
14✔
1084
                mergedVolumes = append(mergedVolumes, volume)
3✔
1085
        }
3✔
1086
        return mergedVolumes
11✔
1087
}
1088

1089
func (s *SynchronizationController) SyncTargetMigrationStatus(ctx context.Context, request *syncv1.VMIStatusRequest) (*syncv1.VMIStatusResponse, error) {
15✔
1090
        if request.VmiStatus == nil || len(request.VmiStatus.VmiStatusJson) == 0 {
16✔
1091
                return &syncv1.VMIStatusResponse{
1✔
1092
                        Message: noTargetStatusErrorMsg,
1✔
1093
                }, fmt.Errorf(noTargetStatusErrorMsg)
1✔
1094
        }
1✔
1095

1096
        migration, err := s.findSourceMigrationFromMigrationID(request.MigrationID)
14✔
1097
        if migration == nil {
15✔
1098
                return &syncv1.VMIStatusResponse{
1✔
1099
                        Message: fmt.Sprintf(targetUnableToLocateVMIMigrationIDErrorMsg, request.MigrationID),
1✔
1100
                }, fmt.Errorf(targetUnableToLocateVMIMigrationIDErrorMsg, request.MigrationID)
1✔
1101
        }
1✔
1102

1103
        key := controller.NamespacedKey(migration.Namespace, migration.Spec.VMIName)
13✔
1104
        obj, exists, err := s.vmiInformer.GetStore().GetByKey(key)
13✔
1105
        if err != nil || !exists {
15✔
1106
                if err == nil {
4✔
1107
                        err = fmt.Errorf(targetUnableToLocateVMIMigrationIDErrorMsgVMI, request.MigrationID, key)
2✔
1108
                }
2✔
1109
                return &syncv1.VMIStatusResponse{
2✔
1110
                        Message: fmt.Sprintf(targetUnableToLocateVMIMigrationIDErrorMsgVMI, request.MigrationID, key),
2✔
1111
                }, err
2✔
1112
        }
1113
        vmi := obj.(*virtv1.VirtualMachineInstance)
11✔
1114
        remoteStatus := &virtv1.VirtualMachineInstanceStatus{}
11✔
1115
        if err := json.Unmarshal(request.VmiStatus.VmiStatusJson, remoteStatus); err != nil {
12✔
1116
                return &syncv1.VMIStatusResponse{
1✔
1117
                        Message: fmt.Sprintf("unable to unmarshal vmistatus for migrationID %s", request.MigrationID),
1✔
1118
                }, err
1✔
1119
        }
1✔
1120
        if remoteStatus.MigrationState == nil {
11✔
1121
                return &syncv1.VMIStatusResponse{
1✔
1122
                        Message: noTargetStatusErrorMsg,
1✔
1123
                }, fmt.Errorf(noTargetStatusErrorMsg)
1✔
1124
        }
1✔
1125
        newVMI := vmi.DeepCopy()
9✔
1126
        if newVMI.Status.MigrationState == nil {
13✔
1127
                newVMI.Status.MigrationState = &virtv1.VirtualMachineInstanceMigrationState{}
4✔
1128
        }
4✔
1129

1130
        // Only update TargetState if this migration is still active and matches the VMI's current migration.
1131
        // This prevents stale updates from a completed decentralized migration from interfering with a new compute migration.
1132
        if migration.IsFinal() {
10✔
1133
                log.Log.Object(migration).Infof("Migration is final, ignoring target state update for VMI %s/%s", vmi.Namespace, vmi.Name)
1✔
1134
                return &syncv1.VMIStatusResponse{
1✔
1135
                        Message: successMessage,
1✔
1136
                }, nil
1✔
1137
        }
1✔
1138

1139
        // Check if the VMI's current migration matches this migration
1140
        if newVMI.Status.MigrationState.MigrationUID != "" && newVMI.Status.MigrationState.MigrationUID != migration.UID {
9✔
1141
                log.Log.Object(migration).Warningf("VMI %s/%s has different migration UID %s, ignoring target state update for migration %s",
1✔
1142
                        vmi.Namespace, vmi.Name, newVMI.Status.MigrationState.MigrationUID, migration.UID)
1✔
1143
                return &syncv1.VMIStatusResponse{
1✔
1144
                        Message: successMessage,
1✔
1145
                }, nil
1✔
1146
        }
1✔
1147

1148
        log.Log.Object(newVMI).V(5).Infof("vmi migration target state: %#v", newVMI.Status.MigrationState.TargetState)
7✔
1149
        log.Log.Object(newVMI).V(5).Infof("remote migration target state: %#v", remoteStatus.MigrationState.TargetState)
7✔
1150
        newVMI.Status.MigrationState.TargetState = remoteStatus.MigrationState.TargetState.DeepCopy()
7✔
1151
        newVMI.Status.MigratedVolumes = getMergedTargetMigratedVolumes(newVMI.Status.MigratedVolumes, remoteStatus.MigratedVolumes)
7✔
1152
        copyLegacyTargetFields(newVMI, remoteStatus.MigrationState)
7✔
1153
        if !apiequality.Semantic.DeepEqual(vmi.Status.MigrationState, newVMI.Status.MigrationState) {
14✔
1154
                if err := s.patchVMI(ctx, vmi, newVMI); err != nil {
8✔
1155
                        return &syncv1.VMIStatusResponse{
1✔
1156
                                Message: fmt.Sprintf("unable to synchronize VMI for migrationID %s", request.MigrationID),
1✔
1157
                        }, err
1✔
1158
                }
1✔
1159
                log.Log.Object(newVMI).With("MigrationID", request.MigrationID).V(5).Info("successfully patched VMI with target state")
6✔
1160
        }
1161
        log.Log.Object(newVMI).V(5).Info("returning success to grpc caller, target")
6✔
1162
        return &syncv1.VMIStatusResponse{
6✔
1163
                Message: successMessage,
6✔
1164
        }, nil
6✔
1165
}
1166

1167
func (s *SynchronizationController) patchVMIConditions(ctx context.Context, origVMI, newVMI *virtv1.VirtualMachineInstance) error {
20✔
1168
        patchSet := patch.New()
20✔
1169
        if !apiequality.Semantic.DeepEqual(origVMI.Status.Conditions, newVMI.Status.Conditions) {
32✔
1170
                patchSet.AddOption(
12✔
1171
                        patch.WithTest("/status/conditions", origVMI.Status.Conditions),
12✔
1172
                        patch.WithReplace("/status/conditions", newVMI.Status.Conditions),
12✔
1173
                )
12✔
1174
        }
12✔
1175
        if !patchSet.IsEmpty() {
32✔
1176
                patchBytes, err := patchSet.GeneratePayload()
12✔
1177
                if err != nil {
12✔
1178
                        return err
×
1179
                }
×
1180
                log.Log.Object(origVMI).V(5).Infof("patch VMI conditions with %s", string(patchBytes))
12✔
1181
                if _, err := s.client.VirtualMachineInstance(origVMI.Namespace).Patch(ctx, origVMI.Name, types.JSONPatchType, patchBytes, metav1.PatchOptions{}); err != nil {
12✔
1182
                        return err
×
1183
                }
×
1184
        }
1185
        return nil
20✔
1186
}
1187

1188
func (s *SynchronizationController) patchVMI(ctx context.Context, origVMI, newVMI *virtv1.VirtualMachineInstance) error {
19✔
1189
        if origVMI.Status.MigrationState != nil &&
19✔
1190
                origVMI.Status.MigrationState.Completed {
19✔
1191
                log.Log.Object(origVMI).V(3).Infof("VMI is completed, skipping patch")
×
1192
                return nil
×
1193
        }
×
1194
        patchSet := patch.New()
19✔
1195

19✔
1196
        if !apiequality.Semantic.DeepEqual(origVMI.Labels, newVMI.Labels) {
19✔
1197
                if len(origVMI.Labels) == 0 {
×
1198
                        patchSet.AddOption(
×
1199
                                patch.WithAdd("/metadata/labels", newVMI.Labels))
×
1200
                } else {
×
1201
                        patchSet.AddOption(
×
1202
                                patch.WithTest("/metadata/labels", origVMI.Labels),
×
1203
                                patch.WithReplace("/metadata/labels", newVMI.Labels),
×
1204
                        )
×
1205
                }
×
1206
        }
1207

1208
        if !apiequality.Semantic.DeepEqual(origVMI.Status.MigrationMethod, newVMI.Status.MigrationMethod) {
19✔
1209
                if origVMI.Status.MigrationMethod == "" {
×
1210
                        patchSet.AddOption(
×
1211
                                patch.WithAdd("/status/migrationMethod", newVMI.Status.MigrationMethod))
×
1212
                } else {
×
1213
                        patchSet.AddOption(
×
1214
                                patch.WithTest("/status/migrationMethod", origVMI.Status.MigrationMethod),
×
1215
                                patch.WithReplace("/status/migrationMethod", newVMI.Status.MigrationMethod),
×
1216
                        )
×
1217
                }
×
1218
        }
1219

1220
        if !apiequality.Semantic.DeepEqual(origVMI.Status.MigrationTransport, newVMI.Status.MigrationTransport) {
20✔
1221
                if origVMI.Status.MigrationTransport == "" {
2✔
1222
                        patchSet.AddOption(
1✔
1223
                                patch.WithAdd("/status/migrationTransport", newVMI.Status.MigrationTransport))
1✔
1224
                } else {
1✔
1225
                        patchSet.AddOption(
×
1226
                                patch.WithTest("/status/migrationTransport", origVMI.Status.MigrationTransport),
×
1227
                                patch.WithReplace("/status/migrationTransport", newVMI.Status.MigrationTransport),
×
1228
                        )
×
1229
                }
×
1230
        }
1231

1232
        if !apiequality.Semantic.DeepEqual(origVMI.Status.MigratedVolumes, newVMI.Status.MigratedVolumes) {
19✔
1233
                if origVMI.Status.MigratedVolumes == nil {
×
1234
                        patchSet.AddOption(
×
1235
                                patch.WithAdd("/status/migratedVolumes", newVMI.Status.MigratedVolumes))
×
1236
                } else {
×
1237
                        patchSet.AddOption(
×
1238
                                patch.WithTest("/status/migratedVolumes", origVMI.Status.MigratedVolumes),
×
1239
                                patch.WithReplace("/status/migratedVolumes", newVMI.Status.MigratedVolumes),
×
1240
                        )
×
1241
                }
×
1242
        }
1243

1244
        if !apiequality.Semantic.DeepEqual(origVMI.Status.MigrationState, newVMI.Status.MigrationState) {
38✔
1245
                if origVMI.Status.MigrationState == nil {
26✔
1246
                        patchSet.AddOption(
7✔
1247
                                patch.WithAdd("/status/migrationState", newVMI.Status.MigrationState))
7✔
1248
                } else {
19✔
1249
                        addMigrationStateFieldPatches(patchSet, origVMI.Status.MigrationState, newVMI.Status.MigrationState)
12✔
1250
                }
12✔
1251
        }
1252

1253
        if !patchSet.IsEmpty() {
38✔
1254
                patchBytes, err := patchSet.GeneratePayload()
19✔
1255
                if err != nil {
19✔
1256
                        return err
×
1257
                }
×
1258
                log.Log.Object(origVMI).V(5).Infof("patch VMI with %s", string(patchBytes))
19✔
1259
                if _, err := s.client.VirtualMachineInstance(origVMI.Namespace).Patch(ctx, origVMI.Name, types.JSONPatchType, patchBytes, metav1.PatchOptions{}); err != nil {
21✔
1260
                        return err
2✔
1261
                }
2✔
1262
        }
1263
        return nil
17✔
1264
}
1265

1266
// addMigrationStateFieldPatches generates individual JSON Patch
1267
// operations for each changed field in MigrationState, rather than a
1268
// single test+replace of the entire object. This prevents spurious
1269
// patch failures when another controller concurrently modifies
1270
// unrelated MigrationState fields (e.g. the virt-handler setting
1271
// ports while the sync controller sets SourceState).
1272
func addMigrationStateFieldPatches(patchSet *patch.PatchSet, origMS, newMS *virtv1.VirtualMachineInstanceMigrationState) {
17✔
1273
        if newMS == nil {
18✔
1274
                patchSet.AddOption(
1✔
1275
                        patch.WithTest("/status/migrationState", origMS),
1✔
1276
                        patch.WithRemove("/status/migrationState"),
1✔
1277
                )
1✔
1278
                return
1✔
1279
        }
1✔
1280

1281
        patchMigrationStateField(patchSet, "startTimestamp", origMS.StartTimestamp, newMS.StartTimestamp)
16✔
1282
        patchMigrationStateField(patchSet, "endTimestamp", origMS.EndTimestamp, newMS.EndTimestamp)
16✔
1283
        patchMigrationStateField(patchSet, "targetNodeDomainReadyTimestamp", origMS.TargetNodeDomainReadyTimestamp, newMS.TargetNodeDomainReadyTimestamp)
16✔
1284
        patchMigrationStateField(patchSet, "targetNodeDomainDetected", origMS.TargetNodeDomainDetected, newMS.TargetNodeDomainDetected)
16✔
1285
        patchMigrationStateField(patchSet, "targetNodeAddress", origMS.TargetNodeAddress, newMS.TargetNodeAddress)
16✔
1286
        patchMigrationStateField(patchSet, "targetDirectMigrationNodePorts", origMS.TargetDirectMigrationNodePorts, newMS.TargetDirectMigrationNodePorts)
16✔
1287
        patchMigrationStateField(patchSet, "targetNode", origMS.TargetNode, newMS.TargetNode)
16✔
1288
        patchMigrationStateField(patchSet, "targetPod", origMS.TargetPod, newMS.TargetPod)
16✔
1289
        patchMigrationStateField(patchSet, "targetAttachmentPodUID", origMS.TargetAttachmentPodUID, newMS.TargetAttachmentPodUID)
16✔
1290
        patchMigrationStateField(patchSet, "sourceNode", origMS.SourceNode, newMS.SourceNode)
16✔
1291
        patchMigrationStateField(patchSet, "sourcePod", origMS.SourcePod, newMS.SourcePod)
16✔
1292
        patchMigrationStateField(patchSet, "completed", origMS.Completed, newMS.Completed)
16✔
1293
        patchMigrationStateField(patchSet, "failed", origMS.Failed, newMS.Failed)
16✔
1294
        patchMigrationStateField(patchSet, "abortRequested", origMS.AbortRequested, newMS.AbortRequested)
16✔
1295
        patchMigrationStateField(patchSet, "abortStatus", origMS.AbortStatus, newMS.AbortStatus)
16✔
1296
        patchMigrationStateField(patchSet, "failureReason", origMS.FailureReason, newMS.FailureReason)
16✔
1297
        patchMigrationStateField(patchSet, "migrationUid", origMS.MigrationUID, newMS.MigrationUID)
16✔
1298
        patchMigrationStateField(patchSet, "mode", origMS.Mode, newMS.Mode)
16✔
1299
        patchMigrationStateField(patchSet, "migrationPolicyName", origMS.MigrationPolicyName, newMS.MigrationPolicyName)
16✔
1300
        patchMigrationStateField(patchSet, "migrationConfiguration", origMS.VMIMConfigurationOptions, newMS.VMIMConfigurationOptions)
16✔
1301
        patchMigrationStateField(patchSet, "targetCPUSet", origMS.TargetCPUSet, newMS.TargetCPUSet)
16✔
1302
        patchMigrationStateField(patchSet, "targetNodeTopology", origMS.TargetNodeTopology, newMS.TargetNodeTopology)
16✔
1303
        patchMigrationStateField(patchSet, "sourcePersistentStatePVCName", origMS.SourcePersistentStatePVCName, newMS.SourcePersistentStatePVCName)
16✔
1304
        patchMigrationStateField(patchSet, "targetPersistentStatePVCName", origMS.TargetPersistentStatePVCName, newMS.TargetPersistentStatePVCName)
16✔
1305
        patchMigrationStateField(patchSet, "sourceState", origMS.SourceState, newMS.SourceState)
16✔
1306
        patchMigrationStateField(patchSet, "targetState", origMS.TargetState, newMS.TargetState)
16✔
1307
        patchMigrationStateField(patchSet, "migrationNetworkType", origMS.MigrationNetworkType, newMS.MigrationNetworkType)
16✔
1308
        patchMigrationStateField(patchSet, "targetMemoryOverhead", origMS.TargetMemoryOverhead, newMS.TargetMemoryOverhead)
16✔
1309
}
1310

1311
// patchMigrationStateField generates an add or test+replace JSON
1312
// Patch operation for a single MigrationState field (all of which are
1313
// omitempty). Unchanged fields are skipped. Zero originals use add
1314
// (field absent from JSON), non-zero use test+replace for conflict
1315
// detection.
1316
func patchMigrationStateField[T any](patchSet *patch.PatchSet, field string, orig, new T) {
448✔
1317
        if apiequality.Semantic.DeepEqual(orig, new) {
868✔
1318
                return
420✔
1319
        }
420✔
1320
        path := "/status/migrationState/" + field
28✔
1321
        var zero T
28✔
1322
        if apiequality.Semantic.DeepEqual(orig, zero) {
44✔
1323
                patchSet.AddOption(patch.WithAdd(path, new))
16✔
1324
        } else {
28✔
1325
                patchSet.AddOption(
12✔
1326
                        patch.WithTest(path, orig),
12✔
1327
                        patch.WithReplace(path, new),
12✔
1328
                )
12✔
1329
        }
12✔
1330
}
1331

1332
func indexByMigrationUID(obj interface{}) ([]string, error) {
81✔
1333
        migration, ok := obj.(*virtv1.VirtualMachineInstanceMigration)
81✔
1334
        if !ok {
82✔
1335
                return nil, nil
1✔
1336
        }
1✔
1337
        return []string{string(migration.UID)}, nil
80✔
1338
}
1339

1340
func indexByActiveVmiName(obj interface{}) ([]string, error) {
81✔
1341
        migration, ok := obj.(*virtv1.VirtualMachineInstanceMigration)
81✔
1342
        if !ok {
82✔
1343
                return nil, nil
1✔
1344
        }
1✔
1345
        return []string{migration.Spec.VMIName}, nil
80✔
1346
}
1347

1348
func indexByTargetMigrationID(obj interface{}) ([]string, error) {
81✔
1349
        migration, ok := obj.(*virtv1.VirtualMachineInstanceMigration)
81✔
1350
        if !ok {
82✔
1351
                return nil, nil
1✔
1352
        }
1✔
1353
        if migration.Spec.Receive != nil {
123✔
1354
                return []string{migration.Spec.Receive.MigrationID}, nil
43✔
1355
        }
43✔
1356
        return []string{}, nil
37✔
1357
}
1358

1359
func indexBySourceMigrationID(obj interface{}) ([]string, error) {
81✔
1360
        migration, ok := obj.(*virtv1.VirtualMachineInstanceMigration)
81✔
1361
        if !ok {
82✔
1362
                return nil, nil
1✔
1363
        }
1✔
1364
        if migration.Spec.SendTo != nil {
117✔
1365
                return []string{migration.Spec.SendTo.MigrationID}, nil
37✔
1366
        }
37✔
1367
        return []string{}, nil
43✔
1368
}
1369

1370
func copyLegacyTargetFields(vmi *virtv1.VirtualMachineInstance, migrationState *virtv1.VirtualMachineInstanceMigrationState) {
7✔
1371
        targetState := migrationState.TargetState
7✔
1372
        vmi.Status.MigrationState.TargetNode = targetState.Node
7✔
1373
        if targetState.AttachmentPodUID != nil {
7✔
1374
                vmi.Status.MigrationState.TargetAttachmentPodUID = *targetState.AttachmentPodUID
×
1375
        }
×
1376
        vmi.Status.MigrationState.TargetCPUSet = targetState.CPUSet
7✔
1377
        vmi.Status.MigrationState.TargetDirectMigrationNodePorts = targetState.DirectMigrationNodePorts
7✔
1378
        if targetState.NodeAddress != nil {
7✔
1379
                vmi.Status.MigrationState.TargetNodeAddress = *targetState.NodeAddress
×
1380
        }
×
1381
        vmi.Status.MigrationState.TargetNodeDomainDetected = targetState.DomainDetected
7✔
1382
        vmi.Status.MigrationState.TargetNodeDomainReadyTimestamp = targetState.DomainReadyTimestamp
7✔
1383
        if targetState.NodeTopology != nil {
7✔
1384
                vmi.Status.MigrationState.TargetNodeTopology = *targetState.NodeTopology
×
1385
        }
×
1386
        if targetState.PersistentStatePVCName != nil {
7✔
1387
                vmi.Status.MigrationState.TargetPersistentStatePVCName = *targetState.PersistentStatePVCName
×
1388
        }
×
1389
        vmi.Status.MigrationState.TargetPod = targetState.Pod
7✔
1390
        copyCommonLegacyFields(vmi.Status.MigrationState, migrationState)
7✔
1391
        vmi.Status.MigrationState.Completed = migrationState.Completed
7✔
1392
        vmi.Status.MigrationState.Failed = migrationState.Failed
7✔
1393
}
1394

1395
func copyLegacySourceFields(vmi *virtv1.VirtualMachineInstance, migrationState *virtv1.VirtualMachineInstanceMigrationState) {
12✔
1396
        vmi.Status.MigrationState.SourceNode = migrationState.SourceState.Node
12✔
1397
        if migrationState.SourceState.PersistentStatePVCName != nil {
12✔
1398
                vmi.Status.MigrationState.SourcePersistentStatePVCName = *migrationState.SourceState.PersistentStatePVCName
×
1399
        }
×
1400
        vmi.Status.MigrationState.SourcePod = migrationState.SourceState.Pod
12✔
1401
        copyCommonLegacyFields(vmi.Status.MigrationState, migrationState)
12✔
1402
        if migrationState.AbortRequested && migrationState.EndTimestamp != nil {
15✔
1403
                vmi.Status.MigrationState.Failed = migrationState.Failed
3✔
1404
                vmi.Status.MigrationState.Completed = migrationState.Completed
3✔
1405
                if migrationState.AbortStatus != "" {
6✔
1406
                        vmi.Status.MigrationState.AbortStatus = migrationState.AbortStatus
3✔
1407
                }
3✔
1408
                if migrationState.FailureReason != "" {
6✔
1409
                        vmi.Status.MigrationState.FailureReason = migrationState.FailureReason
3✔
1410
                }
3✔
1411
        }
1412
}
1413

1414
func copyCommonLegacyFields(targetMigrationState, sourceMigrationState *virtv1.VirtualMachineInstanceMigrationState) {
19✔
1415
        // Copy regular fields.
19✔
1416
        if sourceMigrationState.MigrationPolicyName != nil {
19✔
1417
                targetMigrationState.MigrationPolicyName = sourceMigrationState.MigrationPolicyName
×
1418
        }
×
1419
        if sourceMigrationState.VMIMConfigurationOptions != nil {
19✔
1420
                targetMigrationState.VMIMConfigurationOptions = sourceMigrationState.VMIMConfigurationOptions
×
1421
        }
×
1422
        if sourceMigrationState.StartTimestamp != nil {
19✔
1423
                targetMigrationState.StartTimestamp = sourceMigrationState.StartTimestamp
×
1424
        }
×
1425
        if sourceMigrationState.EndTimestamp != nil {
22✔
1426
                targetMigrationState.EndTimestamp = sourceMigrationState.EndTimestamp
3✔
1427
        }
3✔
1428
}
1429

1430
func (s *SynchronizationController) runConnectionCleanup() {
×
1431
        s.failedCloseConnections.Range(func(k, v interface{}) bool {
×
1432
                retryCount, ok := v.(int)
×
1433
                if !ok {
×
1434
                        log.Log.Warningf("invalid retry count type during connection cleanup: %v", v)
×
1435
                        s.failedCloseConnections.Delete(k)
×
1436
                        return true
×
1437
                }
×
1438
                if retryCount >= maxCloseRetries {
×
1439
                        log.Log.Warningf("connection for migrationID %s failed to close after %d retries, not attempting to close again", k, retryCount)
×
1440
                        s.failedCloseConnections.Delete(k)
×
1441
                }
×
1442
                outboundConnection, ok := k.(*SynchronizationConnection)
×
1443
                if !ok {
×
1444
                        log.Log.Warningf("invalid outbound connection type during connection cleanup: %v", k)
×
1445
                        s.failedCloseConnections.Delete(k)
×
1446
                        return true
×
1447
                }
×
1448
                if err := outboundConnection.Close(); err != nil {
×
1449
                        log.Log.Warningf("unable to close connection for migrationID, trying again: %s, %v", outboundConnection.migrationID, err)
×
1450
                        s.failedCloseConnections.Store(outboundConnection, retryCount+1)
×
1451
                } else {
×
1452
                        s.failedCloseConnections.Delete(k)
×
1453
                }
×
1454
                return true
×
1455
        })
1456
}
1457

1458
func (s *SynchronizationController) CancelMigration(ctx context.Context, request *syncv1.MigrationCancelRequest) (*syncv1.MigrationCancelResponse, error) {
1✔
1459
        migrationUID := request.MigrationUID
1✔
1460

1✔
1461
        migration, err := s.findMigrationFromMigrationIDByIndex("byUID", migrationUID)
1✔
1462
        if err != nil {
1✔
1463
                return &syncv1.MigrationCancelResponse{
×
1464
                        Message: fmt.Sprintf("unable to find migration to cancel for migrationUID %s", migrationUID),
×
1465
                }, err
×
1466
        }
×
1467
        if migration != nil {
1✔
1468
                log.Log.V(2).Object(migration).Infof("found migration to cancel for migrationUID %s", migrationUID)
×
1469
                if err := s.client.VirtualMachineInstanceMigration(migration.Namespace).Delete(ctx, migration.Name, metav1.DeleteOptions{}); err != nil {
×
1470
                        return &syncv1.MigrationCancelResponse{
×
1471
                                Message: fmt.Sprintf("unable to cancel migration for migrationUID %s", migrationUID),
×
1472
                        }, err
×
1473
                }
×
1474
                log.Log.V(2).Object(migration).Infof("successfully deleted migration %s/%s for migrationUID %s", migration.Namespace, migration.Name, migrationUID)
×
1475
        }
1476
        return &syncv1.MigrationCancelResponse{
1✔
1477
                Message: "migration canceled",
1✔
1478
        }, nil
1✔
1479
}
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc