588
ctx context.Context,
589
request *p.PutReplicationTaskToDLQRequest,
591
>
replicationTask := request.TaskInfo
592
>
blob, err := m.serializer.ReplicationTaskInfoToBlob(replicationTask)
593
>
594
>
if err != nil {
595
return err
596
}
597
598
>
_, err = m.DB.InsertIntoReplicationDLQTasks(ctx, []sqlplugin.ReplicationDLQTasksRow{{
execution_tasks.go
599
>
SourceClusterName: request.SourceClusterName,
600
>
ShardID: request.ShardID,
601
>
TaskID: replicationTask.GetTaskId(),
602
>
Data: blob.Data,
603
>
DataEncoding: blob.EncodingType.String(),
604
>
}})
605
>
606
>
// Tasks are immutable. So it's fine if we already persisted it before.
607
>
// This can happen when tasks are retried (ack and cleanup can have lag on source side).
608
>
if err != nil && !m.DB.IsDupEntryError(err) {
609
return serviceerror.NewUnavailablef("Failed to create replication tasks. Error: %v", err)
610
}
611
613
}
614