move SocketService to a separate file
lgz committed
Oct 13, 2017 at 00:17 UTC
ff3d187a39ab5fc6f53f371bbab3921c79d3d86d
1 file changed
+250
python.d/python_modules/bases/FrameworkServices/SocketService.py
new
+250
@@ -0,0 +1,250 @@
1
+# -*- coding: utf-8 -*-
2
+# Description:
3
+# Author: Pawel Krupa (paulfantom)
4
+
5
+import socket
6
+
7
+from bases.FrameworkServices.SimpleService import SimpleService
8
+
9
+
10
+class SocketService(SimpleService):
11
+ def __init__(self, configuration=None, name=None):
12
+ self._sock = None
13
+ self._keep_alive = False
14
+ self.host = 'localhost'
15
+ self.port = None
16
+ self.unix_socket = None
17
+ self.request = ''
18
+ self.__socket_config = None
19
+ self.__empty_request = "".encode()
20
+ SimpleService.__init__(self, configuration=configuration, name=name)
21
+
22
+ def _socket_error(self, message=None):
23
+ if self.unix_socket is not None:
24
+ self.error('unix socket "{socket}": {message}'.format(socket=self.unix_socket,
25
+ message=message))
26
+ else:
27
+ if self.__socket_config is not None:
28
+ af, sock_type, proto, canon_name, sa = self.__socket_config
29
+ self.error('socket to "{address}" port {port}: {message}'.format(address=sa[0],
30
+ port=sa[1],
31
+ message=message))
32
+ else:
33
+ self.error('unknown socket: {0}'.format(message))
34
+
35
+ def _connect2socket(self, res=None):
36
+ """
37
+ Connect to a socket, passing the result of getaddrinfo()
38
+ :return: boolean
39
+ """
40
+ if res is None:
41
+ res = self.__socket_config
42
+ if res is None:
43
+ self.error("Cannot create socket to 'None':")
44
+ return False
45
+
46
+ af, sock_type, proto, canon_name, sa = res
47
+ try:
48
+ self.debug('Creating socket to "{address}", port {port}'.format(address=sa[0], port=sa[1]))
49
+ self._sock = socket.socket(af, sock_type, proto)
50
+ except socket.error as error:
51
+ self.error('Failed to create socket "{address}", port {port}, error: {error}'.format(address=sa[0],
52
+ port=sa[1],
53
+ error=error))
54
+ self._sock = None
55
+ self.__socket_config = None
56
+ return False
57
+
58
+ try:
59
+ self.debug('connecting socket to "{address}", port {port}'.format(address=sa[0], port=sa[1]))
60
+ self._sock.connect(sa)
61
+ except socket.error as error:
62
+ self.error('Failed to connect to "{address}", port {port}, error: {error}'.format(address=sa[0],
63
+ port=sa[1],
64
+ error=error))
65
+ self._disconnect()
66
+ self.__socket_config = None
67
+ return False
68
+
69
+ self.debug('connected to "{address}", port {port}'.format(address=sa[0], port=sa[1]))
70
+ self.__socket_config = res
71
+ return True
72
+
73
+ def _connect2unixsocket(self):
74
+ """
75
+ Connect to a unix socket, given its filename
76
+ :return: boolean
77
+ """
78
+ if self.unix_socket is None:
79
+ self.error("cannot connect to unix socket 'None'")
80
+ return False
81
+
82
+ try:
83
+ self.debug('attempting DGRAM unix socket "{0}"'.format(self.unix_socket))
84
+ self._sock = socket.socket(socket.AF_UNIX, socket.SOCK_DGRAM)
85
+ self._sock.connect(self.unix_socket)
86
+ self.debug('connected DGRAM unix socket "{0}"'.format(self.unix_socket))
87
+ return True
88
+ except socket.error as error:
89
+ self.debug('Failed to connect DGRAM unix socket "{socket}": {error}'.format(socket=self.unix_socket,
90
+ error=error))
91
+
92
+ try:
93
+ self.debug('attempting STREAM unix socket "{0}"'.format(self.unix_socket))
94
+ self._sock = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
95
+ self._sock.connect(self.unix_socket)
96
+ self.debug('connected STREAM unix socket "{0}"'.format(self.unix_socket))
97
+ return True
98
+ except socket.error as error:
99
+ self.debug('Failed to connect STREAM unix socket "{socket}": {error}'.format(socket=self.unix_socket,
100
+ error=error))
101
+ self._sock = None
102
+ return False
103
+
104
+ def _connect(self):
105
+ """
106
+ Recreate socket and connect to it since sockets cannot be reused after closing
107
+ Available configurations are IPv6, IPv4 or UNIX socket
108
+ :return:
109
+ """
110
+ try:
111
+ if self.unix_socket is not None:
112
+ self._connect2unixsocket()
113
+
114
+ else:
115
+ if self.__socket_config is not None:
116
+ self._connect2socket()
117
+ else:
118
+ for res in socket.getaddrinfo(self.host, self.port, socket.AF_UNSPEC, socket.SOCK_STREAM):
119
+ if self._connect2socket(res):
120
+ break
121
+
122
+ except Exception:
123
+ self._sock = None
124
+ self.__socket_config = None
125
+
126
+ if self._sock is not None:
127
+ self._sock.setblocking(0)
128
+ self._sock.settimeout(5)
129
+ self.debug('set socket timeout to: {0}'.format(self._sock.gettimeout()))
130
+
131
+ def _disconnect(self):
132
+ """
133
+ Close socket connection
134
+ :return:
135
+ """
136
+ if self._sock is not None:
137
+ try:
138
+ self.debug('closing socket')
139
+ self._sock.shutdown(2) # 0 - read, 1 - write, 2 - all
140
+ self._sock.close()
141
+ except Exception:
142
+ pass
143
+ self._sock = None
144
+
145
+ def _send(self):
146
+ """
147
+ Send request.
148
+ :return: boolean
149
+ """
150
+ # Send request if it is needed
151
+ if self.request != self.__empty_request:
152
+ try:
153
+ self.debug('sending request: {0}'.format(self.request))
154
+ self._sock.send(self.request)
155
+ except Exception as error:
156
+ self._socket_error('error sending request: {0}'.format(error))
157
+ self._disconnect()
158
+ return False
159
+ return True
160
+
161
+ def _receive(self):
162
+ """
163
+ Receive data from socket
164
+ :return: str
165
+ """
166
+ data = ""
167
+ while True:
168
+ self.debug('receiving response')
169
+ try:
170
+ buf = self._sock.recv(4096)
171
+ except Exception as error:
172
+ self._socket_error('failed to receive response: {0}'.format(error))
173
+ self._disconnect()
174
+ break
175
+
176
+ if buf is None or len(buf) == 0: # handle server disconnect
177
+ if data == "":
178
+ self._socket_error('unexpectedly disconnected')
179
+ else:
180
+ self.debug('server closed the connection')
181
+ self._disconnect()
182
+ break
183
+
184
+ self.debug('received data')
185
+ data += buf.decode('utf-8', 'ignore')
186
+ if self._check_raw_data(data):
187
+ break
188
+
189
+ self.debug('final response: {0}'.format(data))
190
+ return data
191
+
192
+ def _get_raw_data(self):
193
+ """
194
+ Get raw data with low-level "socket" module.
195
+ :return: str
196
+ """
197
+ if self._sock is None:
198
+ self._connect()
199
+ if self._sock is None:
200
+ return None
201
+
202
+ # Send request if it is needed
203
+ if not self._send():
204
+ return None
205
+
206
+ data = self._receive()
207
+
208
+ if not self._keep_alive:
209
+ self._disconnect()
210
+
211
+ return data
212
+
213
+ @staticmethod
214
+ def _check_raw_data(data):
215
+ """
216
+ Check if all data has been gathered from socket
217
+ :param data: str
218
+ :return: boolean
219
+ """
220
+ return bool(data)
221
+
222
+ def _parse_config(self):
223
+ """
224
+ Parse configuration data
225
+ :return: boolean
226
+ """
227
+ try:
228
+ self.unix_socket = str(self.configuration['socket'])
229
+ except (KeyError, TypeError):
230
+ self.debug('No unix socket specified. Trying TCP/IP socket.')
231
+ self.unix_socket = None
232
+ try:
233
+ self.host = str(self.configuration['host'])
234
+ except (KeyError, TypeError):
235
+ self.debug('No host specified. Using: "{0}"'.format(self.host))
236
+ try:
237
+ self.port = int(self.configuration['port'])
238
+ except (KeyError, TypeError):
239
+ self.debug('No port specified. Using: "{0}"'.format(self.port))
240
+
241
+ try:
242
+ self.request = str(self.configuration['request'])
243
+ except (KeyError, TypeError):
244
+ self.debug('No request specified. Using: "{0}"'.format(self.request))
245
+
246
+ self.request = self.request.encode()
247
+
248
+ def check(self):
249
+ self._parse_config()
250
+ return SimpleService.check(self)