2 ==================================================================================
3 Copyright (c) 2019 AT&T Intellectual Property.
4 Copyright (c) 2019 Nokia
6 Licensed under the Apache License, Version 2.0 (the "License");
7 you may not use this file except in compliance with the License.
8 You may obtain a copy of the License at
10 http://www.apache.org/licenses/LICENSE-2.0
12 Unless required by applicable law or agreed to in writing, software
13 distributed under the License is distributed on an "AS IS" BASIS,
14 WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
15 See the License for the specific language governing permissions and
16 limitations under the License.
17 ==================================================================================
25 "github.com/segmentio/ksuid"
30 func (sd *SubscriptionDispatcher) Initialize() {
31 sd.client = &http.Client{}
35 sd.subscriptions = sd.db.RestoreSubscriptions()
38 func (sd *SubscriptionDispatcher) Add(sr SubscriptionReq) (resp SubscriptionResp) {
39 key := ksuid.New().String()
40 resp = SubscriptionResp{key, 0, sr.EventType}
43 sd.subscriptions.Set(key, Subscription{sr, resp})
44 sd.db.StoreSubscriptions(sd.subscriptions)
46 Logger.Info("Sub: New subscription added: key=%s value=%v", key, sr)
50 func (sd *SubscriptionDispatcher) GetAll() (hooks []SubscriptionReq) {
51 hooks = []SubscriptionReq{}
52 for v := range sd.subscriptions.IterBuffered() {
53 hooks = append(hooks, v.Val.(Subscription).req)
59 func (sd *SubscriptionDispatcher) Get(id string) (SubscriptionReq, bool) {
60 if v, found := sd.subscriptions.Get(id); found {
61 Logger.Info("Subscription id=%s found: %v", id, v.(Subscription).req)
63 return v.(Subscription).req, found
65 return SubscriptionReq{}, false
68 func (sd *SubscriptionDispatcher) Delete(id string) (SubscriptionReq, bool) {
69 if v, found := sd.subscriptions.Get(id); found {
70 Logger.Info("Subscription id=%s found: %v ... deleting", id, v.(Subscription).req)
72 sd.subscriptions.Remove(id)
73 sd.db.StoreSubscriptions(sd.subscriptions)
75 return v.(Subscription).req, found
77 return SubscriptionReq{}, false
80 func (sd *SubscriptionDispatcher) Update(id string, sr SubscriptionReq) (SubscriptionReq, bool) {
81 if s, found := sd.subscriptions.Get(id); found {
82 Logger.Info("Subscription id=%s found: %v ... updating", id, s.(Subscription).req)
85 sd.subscriptions.Set(id, Subscription{sr, s.(Subscription).resp})
86 sd.db.StoreSubscriptions(sd.subscriptions)
90 return SubscriptionReq{}, false
93 func (sd *SubscriptionDispatcher) Publish(x Xapp, et EventType) {
94 sd.notifyClients([]Xapp{x}, et)
97 func (sd *SubscriptionDispatcher) notifyClients(xapps []Xapp, et EventType) {
98 if len(xapps) == 0 || len(sd.subscriptions) == 0 {
99 Logger.Info("Nothing to publish [%d:%d]", len(xapps), len(sd.subscriptions))
104 for v := range sd.subscriptions.Iter() {
105 go sd.notify(xapps, et, v.Val.(Subscription), sd.Seq)
109 func (sd *SubscriptionDispatcher) notify(xapps []Xapp, et EventType, s Subscription, seq int) error {
110 xappData, err := json.Marshal(xapps)
112 Logger.Info("json.Marshal failed: %v", err)
116 notif := SubscriptionNotif{Id: s.req.Id, Version: seq, EventType: string(et), XApps: string(xappData)}
117 jsonData, err := json.Marshal(notif)
119 Logger.Info("json.Marshal failed: %v", err)
123 // Execute the request with retry policy
124 return sd.retry(s, func() error {
125 Logger.Info("Posting notification to targetUrl=%s: %v", s.req.TargetUrl, notif)
126 resp, err := http.Post(s.req.TargetUrl, "application/json", bytes.NewBuffer(jsonData))
128 Logger.Info("Posting to subscription failed: %v", err)
132 if resp.StatusCode != http.StatusOK {
133 Logger.Info("Client returned error code: %d", resp.StatusCode)
137 Logger.Info("subscription to '%s' dispatched, response code: %d", s.req.TargetUrl, resp.StatusCode)
142 func (sd *SubscriptionDispatcher) retry(s Subscription, fn func() error) error {
143 if err := fn(); err != nil {
144 // Todo: use exponential backoff, or similar mechanism
145 if s.req.MaxRetries--; s.req.MaxRetries > 0 {
146 time.Sleep(time.Duration(s.req.RetryTimer) * time.Second)
147 return sd.retry(s, fn)
149 sd.subscriptions.Remove(s.req.Id)