Initial version
[ric-plt/xapp-frame.git] / pkg / xapp / rmr.go
1 /*
2 ==================================================================================
3   Copyright (c) 2019 AT&T Intellectual Property.
4   Copyright (c) 2019 Nokia
5
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
9
10        http://www.apache.org/licenses/LICENSE-2.0
11
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 ==================================================================================
18 */
19
20 package xapp
21
22 /*
23 #include <time.h>
24 #include <stdlib.h>
25 #include <stdio.h>
26 #include <string.h>
27 #include <rmr/rmr.h>
28 #include <rmr/RIC_message_types.h>
29
30 void write_bytes_array(unsigned char *dst, void *data, int len) {
31     memcpy((void *)dst, (void *)data, len);
32 }
33
34 #cgo CFLAGS: -I../
35 #cgo LDFLAGS: -lrmr_nng -lnng
36 */
37 import "C"
38
39 import (
40         "github.com/spf13/viper"
41         "strconv"
42         "sync"
43         "time"
44         "unsafe"
45 )
46
47 var RMRCounterOpts = []CounterOpts{
48         {Name: "Transmitted", Help: "The total number of transmited RMR messages"},
49         {Name: "Received", Help: "The total number of received RMR messages"},
50         {Name: "TransmitError", Help: "The total number of RMR transmission errors"},
51         {Name: "ReceiveError", Help: "The total number of RMR receive errors"},
52 }
53
54 // To be removed ...
55 type RMRStatistics struct{}
56
57 type RMRClient struct {
58         context   unsafe.Pointer
59         ready     int
60         wg        sync.WaitGroup
61         mux       sync.Mutex
62         stat      map[string]Counter
63         consumers []MessageConsumer
64 }
65
66 type MessageConsumer interface {
67         Consume(mtype int, sid int, len int, payload []byte) error
68 }
69
70 func NewRMRClient() *RMRClient {
71         r := &RMRClient{}
72         r.consumers = make([]MessageConsumer, 0)
73
74         p := C.CString(viper.GetString("rmr.protPort"))
75         m := C.int(viper.GetInt("rmr.maxSize"))
76         defer C.free(unsafe.Pointer(p))
77
78         r.context = C.rmr_init(p, m, C.int(0))
79         if r.context == nil {
80                 Logger.Fatal("rmrClient: Initializing RMR context failed, bailing out!")
81         }
82
83         return r
84 }
85
86 func (m *RMRClient) Start(c MessageConsumer) {
87         m.RegisterMetrics()
88
89         for {
90                 Logger.Info("rmrClient: Waiting for RMR to be ready ...")
91
92                 if m.ready = int(C.rmr_ready(m.context)); m.ready == 1 {
93                         break
94                 }
95                 time.Sleep(10 * time.Second)
96         }
97         m.wg.Add(viper.GetInt("rmr.numWorkers"))
98
99         if c != nil {
100                 m.consumers = append(m.consumers, c)
101         }
102
103         for w := 0; w < viper.GetInt("rmr.numWorkers"); w++ {
104                 go m.Worker("worker-"+strconv.Itoa(w), 0)
105         }
106
107         m.Wait()
108 }
109
110 func (m *RMRClient) Worker(taskName string, msgSize int) {
111         p := viper.GetString("rmr.protPort")
112         Logger.Info("rmrClient: '%s': receiving messages on [%s]", taskName, p)
113
114         defer m.wg.Done()
115         for {
116                 rxBuffer := C.rmr_rcv_msg(m.context, nil)
117                 if rxBuffer == nil {
118                         m.UpdateStatCounter("ReceiveError")
119                         continue
120                 }
121                 m.UpdateStatCounter("Received")
122
123                 go m.parseMessage(rxBuffer)
124         }
125 }
126
127 func (m *RMRClient) parseMessage(rxBuffer *C.rmr_mbuf_t) {
128         if len(m.consumers) == 0 {
129                 Logger.Info("rmrClient: No message handlers defined, message discarded!")
130                 return
131         }
132
133         for _, c := range m.consumers {
134                 cptr := unsafe.Pointer(rxBuffer.payload)
135                 payload := C.GoBytes(cptr, C.int(rxBuffer.len))
136
137                 err := c.Consume(int(rxBuffer.mtype), int(rxBuffer.sub_id), int(rxBuffer.len), payload)
138                 if err != nil {
139                         Logger.Warn("rmrClient: Consumer returned error: %v", err)
140                 }
141         }
142 }
143
144 func (m *RMRClient) Allocate() *C.rmr_mbuf_t {
145         buf := C.rmr_alloc_msg(m.context, 0)
146         if buf == nil {
147                 Logger.Fatal("rmrClient: Allocating message buffer failed!")
148         }
149
150         return buf
151 }
152
153 func (m *RMRClient) Send(mtype int, sid int, len int, payload []byte) bool {
154         buf := m.Allocate()
155
156         buf.mtype = C.int(mtype)
157         buf.sub_id = C.int(sid)
158         buf.len = C.int(len)
159         datap := C.CBytes(payload)
160         defer C.free(datap)
161
162         C.write_bytes_array(buf.payload, datap, C.int(len))
163
164         return m.SendBuf(buf)
165 }
166
167 func (m *RMRClient) SendBuf(txBuffer *C.rmr_mbuf_t) bool {
168         for i := 0; i < 10; i++ {
169                 txBuffer.state = 0
170                 txBuffer := C.rmr_send_msg(m.context, txBuffer)
171                 if txBuffer == nil {
172                         break
173                 } else if txBuffer.state != C.RMR_OK {
174                         if txBuffer.state != C.RMR_ERR_RETRY {
175                                 time.Sleep(100 * time.Microsecond)
176                                 m.UpdateStatCounter("TransmitError")
177                         }
178                         for j := 0; j < 100 && txBuffer.state == C.RMR_ERR_RETRY; j++ {
179                                 txBuffer = C.rmr_send_msg(m.context, txBuffer)
180                         }
181                 }
182
183                 if txBuffer.state == C.RMR_OK {
184                         m.UpdateStatCounter("Transmitted")
185                         return true
186                 }
187         }
188         m.UpdateStatCounter("TransmitError")
189         return false
190 }
191
192 func (m *RMRClient) UpdateStatCounter(name string) {
193         m.mux.Lock()
194         m.stat[name].Inc()
195         m.mux.Unlock()
196 }
197
198 func (m *RMRClient) RegisterMetrics() {
199         m.stat = Metric.RegisterCounterGroup(RMRCounterOpts, "RMR")
200 }
201
202 // To be removed ...
203 func (m *RMRClient) GetStat() (r RMRStatistics) {
204         return
205 }
206
207 func (m *RMRClient) Wait() {
208         m.wg.Wait()
209 }
210
211 func (m *RMRClient) IsReady() bool {
212         return m.ready != 0
213 }
214
215 func (m *RMRClient) GetRicMessageId(mid string) int {
216         return RICMessageTypes[mid]
217 }