-
Notifications
You must be signed in to change notification settings - Fork 10
/
Copy pathmain.go
74 lines (65 loc) · 2.48 KB
/
main.go
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
package main
import (
"log"
"os"
updatabletimerv1 "github.com/cludden/protoc-gen-go-temporal/gen/example/updatabletimer/v1"
"github.com/urfave/cli/v2"
"go.temporal.io/sdk/client"
tlog "go.temporal.io/sdk/log"
"go.temporal.io/sdk/worker"
"go.temporal.io/sdk/workflow"
"google.golang.org/protobuf/types/known/timestamppb"
)
// UpdatableTimerWorkflow provides a updatabletimerv1.UpdatableTimerWorkflow implementation
type UpdatableTimerWorkflow struct {
*updatabletimerv1.UpdatableTimerWorkflowInput
log tlog.Logger
wakeUpTime *timestamppb.Timestamp
}
// NewUpdatableTimerWorkflow initializes a new updatabletimerv1.UpdatableTimerWorkflow value
func NewUpdatableTimerWorkflow(ctx workflow.Context, input *updatabletimerv1.UpdatableTimerWorkflowInput) (updatabletimerv1.UpdatableTimerWorkflow, error) {
return &UpdatableTimerWorkflow{input, workflow.GetLogger(ctx), input.Req.GetInitialWakeUpTime()}, nil
}
// Execute defines the entrypoint to a UpdatableTimer workflow
func (w *UpdatableTimerWorkflow) Execute(ctx workflow.Context) error {
for timerFired := false; !timerFired && ctx.Err() == nil; {
timerCtx, cancelTimer := workflow.WithCancel(ctx)
timer := workflow.NewTimer(timerCtx, w.wakeUpTime.AsTime().Sub(workflow.Now(ctx)))
w.log.Info("SleepUntil", "WakeUpTime", w.wakeUpTime)
workflow.NewSelector(ctx).
AddFuture(timer, func(f workflow.Future) {
if err := f.Get(timerCtx, nil); err != nil {
w.log.Info("Timer canceled")
} else {
w.log.Info("Timer fired")
timerFired = true
}
}).
AddReceive(w.UpdateWakeUpTime.Channel, func(workflow.ReceiveChannel, bool) {
defer cancelTimer()
w.wakeUpTime = w.UpdateWakeUpTime.ReceiveAsync().GetWakeUpTime()
w.log.Info("WakeUpTime updated", "WakeUpTime", w.wakeUpTime)
}).
Select(ctx)
}
return ctx.Err()
}
// GetWakeUpTime defines the entrypoint to a GetWakeUpTime query
func (w *UpdatableTimerWorkflow) GetWakeUpTime() (*updatabletimerv1.GetWakeUpTimeOutput, error) {
return &updatabletimerv1.GetWakeUpTimeOutput{WakeUpTime: w.wakeUpTime}, nil
}
func main() {
app, err := updatabletimerv1.NewExampleCli(
updatabletimerv1.NewExampleCliOptions().WithWorker(func(cmd *cli.Context, c client.Client) (worker.Worker, error) {
w := worker.New(c, updatabletimerv1.ExampleTaskQueue, worker.Options{})
updatabletimerv1.RegisterUpdatableTimerWorkflow(w, NewUpdatableTimerWorkflow)
return w, nil
}),
)
if err != nil {
log.Fatal(err)
}
if err := app.Run(os.Args); err != nil {
log.Fatal(err)
}
}