-
Notifications
You must be signed in to change notification settings - Fork 21
Expand file tree
/
Copy pathdatabase.py
More file actions
102 lines (79 loc) · 2.87 KB
/
Copy pathdatabase.py
File metadata and controls
102 lines (79 loc) · 2.87 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
from weakref import WeakKeyDictionary
from nameko.extensions import DependencyProvider
from sqlalchemy import create_engine
from sqlalchemy.orm import Session as BaseSession
from sqlalchemy.orm import sessionmaker
DB_URIS_KEY = 'DB_URIS'
DB_ENGINE_OPTIONS_KEY = 'DB_ENGINE_OPTIONS'
DB_SESSION_OPTIONS_KEY = 'DB_SESSION_OPTIONS'
class Session(BaseSession):
def __init__(self, *args, **kwargs):
self.close_on_exit = kwargs.pop('close_on_exit', False)
super(Session, self).__init__(*args, **kwargs)
def __enter__(self):
return self
def __exit__(self, exc_type, exc_val, exc_tb):
try:
if exc_type:
self.rollback()
else:
try:
self.commit()
except Exception:
self.rollback()
raise
finally:
if self.close_on_exit:
self.close()
class DatabaseWrapper(object):
def __init__(self, Session):
self.Session = Session
self._worker_session = None
self._context_sessions = []
def get_session(self, close_on_exit=False):
session = self.Session(close_on_exit=close_on_exit)
self._context_sessions.append(session)
return session
@property
def session(self):
if self._worker_session is None:
self._worker_session = self.Session()
return self._worker_session
def close(self):
if self._worker_session:
self._worker_session.close()
for session in self._context_sessions:
session.close()
class Database(DependencyProvider):
def __init__(
self, declarative_base, session_options=None, engine_options=None
):
self.declarative_base = declarative_base
self.dbs = WeakKeyDictionary()
self.session_options = session_options or {}
self.engine_options = engine_options or {}
def setup(self):
service_name = self.container.service_name
declarative_base_name = self.declarative_base.__name__
uri_key = f'{service_name}:{declarative_base_name}'
db_uris = self.container.config[DB_URIS_KEY]
self.db_uri = db_uris[uri_key].format({
'service_name': service_name,
'declarative_base_name': declarative_base_name,
})
self.engine = create_engine(self.db_uri, **self.engine_options)
self.Session = sessionmaker(
bind=self.engine, class_=Session, **self.session_options)
def stop(self):
self.engine.dispose()
del self.engine
def kill(self):
self.engine.dispose()
del self.engine
def worker_teardown(self, worker_ctx):
db = self.dbs.pop(worker_ctx)
db.close()
def get_dependency(self, worker_ctx):
db = DatabaseWrapper(self.Session)
self.dbs[worker_ctx] = db
return db