|
|
|
@ -26,6 +26,7 @@ import threading |
|
|
|
import copy |
|
|
|
import copy |
|
|
|
import json |
|
|
|
import json |
|
|
|
from typing import TYPE_CHECKING |
|
|
|
from typing import TYPE_CHECKING |
|
|
|
|
|
|
|
import jsonpatch |
|
|
|
|
|
|
|
|
|
|
|
from . import util |
|
|
|
from . import util |
|
|
|
from .util import WalletFileException, profiler |
|
|
|
from .util import WalletFileException, profiler |
|
|
|
@ -80,22 +81,35 @@ def stored_in(name, _type=dict): |
|
|
|
return decorator |
|
|
|
return decorator |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def key_path(path, key): |
|
|
|
|
|
|
|
def to_str(x): |
|
|
|
|
|
|
|
if isinstance(x, int): |
|
|
|
|
|
|
|
return str(int(x)) |
|
|
|
|
|
|
|
else: |
|
|
|
|
|
|
|
assert isinstance(x, str) |
|
|
|
|
|
|
|
return x |
|
|
|
|
|
|
|
return '/' + '/'.join([to_str(x) for x in path + [to_str(key)]]) |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class StoredObject: |
|
|
|
class StoredObject: |
|
|
|
|
|
|
|
|
|
|
|
db = None |
|
|
|
db = None |
|
|
|
|
|
|
|
path = None |
|
|
|
|
|
|
|
|
|
|
|
def __setattr__(self, key, value): |
|
|
|
def __setattr__(self, key, value): |
|
|
|
if self.db: |
|
|
|
if self.db and key not in ['path', 'db'] and not key.startswith('_'): |
|
|
|
self.db.set_modified(True) |
|
|
|
if value != getattr(self, key): |
|
|
|
|
|
|
|
self.db.add_patch({'op': 'replace', 'path': key_path(self.path, key), 'value': value}) |
|
|
|
object.__setattr__(self, key, value) |
|
|
|
object.__setattr__(self, key, value) |
|
|
|
|
|
|
|
|
|
|
|
def set_db(self, db): |
|
|
|
def set_db(self, db, path): |
|
|
|
self.db = db |
|
|
|
self.db = db |
|
|
|
|
|
|
|
self.path = path |
|
|
|
|
|
|
|
|
|
|
|
def to_json(self): |
|
|
|
def to_json(self): |
|
|
|
d = dict(vars(self)) |
|
|
|
d = dict(vars(self)) |
|
|
|
d.pop('db', None) |
|
|
|
d.pop('db', None) |
|
|
|
|
|
|
|
d.pop('path', None) |
|
|
|
# don't expose/store private stuff |
|
|
|
# don't expose/store private stuff |
|
|
|
d = {k: v for k, v in d.items() |
|
|
|
d = {k: v for k, v in d.items() |
|
|
|
if not k.startswith('_')} |
|
|
|
if not k.startswith('_')} |
|
|
|
@ -112,20 +126,22 @@ class StoredDict(dict): |
|
|
|
self.path = path |
|
|
|
self.path = path |
|
|
|
# recursively convert dicts to StoredDict |
|
|
|
# recursively convert dicts to StoredDict |
|
|
|
for k, v in list(data.items()): |
|
|
|
for k, v in list(data.items()): |
|
|
|
self.__setitem__(k, v) |
|
|
|
self.__setitem__(k, v, patch=False) |
|
|
|
|
|
|
|
|
|
|
|
@locked |
|
|
|
@locked |
|
|
|
def __setitem__(self, key, v): |
|
|
|
def __setitem__(self, key, v, patch=True): |
|
|
|
is_new = key not in self |
|
|
|
is_new = key not in self |
|
|
|
# early return to prevent unnecessary disk writes |
|
|
|
# early return to prevent unnecessary disk writes |
|
|
|
if not is_new and self[key] == v: |
|
|
|
if not is_new and patch: |
|
|
|
return |
|
|
|
if self.db and json.dumps(v, cls=self.db.encoder) == json.dumps(self[key], cls=self.db.encoder): |
|
|
|
|
|
|
|
return |
|
|
|
# recursively set db and path |
|
|
|
# recursively set db and path |
|
|
|
if isinstance(v, StoredDict): |
|
|
|
if isinstance(v, StoredDict): |
|
|
|
|
|
|
|
#assert v.db is None |
|
|
|
v.db = self.db |
|
|
|
v.db = self.db |
|
|
|
v.path = self.path + [key] |
|
|
|
v.path = self.path + [key] |
|
|
|
for k, vv in v.items(): |
|
|
|
for k, vv in v.items(): |
|
|
|
v[k] = vv |
|
|
|
v.__setitem__(k, vv, patch=False) |
|
|
|
# recursively convert dict to StoredDict. |
|
|
|
# recursively convert dict to StoredDict. |
|
|
|
# _convert_dict is called breadth-first |
|
|
|
# _convert_dict is called breadth-first |
|
|
|
elif isinstance(v, dict): |
|
|
|
elif isinstance(v, dict): |
|
|
|
@ -139,29 +155,57 @@ class StoredDict(dict): |
|
|
|
v = self.db._convert_value(self.path, key, v) |
|
|
|
v = self.db._convert_value(self.path, key, v) |
|
|
|
# set parent of StoredObject |
|
|
|
# set parent of StoredObject |
|
|
|
if isinstance(v, StoredObject): |
|
|
|
if isinstance(v, StoredObject): |
|
|
|
v.set_db(self.db) |
|
|
|
v.set_db(self.db, self.path + [key]) |
|
|
|
|
|
|
|
# convert lists |
|
|
|
|
|
|
|
if isinstance(v, list): |
|
|
|
|
|
|
|
v = StoredList(v, self.db, self.path + [key]) |
|
|
|
# set item |
|
|
|
# set item |
|
|
|
dict.__setitem__(self, key, v) |
|
|
|
dict.__setitem__(self, key, v) |
|
|
|
if self.db: |
|
|
|
if self.db and patch: |
|
|
|
self.db.set_modified(True) |
|
|
|
op = 'add' if is_new else 'replace' |
|
|
|
|
|
|
|
self.db.add_patch({'op': op, 'path': key_path(self.path, key), 'value': v}) |
|
|
|
|
|
|
|
|
|
|
|
@locked |
|
|
|
@locked |
|
|
|
def __delitem__(self, key): |
|
|
|
def __delitem__(self, key): |
|
|
|
dict.__delitem__(self, key) |
|
|
|
dict.__delitem__(self, key) |
|
|
|
if self.db: |
|
|
|
if self.db: |
|
|
|
self.db.set_modified(True) |
|
|
|
self.db.add_patch({'op': 'remove', 'path': key_path(self.path, key)}) |
|
|
|
|
|
|
|
|
|
|
|
@locked |
|
|
|
@locked |
|
|
|
def pop(self, key, v=_RaiseKeyError): |
|
|
|
def pop(self, key, v=_RaiseKeyError): |
|
|
|
if v is _RaiseKeyError: |
|
|
|
if key not in self: |
|
|
|
r = dict.pop(self, key) |
|
|
|
if v is _RaiseKeyError: |
|
|
|
else: |
|
|
|
raise KeyError(key) |
|
|
|
r = dict.pop(self, key, v) |
|
|
|
else: |
|
|
|
|
|
|
|
return v |
|
|
|
|
|
|
|
r = dict.pop(self, key) |
|
|
|
if self.db: |
|
|
|
if self.db: |
|
|
|
self.db.set_modified(True) |
|
|
|
self.db.add_patch({'op': 'remove', 'path': key_path(self.path, key)}) |
|
|
|
return r |
|
|
|
return r |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class StoredList(list): |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def __init__(self, data, db, path): |
|
|
|
|
|
|
|
list.__init__(self, data) |
|
|
|
|
|
|
|
self.db = db |
|
|
|
|
|
|
|
self.lock = self.db.lock if self.db else threading.RLock() |
|
|
|
|
|
|
|
self.path = path |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
@locked |
|
|
|
|
|
|
|
def append(self, item): |
|
|
|
|
|
|
|
n = len(self) |
|
|
|
|
|
|
|
list.append(self, item) |
|
|
|
|
|
|
|
if self.db: |
|
|
|
|
|
|
|
self.db.add_patch({'op': 'add', 'path': key_path(self.path, '%d'%n), 'value':item}) |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
@locked |
|
|
|
|
|
|
|
def remove(self, item): |
|
|
|
|
|
|
|
n = self.index(item) |
|
|
|
|
|
|
|
list.remove(self, item) |
|
|
|
|
|
|
|
if self.db: |
|
|
|
|
|
|
|
self.db.add_patch({'op': 'remove', 'path': key_path(self.path, '%d'%n)}) |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class JsonDB(Logger): |
|
|
|
class JsonDB(Logger): |
|
|
|
@ -171,34 +215,41 @@ class JsonDB(Logger): |
|
|
|
self.lock = threading.RLock() |
|
|
|
self.lock = threading.RLock() |
|
|
|
self.storage = storage |
|
|
|
self.storage = storage |
|
|
|
self.encoder = encoder |
|
|
|
self.encoder = encoder |
|
|
|
|
|
|
|
self.pending_changes = [] |
|
|
|
self._modified = False |
|
|
|
self._modified = False |
|
|
|
# load data |
|
|
|
# load data |
|
|
|
data = self.load_data(s) |
|
|
|
data = self.load_data(s) |
|
|
|
if upgrader: |
|
|
|
if upgrader: |
|
|
|
data, was_upgraded = upgrader(data) |
|
|
|
data, was_upgraded = upgrader(data) |
|
|
|
else: |
|
|
|
self._modified |= was_upgraded |
|
|
|
was_upgraded = False |
|
|
|
|
|
|
|
# convert to StoredDict |
|
|
|
# convert to StoredDict |
|
|
|
self.data = StoredDict(data, self, []) |
|
|
|
self.data = StoredDict(data, self, []) |
|
|
|
# note: self._modified may have been affected by StoredDict |
|
|
|
|
|
|
|
self._modified = was_upgraded |
|
|
|
|
|
|
|
# write file in case there was a db upgrade |
|
|
|
# write file in case there was a db upgrade |
|
|
|
if self.storage and self.storage.file_exists(): |
|
|
|
if self.storage and self.storage.file_exists(): |
|
|
|
self.write() |
|
|
|
self._write() |
|
|
|
|
|
|
|
|
|
|
|
def load_data(self, s:str) -> dict: |
|
|
|
def load_data(self, s:str) -> dict: |
|
|
|
""" overloaded in wallet_db """ |
|
|
|
""" overloaded in wallet_db """ |
|
|
|
if s == '': |
|
|
|
if s == '': |
|
|
|
return {} |
|
|
|
return {} |
|
|
|
try: |
|
|
|
try: |
|
|
|
data = json.loads(s) |
|
|
|
data = json.loads('[' + s + ']') |
|
|
|
|
|
|
|
data, patches = data[0], data[1:] |
|
|
|
except Exception: |
|
|
|
except Exception: |
|
|
|
if r := self.maybe_load_ast_data(s): |
|
|
|
if r := self.maybe_load_ast_data(s): |
|
|
|
data = r |
|
|
|
data, patches = r, [] |
|
|
|
|
|
|
|
elif r := self.maybe_load_incomplete_data(s): |
|
|
|
|
|
|
|
data, patches = r, [] |
|
|
|
else: |
|
|
|
else: |
|
|
|
raise WalletFileException("Cannot read wallet file. (parsing failed)") |
|
|
|
raise WalletFileException("Cannot read wallet file. (parsing failed)") |
|
|
|
if not isinstance(data, dict): |
|
|
|
if not isinstance(data, dict): |
|
|
|
raise WalletFileException("Malformed wallet file (not dict)") |
|
|
|
raise WalletFileException("Malformed wallet file (not dict)") |
|
|
|
|
|
|
|
if patches: |
|
|
|
|
|
|
|
# apply patches |
|
|
|
|
|
|
|
self.logger.info('found %d patches'%len(patches)) |
|
|
|
|
|
|
|
patch = jsonpatch.JsonPatch(patches) |
|
|
|
|
|
|
|
data = patch.apply(data) |
|
|
|
|
|
|
|
self.set_modified(True) |
|
|
|
return data |
|
|
|
return data |
|
|
|
|
|
|
|
|
|
|
|
def maybe_load_ast_data(self, s): |
|
|
|
def maybe_load_ast_data(self, s): |
|
|
|
@ -220,6 +271,21 @@ class JsonDB(Logger): |
|
|
|
data[key] = value |
|
|
|
data[key] = value |
|
|
|
return data |
|
|
|
return data |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def maybe_load_incomplete_data(self, s): |
|
|
|
|
|
|
|
n = s.count('{') - s.count('}') |
|
|
|
|
|
|
|
i = len(s) |
|
|
|
|
|
|
|
while n > 0 and i > 0: |
|
|
|
|
|
|
|
i = i - 1 |
|
|
|
|
|
|
|
if s[i] == '{': |
|
|
|
|
|
|
|
n = n - 1 |
|
|
|
|
|
|
|
if s[i] == '}': |
|
|
|
|
|
|
|
n = n + 1 |
|
|
|
|
|
|
|
if n == 0: |
|
|
|
|
|
|
|
s = s[0:i] |
|
|
|
|
|
|
|
assert s[-2:] == ',\n' |
|
|
|
|
|
|
|
self.logger.info('found incomplete data {s[i:]}') |
|
|
|
|
|
|
|
return self.load_data(s[0:-2]) |
|
|
|
|
|
|
|
|
|
|
|
def set_modified(self, b): |
|
|
|
def set_modified(self, b): |
|
|
|
with self.lock: |
|
|
|
with self.lock: |
|
|
|
self._modified = b |
|
|
|
self._modified = b |
|
|
|
@ -227,6 +293,11 @@ class JsonDB(Logger): |
|
|
|
def modified(self): |
|
|
|
def modified(self): |
|
|
|
return self._modified |
|
|
|
return self._modified |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
@locked |
|
|
|
|
|
|
|
def add_patch(self, patch): |
|
|
|
|
|
|
|
self.pending_changes.append(json.dumps(patch, cls=self.encoder)) |
|
|
|
|
|
|
|
self.set_modified(True) |
|
|
|
|
|
|
|
|
|
|
|
@locked |
|
|
|
@locked |
|
|
|
def get(self, key, default=None): |
|
|
|
def get(self, key, default=None): |
|
|
|
v = self.data.get(key) |
|
|
|
v = self.data.get(key) |
|
|
|
@ -259,6 +330,12 @@ class JsonDB(Logger): |
|
|
|
self.data[name] = {} |
|
|
|
self.data[name] = {} |
|
|
|
return self.data[name] |
|
|
|
return self.data[name] |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
@locked |
|
|
|
|
|
|
|
def get_stored_item(self, key, default) -> dict: |
|
|
|
|
|
|
|
if key not in self.data: |
|
|
|
|
|
|
|
self.data[key] = default |
|
|
|
|
|
|
|
return self.data[key] |
|
|
|
|
|
|
|
|
|
|
|
@locked |
|
|
|
@locked |
|
|
|
def dump(self, *, human_readable: bool = True) -> str: |
|
|
|
def dump(self, *, human_readable: bool = True) -> str: |
|
|
|
"""Serializes the DB as a string. |
|
|
|
"""Serializes the DB as a string. |
|
|
|
@ -302,10 +379,29 @@ class JsonDB(Logger): |
|
|
|
v = constructor(v) |
|
|
|
v = constructor(v) |
|
|
|
return v |
|
|
|
return v |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
@locked |
|
|
|
def write(self): |
|
|
|
def write(self): |
|
|
|
with self.lock: |
|
|
|
if not self.storage.file_exists()\ |
|
|
|
|
|
|
|
or self.storage.is_encrypted()\ |
|
|
|
|
|
|
|
or self.storage.needs_consolidation(): |
|
|
|
self._write() |
|
|
|
self._write() |
|
|
|
|
|
|
|
else: |
|
|
|
|
|
|
|
self._append_pending_changes() |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
@locked |
|
|
|
|
|
|
|
def _append_pending_changes(self): |
|
|
|
|
|
|
|
if threading.current_thread().daemon: |
|
|
|
|
|
|
|
self.logger.warning('daemon thread cannot write db') |
|
|
|
|
|
|
|
return |
|
|
|
|
|
|
|
if not self.pending_changes: |
|
|
|
|
|
|
|
self.logger.info('no pending changes') |
|
|
|
|
|
|
|
return |
|
|
|
|
|
|
|
self.logger.info(f'appending {len(self.pending_changes)} pending changes') |
|
|
|
|
|
|
|
s = ''.join([',\n' + x for x in self.pending_changes]) |
|
|
|
|
|
|
|
self.storage.append(s) |
|
|
|
|
|
|
|
self.pending_changes = [] |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
@locked |
|
|
|
@profiler |
|
|
|
@profiler |
|
|
|
def _write(self): |
|
|
|
def _write(self): |
|
|
|
if threading.current_thread().daemon: |
|
|
|
if threading.current_thread().daemon: |
|
|
|
@ -315,4 +411,5 @@ class JsonDB(Logger): |
|
|
|
return |
|
|
|
return |
|
|
|
json_str = self.dump(human_readable=not self.storage.is_encrypted()) |
|
|
|
json_str = self.dump(human_readable=not self.storage.is_encrypted()) |
|
|
|
self.storage.write(json_str) |
|
|
|
self.storage.write(json_str) |
|
|
|
|
|
|
|
self.pending_changes = [] |
|
|
|
self.set_modified(False) |
|
|
|
self.set_modified(False) |
|
|
|
|