MCPcopy Create free account
hub / github.com/temporalio/temporal / EnqueueMessage

Method EnqueueMessage

common/persistence/sql/queue_v2.go:45–105  ·  view source on GitHub ↗
(
	ctx context.Context,
	request *persistence.InternalEnqueueMessageRequest,
)

Source from the content-addressed store, hash-verified

43}
44
45func (q *queueV2) EnqueueMessage(
46 ctx context.Context,
47 request *persistence.InternalEnqueueMessageRequest,
48) (*persistence.InternalEnqueueMessageResponse, error) {
49
50 _, err := q.getQueueMetadata(ctx, q.DB, request.QueueType, request.QueueName)
51 if err != nil {
52 return nil, err
53 }
54 tx, err := q.DB.BeginTx(ctx)
55 if err != nil {
56 return nil, serviceerror.NewUnavailablef(
57 "EnqueueMessage failed for queue with type: %v and name: %v. BeginTx operation failed. Error: %v",
58 request.QueueType,
59 request.QueueName,
60 err,
61 )
62 }
63 lastMessageID, ok, err := q.getMaxMessageID(ctx, request.QueueType, request.QueueName, tx)
64 if err != nil {
65 rollBackErr := tx.Rollback()
66 if rollBackErr != nil {
67 q.logger.Error("transaction rollback error", tag.Error(rollBackErr))
68 }
69 return nil, serviceerror.NewUnavailablef(
70 "EnqueueMessage failed for queue with type: %v and name: %v. failed to get last messageId. Error: %v",
71 request.QueueType,
72 request.QueueName,
73 err,
74 )
75 }
76 nextMessageID := int64(persistence.FirstQueueMessageID)
77 if ok {
78 nextMessageID = lastMessageID + 1
79 }
80 _, err = tx.InsertIntoQueueV2Messages(ctx, []sqlplugin.QueueV2MessageRow{
81 newQueueV2Row(request.QueueType, request.QueueName, nextMessageID, request.Blob),
82 })
83 if err != nil {
84 rollBackErr := tx.Rollback()
85 if rollBackErr != nil {
86 q.logger.Error("transaction rollback error", tag.Error(rollBackErr))
87 }
88 return nil, serviceerror.NewUnavailablef(
89 "EnqueueMessage failed for queue with type: %v and name: %v. InsertIntoQueueV2Messages operation failed. Error: %v",
90 request.QueueType,
91 request.QueueName,
92 err,
93 )
94 }
95
96 if err := tx.Commit(); err != nil {
97 return nil, serviceerror.NewUnavailablef(
98 "EnqueueMessage failed for queue with type: %v and name: %v. commit operation failed. Error: %v",
99 request.QueueType,
100 request.QueueName,
101 err,
102 )

Callers

nothing calls this directly

Calls 9

getQueueMetadataMethod · 0.95
getMaxMessageIDMethod · 0.95
ErrorFunction · 0.92
newQueueV2RowFunction · 0.85
BeginTxMethod · 0.65
RollbackMethod · 0.65
ErrorMethod · 0.65
CommitMethod · 0.65

Tested by

no test coverage detected