-
Notifications
You must be signed in to change notification settings - Fork 4
/
writer.go
105 lines (89 loc) · 2.19 KB
/
writer.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
99
100
101
102
103
104
105
package nq
import (
"errors"
"github.com/nats-io/nats.go"
)
type ResultHandlerIFACE interface {
// Get the result of a task in nats kv store
Get(id string) (*TaskMessage, error)
// Set the result of a task in nats kv store
Set(id string, data []byte) error
Watch(id string) (chan *TaskMessage, error)
GetAllKeys(id string, data []byte) ([]string, error)
}
type ResultHandlerNats struct {
kv nats.KeyValue
}
func NewResultHandlerNats(name string, js nats.JetStreamContext) *ResultHandlerNats {
kv, err := js.KeyValue(name)
if errors.Is(err, nats.ErrBucketNotFound) {
// create the bucket
kv, err := js.CreateKeyValue(&nats.KeyValueConfig{
Bucket: defaultKVName,
Description: "used by package for status retention and fetching",
Storage: nats.FileStorage,
})
if err != nil {
// failed to create a kv store
panic(err)
}
return &ResultHandlerNats{
kv: kv,
}
}
if err != nil {
panic(err)
}
return &ResultHandlerNats{
kv: kv,
}
}
func (rn *ResultHandlerNats) Get(id string) (*TaskMessage, error) {
x, err := rn.kv.Get(id)
if err != nil {
if errors.Is(err, nats.ErrKeyNotFound) {
return nil, ErrTaskNotFound
}
return nil, err
}
return DecodeTMFromJSON(x.Value())
}
func (rn *ResultHandlerNats) Set(id string, data []byte) error {
if _, err := rn.kv.Put(id, data); err != nil {
return err
} else {
return nil
}
}
// Get all keys from nats key-value store
func (rn *ResultHandlerNats) GetAllKeys(id string, data []byte) ([]string, error) {
if keys, err := rn.kv.Keys(); err != nil {
return nil, err
} else {
return keys, nil
}
}
func (rn *ResultHandlerNats) Watch(id string) (chan *TaskMessage, error) {
watcher, err := rn.kv.Watch(id)
status := make(chan *TaskMessage)
go func() {
for updated := range watcher.Updates() {
if updated != nil {
// TODO: handle error
x, err := DecodeTMFromJSON(updated.Value())
if err != nil {
close(status)
return
}
// log.Println("decode err", err, x)
status <- x
// check if status is terminal
if x.Status == Failed || x.Status == Completed || x.Status == Deleted || x.Status == Cancelled {
close(status)
return
}
}
}
}()
return status, err
}