stream.where and stream.async_where

This commit is contained in:
Lex Berezhny 2018-05-25 09:54:01 -04:00
parent ece2db08da
commit ca6d5937ba
3 changed files with 24 additions and 3 deletions

View file

@ -105,7 +105,7 @@ class BaseLedger:
if address not in self.addresses:
self.addresses[address] = Address(self.coin_class.address_to_hash160(address))
self.addresses[address].add_transaction(transaction)
self.transactions.setdefault(hexlify(transaction.id), transaction)
self.transactions.setdefault(transaction.id, transaction)
self._on_transaction_controller.add(transaction)
def has_address(self, address):

View file

@ -1,7 +1,7 @@
import six
import logging
from typing import List
from collections import namedtuple
from binascii import hexlify
from torba.basecoin import BaseCoin
from torba.basescript import BaseInputScript, BaseOutputScript
@ -161,7 +161,7 @@ class BaseTransaction:
@property
def id(self):
if self._id is None:
self._id = self.hash[::-1]
self._id = hexlify(self.hash[::-1])
return self._id
@property

View file

@ -1,6 +1,10 @@
import six
from twisted.internet.defer import Deferred, DeferredLock, maybeDeferred, inlineCallbacks
from twisted.python.failure import Failure
if six.PY3:
import asyncio
def execute_serially(f):
_lock = DeferredLock()
@ -124,6 +128,23 @@ class Stream:
def listen(self, on_data, on_error=None, on_done=None):
return self._controller._listen(on_data, on_error, on_done)
def where(self, condition):
deferred = Deferred()
def where_test(value):
if condition(value):
self._cancel_and_callback(subscription, deferred, value)
subscription = self.listen(
where_test,
lambda error, traceback: self._cancel_and_error(subscription, deferred, error, traceback)
)
return deferred
def async_where(self, condition):
return self.where(condition).asFuture(asyncio.get_event_loop())
@property
def first(self):
deferred = Deferred()