118
ctx context.Context,
119
metadata *persistence.InternalQueueMetadata,
121
>
err := q.txExecute(ctx, "UpdateAckLevel", func(tx sqlplugin.Tx) error {
122
>
result, err := tx.UpdateQueueMetadata(ctx, &sqlplugin.QueueMetadataRow{
123
>
QueueType: q.queueType,
124
>
Data: metadata.Blob.Data,
125
>
DataEncoding: metadata.Blob.EncodingType.String(),
126
>
Version: metadata.Version,
127
>
})
128
>
if err != nil {
129
return serviceerror.NewUnavailablef("UpdateAckLevel operation failed. Error %v", err)
130
}
131
>
rowsAffected, err := result.RowsAffected()
queue.go
132
>
if err != nil {
133
return fmt.Errorf("rowsAffected returned error for queue metadata %v: %v", q.queueType, err)
134
}
136
return &persistence.ConditionFailedError{Msg: "UpdateAckLevel operation encountered concurrent write."}
137
}
139
})
140
142
return serviceerror.NewUnavailable(err.Error())
143
}
145
}
146
147
func (q *sqlQueue) GetAckLevels(
148
ctx context.Context,
149
>
) (*persistence.InternalQueueMetadata, error) {
queue.go
150
>
row, err := q.DB.SelectFromQueueMetadata(ctx, sqlplugin.QueueMetadataFilter{
151
>
QueueType: q.queueType,
152
>
})
153
>
if err != nil {
154
return nil, serviceerror.NewUnavailablef("GetAckLevels operation failed. Error %v", err)
155
}
156
157
>
return &persistence.InternalQueueMetadata{
queue.go
158
>
Blob: persistence.NewDataBlob(row.Data, row.DataEncoding),
159
>
Version: row.Version,
160
>
}, nil
161
}
162