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

Method EnqueueMessage

common/persistence/sql/queue.go:48–73  ·  view source on GitHub ↗
(
	ctx context.Context,
	blob *commonpb.DataBlob,
)

Source from the content-addressed store, hash-verified

46}
47
48func (q *sqlQueue) EnqueueMessage(
49 ctx context.Context,
50 blob *commonpb.DataBlob,
51) error {
52 err := q.txExecute(ctx, "EnqueueMessage", func(tx sqlplugin.Tx) error {
53 lastMessageID, err := tx.GetLastEnqueuedMessageIDForUpdate(ctx, q.queueType)
54 switch err {
55 case nil:
56 _, err = tx.InsertIntoMessages(ctx, []sqlplugin.QueueMessageRow{
57 newQueueRow(q.queueType, lastMessageID+1, blob),
58 })
59 return err
60 case sql.ErrNoRows:
61 _, err = tx.InsertIntoMessages(ctx, []sqlplugin.QueueMessageRow{
62 newQueueRow(q.queueType, persistence.EmptyQueueMessageID+1, blob),
63 })
64 return err
65 default:
66 return fmt.Errorf("failed to get last enqueued message id: %v", err)
67 }
68 })
69 if err != nil {
70 return serviceerror.NewUnavailable(err.Error())
71 }
72 return nil
73}
74
75func (q *sqlQueue) ReadMessages(
76 ctx context.Context,

Callers

nothing calls this directly

Calls 6

newQueueRowFunction · 0.85
InsertIntoMessagesMethod · 0.65
ErrorMethod · 0.65
txExecuteMethod · 0.45
ErrorfMethod · 0.45

Tested by

no test coverage detected