115
}
116
118
>
tq1, tq2, tq3 := "tq1", "tq2", "tq3"
119
>
120
>
// set up three task queues
121
>
data := s.makeData(hlc.Zero(12345), 0)
122
>
var applied1, applied2, applied3 bool
123
>
var conflict1, conflict2, conflict3 bool
124
>
for range 3 {
125
>
err := s.taskManager.UpdateTaskQueueUserData(s.ctx, &p.UpdateTaskQueueUserDataRequest{
126
>
NamespaceID: s.namespaceID,
127
>
Updates: map[string]*p.SingleTaskQueueUserDataUpdate{
128
>
tq1: &p.SingleTaskQueueUserDataUpdate{UserData: data, Applied: &applied1, Conflicting: &conflict1},
129
>
tq2: &p.SingleTaskQueueUserDataUpdate{UserData: data, Applied: &applied2, Conflicting: &conflict2},
130
>
tq3: &p.SingleTaskQueueUserDataUpdate{UserData: data, Applied: &applied3, Conflicting: &conflict3},
131
>
},
132
>
})
133
>
s.NoError(err)
134
>
data.Version++
135
>
}
136
>
s.True(applied1)
137
>
s.True(applied2)
138
>
s.True(applied3)
139
>
s.False(conflict1)
140
>
s.False(conflict2)
141
>
s.False(conflict3)
142
>
143
>
// get all and verify
144
>
for _, tq := range []string{tq1, tq2, tq3} {
145
>
res, err := s.taskManager.GetTaskQueueUserData(s.ctx, &p.GetTaskQueueUserDataRequest{
146
>
NamespaceID: s.namespaceID,
147
>
TaskQueue: tq,
148
>
})
149
>
s.NoError(err)
150
>
s.Equal(int64(3), res.UserData.Version)
151
>
s.True(hlc.Equal(data.Data.Clock, res.UserData.Data.Clock))
152
>
}
153
154
// do update where one conflicts
156
>
err := s.taskManager.UpdateTaskQueueUserData(s.ctx, &p.UpdateTaskQueueUserDataRequest{
157
>
NamespaceID: s.namespaceID,
158
>
Updates: map[string]*p.SingleTaskQueueUserDataUpdate{
159
>
tq1: &p.SingleTaskQueueUserDataUpdate{UserData: data, Applied: &applied1, Conflicting: &conflict1},
160
>
tq2: &p.SingleTaskQueueUserDataUpdate{UserData: d4, Applied: &applied2, Conflicting: &conflict2},
161
>
tq3: &p.SingleTaskQueueUserDataUpdate{UserData: data, Applied: &applied3, Conflicting: &conflict3},
162
>
},
163
>
})
164
>
s.Error(err)
165
>
s.True(p.IsConflictErr(err))
166
>
s.False(applied1)
167
>
s.False(applied2)
168
>
s.False(applied3)
169
>
s.False(conflict1)
170
>
s.True(conflict2)
171
>
s.False(conflict3)
172
>
173
>
// verify that none were updated
174
>
for _, tq := range []string{tq1, tq2, tq3} {
175
>
res, err := s.taskManager.GetTaskQueueUserData(s.ctx, &p.GetTaskQueueUserDataRequest{
176
>
NamespaceID: s.namespaceID,
177
>
TaskQueue: tq,
178
>
})
179
>
s.NoError(err)
180
>
s.Equal(int64(3), res.UserData.Version)
181
>
s.True(hlc.Equal(data.Data.Clock, res.UserData.Data.Clock))
182
>
}
183
}
184
185
>
func (s *TaskQueueUserDataSuite) makeData(prev *hlc.Clock, ver int64) *persistencespb.VersionedTaskQueueUserData {
task_queue_user_data.go
186
>
return &persistencespb.VersionedTaskQueueUserData{
187
>
Data: &persistencespb.TaskQueueUserData{
188
>
Clock: hlc.Next(prev, clock.NewRealTimeSource()),
189
>
},
190
>
Version: ver,
191
>
}
192
>
}