protocol.py 16.1 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
from agkyra.syncer import (
28
    syncer, setup, pithos_client, localfs_client, messaging, utils)
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
        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()

52
53
        self.db.close()

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
    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
79
80
        WebSocketProtocol.session_db = self.session_db
        WebSocketProtocol.session_relation = self.session_relation
81
        server = make_server(
82
            '', 0,
83
84
85
86
            server_class=WSGIServer,
            handler_class=WebSocketWSGIRequestHandler,
            app=WebSocketWSGIApplication(handler_cls=WebSocketProtocol))
        server.initialize_websockets_manager()
87
        address = 'ws://%s:%s' % (server.server_name, server.server_port)
88
89
        self.server = server

90
91
92
93
94
95
96
97
        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)

98
99
    def start(self):
        """Start the helper server in a thread"""
100
101
        if getattr(self, 'server', None):
            Thread(target=self.server.serve_forever).start()
102
103
104

    def shutdown(self):
        """Shutdown the server (needs another thread) and join threads"""
105
106
107
108
        if getattr(self, 'server', None):
            t = Thread(target=self.server.shutdown)
            t.start()
            t.join()
109
110


111
112
113
114
class WebSocketProtocol(WebSocket):
    """Helper-side WebSocket protocol for communication with GUI:

    -- INTERRNAL HANDSAKE --
115
    GUI: {"method": "post", "ui_id": <GUI ID>}
116
    HELPER: {"ACCEPTED": 202, "method": "post"}" or
117
        "{"REJECTED": 401, "action": "post ui_id"}
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
152
153
154
155

    -- 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"}
156
157
158
159
160
    HELPER: {
        "can_sync": <boolean>,
        "progress": <int>,
        "paused": <boolean>,
        "action": "get status"} or
161
162
163
        {<ERROR>: <ERROR CODE>, "action": "get status"}
    """

164
    ui_id = None
165
    db, session_db, session_relation = None, None, None
166
167
168
169
    accepted = False
    settings = dict(
        token=None, url=None,
        container=None, directory=None,
170
        exclude=None)
171
172
    status = dict(
        progress=0, synced=0, unsynced=0, paused=True, can_sync=False)
173
174
    file_syncer = None
    cnf = AgkyraConfig()
175
    essentials = ('url', 'token', 'container', 'directory')
176

177
178
179
180
181
182
183
184
185
    def heartbeat(self):
        if not self.db:
            self.db = sqlite3.connect(self.session_db)
        self.db.execute('BEGIN')
        self.db.execute('UPDATE %s SET beat="%s" WHERE ui_id="%s"' % (
            self.session_relation, time.time(), self.ui_id))
        self.db.commit()
        time.sleep(2)

186
187
188
189
    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.
        """
190
        sync = self.cnf.get('global', 'default_sync')
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
        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'
212

213
214
215
216
    def _load_settings(self):
        LOG.debug('Start loading settings')
        sync = self._get_default_sync()
        cloud = self._get_sync_cloud(sync)
217
218

        try:
219
220
221
222
223
224
225
            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
226
227

        for option in ('container', 'directory', 'exclude'):
228
229
230
231
            try:
                self.settings[option] = self.cnf.get_sync(sync, option)
            except KeyError:
                LOG.debug('No %s is set' % option)
232

233
234
        LOG.debug('Finished loading settings')

235
    def _dump_settings(self):
236
        LOG.debug('Saving settings')
237
238
239
240
241
242
243
244
245
246
247
        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']
248
249
250
251
252
253
254
255
256

        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'])
257
        self.cnf.set_cloud(cloud, 'token', self.settings['token'] or '')
258
259
260
        self.cnf.set_sync(sync, 'cloud', cloud)

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

263
264
265
        self.cnf.write()
        LOG.debug('Settings saved')

266
267
268
269
270
    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])

271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
    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])

298
299
300
    def init_sync(self):
        """Initialize syncer"""
        sync = self._get_default_sync()
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315

        kwargs = dict(agkyra_path=AGKYRA_DIR)
        # Get SSL settings
        cloud = self._get_sync_cloud(sync)
        try:
            ignore_ssl = self.cnf.get_cloud(cloud, 'ignore_ssl') in ('on', )
            kwargs['ignore_ssl'] = ignore_ssl
        except KeyError:
            ignore_ssl = None
        if not ignore_ssl:
            try:
                kwargs['ca_certs'] = self.cnf.get_cloud(cloud, 'ca_certs')
            except KeyError:
                pass

316
317
318
        syncer_settings = setup.SyncerSettings(
            self.settings['url'], self.settings['token'],
            self.settings['container'], self.settings['directory'],
319
            **kwargs)
320
321
322
323
        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
324
        self.syncer.initiate_probe()
325

326
327
    # Syncer-related methods
    def get_status(self):
328
329
330
331
332
333
334
        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)
335
336
337
338
339
340
        return self.status

    def get_settings(self):
        return self.settings

    def set_settings(self, new_settings):
341
342
        # Prepare setting save
        could_sync = self.can_sync()
343
344
345
        was_active = False
        if could_sync and not self.syncer.paused:
            was_active = True
346
347
348
349
            self.pause_sync()
        must_reset_syncing = self._essentials_changed(new_settings)

        # save settings
350
351
352
        self.settings = new_settings
        self._dump_settings()

353
354
355
356
357
358
359
        # Restart
        if self.can_sync():
            if must_reset_syncing or not could_sync:
                self.init_sync()
            elif was_active:
                self.start_sync()

360
    def pause_sync(self):
Giorgos Korfiatis's avatar
Giorgos Korfiatis committed
361
        self.syncer.stop_decide()
362
363
        LOG.debug('Wait open syncs to complete')
        self.syncer.wait_sync_threads()
364
365

    def start_sync(self):
366
        self.syncer.start_decide()
367
368
369
370

    # WebSocket connection methods
    def opened(self):
        LOG.debug('Helper: connection established')
371
372
373
        self.heart = utils.StoppableThread()
        self.heart.run_body = self.heartbeat
        self.heart.start()
374
375

    def closed(self, *args):
376
377
378
379
380
381
382
383
384
        """Stop server heart, empty DB and exit"""
        LOG.debug('Stop protocol heart')
        self.heart.stop()
        LOG.debug('Remove session traces')
        self.db = sqlite3.connect(self.session_db)
        self.db.execute('BEGIN')
        self.db.execute('DELETE FROM %s' % self.session_relation)
        self.db.commit()
        self.db.close()
385
386
387
388
389
390
391
392
393
394
395
396
        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':
397
398
                if self.can_sync():
                    self.syncer.stop_all_daemons()
399
400
                    LOG.debug('Wait open syncs to complete')
                    self.syncer.wait_sync_threads()
401
402
403
404
405
406
407
                self.close()
                return
            {
                'start': self.start_sync,
                'pause': self.pause_sync
            }[action]()
            self.send_json({'OK': 200, 'action': 'post %s' % action})
408
        elif r['ui_id'] == self.ui_id:
409
            self._load_settings()
410
            self.accepted = True
411
            self.send_json({'ACCEPTED': 202, 'action': 'post ui_id'})
412
413
414
            if self.can_sync():
                self.init_sync()
                self.pause_sync()
415
        else:
416
            action = r.get('path', 'ui_id')
417
418
419
420
421
            self.send_json({'REJECTED': 401, 'action': 'post %s' % action})
            self.terminate()

    def _put(self, r):
        """Handle PUT requests"""
422
        if self.accepted:
423
424
425
426
427
            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)
428
429
430
431
        else:
            action = r['path']
            self.send_json({'UNAUTHORIZED': 401, 'action': 'put %s' % action})
            self.terminate()
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469

    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()