protocol.py 14.7 KB
Newer Older
Giorgos Korfiatis's avatar
Giorgos Korfiatis committed
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
# Copyright (C) 2015 GRNET S.A.
#
# This program is free software: you can redistribute it and/or modify
# it under the terms of the GNU General Public License as published by
# the Free Software Foundation, either version 3 of the License, or
# (at your option) any later version.
#
# This program is distributed in the hope that it will be useful,
# but WITHOUT ANY WARRANTY; without even the implied warranty of
# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE.  See the
# GNU General Public License for more details.
#
# You should have received a copy of the GNU General Public License
# along with this program.  If not, see <http://www.gnu.org/licenses/>.

16
from wsgiref.simple_server import make_server
17
from ws4py.websocket import WebSocket
18
19
20
21
from ws4py.server.wsgiutils import WebSocketWSGIApplication
from ws4py.server.wsgirefserver import WSGIServer, WebSocketWSGIRequestHandler
from hashlib import sha1
from threading import Thread
22
23
import sqlite3
import time
24
import os
25
26
import json
import logging
27
28
from agkyra.syncer import (
    syncer, setup, pithos_client, localfs_client, messaging)
29
from agkyra.config import AgkyraConfig, AGKYRA_DIR
30
31
32
33
34


LOG = logging.getLogger(__name__)


35
class SessionHelper(object):
36
37
38
    """Agkyra Helper Server sets a WebSocket server with the Helper protocol
    It also provided methods for running and killing the Helper server
    """
39
    session_timeout = 20
40

41
    def __init__(self, **kwargs):
42
        """Setup the helper server"""
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
        self.session_db = kwargs.get(
            'session_db', os.path.join(AGKYRA_DIR, 'session.db'))
        self.session_relation = kwargs.get('session_relation', 'heart')

        LOG.debug('Connect to db')
        self.db = sqlite3.connect(self.session_db)
        self._init_db_relation()
        self.session = self._load_active_session() or self._create_session()

    def _init_db_relation(self):
        self.db.execute('BEGIN')
        self.db.execute(
            'CREATE TABLE IF NOT EXISTS %s ('
            'ui_id VARCHAR(256), address text, beat VARCHAR(32)'
            ')' % self.session_relation)
        self.db.commit()

    def _load_active_session(self):
        """Load a session from db"""
        r = self.db.execute('SELECT * FROM %s' % self.session_relation)
        sessions = r.fetchall()
        if sessions:
            last = sessions[-1]
            now, last_beat = time.time(), float(last[2])
            if now - last_beat < self.session_timeout:
                # Found an active session
                return dict(ui_id=last[0], address=last[1])
        return None

    def _create_session(self):
        """Create session credentials"""
        ui_id = sha1(os.urandom(128)).hexdigest()

        WebSocketProtocol.ui_id = ui_id
77
        server = make_server(
78
            '', 0,
79
80
81
82
            server_class=WSGIServer,
            handler_class=WebSocketWSGIRequestHandler,
            app=WebSocketWSGIApplication(handler_cls=WebSocketProtocol))
        server.initialize_websockets_manager()
83
        address = 'ws://%s:%s' % (server.server_name, server.server_port)
84
85
        self.server = server

86
87
88
89
90
91
92
93
        self.db.execute('BEGIN')
        self.db.execute('DELETE FROM %s' % self.session_relation)
        self.db.execute('INSERT INTO %s VALUES ("%s", "%s", "%s")' % (
            self.session_relation, ui_id, address, time.time()))
        self.db.commit()

        return dict(ui_id=ui_id, address=address)

94
95
    def start(self):
        """Start the helper server in a thread"""
96
97
        if getattr(self, 'server', None):
            Thread(target=self.server.serve_forever).start()
98
99
100

    def shutdown(self):
        """Shutdown the server (needs another thread) and join threads"""
101
102
103
104
        if getattr(self, 'server', None):
            t = Thread(target=self.server.shutdown)
            t.start()
            t.join()
105
106


107
108
109
110
class WebSocketProtocol(WebSocket):
    """Helper-side WebSocket protocol for communication with GUI:

    -- INTERRNAL HANDSAKE --
111
    GUI: {"method": "post", "ui_id": <GUI ID>}
112
    HELPER: {"ACCEPTED": 202, "method": "post"}" or
113
        "{"REJECTED": 401, "action": "post ui_id"}
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
148
149
150
151

    -- SHUT DOWN --
    GUI: {"method": "post", "path": "shutdown"}

    -- PAUSE --
    GUI: {"method": "post", "path": "pause"}
    HELPER: {"OK": 200, "action": "post pause"} or error

    -- start --
    GUI: {"method": "post", "path": "start"}
    HELPER: {"OK": 200, "action": "post start"} or error

    -- GET SETTINGS --
    GUI: {"method": "get", "path": "settings"}
    HELPER:
        {
            "action": "get settings",
            "token": <user token>,
            "url": <auth url>,
            "container": <container>,
            "directory": <local directory>,
            "exclude": <file path>
        } or {<ERROR>: <ERROR CODE>}

    -- PUT SETTINGS --
    GUI: {
            "method": "put", "path": "settings",
            "token": <user token>,
            "url": <auth url>,
            "container": <container>,
            "directory": <local directory>,
            "exclude": <file path>
        }
    HELPER: {"CREATED": 201, "action": "put settings",} or
        {<ERROR>: <ERROR CODE>, "action": "get settings",}

    -- GET STATUS --
    GUI: {"method": "get", "path": "status"}
152
153
154
155
156
    HELPER: {
        "can_sync": <boolean>,
        "progress": <int>,
        "paused": <boolean>,
        "action": "get status"} or
157
158
159
        {<ERROR>: <ERROR CODE>, "action": "get status"}
    """

160
    ui_id = None
161
162
163
164
    accepted = False
    settings = dict(
        token=None, url=None,
        container=None, directory=None,
165
        exclude=None)
166
167
    status = dict(
        progress=0, synced=0, unsynced=0, paused=True, can_sync=False)
168
169
    file_syncer = None
    cnf = AgkyraConfig()
170
    essentials = ('url', 'token', 'container', 'directory')
171

172
173
174
175
    def _get_default_sync(self):
        """Get global.default_sync or pick the first sync as default
        If there are no syncs, create a 'default' sync.
        """
176
        sync = self.cnf.get('global', 'default_sync')
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
        if not sync:
            for sync in self.cnf.keys('sync'):
                break
            self.cnf.set('global', 'default_sync', sync or 'default')
        return sync or 'default'

    def _get_sync_cloud(self, sync):
        """Get the <sync>.cloud or pick the first cloud and use it
        In case of cloud picking, set the cloud as the <sync>.cloud for future
        sessions.
        If no clouds are found, create a 'default' cloud, with an empty url.
        """
        try:
            cloud = self.cnf.get_sync(sync, 'cloud')
        except KeyError:
            cloud = None
        if not cloud:
            for cloud in self.cnf.keys('cloud'):
                break
            self.cnf.set_sync(sync, 'cloud', cloud or 'default')
        return cloud or 'default'
198

199
200
201
202
    def _load_settings(self):
        LOG.debug('Start loading settings')
        sync = self._get_default_sync()
        cloud = self._get_sync_cloud(sync)
203
204

        try:
205
206
207
208
209
210
211
            self.settings['url'] = self.cnf.get_cloud(cloud, 'url')
        except Exception:
            self.settings['url'] = None
        try:
            self.settings['token'] = self.cnf.get_cloud(cloud, 'token')
        except Exception:
            self.settings['url'] = None
212
213

        for option in ('container', 'directory', 'exclude'):
214
215
216
217
            try:
                self.settings[option] = self.cnf.get_sync(sync, option)
            except KeyError:
                LOG.debug('No %s is set' % option)
218

219
220
        LOG.debug('Finished loading settings')

221
    def _dump_settings(self):
222
        LOG.debug('Saving settings')
223
224
225
226
227
228
229
230
231
232
233
        if not self.settings.get('url', None):
            LOG.debug('No settings to save')
            return

        sync = self._get_default_sync()
        cloud = self._get_sync_cloud(sync)

        try:
            old_url = self.cnf.get_cloud(cloud, 'url') or ''
        except KeyError:
            old_url = self.settings['url']
234
235
236
237
238
239
240
241
242

        while old_url != self.settings['url']:
            cloud = '%s_%s' % (cloud, sync)
            try:
                self.cnf.get_cloud(cloud, 'url')
            except KeyError:
                break

        self.cnf.set_cloud(cloud, 'url', self.settings['url'])
243
        self.cnf.set_cloud(cloud, 'token', self.settings['token'] or '')
244
245
246
        self.cnf.set_sync(sync, 'cloud', cloud)

        for option in ('directory', 'container', 'exclude'):
247
            self.cnf.set_sync(sync, option, self.settings[option] or '')
248

249
250
251
        self.cnf.write()
        LOG.debug('Settings saved')

252
253
254
255
256
    def _essentials_changed(self, new_settings):
        """Check if essential settings have changed in new_settings"""
        return all([
            self.settings[e] == self.settings[e] for e in self.essentials])

257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
    def _update_statistics(self):
        """Update statistics by consuming and understanding syncer messages"""
        if self.can_sync():
            msg = self.syncer.get_next_message()
            if not msg:
                if self.status['unsynced'] == self.status['synced']:
                    self.status['unsynced'] = 0
                    self.status['synced'] = 0
            while (msg):
                if isinstance(msg, messaging.SyncMessage):
                    LOG.info('Start syncing "%s"' % msg.objname)
                    self.status['unsynced'] += 1
                elif isinstance(msg, messaging.AckSyncMessage):
                    LOG.info('Finished syncing "%s"' % msg.objname)
                    self.status['synced'] += 1
                elif isinstance(msg, messaging.CollisionMessage):
                    LOG.info('Collision for "%s"' % msg.objname)
                elif isinstance(msg, messaging.ConflictStashMessage):
                    LOG.info('Conflict for "%s"' % msg.objname)
                else:
                    LOG.debug('Consumed msg %s' % msg)
                msg = self.syncer.get_next_message()

    def can_sync(self):
        """Check if settings are enough to setup a syncing proccess"""
        return all([self.settings[e] for e in self.essentials])

284
285
286
287
288
289
290
    def init_sync(self):
        """Initialize syncer"""
        sync = self._get_default_sync()
        syncer_settings = setup.SyncerSettings(
            sync,
            self.settings['url'], self.settings['token'],
            self.settings['container'], self.settings['directory'],
291
            agkyra_path=AGKYRA_DIR,
292
293
294
295
296
            ignore_ssl=True)
        master = pithos_client.PithosFileClient(syncer_settings)
        slave = localfs_client.LocalfsFileClient(syncer_settings)
        self.syncer = syncer.FileSyncer(syncer_settings, master, slave)
        self.syncer_settings = syncer_settings
Giorgos Korfiatis's avatar
Giorgos Korfiatis committed
297
        self.syncer.initiate_probe()
298

299
300
    # Syncer-related methods
    def get_status(self):
301
302
303
304
305
306
307
        if self.can_sync():
            self._update_statistics()
            self.status['paused'] = self.syncer.paused
            self.status['can_sync'] = self.can_sync()
        else:
            self.status = dict(
                progress=0, synced=0, unsynced=0, paused=True, can_sync=False)
308
309
310
311
312
313
        return self.status

    def get_settings(self):
        return self.settings

    def set_settings(self, new_settings):
314
315
        # Prepare setting save
        could_sync = self.can_sync()
316
317
318
        was_active = False
        if could_sync and not self.syncer.paused:
            was_active = True
319
320
321
322
            self.pause_sync()
        must_reset_syncing = self._essentials_changed(new_settings)

        # save settings
323
324
325
        self.settings = new_settings
        self._dump_settings()

326
327
328
329
330
331
332
        # Restart
        if self.can_sync():
            if must_reset_syncing or not could_sync:
                self.init_sync()
            elif was_active:
                self.start_sync()

333
    def pause_sync(self):
Giorgos Korfiatis's avatar
Giorgos Korfiatis committed
334
        self.syncer.stop_decide()
335
336
        LOG.debug('Wait open syncs to complete')
        self.syncer.wait_sync_threads()
337
338

    def start_sync(self):
339
        self.syncer.start_decide()
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357

    # WebSocket connection methods
    def opened(self):
        LOG.debug('Helper: connection established')

    def closed(self, *args):
        LOG.debug('Helper: connection closed')

    def send_json(self, msg):
        LOG.debug('send: %s' % msg)
        self.send(json.dumps(msg))

    # Protocol handling methods
    def _post(self, r):
        """Handle POST requests"""
        if self.accepted:
            action = r['path']
            if action == 'shutdown':
358
359
                if self.can_sync():
                    self.syncer.stop_all_daemons()
360
361
                    LOG.debug('Wait open syncs to complete')
                    self.syncer.wait_sync_threads()
362
363
364
365
366
367
368
                self.close()
                return
            {
                'start': self.start_sync,
                'pause': self.pause_sync
            }[action]()
            self.send_json({'OK': 200, 'action': 'post %s' % action})
369
        elif r['ui_id'] == self.ui_id:
370
            self._load_settings()
371
            self.accepted = True
372
            self.send_json({'ACCEPTED': 202, 'action': 'post ui_id'})
373
374
375
            if self.can_sync():
                self.init_sync()
                self.pause_sync()
376
        else:
377
            action = r.get('path', 'ui_id')
378
379
380
381
382
            self.send_json({'REJECTED': 401, 'action': 'post %s' % action})
            self.terminate()

    def _put(self, r):
        """Handle PUT requests"""
383
        if self.accepted:
384
385
386
387
388
            LOG.debug('put %s' % r)
            action = r.pop('path')
            self.set_settings(r)
            r.update({'CREATED': 201, 'action': 'put %s' % action})
            self.send_json(r)
389
390
391
392
        else:
            action = r['path']
            self.send_json({'UNAUTHORIZED': 401, 'action': 'put %s' % action})
            self.terminate()
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430

    def _get(self, r):
        """Handle GET requests"""
        action = r.pop('path')
        if not self.accepted:
            self.send_json({'UNAUTHORIZED': 401, 'action': 'get %s' % action})
            self.terminate()
        else:
            data = {
                'settings': self.get_settings,
                'status': self.get_status,
            }[action]()
            data['action'] = 'get %s' % action
            self.send_json(data)

    def received_message(self, message):
        """Route requests to corresponding handling methods"""
        LOG.debug('recv: %s' % message)
        try:
            r = json.loads('%s' % message)
        except ValueError as ve:
            self.send_json({'BAD REQUEST': 400})
            LOG.error('JSON ERROR: %s' % ve)
            return
        try:
            method = r.pop('method')
            {
                'post': self._post,
                'put': self._put,
                'get': self._get
            }[method](r)
        except KeyError as ke:
            self.send_json({'BAD REQUEST': 400})
            LOG.error('KEY ERROR: %s' % ke)
        except Exception as e:
            self.send_json({'INTERNAL ERROR': 500})
            LOG.error('EXCEPTION: %s' % e)
            self.terminate()