Repository navigation
Expand file tree
/
Copy pathexample_test.go
More file actions
109 lines (86 loc) · 4.01 KB
/
Copy pathexample_test.go
File metadata and controls
109 lines (86 loc) · 4.01 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
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
package workmanager_test
import (
"context"
"fmt"
"time"
wm "github.com/tr1v3r/workmanager"
)
func ExampleWorkerManager_newInstance() {
mgr := wm.NewWorkerManager(context.Background())
// register worker by workerbuilder with name
mgr.RegisterWorker(wm.DummyWorkerA, wm.DummyBuilder(wm.DummyWorkerA))
mgr.RegisterWorker(wm.DummyWorkerB, wm.DummyBuilder(wm.DummyWorkerB))
// register step, specify from which step to which step
mgr.RegisterStep(wm.StepA, wm.DummyStepRunner, wm.StepB)
mgr.RegisterStep(wm.StepB, wm.DummyStepRunner)
// register hooks
mgr.RegisterBeforeCallbacks(wm.StepA, func(_ context.Context, t ...wm.WorkTarget) []wm.WorkTarget {
fmt.Printf("[%s] before callback got target: %+v\n", wm.StepA, t[0])
return t
})
mgr.RegisterAfterCallbacks(wm.StepA, func(_ context.Context, t ...wm.WorkTarget) []wm.WorkTarget {
fmt.Printf("[%s] after callback got target: %+v\n", wm.StepA, t[0])
return t
})
mgr.RegisterAfterCallbacks(wm.StepB, func(_ context.Context, t ...wm.WorkTarget) []wm.WorkTarget {
mgr.FinishTask(t[0].Token())
return t
})
mgr.SetPipe(wm.StepA, wm.PipeChSize(8))
// start serve
mgr.Serve(wm.StepA, wm.StepB)
task := wm.NewTask(context.Background())
task.(*wm.Task).SetToken("example_token_123")
mgr.AddTask(task)
err := mgr.Recv(wm.StepA, &wm.DummyTestTarget{DummyTarget: wm.DummyTarget{TaskToken: task.Token()}, Step: wm.StepA})
if err != nil {
fmt.Printf("send target fail: %s", err)
}
for c := time.Tick(100 * time.Millisecond); !task.IsFinished(); <-c {
}
resultTask := mgr.GetTask(task.Token()).(*wm.Task)
fmt.Printf("task final status: { token: %s, finished: %t }", resultTask.Token(), resultTask.IsFinished())
// Output:
// [step_a] before callback got target: &{DummyTarget:{TaskToken:example_token_123} Step:step_a Remark: Count:0}
// [step_a] after callback got target: &{DummyTarget:{TaskToken:example_token_123} Step:step_b Remark: Count:0}
// [step_a] got result: &{DummyTarget:{TaskToken:example_token_123} Step:step_b Remark: Count:0}
// [step_b] got result: &{DummyTarget:{TaskToken:example_token_123} Step:step_a Remark: Count:0}
// task final status: { token: example_token_123, finished: true }
}
func ExampleWorkerManager_singleton() {
wm.RegisterWorker(wm.DummyWorkerA, wm.DummyBuilder(wm.DummyWorkerA))
wm.RegisterWorker(wm.DummyWorkerB, wm.DummyBuilder(wm.DummyWorkerB))
wm.RegisterStep(wm.StepA, wm.DummyStepRunner, wm.StepB)
wm.RegisterStep(wm.StepB, wm.DummyStepRunner)
wm.RegisterBeforeCallbacks(wm.StepA, func(_ context.Context, t ...wm.WorkTarget) []wm.WorkTarget {
fmt.Printf("[%s] before callback got target: %+v\n", wm.StepA, t[0])
return t
})
wm.RegisterAfterCallbacks(wm.StepA, func(_ context.Context, t ...wm.WorkTarget) []wm.WorkTarget {
fmt.Printf("[%s] after callback got target: %+v\n", wm.StepA, t[0])
return t
})
wm.RegisterAfterCallbacks(wm.StepB, func(_ context.Context, t ...wm.WorkTarget) []wm.WorkTarget {
wm.FinishTask(t[0].Token())
return t
})
wm.SetPipe(wm.StepA, wm.PipeChSize(8))
wm.Serve(wm.StepA, wm.StepB)
task := wm.NewTask(context.Background())
task.(*wm.Task).SetToken("example_token_123")
wm.AddTask(task)
err := wm.Recv(wm.StepA, &wm.DummyTestTarget{DummyTarget: wm.DummyTarget{TaskToken: task.Token()}, Step: wm.StepA})
if err != nil {
fmt.Printf("send target fail: %s", err)
}
for c := time.Tick(100 * time.Millisecond); !task.IsFinished(); <-c {
}
resultTask := wm.GetTask(task.Token()).(*wm.Task)
fmt.Printf("task final status: { token: %s, finished: %t }", resultTask.Token(), resultTask.IsFinished())
// Output:
// [step_a] before callback got target: &{DummyTarget:{TaskToken:example_token_123} Step:step_a Remark: Count:0}
// [step_a] after callback got target: &{DummyTarget:{TaskToken:example_token_123} Step:step_b Remark: Count:0}
// [step_a] got result: &{DummyTarget:{TaskToken:example_token_123} Step:step_b Remark: Count:0}
// [step_b] got result: &{DummyTarget:{TaskToken:example_token_123} Step:step_a Remark: Count:0}
// task final status: { token: example_token_123, finished: true }
}