-
Notifications
You must be signed in to change notification settings - Fork 94
Expand file tree
/
Copy pathactivity.go
More file actions
80 lines (66 loc) · 2.19 KB
/
Copy pathactivity.go
File metadata and controls
80 lines (66 loc) · 2.19 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
package valkey
import (
"context"
"fmt"
"github.com/cschleiden/go-workflows/backend"
"github.com/cschleiden/go-workflows/backend/history"
"github.com/cschleiden/go-workflows/workflow"
)
func (vb *valkeyBackend) PrepareActivityQueues(ctx context.Context, queues []workflow.Queue) error {
return vb.activityQueue.Prepare(ctx, vb.client, queues)
}
func (vb *valkeyBackend) GetActivityTask(ctx context.Context, queues []workflow.Queue) (*backend.ActivityTask, error) {
activityTask, err := vb.activityQueue.Dequeue(ctx, vb.client, queues, vb.options.ActivityLockTimeout, vb.options.BlockTimeout)
if err != nil {
return nil, err
}
if activityTask == nil {
return nil, nil
}
return &backend.ActivityTask{
WorkflowInstance: activityTask.Data.Instance,
Queue: workflow.Queue(activityTask.Data.Queue),
ID: activityTask.TaskID,
ActivityID: activityTask.Data.ID,
Event: activityTask.Data.Event,
}, nil
}
func (vb *valkeyBackend) ExtendActivityTask(ctx context.Context, task *backend.ActivityTask) error {
if err := vb.activityQueue.Extend(ctx, vb.client, task.Queue, task.ID); err != nil {
return err
}
return nil
}
func (vb *valkeyBackend) CompleteActivityTask(ctx context.Context, task *backend.ActivityTask, result *history.Event) error {
instance, err := readInstance(ctx, vb.client, vb.keys.instanceKey(task.WorkflowInstance))
if err != nil {
return err
}
eventData, payload, err := marshalEvent(result)
if err != nil {
return err
}
activityQueueKeys := vb.activityQueue.Keys(task.Queue)
workflowQueueKeys := vb.workflowQueue.Keys(workflow.Queue(instance.Queue))
err = vb.completeActivityTaskScript.Exec(ctx, vb.client, []string{
activityQueueKeys.SetKey,
activityQueueKeys.StreamKey,
vb.keys.pendingEventsKey(task.WorkflowInstance),
vb.keys.payloadKey(task.WorkflowInstance),
vb.workflowQueue.queueSetKey,
workflowQueueKeys.SetKey,
workflowQueueKeys.StreamKey,
}, []string{
task.ID,
vb.activityQueue.groupName,
result.ID,
eventData,
payload,
vb.workflowQueue.groupName,
instanceSegment(task.WorkflowInstance),
}).Error()
if err != nil {
return fmt.Errorf("completing activity task: %w", err)
}
return nil
}