first commit
This commit is contained in:
+124
@@ -0,0 +1,124 @@
|
||||
'''
|
||||
Bitcoin base58 encoding and decoding.
|
||||
|
||||
Based on https://bitcointalk.org/index.php?topic=1026.0 (public domain)
|
||||
'''
|
||||
import hashlib
|
||||
|
||||
|
||||
# for compatibility with following code...
|
||||
class SHA256(object):
|
||||
new = hashlib.sha256
|
||||
|
||||
|
||||
if str != bytes:
|
||||
# Python 3.x
|
||||
def ord(c):
|
||||
return c
|
||||
|
||||
def chr(n):
|
||||
return bytes((n,))
|
||||
|
||||
|
||||
__b58chars = '123456789ABCDEFGHJKLMNPQRSTUVWXYZabcdefghijkmnopqrstuvwxyz'
|
||||
__b58base = len(__b58chars)
|
||||
b58chars = __b58chars
|
||||
|
||||
|
||||
def b58encode(v):
|
||||
""" encode v, which is a string of bytes, to base58.
|
||||
"""
|
||||
long_value = 0
|
||||
for (i, c) in enumerate(v[::-1]):
|
||||
long_value += (256**i) * ord(c)
|
||||
|
||||
result = ''
|
||||
while long_value >= __b58base:
|
||||
div, mod = divmod(long_value, __b58base)
|
||||
result = __b58chars[mod] + result
|
||||
long_value = div
|
||||
result = __b58chars[long_value] + result
|
||||
|
||||
# Bitcoin does a little leading-zero-compression:
|
||||
# leading 0-bytes in the input become leading-1s
|
||||
nPad = 0
|
||||
for c in v:
|
||||
if c == '\0':
|
||||
nPad += 1
|
||||
else:
|
||||
break
|
||||
|
||||
return (__b58chars[0] * nPad) + result
|
||||
|
||||
|
||||
def b58decode(v, length=None):
|
||||
""" decode v into a string of len bytes
|
||||
"""
|
||||
long_value = 0
|
||||
for (i, c) in enumerate(v[::-1]):
|
||||
long_value += __b58chars.find(c) * (__b58base**i)
|
||||
|
||||
result = bytes()
|
||||
while long_value >= 256:
|
||||
div, mod = divmod(long_value, 256)
|
||||
result = chr(mod) + result
|
||||
long_value = div
|
||||
result = chr(long_value) + result
|
||||
|
||||
nPad = 0
|
||||
for c in v:
|
||||
if c == __b58chars[0]:
|
||||
nPad += 1
|
||||
else:
|
||||
break
|
||||
|
||||
result = chr(0) * nPad + result
|
||||
|
||||
if length is not None and len(result) != length:
|
||||
return None
|
||||
|
||||
return result
|
||||
|
||||
|
||||
def checksum(v):
|
||||
"""Return 32-bit checksum based on SHA256"""
|
||||
return SHA256.new(SHA256.new(v).digest()).digest()[0:4]
|
||||
|
||||
|
||||
def b58encode_chk(v):
|
||||
"""b58encode a string, with 32-bit checksum"""
|
||||
return b58encode(v + checksum(v))
|
||||
|
||||
|
||||
def b58decode_chk(v):
|
||||
"""decode a base58 string, check and remove checksum"""
|
||||
result = b58decode(v)
|
||||
|
||||
if result is None:
|
||||
return None
|
||||
|
||||
h3 = checksum(result[:-4])
|
||||
|
||||
if result[-4:] == checksum(result[:-4]):
|
||||
return result[:-4]
|
||||
else:
|
||||
return None
|
||||
|
||||
|
||||
def get_bcaddress_version(strAddress):
|
||||
""" Returns None if strAddress is invalid. Otherwise returns integer version of address. """
|
||||
addr = b58decode_chk(strAddress)
|
||||
if addr is None or len(addr) != 21:
|
||||
return None
|
||||
version = addr[0]
|
||||
return ord(version)
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
# Test case (from http://gitorious.org/bitcoin/python-base58.git)
|
||||
assert get_bcaddress_version('15VjRaDX9zpbA8LVnbrCAFzrVzN7ixHNsC') is 0
|
||||
_ohai = 'o hai'.encode('ascii')
|
||||
_tmp = b58encode(_ohai)
|
||||
assert _tmp == 'DYB3oMS'
|
||||
assert b58decode(_tmp, 5) == _ohai
|
||||
print("Tests passed")
|
||||
+110
@@ -0,0 +1,110 @@
|
||||
"""
|
||||
Set up defaults and read sentinel.conf
|
||||
"""
|
||||
import sys
|
||||
import os
|
||||
from sib_config import SibcoinConfig
|
||||
|
||||
default_sentinel_config = os.path.normpath(
|
||||
os.path.join(os.path.dirname(__file__), '../sentinel.conf')
|
||||
)
|
||||
|
||||
debug_enabled = os.environ.get('SENTINEL_DEBUG', False)
|
||||
|
||||
sentinel_config_file = os.environ.get('SENTINEL_CONFIG', default_sentinel_config)
|
||||
sentinel_cfg = SibcoinConfig.tokenize(sentinel_config_file)
|
||||
sentinel_version = "1.3.0"
|
||||
min_dashd_proto_version_with_sentinel_ping = 70208
|
||||
|
||||
|
||||
def get_dash_conf():
|
||||
if sys.platform == 'win32':
|
||||
dash_conf = os.path.join(os.getenv('APPDATA'), "DashCore/dash.conf")
|
||||
else:
|
||||
home = os.environ.get('HOME')
|
||||
|
||||
dash_conf = os.path.join(home, ".dashcore/dash.conf")
|
||||
if sys.platform == 'darwin':
|
||||
dash_conf = os.path.join(home, "Library/Application Support/DashCore/dash.conf")
|
||||
|
||||
dash_conf = sentinel_cfg.get('dash_conf', dash_conf)
|
||||
|
||||
return dash_conf
|
||||
|
||||
def get_sibcoin_conf():
|
||||
if sys.platform == 'win32':
|
||||
sibcoin_conf = os.path.join(os.getenv('APPDATA'), "Sibcoin/sibcoin.conf")
|
||||
else:
|
||||
home = os.environ.get('HOME')
|
||||
|
||||
sibcoin_conf = os.path.join(home, ".sibcoin/sibcoin.conf")
|
||||
if sys.platform == 'darwin':
|
||||
sibcoin_conf = os.path.join(home, "Library/Application Support/Sibcoin/sibcoin.conf")
|
||||
|
||||
sibcoin_conf = sentinel_cfg.get('sibcoin_conf', sibcoin_conf)
|
||||
|
||||
return sibcoin_conf
|
||||
|
||||
|
||||
def get_network():
|
||||
return sentinel_cfg.get('network', 'mainnet')
|
||||
|
||||
|
||||
def get_rpchost():
|
||||
return sentinel_cfg.get('rpchost', '127.0.0.1')
|
||||
|
||||
|
||||
def sqlite_test_db_name(sqlite_file_path):
|
||||
(root, ext) = os.path.splitext(sqlite_file_path)
|
||||
test_sqlite_file_path = root + '_test' + ext
|
||||
return test_sqlite_file_path
|
||||
|
||||
|
||||
def get_db_conn():
|
||||
import peewee
|
||||
env = os.environ.get('SENTINEL_ENV', 'production')
|
||||
|
||||
# default values should be used unless you need a different config for development
|
||||
db_host = sentinel_cfg.get('db_host', '127.0.0.1')
|
||||
db_port = sentinel_cfg.get('db_port', None)
|
||||
db_name = sentinel_cfg.get('db_name', 'sentinel')
|
||||
db_user = sentinel_cfg.get('db_user', 'sentinel')
|
||||
db_password = sentinel_cfg.get('db_password', 'sentinel')
|
||||
db_charset = sentinel_cfg.get('db_charset', 'utf8mb4')
|
||||
db_driver = sentinel_cfg.get('db_driver', 'sqlite')
|
||||
|
||||
if (env == 'test'):
|
||||
if db_driver == 'sqlite':
|
||||
db_name = sqlite_test_db_name(db_name)
|
||||
else:
|
||||
db_name = "%s_test" % db_name
|
||||
|
||||
peewee_drivers = {
|
||||
'mysql': peewee.MySQLDatabase,
|
||||
'postgres': peewee.PostgresqlDatabase,
|
||||
'sqlite': peewee.SqliteDatabase,
|
||||
}
|
||||
driver = peewee_drivers.get(db_driver)
|
||||
|
||||
dbpfn = 'passwd' if db_driver == 'mysql' else 'password'
|
||||
db_conn = {
|
||||
'host': db_host,
|
||||
'user': db_user,
|
||||
dbpfn: db_password,
|
||||
}
|
||||
if db_port:
|
||||
db_conn['port'] = int(db_port)
|
||||
|
||||
if driver == peewee.SqliteDatabase:
|
||||
db_conn = {}
|
||||
|
||||
db = driver(db_name, **db_conn)
|
||||
|
||||
return db
|
||||
|
||||
|
||||
#dash_conf = get_dash_conf()
|
||||
#sibcoin_conf = get_sibcoin_conf()
|
||||
#network = get_network()
|
||||
#rpc_host = get_rpchost()
|
||||
#db = get_db_conn()
|
||||
@@ -0,0 +1,4 @@
|
||||
# for constants which need to be accessed by various parts of Sentinel
|
||||
|
||||
# skip proposals on superblock creation if the SB isn't within the fudge window
|
||||
SUPERBLOCK_FUDGE_WINDOW = 60 * 60 * 2
|
||||
@@ -0,0 +1,59 @@
|
||||
import sys
|
||||
import os
|
||||
import io
|
||||
import re
|
||||
sys.path.append(os.path.join(os.path.dirname(__file__), '..'))
|
||||
sys.path.append(os.path.join(os.path.dirname(__file__), '..', 'lib'))
|
||||
from misc import printdbg
|
||||
|
||||
|
||||
class DashConfig():
|
||||
|
||||
@classmethod
|
||||
def slurp_config_file(self, filename):
|
||||
# read dash.conf config but skip commented lines
|
||||
f = io.open(filename)
|
||||
lines = []
|
||||
for line in f:
|
||||
if re.match(r'^\s*#', line):
|
||||
continue
|
||||
lines.append(line)
|
||||
f.close()
|
||||
|
||||
# data is dash.conf without commented lines
|
||||
data = ''.join(lines)
|
||||
|
||||
return data
|
||||
|
||||
@classmethod
|
||||
def get_rpc_creds(self, data, network='mainnet'):
|
||||
# get rpc info from dash.conf
|
||||
match = re.findall(r'rpc(user|password|port)=(.*?)$', data, re.MULTILINE)
|
||||
|
||||
# python >= 2.7
|
||||
creds = {key: value for (key, value) in match}
|
||||
|
||||
# standard Dash defaults...
|
||||
default_port = 9998 if (network == 'mainnet') else 19998
|
||||
|
||||
# use default port for network if not specified in dash.conf
|
||||
if not ('port' in creds):
|
||||
creds[u'port'] = default_port
|
||||
|
||||
# convert to an int if taken from dash.conf
|
||||
creds[u'port'] = int(creds[u'port'])
|
||||
|
||||
# return a dictionary with RPC credential key, value pairs
|
||||
return creds
|
||||
|
||||
@classmethod
|
||||
def tokenize(self, filename):
|
||||
tokens = {}
|
||||
try:
|
||||
data = self.slurp_config_file(filename)
|
||||
match = re.findall(r'(.*?)=(.*?)$', data, re.MULTILINE)
|
||||
tokens = {key: value for (key, value) in match}
|
||||
except IOError as e:
|
||||
printdbg("[warning] error reading config file: %s" % e)
|
||||
|
||||
return tokens
|
||||
+226
@@ -0,0 +1,226 @@
|
||||
"""
|
||||
dashd JSONRPC interface
|
||||
"""
|
||||
import sys
|
||||
import os
|
||||
sys.path.append(os.path.join(os.path.dirname(__file__), '..'))
|
||||
sys.path.append(os.path.join(os.path.dirname(__file__), '..', 'lib'))
|
||||
import config
|
||||
import base58
|
||||
from bitcoinrpc.authproxy import AuthServiceProxy, JSONRPCException
|
||||
from masternode import Masternode
|
||||
from decimal import Decimal
|
||||
import time
|
||||
|
||||
|
||||
class DashDaemon():
|
||||
def __init__(self, **kwargs):
|
||||
host = kwargs.get('host', '127.0.0.1')
|
||||
user = kwargs.get('user')
|
||||
password = kwargs.get('password')
|
||||
port = kwargs.get('port')
|
||||
|
||||
self.creds = (user, password, host, port)
|
||||
|
||||
# memoize calls to some dashd methods
|
||||
self.governance_info = None
|
||||
self.gobject_votes = {}
|
||||
|
||||
@property
|
||||
def rpc_connection(self):
|
||||
return AuthServiceProxy("http://{0}:{1}@{2}:{3}".format(*self.creds))
|
||||
|
||||
@classmethod
|
||||
def from_dash_conf(self, dash_dot_conf):
|
||||
from dash_config import DashConfig
|
||||
config_text = DashConfig.slurp_config_file(dash_dot_conf)
|
||||
creds = DashConfig.get_rpc_creds(config_text, config.network)
|
||||
|
||||
creds[u'host'] = config.rpc_host
|
||||
|
||||
return self(**creds)
|
||||
|
||||
def rpc_command(self, *params):
|
||||
return self.rpc_connection.__getattr__(params[0])(*params[1:])
|
||||
|
||||
# common RPC convenience methods
|
||||
|
||||
def get_masternodes(self):
|
||||
mnlist = self.rpc_command('masternodelist', 'full')
|
||||
return [Masternode(k, v) for (k, v) in mnlist.items()]
|
||||
|
||||
def get_current_masternode_vin(self):
|
||||
from dashlib import parse_masternode_status_vin
|
||||
|
||||
my_vin = None
|
||||
|
||||
try:
|
||||
status = self.rpc_command('masternode', 'status')
|
||||
mn_outpoint = status.get('outpoint') or status.get('vin')
|
||||
my_vin = parse_masternode_status_vin(mn_outpoint)
|
||||
except JSONRPCException as e:
|
||||
pass
|
||||
|
||||
return my_vin
|
||||
|
||||
def governance_quorum(self):
|
||||
# TODO: expensive call, so memoize this
|
||||
total_masternodes = self.rpc_command('masternode', 'count', 'enabled')
|
||||
min_quorum = self.govinfo['governanceminquorum']
|
||||
|
||||
# the minimum quorum is calculated based on the number of masternodes
|
||||
quorum = max(min_quorum, (total_masternodes // 10))
|
||||
return quorum
|
||||
|
||||
@property
|
||||
def govinfo(self):
|
||||
if (not self.governance_info):
|
||||
self.governance_info = self.rpc_command('getgovernanceinfo')
|
||||
return self.governance_info
|
||||
|
||||
# governance info convenience methods
|
||||
def superblockcycle(self):
|
||||
return self.govinfo['superblockcycle']
|
||||
|
||||
def last_superblock_height(self):
|
||||
height = self.rpc_command('getblockcount')
|
||||
cycle = self.superblockcycle()
|
||||
return cycle * (height // cycle)
|
||||
|
||||
def next_superblock_height(self):
|
||||
return self.last_superblock_height() + self.superblockcycle()
|
||||
|
||||
def is_masternode(self):
|
||||
return not (self.get_current_masternode_vin() is None)
|
||||
|
||||
def is_synced(self):
|
||||
mnsync_status = self.rpc_command('mnsync', 'status')
|
||||
synced = (mnsync_status['IsBlockchainSynced'] and
|
||||
mnsync_status['IsMasternodeListSynced'] and
|
||||
mnsync_status['IsWinnersListSynced'] and
|
||||
mnsync_status['IsSynced'] and
|
||||
not mnsync_status['IsFailed'])
|
||||
return synced
|
||||
|
||||
def current_block_hash(self):
|
||||
height = self.rpc_command('getblockcount')
|
||||
block_hash = self.rpc_command('getblockhash', height)
|
||||
return block_hash
|
||||
|
||||
def get_superblock_budget_allocation(self, height=None):
|
||||
if height is None:
|
||||
height = self.rpc_command('getblockcount')
|
||||
return Decimal(self.rpc_command('getsuperblockbudget', height))
|
||||
|
||||
def next_superblock_max_budget(self):
|
||||
cycle = self.superblockcycle()
|
||||
current_block_height = self.rpc_command('getblockcount')
|
||||
|
||||
last_superblock_height = (current_block_height // cycle) * cycle
|
||||
next_superblock_height = last_superblock_height + cycle
|
||||
|
||||
last_allocation = self.get_superblock_budget_allocation(last_superblock_height)
|
||||
next_allocation = self.get_superblock_budget_allocation(next_superblock_height)
|
||||
|
||||
next_superblock_max_budget = next_allocation
|
||||
|
||||
return next_superblock_max_budget
|
||||
|
||||
# "my" votes refers to the current running masternode
|
||||
# memoized on a per-run, per-object_hash basis
|
||||
def get_my_gobject_votes(self, object_hash):
|
||||
import dashlib
|
||||
if not self.gobject_votes.get(object_hash):
|
||||
my_vin = self.get_current_masternode_vin()
|
||||
# if we can't get MN vin from output of `masternode status`,
|
||||
# return an empty list
|
||||
if not my_vin:
|
||||
return []
|
||||
|
||||
(txid, vout_index) = my_vin.split('-')
|
||||
|
||||
cmd = ['gobject', 'getcurrentvotes', object_hash, txid, vout_index]
|
||||
raw_votes = self.rpc_command(*cmd)
|
||||
self.gobject_votes[object_hash] = dashlib.parse_raw_votes(raw_votes)
|
||||
|
||||
return self.gobject_votes[object_hash]
|
||||
|
||||
def is_govobj_maturity_phase(self):
|
||||
# 3-day period for govobj maturity
|
||||
maturity_phase_delta = 1662 # ~(60*24*3)/2.6
|
||||
if config.network == 'testnet':
|
||||
maturity_phase_delta = 24 # testnet
|
||||
|
||||
event_block_height = self.next_superblock_height()
|
||||
maturity_phase_start_block = event_block_height - maturity_phase_delta
|
||||
|
||||
current_height = self.rpc_command('getblockcount')
|
||||
event_block_height = self.next_superblock_height()
|
||||
|
||||
# print "current_height = %d" % current_height
|
||||
# print "event_block_height = %d" % event_block_height
|
||||
# print "maturity_phase_delta = %d" % maturity_phase_delta
|
||||
# print "maturity_phase_start_block = %d" % maturity_phase_start_block
|
||||
|
||||
return (current_height >= maturity_phase_start_block)
|
||||
|
||||
def we_are_the_winner(self):
|
||||
import dashlib
|
||||
# find the elected MN vin for superblock creation...
|
||||
current_block_hash = self.current_block_hash()
|
||||
mn_list = self.get_masternodes()
|
||||
winner = dashlib.elect_mn(block_hash=current_block_hash, mnlist=mn_list)
|
||||
my_vin = self.get_current_masternode_vin()
|
||||
|
||||
# print "current_block_hash: [%s]" % current_block_hash
|
||||
# print "MN election winner: [%s]" % winner
|
||||
# print "current masternode VIN: [%s]" % my_vin
|
||||
|
||||
return (winner == my_vin)
|
||||
|
||||
def estimate_block_time(self, height):
|
||||
import dashlib
|
||||
"""
|
||||
Called by block_height_to_epoch if block height is in the future.
|
||||
Call `block_height_to_epoch` instead of this method.
|
||||
|
||||
DO NOT CALL DIRECTLY if you don't want a "Oh Noes." exception.
|
||||
"""
|
||||
current_block_height = self.rpc_command('getblockcount')
|
||||
diff = height - current_block_height
|
||||
|
||||
if (diff < 0):
|
||||
raise Exception("Oh Noes.")
|
||||
|
||||
future_seconds = dashlib.blocks_to_seconds(diff)
|
||||
estimated_epoch = int(time.time() + future_seconds)
|
||||
|
||||
return estimated_epoch
|
||||
|
||||
def block_height_to_epoch(self, height):
|
||||
"""
|
||||
Get the epoch for a given block height, or estimate it if the block hasn't
|
||||
been mined yet. Call this method instead of `estimate_block_time`.
|
||||
"""
|
||||
epoch = -1
|
||||
|
||||
try:
|
||||
bhash = self.rpc_command('getblockhash', height)
|
||||
block = self.rpc_command('getblock', bhash)
|
||||
epoch = block['time']
|
||||
except JSONRPCException as e:
|
||||
if e.message == 'Block height out of range':
|
||||
epoch = self.estimate_block_time(height)
|
||||
else:
|
||||
print("error: %s" % e)
|
||||
raise e
|
||||
|
||||
return epoch
|
||||
|
||||
@property
|
||||
def has_sentinel_ping(self):
|
||||
getinfo = self.rpc_command('getinfo')
|
||||
return (getinfo['protocolversion'] >= config.min_dashd_proto_version_with_sentinel_ping)
|
||||
|
||||
def ping(self):
|
||||
self.rpc_command('sentinelping', config.sentinel_version)
|
||||
+304
@@ -0,0 +1,304 @@
|
||||
import sys
|
||||
import os
|
||||
sys.path.append(os.path.join(os.path.dirname(__file__), '..'))
|
||||
sys.path.append(os.path.join(os.path.dirname(__file__), '..', 'lib'))
|
||||
import base58
|
||||
import hashlib
|
||||
import re
|
||||
from decimal import Decimal
|
||||
import simplejson
|
||||
import binascii
|
||||
from misc import printdbg, epoch2str
|
||||
import time
|
||||
|
||||
|
||||
def is_valid_dash_address(address, network='mainnet'):
|
||||
raise RuntimeWarning('This method should not be used with sibcoin')
|
||||
# Only public key addresses are allowed
|
||||
# A valid address is a RIPEMD-160 hash which contains 20 bytes
|
||||
# Prior to base58 encoding 1 version byte is prepended and
|
||||
# 4 checksum bytes are appended so the total number of
|
||||
# base58 encoded bytes should be 25. This means the number of characters
|
||||
# in the encoding should be about 34 ( 25 * log2( 256 ) / log2( 58 ) ).
|
||||
dash_version = 140 if network == 'testnet' else 76
|
||||
|
||||
# Check length (This is important because the base58 library has problems
|
||||
# with long addresses (which are invalid anyway).
|
||||
if ((len(address) < 26) or (len(address) > 35)):
|
||||
return False
|
||||
|
||||
address_version = None
|
||||
|
||||
try:
|
||||
decoded = base58.b58decode_chk(address)
|
||||
address_version = ord(decoded[0:1])
|
||||
except:
|
||||
# rescue from exception, not a valid Dash address
|
||||
return False
|
||||
|
||||
if (address_version != dash_version):
|
||||
return False
|
||||
|
||||
return True
|
||||
|
||||
def is_valid_sibcoin_address(address, network='mainnet'):
|
||||
# Only public key addresses are allowed
|
||||
# A valid address is a RIPEMD-160 hash which contains 20 bytes
|
||||
# Prior to base58 encoding 1 version byte is prepended and
|
||||
# 4 checksum bytes are appended so the total number of
|
||||
# base58 encoded bytes should be 25. This means the number of characters
|
||||
# in the encoding should be about 34 ( 25 * log2( 256 ) / log2( 58 ) ).
|
||||
dash_version = 125 if network == 'testnet' else 63
|
||||
|
||||
# Check length (This is important because the base58 library has problems
|
||||
# with long addresses (which are invalid anyway).
|
||||
if ((len(address) < 26) or (len(address) > 35)):
|
||||
return False
|
||||
|
||||
address_version = None
|
||||
|
||||
try:
|
||||
decoded = base58.b58decode_chk(address)
|
||||
address_version = ord(decoded[0:1])
|
||||
except:
|
||||
# rescue from exception, not a valid Dash address
|
||||
return False
|
||||
|
||||
if (address_version != dash_version):
|
||||
return False
|
||||
|
||||
return True
|
||||
|
||||
def is_valid_address(address, network='mainnet'):
|
||||
return is_valid_sibcoin_address(address, network)
|
||||
|
||||
|
||||
def hashit(data):
|
||||
return int(hashlib.sha256(data.encode('utf-8')).hexdigest(), 16)
|
||||
|
||||
|
||||
# returns the masternode VIN of the elected winner
|
||||
def elect_mn(**kwargs):
|
||||
current_block_hash = kwargs['block_hash']
|
||||
mn_list = kwargs['mnlist']
|
||||
|
||||
# filter only enabled MNs
|
||||
enabled = [mn for mn in mn_list if mn.status == 'ENABLED']
|
||||
|
||||
block_hash_hash = hashit(current_block_hash)
|
||||
|
||||
candidates = []
|
||||
for mn in enabled:
|
||||
mn_vin_hash = hashit(mn.vin)
|
||||
diff = mn_vin_hash - block_hash_hash
|
||||
absdiff = abs(diff)
|
||||
candidates.append({'vin': mn.vin, 'diff': absdiff})
|
||||
|
||||
candidates.sort(key=lambda k: k['diff'])
|
||||
|
||||
try:
|
||||
winner = candidates[0]['vin']
|
||||
except:
|
||||
winner = None
|
||||
|
||||
return winner
|
||||
|
||||
|
||||
def parse_masternode_status_vin(status_vin_string):
|
||||
status_vin_string_regex = re.compile(r'CTxIn\(COutPoint\(([0-9a-zA-Z]+),\s*(\d+)\),')
|
||||
|
||||
m = status_vin_string_regex.match(status_vin_string)
|
||||
|
||||
# To Support additional format of string return from masternode status rpc.
|
||||
if m is None:
|
||||
status_output_string_regex = re.compile(r'([0-9a-zA-Z]+)-(\d+)')
|
||||
m = status_output_string_regex.match(status_vin_string)
|
||||
|
||||
txid = m.group(1)
|
||||
index = m.group(2)
|
||||
|
||||
vin = txid + '-' + index
|
||||
if (txid == '0000000000000000000000000000000000000000000000000000000000000000'):
|
||||
vin = None
|
||||
|
||||
return vin
|
||||
|
||||
|
||||
def create_superblock(proposals, event_block_height, budget_max, sb_epoch_time):
|
||||
from models import Superblock, GovernanceObject, Proposal
|
||||
from constants import SUPERBLOCK_FUDGE_WINDOW
|
||||
import copy
|
||||
|
||||
# don't create an empty superblock
|
||||
if (len(proposals) == 0):
|
||||
printdbg("No proposals, cannot create an empty superblock.")
|
||||
return None
|
||||
|
||||
budget_allocated = Decimal(0)
|
||||
fudge = SUPERBLOCK_FUDGE_WINDOW # fudge-factor to allow for slightly incorrect estimates
|
||||
|
||||
payments_list = []
|
||||
|
||||
for proposal in proposals:
|
||||
fmt_string = "name: %s, rank: %4d, hash: %s, amount: %s <= %s"
|
||||
|
||||
# skip proposals that are too expensive...
|
||||
if (budget_allocated + proposal.payment_amount) > budget_max:
|
||||
printdbg(
|
||||
fmt_string % (
|
||||
proposal.name,
|
||||
proposal.rank,
|
||||
proposal.object_hash,
|
||||
proposal.payment_amount,
|
||||
"skipped (blows the budget)",
|
||||
)
|
||||
)
|
||||
continue
|
||||
|
||||
# skip proposals if the SB isn't within the Proposal time window...
|
||||
window_start = proposal.start_epoch - fudge
|
||||
window_end = proposal.end_epoch + fudge
|
||||
|
||||
printdbg("\twindow_start: %s" % epoch2str(window_start))
|
||||
printdbg("\twindow_end: %s" % epoch2str(window_end))
|
||||
printdbg("\tsb_epoch_time: %s" % epoch2str(sb_epoch_time))
|
||||
|
||||
if (sb_epoch_time < window_start or sb_epoch_time > window_end):
|
||||
printdbg(
|
||||
fmt_string % (
|
||||
proposal.name,
|
||||
proposal.rank,
|
||||
proposal.object_hash,
|
||||
proposal.payment_amount,
|
||||
"skipped (SB time is outside of Proposal window)",
|
||||
)
|
||||
)
|
||||
continue
|
||||
|
||||
printdbg(
|
||||
fmt_string % (
|
||||
proposal.name,
|
||||
proposal.rank,
|
||||
proposal.object_hash,
|
||||
proposal.payment_amount,
|
||||
"adding",
|
||||
)
|
||||
)
|
||||
|
||||
payment = {
|
||||
'address': proposal.payment_address,
|
||||
'amount': "{0:.8f}".format(proposal.payment_amount),
|
||||
'proposal': "{}".format(proposal.object_hash)
|
||||
}
|
||||
|
||||
temp_payments_list = copy.deepcopy(payments_list)
|
||||
temp_payments_list.append(payment)
|
||||
|
||||
# calculate size of proposed Superblock
|
||||
sb_temp = Superblock(
|
||||
event_block_height=event_block_height,
|
||||
payment_addresses='|'.join([pd['address'] for pd in temp_payments_list]),
|
||||
payment_amounts='|'.join([pd['amount'] for pd in temp_payments_list]),
|
||||
proposal_hashes='|'.join([pd['proposal'] for pd in temp_payments_list])
|
||||
)
|
||||
proposed_sb_size = len(sb_temp.serialise())
|
||||
|
||||
# add proposal and keep track of total budget allocation
|
||||
budget_allocated += proposal.payment_amount
|
||||
payments_list.append(payment)
|
||||
|
||||
# don't create an empty superblock
|
||||
if not payments_list:
|
||||
printdbg("No proposals made the cut!")
|
||||
return None
|
||||
|
||||
# 'payments' now contains all the proposals for inclusion in the
|
||||
# Superblock, but needs to be sorted by proposal hash descending
|
||||
payments_list.sort(key=lambda k: k['proposal'], reverse=True)
|
||||
|
||||
sb = Superblock(
|
||||
event_block_height=event_block_height,
|
||||
payment_addresses='|'.join([pd['address'] for pd in payments_list]),
|
||||
payment_amounts='|'.join([pd['amount'] for pd in payments_list]),
|
||||
proposal_hashes='|'.join([pd['proposal'] for pd in payments_list]),
|
||||
)
|
||||
printdbg("generated superblock: %s" % sb.__dict__)
|
||||
|
||||
return sb
|
||||
|
||||
|
||||
# convenience
|
||||
def deserialise(hexdata):
|
||||
json = binascii.unhexlify(hexdata)
|
||||
obj = simplejson.loads(json, use_decimal=True)
|
||||
return obj
|
||||
|
||||
|
||||
def serialise(dikt):
|
||||
json = simplejson.dumps(dikt, sort_keys=True, use_decimal=True)
|
||||
hexdata = binascii.hexlify(json.encode('utf-8')).decode('utf-8')
|
||||
return hexdata
|
||||
|
||||
|
||||
def did_we_vote(output):
|
||||
from bitcoinrpc.authproxy import JSONRPCException
|
||||
|
||||
# sentinel
|
||||
voted = False
|
||||
err_msg = ''
|
||||
|
||||
try:
|
||||
detail = output.get('detail').get('sibcoin.conf')
|
||||
result = detail.get('result')
|
||||
if 'errorMessage' in detail:
|
||||
err_msg = detail.get('errorMessage')
|
||||
except JSONRPCException as e:
|
||||
result = 'failed'
|
||||
err_msg = e.message
|
||||
|
||||
# success, failed
|
||||
printdbg("result = [%s]" % result)
|
||||
if err_msg:
|
||||
printdbg("err_msg = [%s]" % err_msg)
|
||||
|
||||
voted = False
|
||||
if result == 'success':
|
||||
voted = True
|
||||
|
||||
# in case we spin up a new instance or server, but have already voted
|
||||
# on the network and network has recorded those votes
|
||||
m_old = re.match(r'^time between votes is too soon', err_msg)
|
||||
m_new = re.search(r'Masternode voting too often', err_msg, re.M)
|
||||
|
||||
if result == 'failed' and (m_old or m_new):
|
||||
printdbg("DEBUG: Voting too often, need to sync w/network")
|
||||
voted = False
|
||||
|
||||
return voted
|
||||
|
||||
|
||||
def parse_raw_votes(raw_votes):
|
||||
votes = []
|
||||
for v in list(raw_votes.values()):
|
||||
(outpoint, ntime, outcome, signal) = v.split(':')
|
||||
signal = signal.lower()
|
||||
outcome = outcome.lower()
|
||||
|
||||
mn_collateral_outpoint = parse_masternode_status_vin(outpoint)
|
||||
v = {
|
||||
'mn_collateral_outpoint': mn_collateral_outpoint,
|
||||
'signal': signal,
|
||||
'outcome': outcome,
|
||||
'ntime': ntime,
|
||||
}
|
||||
votes.append(v)
|
||||
|
||||
return votes
|
||||
|
||||
|
||||
def blocks_to_seconds(blocks):
|
||||
"""
|
||||
Return the estimated number of seconds which will transpire for a given
|
||||
number of blocks.
|
||||
"""
|
||||
return blocks * 2.62 * 60
|
||||
@@ -0,0 +1,33 @@
|
||||
import simplejson
|
||||
|
||||
|
||||
def valid_json(input):
|
||||
""" Return true/false depending on whether input is valid JSON """
|
||||
is_valid = False
|
||||
try:
|
||||
simplejson.loads(input)
|
||||
is_valid = True
|
||||
except:
|
||||
pass
|
||||
|
||||
return is_valid
|
||||
|
||||
|
||||
def extract_object(json_input):
|
||||
"""
|
||||
Given either an old-style or new-style Proposal JSON string, extract the
|
||||
actual object used (ignore old-style multi-dimensional array and unused
|
||||
string for object type)
|
||||
"""
|
||||
if not valid_json(json_input):
|
||||
raise Exception("Invalid JSON input.")
|
||||
|
||||
obj = simplejson.loads(json_input, use_decimal=True)
|
||||
|
||||
if (isinstance(obj, list) and
|
||||
isinstance(obj[0], list) and
|
||||
(isinstance(obj[0][0], str) or (isinstance(obj[0][0], unicode))) and
|
||||
isinstance(obj[0][1], dict)):
|
||||
obj = obj[0][1]
|
||||
|
||||
return obj
|
||||
@@ -0,0 +1,92 @@
|
||||
import os
|
||||
import sys
|
||||
sys.path.append(os.path.join(os.path.dirname(__file__), '..', 'lib'))
|
||||
import models
|
||||
from bitcoinrpc.authproxy import JSONRPCException
|
||||
import misc
|
||||
import re
|
||||
from misc import printdbg
|
||||
import time
|
||||
|
||||
|
||||
# mixin for GovObj composed classes like proposal and superblock, etc.
|
||||
class GovernanceClass(object):
|
||||
only_masternode_can_submit = False
|
||||
|
||||
# lazy
|
||||
@property
|
||||
def go(self):
|
||||
return self.governance_object
|
||||
|
||||
# pass thru to GovernanceObject#vote
|
||||
def vote(self, dashd, signal, outcome):
|
||||
return self.go.vote(dashd, signal, outcome)
|
||||
|
||||
# pass thru to GovernanceObject#voted_on
|
||||
def voted_on(self, **kwargs):
|
||||
return self.go.voted_on(**kwargs)
|
||||
|
||||
def vote_validity(self, dashd):
|
||||
if self.is_valid():
|
||||
printdbg("Voting valid! %s: %d" % (self.__class__.__name__, self.id))
|
||||
self.vote(dashd, models.VoteSignals.valid, models.VoteOutcomes.yes)
|
||||
else:
|
||||
printdbg("Voting INVALID! %s: %d" % (self.__class__.__name__, self.id))
|
||||
self.vote(dashd, models.VoteSignals.valid, models.VoteOutcomes.no)
|
||||
|
||||
def get_submit_command(self):
|
||||
obj_data = self.serialise()
|
||||
|
||||
# new objects won't have parent_hash, revision, etc...
|
||||
cmd = ['gobject', 'submit', '0', '1', str(int(time.time())), obj_data]
|
||||
|
||||
# some objects don't have a collateral tx to submit
|
||||
if not self.only_masternode_can_submit:
|
||||
cmd.append(go.object_fee_tx)
|
||||
|
||||
return cmd
|
||||
|
||||
def submit(self, dashd):
|
||||
# don't attempt to submit a superblock unless a masternode
|
||||
# note: will probably re-factor this, this has code smell
|
||||
if (self.only_masternode_can_submit and not dashd.is_masternode()):
|
||||
print("Not a masternode. Only masternodes may submit these objects")
|
||||
return
|
||||
|
||||
try:
|
||||
object_hash = dashd.rpc_command(*self.get_submit_command())
|
||||
printdbg("Submitted: [%s]" % object_hash)
|
||||
except JSONRPCException as e:
|
||||
print("Unable to submit: %s" % e.message)
|
||||
|
||||
def serialise(self):
|
||||
import binascii
|
||||
import simplejson
|
||||
|
||||
return binascii.hexlify(simplejson.dumps(self.get_dict(), sort_keys=True).encode('utf-8')).decode('utf-8')
|
||||
|
||||
@classmethod
|
||||
def serialisable_fields(self):
|
||||
# Python is so not very elegant...
|
||||
pk_column = self._meta.primary_key.db_column
|
||||
fk_columns = [fk.db_column for fk in self._meta.rel.values()]
|
||||
do_not_use = [pk_column]
|
||||
do_not_use.extend(fk_columns)
|
||||
do_not_use.append('object_hash')
|
||||
fields_to_serialise = list(self._meta.columns.keys())
|
||||
|
||||
for field in do_not_use:
|
||||
if field in fields_to_serialise:
|
||||
fields_to_serialise.remove(field)
|
||||
|
||||
return fields_to_serialise
|
||||
|
||||
def get_dict(self):
|
||||
dikt = {}
|
||||
|
||||
for field_name in self.serialisable_fields():
|
||||
dikt[field_name] = getattr(self, field_name)
|
||||
|
||||
dikt['type'] = getattr(self, 'govobj_type')
|
||||
|
||||
return dikt
|
||||
+147
@@ -0,0 +1,147 @@
|
||||
import sys
|
||||
import os
|
||||
sys.path.append(os.path.join(os.path.dirname(__file__), '..'))
|
||||
sys.path.append(os.path.join(os.path.dirname(__file__), '..', 'lib'))
|
||||
import argparse
|
||||
import config
|
||||
|
||||
def is_valid_python_version():
|
||||
version_valid = False
|
||||
|
||||
ver = sys.version_info
|
||||
if (2 == ver.major) and (7 <= ver.minor):
|
||||
version_valid = True
|
||||
|
||||
if (3 == ver.major) and (4 <= ver.minor):
|
||||
version_valid = True
|
||||
|
||||
return version_valid
|
||||
|
||||
|
||||
def python_short_ver_str():
|
||||
ver = sys.version_info
|
||||
return "%s.%s" % (ver.major, ver.minor)
|
||||
|
||||
|
||||
def are_deps_installed():
|
||||
installed = False
|
||||
|
||||
try:
|
||||
import peewee
|
||||
import bitcoinrpc.authproxy
|
||||
import simplejson
|
||||
installed = True
|
||||
except ImportError as e:
|
||||
print("[error]: Missing dependencies")
|
||||
|
||||
return installed
|
||||
|
||||
|
||||
def is_database_correctly_configured():
|
||||
import peewee
|
||||
import config
|
||||
|
||||
configured = False
|
||||
|
||||
cannot_connect_message = "Cannot connect to database. Please ensure database service is running and user access is properly configured in 'sentinel.conf'."
|
||||
|
||||
try:
|
||||
db = config.db
|
||||
db.connect()
|
||||
configured = True
|
||||
except (peewee.ImproperlyConfigured, peewee.OperationalError, ImportError) as e:
|
||||
print("[error]: %s" % e)
|
||||
print(cannot_connect_message)
|
||||
sys.exit(1)
|
||||
|
||||
return configured
|
||||
|
||||
|
||||
def has_sibcoin_conf():
|
||||
import config
|
||||
import io
|
||||
|
||||
valid_sibcoin_conf = False
|
||||
|
||||
# ensure dash_conf exists & readable
|
||||
#
|
||||
# if not, print a message stating that Dash Core must be installed and
|
||||
# configured, including JSONRPC access in dash.conf
|
||||
try:
|
||||
f = io.open(config.sibcoin_conf)
|
||||
valid_sibcoin_conf = True
|
||||
except IOError as e:
|
||||
print(e)
|
||||
|
||||
return valid_sibcoin_conf
|
||||
|
||||
def process_args():
|
||||
|
||||
parser = argparse.ArgumentParser()
|
||||
parser.add_argument('-b', '--bypass-scheduler',
|
||||
action='store_true',
|
||||
help='Bypass scheduler and sync/vote immediately',
|
||||
dest='bypass')
|
||||
parser.add_argument('-c', '--config',
|
||||
help='Path to sentinel.conf (default: ../sentinel.conf)',
|
||||
dest='config')
|
||||
parser.add_argument('-d', '--debug',
|
||||
action='store_true',
|
||||
help='Enable debug mode',
|
||||
dest='debug')
|
||||
args, unknown = parser.parse_known_args()
|
||||
|
||||
return args
|
||||
|
||||
initmodule = sys.modules[__name__]
|
||||
initmodule.options = False
|
||||
|
||||
# === begin main
|
||||
|
||||
|
||||
def main():
|
||||
|
||||
options = process_args()
|
||||
|
||||
if options.config:
|
||||
config.sentinel_config_file = options.config
|
||||
|
||||
# register a handler if SENTINEL_DEBUG is set
|
||||
if os.environ.get('SENTINEL_DEBUG', None) or options.debug:
|
||||
config.debug_enabled = True
|
||||
import logging
|
||||
logger = logging.getLogger('peewee')
|
||||
logger.setLevel(logging.DEBUG)
|
||||
logger.addHandler(logging.StreamHandler())
|
||||
|
||||
initmodule.options = options
|
||||
|
||||
from sib_config import SibcoinConfig
|
||||
config.sentinel_cfg = SibcoinConfig.tokenize(config.sentinel_config_file)
|
||||
|
||||
config.sibcoin_conf = config.get_sibcoin_conf()
|
||||
config.network = config.get_network()
|
||||
config.rpc_host = config.get_rpchost()
|
||||
config.db = config.get_db_conn()
|
||||
|
||||
install_instructions = "\tpip install -r requirements.txt"
|
||||
|
||||
if not is_valid_python_version():
|
||||
print("Python %s is not supported" % python_short_ver_str())
|
||||
sys.exit(1)
|
||||
|
||||
if not are_deps_installed():
|
||||
print("Please ensure all dependencies are installed:")
|
||||
print(install_instructions)
|
||||
sys.exit(1)
|
||||
|
||||
if not is_database_correctly_configured():
|
||||
print("Please ensure correct database configuration.")
|
||||
sys.exit(1)
|
||||
|
||||
if not has_sibcoin_conf():
|
||||
print("Sibcoin Core must be installed and configured, including JSONRPC access in sibcoin.conf")
|
||||
sys.exit(1)
|
||||
|
||||
|
||||
main()
|
||||
@@ -0,0 +1,41 @@
|
||||
# basically just parse & make it easier to access the MN data from the output of
|
||||
# "masternodelist full"
|
||||
|
||||
|
||||
class Masternode():
|
||||
def __init__(self, collateral, mnstring):
|
||||
(txid, vout_index) = self.parse_collateral_string(collateral)
|
||||
self.txid = txid
|
||||
self.vout_index = int(vout_index)
|
||||
|
||||
(status, protocol, address, ip_port, lastseen, activeseconds, lastpaid) = self.parse_mn_string(mnstring)
|
||||
self.status = status
|
||||
self.protocol = int(protocol)
|
||||
self.address = address
|
||||
|
||||
# TODO: break this out... take ipv6 into account
|
||||
self.ip_port = ip_port
|
||||
|
||||
self.lastseen = int(lastseen)
|
||||
self.activeseconds = int(activeseconds)
|
||||
self.lastpaid = int(lastpaid)
|
||||
|
||||
@classmethod
|
||||
def parse_collateral_string(self, collateral):
|
||||
(txid, index) = collateral.split('-')
|
||||
return (txid, index)
|
||||
|
||||
@classmethod
|
||||
def parse_mn_string(self, mn_full_out):
|
||||
# trim whitespace
|
||||
# mn_full_out = mn_full_out.strip()
|
||||
|
||||
(status, protocol, address, lastseen, activeseconds, lastpaid,
|
||||
lastpaidblock, ip_port) = mn_full_out.split()
|
||||
|
||||
# status protocol pubkey IP lastseen activeseconds lastpaid
|
||||
return (status, protocol, address, ip_port, lastseen, activeseconds, lastpaid)
|
||||
|
||||
@property
|
||||
def vin(self):
|
||||
return self.txid + '-' + str(self.vout_index)
|
||||
+52
@@ -0,0 +1,52 @@
|
||||
import time
|
||||
from datetime import datetime
|
||||
import re
|
||||
import sys
|
||||
import os
|
||||
import config
|
||||
|
||||
|
||||
def is_numeric(strin):
|
||||
import decimal
|
||||
|
||||
strin = str(strin)
|
||||
|
||||
# Decimal allows spaces in input, but we don't
|
||||
if strin.strip() != strin:
|
||||
return False
|
||||
try:
|
||||
value = decimal.Decimal(strin)
|
||||
except decimal.InvalidOperation as e:
|
||||
return False
|
||||
|
||||
return True
|
||||
|
||||
|
||||
def printdbg(str):
|
||||
ts = time.strftime('%Y-%m-%d %H:%M:%S', time.gmtime(now()))
|
||||
logstr = "{} {}".format(ts, str)
|
||||
if config.debug_enabled:
|
||||
print(logstr)
|
||||
|
||||
sys.stdout.flush()
|
||||
|
||||
|
||||
def is_hash(s):
|
||||
m = re.match('^[a-f0-9]{64}$', s)
|
||||
return m is not None
|
||||
|
||||
|
||||
def now():
|
||||
return int(time.time())
|
||||
|
||||
|
||||
def epoch2str(epoch):
|
||||
return datetime.utcfromtimestamp(epoch).strftime("%Y-%m-%d %H:%M:%S")
|
||||
|
||||
|
||||
class Bunch(object):
|
||||
def __init__(self, **kwargs):
|
||||
self.__dict__.update(kwargs)
|
||||
|
||||
def get(self, name):
|
||||
return self.__dict__.get(name, None)
|
||||
+768
@@ -0,0 +1,768 @@
|
||||
import sys
|
||||
import os
|
||||
sys.path.append(os.path.join(os.path.dirname(__file__), '..'))
|
||||
sys.path.append(os.path.join(os.path.dirname(__file__), '..', 'lib'))
|
||||
import init
|
||||
import time
|
||||
import datetime
|
||||
import re
|
||||
import simplejson
|
||||
from peewee import IntegerField, CharField, TextField, ForeignKeyField, DecimalField, DateTimeField
|
||||
import peewee
|
||||
import playhouse.signals
|
||||
import misc
|
||||
import dashd
|
||||
from misc import (printdbg, is_numeric)
|
||||
import config
|
||||
from bitcoinrpc.authproxy import JSONRPCException
|
||||
try:
|
||||
import urllib.parse as urlparse
|
||||
except ImportError:
|
||||
import urlparse
|
||||
|
||||
# our mixin
|
||||
from governance_class import GovernanceClass
|
||||
|
||||
db = config.db
|
||||
db.connect()
|
||||
|
||||
|
||||
# TODO: lookup table?
|
||||
DASHD_GOVOBJ_TYPES = {
|
||||
'proposal': 1,
|
||||
'superblock': 2,
|
||||
}
|
||||
GOVOBJ_TYPE_STRINGS = {
|
||||
1: 'proposal',
|
||||
2: 'trigger', # it should be trigger here, not superblock
|
||||
}
|
||||
|
||||
# schema version follows format 'YYYYMMDD-NUM'.
|
||||
#
|
||||
# YYYYMMDD is the 4-digit year, 2-digit month and 2-digit day the schema
|
||||
# changes were added.
|
||||
#
|
||||
# NUM is a numerical version of changes for that specific date. If the date
|
||||
# changes, the NUM resets to 1.
|
||||
SCHEMA_VERSION = '20170111-1'
|
||||
|
||||
# === models ===
|
||||
|
||||
|
||||
class BaseModel(playhouse.signals.Model):
|
||||
class Meta:
|
||||
database = db
|
||||
|
||||
@classmethod
|
||||
def is_database_connected(self):
|
||||
return not db.is_closed()
|
||||
|
||||
|
||||
class GovernanceObject(BaseModel):
|
||||
parent_id = IntegerField(default=0)
|
||||
object_creation_time = IntegerField(default=int(time.time()))
|
||||
object_hash = CharField(max_length=64)
|
||||
object_parent_hash = CharField(default='0')
|
||||
object_type = IntegerField(default=0)
|
||||
object_revision = IntegerField(default=1)
|
||||
object_fee_tx = CharField(default='')
|
||||
yes_count = IntegerField(default=0)
|
||||
no_count = IntegerField(default=0)
|
||||
abstain_count = IntegerField(default=0)
|
||||
absolute_yes_count = IntegerField(default=0)
|
||||
|
||||
class Meta:
|
||||
db_table = 'governance_objects'
|
||||
|
||||
# sync dashd gobject list with our local relational DB backend
|
||||
@classmethod
|
||||
def sync(self, dashd):
|
||||
golist = dashd.rpc_command('gobject', 'list')
|
||||
|
||||
# objects which are removed from the network should be removed from the DB
|
||||
try:
|
||||
for purged in self.purged_network_objects(list(golist.keys())):
|
||||
# SOMEDAY: possible archive step here
|
||||
purged.delete_instance(recursive=True, delete_nullable=True)
|
||||
except Exception as e:
|
||||
printdbg("Got an error while purging: %s" % e)
|
||||
|
||||
for item in golist.values():
|
||||
try:
|
||||
(go, subobj) = self.import_gobject_from_dashd(dashd, item)
|
||||
except Exception as e:
|
||||
printdbg("Got an error upon import: %s" % e)
|
||||
|
||||
@classmethod
|
||||
def purged_network_objects(self, network_object_hashes):
|
||||
query = self.select()
|
||||
if network_object_hashes:
|
||||
query = query.where(~(self.object_hash << network_object_hashes))
|
||||
return query
|
||||
|
||||
@classmethod
|
||||
def import_gobject_from_dashd(self, dashd, rec):
|
||||
import decimal
|
||||
import dashlib
|
||||
import binascii
|
||||
import gobject_json
|
||||
|
||||
object_hash = rec['Hash']
|
||||
|
||||
gobj_dict = {
|
||||
'object_hash': object_hash,
|
||||
'object_fee_tx': rec['CollateralHash'],
|
||||
'absolute_yes_count': rec['AbsoluteYesCount'],
|
||||
'abstain_count': rec['AbstainCount'],
|
||||
'yes_count': rec['YesCount'],
|
||||
'no_count': rec['NoCount'],
|
||||
}
|
||||
|
||||
# deserialise and extract object
|
||||
json_str = binascii.unhexlify(rec['DataHex']).decode('utf-8')
|
||||
dikt = gobject_json.extract_object(json_str)
|
||||
|
||||
subobj = None
|
||||
|
||||
type_class_map = {
|
||||
1: Proposal,
|
||||
2: Superblock,
|
||||
}
|
||||
subclass = type_class_map[dikt['type']]
|
||||
|
||||
# set object_type in govobj table
|
||||
gobj_dict['object_type'] = subclass.govobj_type
|
||||
|
||||
# exclude any invalid model data from dashd...
|
||||
valid_keys = subclass.serialisable_fields()
|
||||
subdikt = {k: dikt[k] for k in valid_keys if k in dikt}
|
||||
|
||||
# get/create, then sync vote counts from dashd, with every run
|
||||
govobj, created = self.get_or_create(object_hash=object_hash, defaults=gobj_dict)
|
||||
if created:
|
||||
printdbg("govobj created = %s" % created)
|
||||
count = govobj.update(**gobj_dict).where(self.id == govobj.id).execute()
|
||||
if count:
|
||||
printdbg("govobj updated = %d" % count)
|
||||
subdikt['governance_object'] = govobj
|
||||
|
||||
# get/create, then sync payment amounts, etc. from dashd - Dashd is the master
|
||||
try:
|
||||
newdikt = subdikt.copy()
|
||||
newdikt['object_hash'] = object_hash
|
||||
if subclass(**newdikt).is_valid() is False:
|
||||
govobj.vote_delete(dashd)
|
||||
return (govobj, None)
|
||||
|
||||
subobj, created = subclass.get_or_create(object_hash=object_hash, defaults=subdikt)
|
||||
|
||||
except Exception as e:
|
||||
# in this case, vote as delete, and log the vote in the DB
|
||||
printdbg("Got invalid object from dashd! %s" % e)
|
||||
govobj.vote_delete(dashd)
|
||||
return (govobj, None)
|
||||
|
||||
if created:
|
||||
printdbg("subobj created = %s" % created)
|
||||
count = subobj.update(**subdikt).where(subclass.id == subobj.id).execute()
|
||||
if count:
|
||||
printdbg("subobj updated = %d" % count)
|
||||
|
||||
# ATM, returns a tuple w/gov attributes and the govobj
|
||||
return (govobj, subobj)
|
||||
|
||||
def vote_delete(self, dashd):
|
||||
if not self.voted_on(signal=VoteSignals.delete, outcome=VoteOutcomes.yes):
|
||||
self.vote(dashd, VoteSignals.delete, VoteOutcomes.yes)
|
||||
return
|
||||
|
||||
def get_vote_command(self, signal, outcome):
|
||||
cmd = ['gobject', 'vote-conf', self.object_hash,
|
||||
signal.name, outcome.name]
|
||||
return cmd
|
||||
|
||||
def vote(self, dashd, signal, outcome):
|
||||
import dashlib
|
||||
|
||||
# At this point, will probably never reach here. But doesn't hurt to
|
||||
# have an extra check just in case objects get out of sync (people will
|
||||
# muck with the DB).
|
||||
if (self.object_hash == '0' or not misc.is_hash(self.object_hash)):
|
||||
printdbg("No governance object hash, nothing to vote on.")
|
||||
return
|
||||
|
||||
# have I already voted on this gobject with this particular signal and outcome?
|
||||
if self.voted_on(signal=signal):
|
||||
printdbg("Found a vote for this gobject/signal...")
|
||||
vote = self.votes.where(Vote.signal == signal)[0]
|
||||
|
||||
# if the outcome is the same, move on, nothing more to do
|
||||
if vote.outcome == outcome:
|
||||
# move on.
|
||||
printdbg("Already voted for this same gobject/signal/outcome, no need to re-vote.")
|
||||
return
|
||||
else:
|
||||
printdbg("Found a STALE vote for this gobject/signal, deleting so that we can re-vote.")
|
||||
vote.delete_instance()
|
||||
|
||||
else:
|
||||
printdbg("Haven't voted on this gobject/signal yet...")
|
||||
|
||||
# now ... vote!
|
||||
|
||||
vote_command = self.get_vote_command(signal, outcome)
|
||||
printdbg(' '.join(vote_command))
|
||||
output = dashd.rpc_command(*vote_command)
|
||||
|
||||
# extract vote output parsing to external lib
|
||||
voted = dashlib.did_we_vote(output)
|
||||
|
||||
if voted:
|
||||
printdbg('VOTE success, saving Vote object to database')
|
||||
Vote(governance_object=self, signal=signal, outcome=outcome,
|
||||
object_hash=self.object_hash).save()
|
||||
else:
|
||||
printdbg('VOTE failed, trying to sync with network vote')
|
||||
self.sync_network_vote(dashd, signal)
|
||||
|
||||
def sync_network_vote(self, dashd, signal):
|
||||
printdbg('\tSyncing network vote for object %s with signal %s' % (self.object_hash, signal.name))
|
||||
vote_info = dashd.get_my_gobject_votes(self.object_hash)
|
||||
for vdikt in vote_info:
|
||||
if vdikt['signal'] != signal.name:
|
||||
continue
|
||||
|
||||
# ensure valid outcome
|
||||
outcome = VoteOutcomes.get(vdikt['outcome'])
|
||||
if not outcome:
|
||||
continue
|
||||
|
||||
printdbg('\tFound a matching valid vote on the network, outcome = %s' % vdikt['outcome'])
|
||||
Vote(governance_object=self, signal=signal, outcome=outcome,
|
||||
object_hash=self.object_hash).save()
|
||||
|
||||
def voted_on(self, **kwargs):
|
||||
signal = kwargs.get('signal', None)
|
||||
outcome = kwargs.get('outcome', None)
|
||||
|
||||
query = self.votes
|
||||
|
||||
if signal:
|
||||
query = query.where(Vote.signal == signal)
|
||||
|
||||
if outcome:
|
||||
query = query.where(Vote.outcome == outcome)
|
||||
|
||||
count = query.count()
|
||||
return count
|
||||
|
||||
|
||||
class Setting(BaseModel):
|
||||
name = CharField(default='')
|
||||
value = CharField(default='')
|
||||
created_at = DateTimeField(default=datetime.datetime.utcnow())
|
||||
updated_at = DateTimeField(default=datetime.datetime.utcnow())
|
||||
|
||||
class Meta:
|
||||
db_table = 'settings'
|
||||
|
||||
|
||||
class Proposal(GovernanceClass, BaseModel):
|
||||
governance_object = ForeignKeyField(GovernanceObject, related_name='proposals', on_delete='CASCADE', on_update='CASCADE')
|
||||
name = CharField(default='', max_length=40)
|
||||
url = CharField(default='')
|
||||
start_epoch = IntegerField()
|
||||
end_epoch = IntegerField()
|
||||
payment_address = CharField(max_length=36)
|
||||
payment_amount = DecimalField(max_digits=16, decimal_places=8)
|
||||
object_hash = CharField(max_length=64)
|
||||
|
||||
# src/governance-validators.cpp
|
||||
MAX_DATA_SIZE = 512
|
||||
|
||||
govobj_type = DASHD_GOVOBJ_TYPES['proposal']
|
||||
|
||||
class Meta:
|
||||
db_table = 'proposals'
|
||||
|
||||
def is_valid(self):
|
||||
import dashlib
|
||||
|
||||
printdbg("In Proposal#is_valid, for Proposal: %s" % self.__dict__)
|
||||
|
||||
try:
|
||||
# proposal name exists and is not null/whitespace
|
||||
if (len(self.name.strip()) == 0):
|
||||
printdbg("\tInvalid Proposal name [%s], returning False" % self.name)
|
||||
return False
|
||||
|
||||
# proposal name is normalized (something like "[a-zA-Z0-9-_]+")
|
||||
if not re.match(r'^[-_a-zA-Z0-9]+$', self.name):
|
||||
printdbg("\tInvalid Proposal name [%s] (does not match regex), returning False" % self.name)
|
||||
return False
|
||||
|
||||
# end date < start date
|
||||
if (self.end_epoch <= self.start_epoch):
|
||||
printdbg("\tProposal end_epoch [%s] <= start_epoch [%s] , returning False" % (self.end_epoch, self.start_epoch))
|
||||
return False
|
||||
|
||||
# amount must be numeric
|
||||
if misc.is_numeric(self.payment_amount) is False:
|
||||
printdbg("\tProposal amount [%s] is not valid, returning False" % self.payment_amount)
|
||||
return False
|
||||
|
||||
# amount can't be negative or 0
|
||||
if (float(self.payment_amount) <= 0):
|
||||
printdbg("\tProposal amount [%s] is negative or zero, returning False" % self.payment_amount)
|
||||
return False
|
||||
|
||||
# payment address is valid base58 dash addr, non-multisig
|
||||
if not dashlib.is_valid_address(self.payment_address, config.network):
|
||||
printdbg("\tPayment address [%s] not a valid Dash address for network [%s], returning False" % (self.payment_address, config.network))
|
||||
return False
|
||||
|
||||
# URL
|
||||
if (len(self.url.strip()) < 4):
|
||||
printdbg("\tProposal URL [%s] too short, returning False" % self.url)
|
||||
return False
|
||||
|
||||
# proposal URL has any whitespace
|
||||
if (re.search(r'\s', self.url)):
|
||||
printdbg("\tProposal URL [%s] has whitespace, returning False" % self.name)
|
||||
return False
|
||||
|
||||
# Dash Core restricts proposals to 512 bytes max
|
||||
if len(self.serialise()) > (self.MAX_DATA_SIZE * 2):
|
||||
printdbg("\tProposal [%s] is too big, returning False" % self.name)
|
||||
return False
|
||||
|
||||
try:
|
||||
parsed = urlparse.urlparse(self.url)
|
||||
except Exception as e:
|
||||
printdbg("\tUnable to parse Proposal URL, marking invalid: %s" % e)
|
||||
return False
|
||||
|
||||
except Exception as e:
|
||||
printdbg("Unable to validate in Proposal#is_valid, marking invalid: %s" % e.message)
|
||||
return False
|
||||
|
||||
printdbg("Leaving Proposal#is_valid, Valid = True")
|
||||
return True
|
||||
|
||||
def is_expired(self, superblockcycle=None):
|
||||
from constants import SUPERBLOCK_FUDGE_WINDOW
|
||||
import dashlib
|
||||
|
||||
if not superblockcycle:
|
||||
raise Exception("Required field superblockcycle missing.")
|
||||
|
||||
printdbg("In Proposal#is_expired, for Proposal: %s" % self.__dict__)
|
||||
now = misc.now()
|
||||
printdbg("\tnow = %s" % now)
|
||||
|
||||
# half the SB cycle, converted to seconds
|
||||
# add the fudge_window in seconds, defined elsewhere in Sentinel
|
||||
expiration_window_seconds = int(
|
||||
(dashlib.blocks_to_seconds(superblockcycle) / 2) +
|
||||
SUPERBLOCK_FUDGE_WINDOW
|
||||
)
|
||||
printdbg("\texpiration_window_seconds = %s" % expiration_window_seconds)
|
||||
|
||||
# "fully expires" adds the expiration window to end time to ensure a
|
||||
# valid proposal isn't excluded from SB by cutting it too close
|
||||
fully_expires_at = self.end_epoch + expiration_window_seconds
|
||||
printdbg("\tfully_expires_at = %s" % fully_expires_at)
|
||||
|
||||
if (fully_expires_at < now):
|
||||
printdbg("\tProposal end_epoch [%s] < now [%s] , returning True" % (self.end_epoch, now))
|
||||
return True
|
||||
|
||||
printdbg("Leaving Proposal#is_expired, Expired = False")
|
||||
return False
|
||||
|
||||
@classmethod
|
||||
def approved_and_ranked(self, proposal_quorum, next_superblock_max_budget):
|
||||
# return all approved proposals, in order of descending vote count
|
||||
#
|
||||
# we need a secondary 'order by' in case of a tie on vote count, since
|
||||
# superblocks must be deterministic
|
||||
query = (self
|
||||
.select(self, GovernanceObject) # Note that we are selecting both models.
|
||||
.join(GovernanceObject)
|
||||
.where(GovernanceObject.absolute_yes_count > proposal_quorum)
|
||||
.order_by(GovernanceObject.absolute_yes_count.desc(), GovernanceObject.object_hash.desc())
|
||||
)
|
||||
|
||||
ranked = []
|
||||
for proposal in query:
|
||||
proposal.max_budget = next_superblock_max_budget
|
||||
if proposal.is_valid():
|
||||
ranked.append(proposal)
|
||||
|
||||
return ranked
|
||||
|
||||
@classmethod
|
||||
def expired(self, superblockcycle=None):
|
||||
if not superblockcycle:
|
||||
raise Exception("Required field superblockcycle missing.")
|
||||
|
||||
expired = []
|
||||
|
||||
for proposal in self.select():
|
||||
if proposal.is_expired(superblockcycle):
|
||||
expired.append(proposal)
|
||||
|
||||
return expired
|
||||
|
||||
@property
|
||||
def rank(self):
|
||||
rank = 0
|
||||
if self.governance_object:
|
||||
rank = self.governance_object.absolute_yes_count
|
||||
return rank
|
||||
|
||||
|
||||
class Superblock(BaseModel, GovernanceClass):
|
||||
governance_object = ForeignKeyField(GovernanceObject, related_name='superblocks', on_delete='CASCADE', on_update='CASCADE')
|
||||
event_block_height = IntegerField()
|
||||
payment_addresses = TextField()
|
||||
payment_amounts = TextField()
|
||||
proposal_hashes = TextField(default='')
|
||||
sb_hash = CharField()
|
||||
object_hash = CharField(max_length=64)
|
||||
|
||||
govobj_type = DASHD_GOVOBJ_TYPES['superblock']
|
||||
only_masternode_can_submit = True
|
||||
|
||||
class Meta:
|
||||
db_table = 'superblocks'
|
||||
|
||||
def is_valid(self):
|
||||
import dashlib
|
||||
import decimal
|
||||
|
||||
printdbg("In Superblock#is_valid, for SB: %s" % self.__dict__)
|
||||
|
||||
# it's a string from the DB...
|
||||
addresses = self.payment_addresses.split('|')
|
||||
for addr in addresses:
|
||||
if not dashlib.is_valid_address(addr, config.network):
|
||||
printdbg("\tInvalid address [%s], returning False" % addr)
|
||||
return False
|
||||
|
||||
amounts = self.payment_amounts.split('|')
|
||||
for amt in amounts:
|
||||
if not misc.is_numeric(amt):
|
||||
printdbg("\tAmount [%s] is not numeric, returning False" % amt)
|
||||
return False
|
||||
|
||||
# no negative or zero amounts allowed
|
||||
damt = decimal.Decimal(amt)
|
||||
if not damt > 0:
|
||||
printdbg("\tAmount [%s] is zero or negative, returning False" % damt)
|
||||
return False
|
||||
|
||||
# verify proposal hashes correctly formatted...
|
||||
if len(self.proposal_hashes) > 0:
|
||||
hashes = self.proposal_hashes.split('|')
|
||||
for object_hash in hashes:
|
||||
if not misc.is_hash(object_hash):
|
||||
printdbg("\tInvalid proposal hash [%s], returning False" % object_hash)
|
||||
return False
|
||||
|
||||
# ensure number of payment addresses matches number of payments
|
||||
if len(addresses) != len(amounts):
|
||||
printdbg("\tNumber of payment addresses [%s] != number of payment amounts [%s], returning False" % (len(addresses), len(amounts)))
|
||||
return False
|
||||
|
||||
printdbg("Leaving Superblock#is_valid, Valid = True")
|
||||
return True
|
||||
|
||||
def hash(self):
|
||||
import dashlib
|
||||
return dashlib.hashit(self.serialise())
|
||||
|
||||
def hex_hash(self):
|
||||
return "%x" % self.hash()
|
||||
|
||||
# workaround for now, b/c we must uniquely ID a superblock with the hash,
|
||||
# in case of differing superblocks
|
||||
#
|
||||
# this prevents sb_hash from being added to the serialised fields
|
||||
@classmethod
|
||||
def serialisable_fields(self):
|
||||
return [
|
||||
'event_block_height',
|
||||
'payment_addresses',
|
||||
'payment_amounts',
|
||||
'proposal_hashes'
|
||||
]
|
||||
|
||||
# has this masternode voted to fund *any* superblocks at the given
|
||||
# event_block_height?
|
||||
@classmethod
|
||||
def is_voted_funding(self, ebh):
|
||||
count = (self.select()
|
||||
.where(self.event_block_height == ebh)
|
||||
.join(GovernanceObject)
|
||||
.join(Vote)
|
||||
.join(Signal)
|
||||
.switch(Vote) # switch join query context back to Vote
|
||||
.join(Outcome)
|
||||
.where(Vote.signal == VoteSignals.funding)
|
||||
.where(Vote.outcome == VoteOutcomes.yes)
|
||||
.count())
|
||||
return count
|
||||
|
||||
@classmethod
|
||||
def latest(self):
|
||||
try:
|
||||
obj = self.select().order_by(self.event_block_height).desc().limit(1)[0]
|
||||
except IndexError as e:
|
||||
obj = None
|
||||
return obj
|
||||
|
||||
@classmethod
|
||||
def at_height(self, ebh):
|
||||
query = (self.select().where(self.event_block_height == ebh))
|
||||
return query
|
||||
|
||||
@classmethod
|
||||
def find_highest_deterministic(self, sb_hash):
|
||||
# highest block hash wins
|
||||
query = (self.select()
|
||||
.where(self.sb_hash == sb_hash)
|
||||
.order_by(self.object_hash.desc()))
|
||||
try:
|
||||
obj = query.limit(1)[0]
|
||||
except IndexError as e:
|
||||
obj = None
|
||||
return obj
|
||||
|
||||
|
||||
# ok, this is an awkward way to implement these...
|
||||
# "hook" into the Superblock model and run this code just before any save()
|
||||
from playhouse.signals import pre_save
|
||||
|
||||
|
||||
@pre_save(sender=Superblock)
|
||||
def on_save_handler(model_class, instance, created):
|
||||
instance.sb_hash = instance.hex_hash()
|
||||
|
||||
|
||||
class Signal(BaseModel):
|
||||
name = CharField(unique=True)
|
||||
created_at = DateTimeField(default=datetime.datetime.utcnow())
|
||||
updated_at = DateTimeField(default=datetime.datetime.utcnow())
|
||||
|
||||
class Meta:
|
||||
db_table = 'signals'
|
||||
|
||||
|
||||
class Outcome(BaseModel):
|
||||
name = CharField(unique=True)
|
||||
created_at = DateTimeField(default=datetime.datetime.utcnow())
|
||||
updated_at = DateTimeField(default=datetime.datetime.utcnow())
|
||||
|
||||
class Meta:
|
||||
db_table = 'outcomes'
|
||||
|
||||
|
||||
class Vote(BaseModel):
|
||||
governance_object = ForeignKeyField(GovernanceObject, related_name='votes', on_delete='CASCADE', on_update='CASCADE')
|
||||
signal = ForeignKeyField(Signal, related_name='votes', on_delete='CASCADE', on_update='CASCADE')
|
||||
outcome = ForeignKeyField(Outcome, related_name='votes', on_delete='CASCADE', on_update='CASCADE')
|
||||
voted_at = DateTimeField(default=datetime.datetime.utcnow())
|
||||
created_at = DateTimeField(default=datetime.datetime.utcnow())
|
||||
updated_at = DateTimeField(default=datetime.datetime.utcnow())
|
||||
object_hash = CharField(max_length=64)
|
||||
|
||||
class Meta:
|
||||
db_table = 'votes'
|
||||
|
||||
|
||||
class Transient(object):
|
||||
|
||||
def __init__(self, **kwargs):
|
||||
for key in ['created_at', 'timeout', 'value']:
|
||||
self.__setattr__(key, kwargs.get(key))
|
||||
|
||||
def is_expired(self):
|
||||
return (self.created_at + self.timeout) < misc.now()
|
||||
|
||||
@classmethod
|
||||
def deserialise(self, json):
|
||||
try:
|
||||
dikt = simplejson.loads(json)
|
||||
# a no-op, but this tells us what exception to expect
|
||||
except simplejson.scanner.JSONDecodeError as e:
|
||||
raise e
|
||||
|
||||
lizt = [dikt.get(key, None) for key in ['timeout', 'value']]
|
||||
lizt = list(set(lizt))
|
||||
if None in lizt:
|
||||
printdbg("Not all fields required for transient -- moving along.")
|
||||
raise Exception("Required fields not present for transient.")
|
||||
|
||||
return dikt
|
||||
|
||||
@classmethod
|
||||
def from_setting(self, setting):
|
||||
dikt = Transient.deserialise(setting.value)
|
||||
dikt['created_at'] = int((setting.created_at - datetime.datetime.utcfromtimestamp(0)).total_seconds())
|
||||
return Transient(**dikt)
|
||||
|
||||
@classmethod
|
||||
def cleanup(self):
|
||||
for s in Setting.select().where(Setting.name.startswith('__transient_')):
|
||||
try:
|
||||
t = Transient.from_setting(s)
|
||||
except:
|
||||
continue
|
||||
|
||||
if t.is_expired():
|
||||
s.delete_instance()
|
||||
|
||||
@classmethod
|
||||
def get(self, name):
|
||||
setting_name = "__transient_%s" % (name)
|
||||
|
||||
try:
|
||||
the_setting = Setting.get(Setting.name == setting_name)
|
||||
t = Transient.from_setting(the_setting)
|
||||
except Setting.DoesNotExist as e:
|
||||
return False
|
||||
|
||||
if t.is_expired():
|
||||
the_setting.delete_instance()
|
||||
return False
|
||||
else:
|
||||
return t.value
|
||||
|
||||
@classmethod
|
||||
def set(self, name, value, timeout):
|
||||
setting_name = "__transient_%s" % (name)
|
||||
setting_dikt = {
|
||||
'value': simplejson.dumps({
|
||||
'value': value,
|
||||
'timeout': timeout,
|
||||
}),
|
||||
}
|
||||
setting, created = Setting.get_or_create(name=setting_name, defaults=setting_dikt)
|
||||
return setting
|
||||
|
||||
@classmethod
|
||||
def delete(self, name):
|
||||
setting_name = "__transient_%s" % (name)
|
||||
try:
|
||||
s = Setting.get(Setting.name == setting_name)
|
||||
except Setting.DoesNotExist as e:
|
||||
return False
|
||||
return s.delete_instance()
|
||||
|
||||
# === /models ===
|
||||
|
||||
|
||||
def load_db_seeds():
|
||||
rows_created = 0
|
||||
|
||||
for name in ['funding', 'valid', 'delete']:
|
||||
(obj, created) = Signal.get_or_create(name=name)
|
||||
if created:
|
||||
rows_created = rows_created + 1
|
||||
|
||||
for name in ['yes', 'no', 'abstain']:
|
||||
(obj, created) = Outcome.get_or_create(name=name)
|
||||
if created:
|
||||
rows_created = rows_created + 1
|
||||
|
||||
return rows_created
|
||||
|
||||
|
||||
def db_models():
|
||||
""" Return a list of Sentinel DB models. """
|
||||
models = [
|
||||
GovernanceObject,
|
||||
Setting,
|
||||
Proposal,
|
||||
Superblock,
|
||||
Signal,
|
||||
Outcome,
|
||||
Vote
|
||||
]
|
||||
return models
|
||||
|
||||
|
||||
def check_db_sane():
|
||||
""" Ensure DB tables exist, create them if they don't. """
|
||||
check_db_schema_version()
|
||||
|
||||
missing_table_models = []
|
||||
|
||||
for model in db_models():
|
||||
if not getattr(model, 'table_exists')():
|
||||
missing_table_models.append(model)
|
||||
printdbg("[warning]: Table for %s (%s) doesn't exist in DB." % (model, model._meta.db_table))
|
||||
|
||||
if missing_table_models:
|
||||
printdbg("[warning]: Missing database tables. Auto-creating tables.")
|
||||
try:
|
||||
db.create_tables(missing_table_models, safe=True)
|
||||
except (peewee.InternalError, peewee.OperationalError, peewee.ProgrammingError) as e:
|
||||
print("[error] Could not create tables: %s" % e)
|
||||
|
||||
update_schema_version()
|
||||
purge_invalid_amounts()
|
||||
|
||||
|
||||
def check_db_schema_version():
|
||||
""" Ensure DB schema is correct version. Drop tables if not. """
|
||||
db_schema_version = None
|
||||
|
||||
try:
|
||||
db_schema_version = Setting.get(Setting.name == 'DB_SCHEMA_VERSION').value
|
||||
except (peewee.OperationalError, peewee.DoesNotExist, peewee.ProgrammingError) as e:
|
||||
printdbg("[info]: Can't get DB_SCHEMA_VERSION...")
|
||||
|
||||
printdbg("[info]: SCHEMA_VERSION (code) = [%s]" % SCHEMA_VERSION)
|
||||
printdbg("[info]: DB_SCHEMA_VERSION = [%s]" % db_schema_version)
|
||||
if (SCHEMA_VERSION != db_schema_version):
|
||||
printdbg("[info]: Schema version mis-match. Syncing tables.")
|
||||
try:
|
||||
existing_table_names = db.get_tables()
|
||||
existing_models = [m for m in db_models() if m._meta.db_table in existing_table_names]
|
||||
if (existing_models):
|
||||
printdbg("[info]: Dropping tables...")
|
||||
db.drop_tables(existing_models, safe=False, cascade=False)
|
||||
except (peewee.InternalError, peewee.OperationalError, peewee.ProgrammingError) as e:
|
||||
print("[error] Could not drop tables: %s" % e)
|
||||
|
||||
|
||||
def update_schema_version():
|
||||
schema_version_setting, created = Setting.get_or_create(name='DB_SCHEMA_VERSION', defaults={'value': SCHEMA_VERSION})
|
||||
if (schema_version_setting.value != SCHEMA_VERSION):
|
||||
schema_version_setting.save()
|
||||
return
|
||||
|
||||
|
||||
def purge_invalid_amounts():
|
||||
result_set = Proposal.select(
|
||||
Proposal.id,
|
||||
Proposal.governance_object
|
||||
).where(Proposal.payment_amount.contains(','))
|
||||
|
||||
for proposal in result_set:
|
||||
gobject = GovernanceObject.get(
|
||||
GovernanceObject.id == proposal.governance_object_id
|
||||
)
|
||||
printdbg("[info]: Pruning governance object w/invalid amount: %s" % gobject.object_hash)
|
||||
gobject.delete_instance(recursive=True, delete_nullable=True)
|
||||
|
||||
|
||||
# sanity checks...
|
||||
check_db_sane() # ensure tables exist
|
||||
load_db_seeds() # ensure seed data loaded
|
||||
|
||||
# convenience accessors
|
||||
VoteSignals = misc.Bunch(**{sig.name: sig for sig in Signal.select()})
|
||||
VoteOutcomes = misc.Bunch(**{out.name: out for out in Outcome.select()})
|
||||
@@ -0,0 +1,50 @@
|
||||
import sys
|
||||
import os
|
||||
sys.path.append(os.path.normpath(os.path.join(os.path.dirname(__file__), '../lib')))
|
||||
import init
|
||||
import misc
|
||||
from models import Transient
|
||||
from misc import printdbg
|
||||
import time
|
||||
import random
|
||||
|
||||
|
||||
class Scheduler(object):
|
||||
transient_key_scheduled = 'NEXT_SENTINEL_CHECK_AT'
|
||||
random_interval_max = 1200
|
||||
|
||||
@classmethod
|
||||
def is_run_time(self):
|
||||
next_run_time = Transient.get(self.transient_key_scheduled) or 0
|
||||
now = misc.now()
|
||||
|
||||
printdbg("current_time = %d" % now)
|
||||
printdbg("next_run_time = %d" % next_run_time)
|
||||
|
||||
return now >= next_run_time
|
||||
|
||||
@classmethod
|
||||
def clear_schedule(self):
|
||||
Transient.delete(self.transient_key_scheduled)
|
||||
|
||||
@classmethod
|
||||
def schedule_next_run(self, random_interval=None):
|
||||
if not random_interval:
|
||||
random_interval = self.random_interval_max
|
||||
|
||||
next_run_at = misc.now() + random.randint(1, random_interval)
|
||||
printdbg("scheduling next sentinel run for %d" % next_run_at)
|
||||
Transient.set(self.transient_key_scheduled, next_run_at,
|
||||
next_run_at)
|
||||
|
||||
@classmethod
|
||||
def delay(self, delay_in_seconds=None):
|
||||
if not delay_in_seconds:
|
||||
delay_in_seconds = random.randint(0, 60)
|
||||
|
||||
# do not delay longer than 60 seconds
|
||||
# in case an int > 60 given as argument
|
||||
delay_in_seconds = delay_in_seconds % 60
|
||||
|
||||
printdbg("Delay of [%d] seconds for cron minute offset" % delay_in_seconds)
|
||||
time.sleep(delay_in_seconds)
|
||||
@@ -0,0 +1,46 @@
|
||||
import sys
|
||||
import os
|
||||
import io
|
||||
import re
|
||||
sys.path.append(os.path.join(os.path.dirname(__file__), '..'))
|
||||
sys.path.append(os.path.join(os.path.dirname(__file__), '..', 'lib'))
|
||||
from misc import printdbg
|
||||
from dash_config import DashConfig
|
||||
|
||||
|
||||
class SibcoinConfig(DashConfig):
|
||||
|
||||
@classmethod
|
||||
def get_rpc_creds(self, data, network='mainnet'):
|
||||
# get rpc info from dash.conf
|
||||
match = re.findall(r'rpc(user|password|port)=(.*?)$', data, re.MULTILINE)
|
||||
|
||||
# python >= 2.7
|
||||
creds = {key: value for (key, value) in match}
|
||||
|
||||
# standard Dash defaults...
|
||||
default_port = 1944 if (network == 'mainnet') else 11944
|
||||
|
||||
# use default port for network if not specified in dash.conf
|
||||
if not ('port' in creds):
|
||||
creds[u'port'] = default_port
|
||||
|
||||
# convert to an int if taken from dash.conf
|
||||
creds[u'port'] = int(creds[u'port'])
|
||||
|
||||
# return a dictionary with RPC credential key, value pairs
|
||||
return creds
|
||||
|
||||
@classmethod
|
||||
def tokenize(self, filename, throw_exception=False):
|
||||
tokens = {}
|
||||
try:
|
||||
data = self.slurp_config_file(filename)
|
||||
match = re.findall(r'(.*?)=(.*?)$', data, re.MULTILINE)
|
||||
tokens = {key: value for (key, value) in match}
|
||||
except IOError as e:
|
||||
printdbg("[warning] error reading config file: %s" % e)
|
||||
if throw_exception:
|
||||
raise e
|
||||
|
||||
return tokens
|
||||
@@ -0,0 +1,33 @@
|
||||
"""
|
||||
dashd JSONRPC interface
|
||||
"""
|
||||
import sys
|
||||
import os
|
||||
sys.path.append(os.path.join(os.path.dirname(__file__), '..'))
|
||||
sys.path.append(os.path.join(os.path.dirname(__file__), '..', 'lib'))
|
||||
import config
|
||||
import base58
|
||||
from bitcoinrpc.authproxy import AuthServiceProxy, JSONRPCException
|
||||
from masternode import Masternode
|
||||
from decimal import Decimal
|
||||
import time
|
||||
from dashd import DashDaemon
|
||||
|
||||
|
||||
class SibcoinDaemon(DashDaemon):
|
||||
|
||||
@classmethod
|
||||
def from_sibcoin_conf(self, sibcoin_dot_conf):
|
||||
from sib_config import SibcoinConfig
|
||||
config_text = SibcoinConfig.slurp_config_file(sibcoin_dot_conf)
|
||||
creds = SibcoinConfig.get_rpc_creds(config_text, config.network)
|
||||
|
||||
creds[u'host'] = config.rpc_host
|
||||
|
||||
return self(**creds)
|
||||
|
||||
@classmethod
|
||||
def from_dash_conf(self, dash_dot_conf):
|
||||
raise RuntimeWarning('This method should not be used with sibcoin')
|
||||
|
||||
|
||||
Reference in New Issue
Block a user