-
Notifications
You must be signed in to change notification settings - Fork 0
/
context.go
98 lines (84 loc) · 2.17 KB
/
context.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
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
package master
import (
"bytes"
"io"
"net/http"
"reflect"
"sync"
"time"
"github.com/benmizrahi/gobig/internal/protos"
"github.com/benmizrahi/gobig/internal/worker"
"github.com/golang/protobuf/proto"
log "github.com/sirupsen/logrus"
)
type Context struct {
IsLocal bool
Workers map[string]string
//private
minWorkers int
masterPath string
}
func NewContext(isLocal bool, minWorkers int, masterPath string) *Context {
context := &Context{
IsLocal: isLocal,
Workers: map[string]string{},
minWorkers: minWorkers,
masterPath: masterPath,
}
return context
}
func (c *Context) InitContext() {
// start all workers
c.handleWorkers(c.minWorkers, c.IsLocal, c.masterPath)
for len(c.Workers) != c.minWorkers {
log.Info("gobig Master, wating for %d workers to register..", c.minWorkers)
time.Sleep(1 * time.Second)
}
log.Info("gobig Master, all workers are ready")
}
func (c *Context) sendAyncTaskToWorker(worker string, partition *protos.IPartition) *protos.IPartitionResult {
body, err := proto.Marshal(partition)
if err != nil {
log.Fatal("error:", err)
}
res, err := http.Post(worker+"/api/v1/tasks", "application/protobuf", bytes.NewReader(body))
if err != nil {
log.Fatal(err)
}
buf, err := io.ReadAll(res.Body)
if err != nil {
log.Fatal(err)
}
result := protos.IPartitionResult{}
err = proto.Unmarshal(buf, &result)
if err != nil {
log.Fatal(err)
}
return &result
}
func (c *Context) DoAction(plan []*protos.IPartition) []*protos.IPartitionResult {
//TODO publish actions to workers
var wg sync.WaitGroup
allPartitionResults := []*protos.IPartitionResult{}
keys := reflect.ValueOf(c.Workers).MapKeys()
for index, partition := range plan {
wg.Add(1)
num := index % len(keys)
worker := c.Workers[keys[num].String()]
go func() {
defer wg.Done()
allPartitionResults = append(allPartitionResults, c.sendAyncTaskToWorker(worker, partition))
}()
}
wg.Wait()
return allPartitionResults
}
func (c *Context) handleWorkers(minWorkers int, isLocal bool, masterPath string) {
if isLocal {
for i := 0; i < minWorkers; i++ {
worker.NewWorker("localhost", 8080+i, masterPath)
}
} else {
//TODO: implement GKE based orchstrations
}
}