Make the verifier a thread job instead of a thread
The verifier, like the synchronizer, now runs as part of the network proxy thread.
This commit is contained in:
parent
01491dd1d0
commit
b64c42b1eb
|
@ -66,7 +66,7 @@ class NetworkProxy(util.DaemonThread):
|
||||||
|
|
||||||
def run(self):
|
def run(self):
|
||||||
while self.is_running():
|
while self.is_running():
|
||||||
self.run_jobs() # Synchronizer, for now
|
self.run_jobs() # Synchronizer and Verifier
|
||||||
try:
|
try:
|
||||||
response = self.pipe.get()
|
response = self.pipe.get()
|
||||||
except util.timeout:
|
except util.timeout:
|
||||||
|
@ -185,9 +185,6 @@ class NetworkProxy(util.DaemonThread):
|
||||||
def get_interfaces(self):
|
def get_interfaces(self):
|
||||||
return self.interfaces
|
return self.interfaces
|
||||||
|
|
||||||
def get_header(self, height):
|
|
||||||
return self.synchronous_get([('network.get_header', [height])])[0]
|
|
||||||
|
|
||||||
def get_local_height(self):
|
def get_local_height(self):
|
||||||
return self.blockchain_height
|
return self.blockchain_height
|
||||||
|
|
||||||
|
|
|
@ -17,63 +17,52 @@
|
||||||
# along with this program. If not, see <http://www.gnu.org/licenses/>.
|
# along with this program. If not, see <http://www.gnu.org/licenses/>.
|
||||||
|
|
||||||
|
|
||||||
import threading
|
from util import ThreadJob
|
||||||
import Queue
|
from functools import partial
|
||||||
|
|
||||||
|
|
||||||
import util
|
|
||||||
from bitcoin import *
|
from bitcoin import *
|
||||||
|
|
||||||
|
|
||||||
class SPV(util.DaemonThread):
|
class SPV(ThreadJob):
|
||||||
""" Simple Payment Verification """
|
""" Simple Payment Verification """
|
||||||
|
|
||||||
def __init__(self, network, wallet):
|
def __init__(self, network, wallet):
|
||||||
util.DaemonThread.__init__(self)
|
|
||||||
self.wallet = wallet
|
self.wallet = wallet
|
||||||
self.network = network
|
self.network = network
|
||||||
self.merkle_roots = {} # hashed by me
|
self.merkle_roots = {} # hashed by me
|
||||||
self.queue = Queue.Queue()
|
self.requested_merkle = set()
|
||||||
|
|
||||||
def run(self):
|
def run(self):
|
||||||
requested_merkle = set()
|
unverified = self.wallet.get_unverified_txs()
|
||||||
while self.is_running():
|
for (tx_hash, tx_height) in unverified:
|
||||||
unverified = self.wallet.get_unverified_txs()
|
if tx_hash not in self.merkle_roots and tx_hash not in self.requested_merkle:
|
||||||
for (tx_hash, tx_height) in unverified:
|
request = ('blockchain.transaction.get_merkle',
|
||||||
if tx_hash not in self.merkle_roots and tx_hash not in requested_merkle:
|
[tx_hash, tx_height])
|
||||||
if self.network.send([ ('blockchain.transaction.get_merkle',[tx_hash, tx_height]) ], self.queue.put):
|
if self.network.send([request], self.merkle_response):
|
||||||
self.print_error('requested merkle', tx_hash)
|
self.print_error('requested merkle', tx_hash)
|
||||||
requested_merkle.add(tx_hash)
|
self.requested_merkle.add(tx_hash)
|
||||||
try:
|
|
||||||
r = self.queue.get(timeout=0.1)
|
|
||||||
except Queue.Empty:
|
|
||||||
continue
|
|
||||||
if not r:
|
|
||||||
continue
|
|
||||||
|
|
||||||
if r.get('error'):
|
def merkle_response(self, r):
|
||||||
self.print_error('Verifier received an error:', r)
|
if r.get('error'):
|
||||||
continue
|
self.print_error('received an error:', r)
|
||||||
|
return
|
||||||
|
|
||||||
# 3. handle response
|
params = r['params']
|
||||||
method = r['method']
|
result = r['result']
|
||||||
params = r['params']
|
|
||||||
result = r['result']
|
|
||||||
|
|
||||||
if method == 'blockchain.transaction.get_merkle':
|
# Get the header asynchronously - as a thread job we cannot block
|
||||||
tx_hash = params[0]
|
tx_hash = params[0]
|
||||||
self.verify_merkle(tx_hash, result)
|
request = ('network.get_header',[result.get('block_height')])
|
||||||
|
self.network.send([request], partial(self.verify, tx_hash, result))
|
||||||
|
|
||||||
self.print_error("stopped")
|
def verify(self, tx_hash, merkle, header):
|
||||||
|
'''Verify the hash of the server-provided merkle branch to a
|
||||||
|
transaction matches the merkle root of its block
|
||||||
def verify_merkle(self, tx_hash, result):
|
'''
|
||||||
tx_height = result.get('block_height')
|
tx_height = merkle.get('block_height')
|
||||||
pos = result.get('pos')
|
pos = merkle.get('pos')
|
||||||
merkle_root = self.hash_merkle_root(result['merkle'], tx_hash, pos)
|
merkle_root = self.hash_merkle_root(merkle['merkle'], tx_hash, pos)
|
||||||
header = self.network.get_header(tx_height)
|
header = header.get('result')
|
||||||
if not header: return
|
if not header or header.get('merkle_root') != merkle_root:
|
||||||
if header.get('merkle_root') != merkle_root:
|
|
||||||
self.print_error("merkle verification failed for", tx_hash)
|
self.print_error("merkle verification failed for", tx_hash)
|
||||||
return
|
return
|
||||||
|
|
||||||
|
|
|
@ -1086,9 +1086,7 @@ class Abstract_Wallet(object):
|
||||||
return True
|
return True
|
||||||
return False
|
return False
|
||||||
|
|
||||||
def set_verifier(self, verifier):
|
def prepare_for_verifier(self):
|
||||||
self.verifier = verifier
|
|
||||||
|
|
||||||
# review transactions that are in the history
|
# review transactions that are in the history
|
||||||
for addr, hist in self.history.items():
|
for addr, hist in self.history.items():
|
||||||
for tx_hash, tx_height in hist:
|
for tx_hash, tx_height in hist:
|
||||||
|
@ -1107,9 +1105,9 @@ class Abstract_Wallet(object):
|
||||||
from verifier import SPV
|
from verifier import SPV
|
||||||
self.network = network
|
self.network = network
|
||||||
if self.network is not None:
|
if self.network is not None:
|
||||||
|
self.prepare_for_verifier()
|
||||||
self.verifier = SPV(self.network, self)
|
self.verifier = SPV(self.network, self)
|
||||||
self.verifier.start()
|
network.add_job(self.verifier)
|
||||||
self.set_verifier(self.verifier)
|
|
||||||
self.synchronizer = Synchronizer(self, network)
|
self.synchronizer = Synchronizer(self, network)
|
||||||
network.add_job(self.synchronizer)
|
network.add_job(self.synchronizer)
|
||||||
else:
|
else:
|
||||||
|
@ -1118,9 +1116,10 @@ class Abstract_Wallet(object):
|
||||||
|
|
||||||
def stop_threads(self):
|
def stop_threads(self):
|
||||||
if self.network:
|
if self.network:
|
||||||
self.verifier.stop()
|
|
||||||
self.network.remove_job(self.synchronizer)
|
self.network.remove_job(self.synchronizer)
|
||||||
|
self.network.remove_job(self.verifier)
|
||||||
self.synchronizer = None
|
self.synchronizer = None
|
||||||
|
self.verifier = None
|
||||||
self.storage.put('stored_height', self.get_local_height(), True)
|
self.storage.put('stored_height', self.get_local_height(), True)
|
||||||
|
|
||||||
def restore(self, cb):
|
def restore(self, cb):
|
||||||
|
|
Loading…
Reference in New Issue