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

Method UpdateTaskQueue

common/persistence/sql/task_queues.go:86–128  ·  view source on GitHub ↗
(
	ctx context.Context,
	request *persistence.InternalUpdateTaskQueueRequest,
)

Source from the content-addressed store, hash-verified

84}
85
86func (m *taskQueueStore) UpdateTaskQueue(
87 ctx context.Context,
88 request *persistence.InternalUpdateTaskQueueRequest,
89) (*persistence.UpdateTaskQueueResponse, error) {
90 nidBytes, err := primitives.ParseUUID(request.NamespaceID)
91 if err != nil {
92 return nil, serviceerror.NewInternal(err.Error())
93 }
94
95 tqId, tqHash := taskQueueIdAndHash(nidBytes, request.TaskQueue, request.TaskType, persistence.SubqueueZero)
96 var resp *persistence.UpdateTaskQueueResponse
97 err = m.txExecute(ctx, "UpdateTaskQueue", func(tx sqlplugin.Tx) error {
98 if err := lockTaskQueue(ctx,
99 tx,
100 tqHash,
101 tqId,
102 request.PrevRangeID,
103 m.version,
104 ); err != nil {
105 return err
106 }
107 result, err := tx.UpdateTaskQueues(ctx, &sqlplugin.TaskQueuesRow{
108 RangeHash: tqHash,
109 TaskQueueID: tqId,
110 RangeID: request.RangeID,
111 Data: request.TaskQueueInfo.Data,
112 DataEncoding: request.TaskQueueInfo.EncodingType.String(),
113 }, m.version)
114 if err != nil {
115 return err
116 }
117 rowsAffected, err := result.RowsAffected()
118 if err != nil {
119 return err
120 }
121 if rowsAffected != 1 {
122 return fmt.Errorf("%v rows were affected instead of 1", rowsAffected)
123 }
124 resp = &persistence.UpdateTaskQueueResponse{}
125 return nil
126 })
127 return resp, err
128}
129
130func (m *taskQueueStore) ListTaskQueue(
131 ctx context.Context,

Callers

nothing calls this directly

Calls 8

ParseUUIDFunction · 0.92
taskQueueIdAndHashFunction · 0.85
lockTaskQueueFunction · 0.85
ErrorMethod · 0.65
UpdateTaskQueuesMethod · 0.65
StringMethod · 0.65
txExecuteMethod · 0.45
ErrorfMethod · 0.45

Tested by

no test coverage detected