-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathremote_sqlitedict.py
More file actions
110 lines (83 loc) · 3.3 KB
/
Copy pathremote_sqlitedict.py
File metadata and controls
110 lines (83 loc) · 3.3 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
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
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
import json
import os
import sys
import rpyc
from rpyc.core.protocol import PingError
from rpyc.utils.server import ThreadedServer
from sqlitedict import SqliteDict
# prevent JSON serialization issues by specifying and encoder
# see: https://github.com/tomerfiliba/rpyc/issues/393#issuecomment-662901702
def json_dumps(obj):
return json.dumps(obj, indent=0)
DEF_PORT = 18753
PING_TIMEOUT = 1
class RemoteSQLiteDictConnector(object):
def __init__(self, host, port, db_name, autocommit):
self._host = host
self._port = port
self._db_name = db_name
self._autocommit = autocommit
self._connection = None
def __enter__(self):
make_connection = False
if self._connection is not None:
try:
self._connection.ping(timeout=PING_TIMEOUT)
except (PingError, EOFError):
make_connection = True
else:
make_connection = True
if make_connection:
self._connection = rpyc.connect(self._host, self._port, config={
'allow_public_attrs': True,
'allow_all_attrs': True,
})
return self._connection.root.proxy_sqlitedict(
self._db_name, autocommit=self._autocommit)
def __exit__(self, exc_type, exc_val, exc_tb):
self._connection.close()
self._connection = None
def get_sqlitedict_connector(host, port, db_name, autocommit=False):
return RemoteSQLiteDictConnector(host, port, db_name, autocommit)
def start_server(port, db_root, single_db):
# define new service with DB_ROOT set based on the parameter
class SQLiteDictService(rpyc.Service):
DB_ROOT = db_root
def __init__(self):
self._instance = None
def exposed_proxy_sqlitedict(self, db_name, **kwargs):
if single_db:
db_path = os.path.join(self.DB_ROOT, '_root.sqlite')
else:
db_path = os.path.join(self.DB_ROOT, db_name + '.sqlite')
self._instance = SqliteDict(
db_path, tablename=db_name, encode=json_dumps,
decode=json.loads, **kwargs)
return self._instance
def on_disconnect(self, conn):
if self._instance is not None:
self._instance.close()
ThreadedServer(
SQLiteDictService, port=port,
protocol_config={
'allow_public_attrs': True,
'allow_all_attrs': True,
}
).start()
if __name__ == '__main__':
import argparse
parser = argparse.ArgumentParser()
parser.add_argument('--single-db', '-s', action='store_true',
dest='single_db',
help='Save all data in a single database file')
parser.add_argument('--directory', '-d', default=os.getcwd(),
help='Specify alternative directory '
'[default:current directory]')
parser.add_argument('port', action='store',
default=DEF_PORT, type=int,
nargs='?',
help=f'Specify alternate port [default: {DEF_PORT}]')
args = parser.parse_args()
sys.stdout.write(f'Server starting on port {args.port}...')
sys.stdout.flush()
start_server(args.port, args.directory, args.single_db)