1617
1618
// syncUnversionedRamp should not be called in async mode
1619
>
func (d *WorkflowRunner) syncUnversionedRamp(ctx workflow.Context, versionUpdateArgs *deploymentspb.SyncVersionStateUpdateArgs) error {
workflow.go
1620
>
activityCtx := workflow.WithActivityOptions(ctx, defaultActivityOptions)
1621
>
1622
>
// DescribeVersion activity to get all the task queues in the current version, or the ramping version if current is nil
1623
>
version := d.State.RoutingConfig.CurrentVersion //nolint:staticcheck // SA1019: worker versioning v0.31
1624
>
if version == worker_versioning.UnversionedVersionId {
1625
version = d.State.RoutingConfig.RampingVersion //nolint:staticcheck // SA1019: worker versioning v0.31
1626
}
1627
1628
>
if d.rampingVersionStringUnversioned(version) {
workflow.go
1629
return nil
1630
}
1631
1632
>
var res deploymentspb.DescribeVersionFromWorkerDeploymentActivityResult
workflow.go
1633
>
err := workflow.ExecuteActivity(
1634
>
activityCtx,
1635
>
d.a.DescribeVersionFromWorkerDeployment,
1636
>
&deploymentspb.DescribeVersionFromWorkerDeploymentActivityArgs{
1637
>
Version: version,
1638
>
}).Get(ctx, &res)
1639
>
if err != nil {
1640
return err
1641
}
1642
1643
// send in the task-queue families in batches of syncBatchSize
1644
>
batches := make([][]*deploymentspb.SyncDeploymentVersionUserDataRequest_SyncUserData, 0)
workflow.go
1645
>
syncReqs := make([]*deploymentspb.SyncDeploymentVersionUserDataRequest_SyncUserData, 0)
1646
>
1647
>
// Grouping by task-queue name
1648
>
taskQueuesByName := make(map[string][]enumspb.TaskQueueType)
1649
>
for _, tq := range res.GetTaskQueueInfos() {
1650
>
taskQueuesByName[tq.GetName()] = append(taskQueuesByName[tq.GetName()], tq.GetType())
1651
>
}
1652
1653
>
for _, tqName := range workflow.DeterministicKeys(taskQueuesByName) {
workflow.go
1654
>
tqTypes := taskQueuesByName[tqName]
1655
>
sync := &deploymentspb.SyncDeploymentVersionUserDataRequest_SyncUserData{
1656
>
Name: tqName,
1657
>
Types: tqTypes,
1658
>
Data: &deploymentspb.DeploymentVersionData{
1659
>
Version: nil,
1660
>
RoutingUpdateTime: versionUpdateArgs.RoutingUpdateTime,
1661
>
RampingSinceTime: versionUpdateArgs.RampingSinceTime,
1662
>
RampPercentage: versionUpdateArgs.RampPercentage,
1663
>
},
1664
>
}
1665
>
syncReqs = append(syncReqs, sync)
1666
>
1667
>
if len(syncReqs) == int(d.State.SyncBatchSize) {
1668
batches = append(batches, syncReqs)
1669
syncReqs = make([]*deploymentspb.SyncDeploymentVersionUserDataRequest_SyncUserData, 0) // reset the syncReq.Sync slice for the next batch
1670
}
1671
}
1673
>
batches = append(batches, syncReqs)
1674
>
}
1675
1676
// calling SyncDeploymentVersionUserData for each batch
1678
>
var syncRes deploymentspb.SyncDeploymentVersionUserDataResponse
1679
>
1680
>
err = workflow.ExecuteActivity(activityCtx, d.a.SyncDeploymentVersionUserDataFromWorkerDeployment, &deploymentspb.SyncDeploymentVersionUserDataRequest{
1681
>
DeploymentName: d.DeploymentName,
1682
>
Version: nil,
1683
>
ForgetVersion: false,
1684
>
Sync: batch,
1685
>
}).Get(ctx, &syncRes)
1686
>
if err != nil {
1687
// TODO (Shivam): Compensation functions required to roll back the local state + activity changes.
1688
return err
1689
}
1690
1691
>
if len(syncRes.TaskQueueMaxVersions) > 0 {
workflow.go
1692
>
// wait for propagation
1693
>
err = workflow.ExecuteActivity(
1694
>
activityCtx,
1695
>
d.a.CheckUnversionedRampUserDataPropagation,
1696
>
&deploymentspb.CheckWorkerDeploymentUserDataPropagationRequest{
1697
>
TaskQueueMaxVersions: syncRes.TaskQueueMaxVersions,
1698
>
}).Get(ctx, nil)
1699
>
if err != nil {
1700
// TODO (Shivam): Compensation functions required to roll back the local state + activity changes.
1701
return err