diff --git a/pkg/services/cloudmigration/cloudmigrationimpl/cloudmigration.go b/pkg/services/cloudmigration/cloudmigrationimpl/cloudmigration.go index 6e121032a6d..933b49358a4 100644 --- a/pkg/services/cloudmigration/cloudmigrationimpl/cloudmigration.go +++ b/pkg/services/cloudmigration/cloudmigrationimpl/cloudmigration.go @@ -524,9 +524,11 @@ func (s *Service) GetSnapshot(ctx context.Context, query cloudmigration.GetSnaps return nil, fmt.Errorf("fetching session for uid %s: %w", sessionUid, err) } + // Ask GMS for snapshot status while the source of truth is in the cloud if snapshot.ShouldQueryGMS() { - // ask GMS for status if it's in the cloud - snapshotMeta, err := s.gmsClient.GetSnapshotStatus(ctx, *session, *snapshot) + // Calculate offset based on how many results we currently have responses for + pending := snapshot.StatsRollup.CountsByStatus[cloudmigration.ItemStatusPending] + snapshotMeta, err := s.gmsClient.GetSnapshotStatus(ctx, *session, *snapshot, snapshot.StatsRollup.Total-pending) if err != nil { return snapshot, fmt.Errorf("error fetching snapshot status from GMS: sessionUid: %s, snapshotUid: %s", sessionUid, snapshotUid) } diff --git a/pkg/services/cloudmigration/cloudmigrationimpl/cloudmigration_test.go b/pkg/services/cloudmigration/cloudmigrationimpl/cloudmigration_test.go index 74a1b56b7c7..3a0458ff4b6 100644 --- a/pkg/services/cloudmigration/cloudmigrationimpl/cloudmigration_test.go +++ b/pkg/services/cloudmigration/cloudmigrationimpl/cloudmigration_test.go @@ -441,7 +441,7 @@ func (m *gmsClientMock) StartSnapshot(_ context.Context, _ cloudmigration.CloudM return nil, nil } -func (m *gmsClientMock) GetSnapshotStatus(_ context.Context, _ cloudmigration.CloudMigrationSession, _ cloudmigration.CloudMigrationSnapshot) (*cloudmigration.GetSnapshotStatusResponse, error) { +func (m *gmsClientMock) GetSnapshotStatus(_ context.Context, _ cloudmigration.CloudMigrationSession, _ cloudmigration.CloudMigrationSnapshot, _ int) (*cloudmigration.GetSnapshotStatusResponse, error) { m.getStatusCalled++ return m.getSnapshotResponse, nil } diff --git a/pkg/services/cloudmigration/cloudmigrationimpl/xorm_store.go b/pkg/services/cloudmigration/cloudmigrationimpl/xorm_store.go index 37e24d1272b..4e45d7ce10f 100644 --- a/pkg/services/cloudmigration/cloudmigrationimpl/xorm_store.go +++ b/pkg/services/cloudmigration/cloudmigrationimpl/xorm_store.go @@ -339,7 +339,13 @@ func (ss *sqlStore) GetSnapshotResourceStats(ctx context.Context, snapshotUid st Count int `json:"count"` Status string `json:"status"` }, 0) + total := 0 err := ss.db.WithDbSession(ctx, func(sess *sqlstore.DBSession) error { + if t, err := sess.Count(cloudmigration.CloudMigrationResource{SnapshotUID: snapshotUid}); err != nil { + return err + } else { + total = int(t) + } sess.Select("count(uid) as 'count', resource_type as 'type'"). Table(tableName). GroupBy("type"). @@ -360,6 +366,7 @@ func (ss *sqlStore) GetSnapshotResourceStats(ctx context.Context, snapshotUid st stats := &cloudmigration.SnapshotResourceStats{ CountsByType: make(map[cloudmigration.MigrateDataType]int, len(typeCounts)), CountsByStatus: make(map[cloudmigration.ItemStatus]int, len(statusCounts)), + Total: total, } for _, c := range typeCounts { stats.CountsByType[cloudmigration.MigrateDataType(c.Type)] = c.Count diff --git a/pkg/services/cloudmigration/cloudmigrationimpl/xorm_store_test.go b/pkg/services/cloudmigration/cloudmigrationimpl/xorm_store_test.go index 6f257b2c4eb..a019c3b24e8 100644 --- a/pkg/services/cloudmigration/cloudmigrationimpl/xorm_store_test.go +++ b/pkg/services/cloudmigration/cloudmigrationimpl/xorm_store_test.go @@ -257,6 +257,7 @@ func Test_SnapshotResources(t *testing.T) { cloudmigration.ItemStatusOK: 3, cloudmigration.ItemStatusPending: 1, }, stats.CountsByStatus) + assert.Equal(t, 4, stats.Total) // delete snapshot resources err = s.DeleteSnapshotResources(ctx, "poiuy") diff --git a/pkg/services/cloudmigration/gmsclient/client.go b/pkg/services/cloudmigration/gmsclient/client.go index 2890d4fbf1c..e6bc631c6d2 100644 --- a/pkg/services/cloudmigration/gmsclient/client.go +++ b/pkg/services/cloudmigration/gmsclient/client.go @@ -10,7 +10,7 @@ type Client interface { ValidateKey(context.Context, cloudmigration.CloudMigrationSession) error MigrateData(context.Context, cloudmigration.CloudMigrationSession, cloudmigration.MigrateDataRequest) (*cloudmigration.MigrateDataResponse, error) StartSnapshot(context.Context, cloudmigration.CloudMigrationSession) (*cloudmigration.StartSnapshotResponse, error) - GetSnapshotStatus(context.Context, cloudmigration.CloudMigrationSession, cloudmigration.CloudMigrationSnapshot) (*cloudmigration.GetSnapshotStatusResponse, error) + GetSnapshotStatus(context.Context, cloudmigration.CloudMigrationSession, cloudmigration.CloudMigrationSnapshot, int) (*cloudmigration.GetSnapshotStatusResponse, error) } const logPrefix = "cloudmigration.gmsclient" diff --git a/pkg/services/cloudmigration/gmsclient/gms_client.go b/pkg/services/cloudmigration/gmsclient/gms_client.go index eca425ece8f..b73254eed29 100644 --- a/pkg/services/cloudmigration/gmsclient/gms_client.go +++ b/pkg/services/cloudmigration/gmsclient/gms_client.go @@ -162,12 +162,12 @@ func (c *gmsClientImpl) StartSnapshot(ctx context.Context, session cloudmigratio return &result, nil } -func (c *gmsClientImpl) GetSnapshotStatus(ctx context.Context, session cloudmigration.CloudMigrationSession, snapshot cloudmigration.CloudMigrationSnapshot) (*cloudmigration.GetSnapshotStatusResponse, error) { +func (c *gmsClientImpl) GetSnapshotStatus(ctx context.Context, session cloudmigration.CloudMigrationSession, snapshot cloudmigration.CloudMigrationSnapshot, offset int) (*cloudmigration.GetSnapshotStatusResponse, error) { c.getStatusMux.Lock() defer c.getStatusMux.Unlock() logger := c.log.FromContext(ctx) - path := fmt.Sprintf("%s/api/v1/status/%s/status", c.buildBasePath(session.ClusterSlug), snapshot.GMSSnapshotUID) + path := fmt.Sprintf("%s/api/v1/status/%s/status?offset=%d", c.buildBasePath(session.ClusterSlug), snapshot.GMSSnapshotUID, offset) // Send the request to gms with the associated auth token req, err := http.NewRequest(http.MethodGet, path, nil) diff --git a/pkg/services/cloudmigration/gmsclient/inmemory_client.go b/pkg/services/cloudmigration/gmsclient/inmemory_client.go index d2635ccfebe..e84780a7f2f 100644 --- a/pkg/services/cloudmigration/gmsclient/inmemory_client.go +++ b/pkg/services/cloudmigration/gmsclient/inmemory_client.go @@ -67,7 +67,7 @@ func (c *memoryClientImpl) StartSnapshot(context.Context, cloudmigration.CloudMi return c.snapshot, nil } -func (c *memoryClientImpl) GetSnapshotStatus(ctx context.Context, session cloudmigration.CloudMigrationSession, snapshot cloudmigration.CloudMigrationSnapshot) (*cloudmigration.GetSnapshotStatusResponse, error) { +func (c *memoryClientImpl) GetSnapshotStatus(ctx context.Context, session cloudmigration.CloudMigrationSession, snapshot cloudmigration.CloudMigrationSnapshot, offset int) (*cloudmigration.GetSnapshotStatusResponse, error) { gmsResp := &cloudmigration.GetSnapshotStatusResponse{ State: cloudmigration.SnapshotStateFinished, Results: []cloudmigration.CloudMigrationResource{ diff --git a/pkg/services/cloudmigration/model.go b/pkg/services/cloudmigration/model.go index 6018f2b167c..d26b4877b85 100644 --- a/pkg/services/cloudmigration/model.go +++ b/pkg/services/cloudmigration/model.go @@ -91,12 +91,12 @@ const ( ItemStatusOK ItemStatus = "OK" ItemStatusError ItemStatus = "ERROR" ItemStatusPending ItemStatus = "PENDING" - ItemStatusUnknown ItemStatus = "UNKNOWN" ) type SnapshotResourceStats struct { CountsByType map[MigrateDataType]int CountsByStatus map[ItemStatus]int + Total int } // Deprecated, use GetSnapshotResult for the async workflow