137
}
138
140
>
numBufferedEvents := 20
141
>
142
>
shardID := rand.Int31()
143
>
namespaceID := primitives.NewUUID()
144
>
workflowID := shuffle.String(testHistoryExecutionWorkflowID)
145
>
runID := primitives.NewUUID()
146
>
147
>
var buffers []sqlplugin.BufferedEventsRow
148
>
for range numBufferedEvents {
149
>
buffer := s.newRandomExecutionBufferRow(shardID, namespaceID, workflowID, runID)
150
>
buffers = append(buffers, buffer)
151
>
}
152
>
result, err := s.store.InsertIntoBufferedEvents(newExecutionContext(), buffers)
153
>
s.NoError(err)
154
>
rowsAffected, err := result.RowsAffected()
155
>
s.NoError(err)
156
>
s.Equal(numBufferedEvents, int(rowsAffected))
157
>
158
>
filter := sqlplugin.BufferedEventsFilter{
159
>
ShardID: shardID,
160
>
NamespaceID: namespaceID,
161
>
WorkflowID: workflowID,
162
>
RunID: runID,
163
>
}
164
>
result, err = s.store.DeleteFromBufferedEvents(newExecutionContext(), filter)
165
>
s.NoError(err)
166
>
rowsAffected, err = result.RowsAffected()
167
>
s.NoError(err)
168
>
s.Equal(numBufferedEvents, int(rowsAffected))
169
>
170
>
rows, err := s.store.SelectFromBufferedEvents(newExecutionContext(), filter)
171
>
s.NoError(err)
172
>
s.Equal([]sqlplugin.BufferedEventsRow(nil), rows)
173
}
174