-
Notifications
You must be signed in to change notification settings - Fork 0
/
app-sync.go
147 lines (122 loc) · 3.62 KB
/
app-sync.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
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
package publisher
import (
"encoding"
"encoding/json"
"fmt"
"net/http"
"time"
"github.com/warp-contracts/syncer/src/utils/config"
"github.com/warp-contracts/syncer/src/utils/monitoring"
"github.com/warp-contracts/syncer/src/utils/task"
"github.com/cenkalti/backoff/v4"
appsync "github.com/sony/appsync-client-go"
"github.com/sony/appsync-client-go/graphql"
)
const mutation = `mutation Publish($data: AWSJSON!, $name: String!) {
publish(data: $data, name: $name) {
data
name
}
}`
// Forwards messages to Redis
type AppSyncPublisher[In encoding.BinaryMarshaler] struct {
*task.Task
monitor monitoring.Monitor
client *appsync.Client
channelName string
input chan In
}
type Args struct {
Name string `json:"name"`
Data string `json:"data"`
}
func NewAppSyncPublisher[In encoding.BinaryMarshaler](config *config.Config, name string) (self *AppSyncPublisher[In]) {
self = new(AppSyncPublisher[In])
self.Task = task.NewTask(config, name).
WithSubtaskFunc(self.run).
WithWorkerPool(config.AppSync.MaxWorkers, config.AppSync.MaxQueueSize)
// Init AppSync client
gqlClient := graphql.NewClient(config.AppSync.Url,
graphql.WithCredential(config.AppSync.Token),
graphql.WithTimeout(time.Second*30),
)
self.client = appsync.NewClient(appsync.NewGraphQLClient(gqlClient))
return
}
func (self *AppSyncPublisher[In]) WithInputChannel(v chan In) *AppSyncPublisher[In] {
self.input = v
return self
}
func (self *AppSyncPublisher[In]) WithChannelName(v string) *AppSyncPublisher[In] {
self.channelName = v
return self
}
func (self *AppSyncPublisher[In]) WithMonitor(monitor monitoring.Monitor) *AppSyncPublisher[In] {
self.monitor = monitor
return self
}
func (self *AppSyncPublisher[In]) publish(data []byte) (err error) {
// Serialize args
args := Args{
Name: self.channelName,
Data: string(data),
}
argsBuf, err := json.Marshal(args)
if err != nil {
return
}
variables := json.RawMessage(argsBuf)
// Perform the request
response, err := self.client.Post(graphql.PostRequest{
Query: mutation,
Variables: &variables,
})
if err != nil {
return err
}
// Check response
if response.StatusCode != nil && *response.StatusCode != http.StatusOK {
err = fmt.Errorf("appsync publish failed with status %d", *response.StatusCode)
if *response.StatusCode > 399 && *response.StatusCode < 500 {
// Something's wrong with client configuration, don't retry
err = backoff.Permanent(err)
}
return
}
return nil
}
func (self *AppSyncPublisher[In]) run() (err error) {
for data := range self.input {
data := data
self.SubmitToWorker(func() {
self.Log.Debug("App sync publish...")
defer self.Log.Debug("...App sync publish done")
// Serialize to JSON
jsonData, err := data.MarshalBinary()
if err != nil {
self.Log.WithError(err).Error("Failed to marshal to json")
return
}
// Retry on failure with exponential backoff
err = task.NewRetry().
WithContext(self.Ctx).
WithMaxElapsedTime(self.Config.AppSync.BackoffMaxElapsedTime).
WithMaxInterval(self.Config.AppSync.BackoffMaxInterval).
WithOnError(func(err error, isDurationeAcceptable bool) error {
self.Log.WithError(err).Error("Appsync publish failed, retrying")
self.monitor.GetReport().AppSyncPublisher.Errors.Publish.Inc()
return err
}).
Run(func() error {
return self.publish(jsonData)
})
if err != nil {
self.Log.WithError(err).Error("Failed to publish to appsync after retries")
self.monitor.GetReport().AppSyncPublisher.Errors.PersistentFailure.Inc()
return
}
self.monitor.GetReport().AppSyncPublisher.State.MessagesPublished.Inc()
})
}
return nil
}