mirror of
https://github.com/cypherstack/stack_wallet.git
synced 2025-01-25 03:35:59 +00:00
INCOMPLETE: WIP use streams instead of change notifier for electrumx socket subscriptions
This commit is contained in:
parent
1b81af1e7e
commit
2be13a89c5
2 changed files with 25 additions and 24 deletions
|
@ -12,17 +12,16 @@ import 'dart:async';
|
||||||
import 'dart:convert';
|
import 'dart:convert';
|
||||||
import 'dart:io';
|
import 'dart:io';
|
||||||
|
|
||||||
import 'package:flutter/foundation.dart';
|
|
||||||
import 'package:stackwallet/electrumx_rpc/electrumx_client.dart';
|
import 'package:stackwallet/electrumx_rpc/electrumx_client.dart';
|
||||||
import 'package:stackwallet/utilities/logger.dart';
|
import 'package:stackwallet/utilities/logger.dart';
|
||||||
|
|
||||||
class ElectrumXSubscription with ChangeNotifier {
|
class ElectrumXSubscription {
|
||||||
dynamic _response;
|
final StreamController<dynamic> _controller =
|
||||||
dynamic get response => _response;
|
StreamController(); // TODO controller params
|
||||||
set response(dynamic newData) {
|
|
||||||
_response = newData;
|
Stream<dynamic> get responseStream => _controller.stream;
|
||||||
notifyListeners();
|
|
||||||
}
|
void addToStream(dynamic data) => _controller.add(data);
|
||||||
}
|
}
|
||||||
|
|
||||||
class SocketTask {
|
class SocketTask {
|
||||||
|
@ -192,13 +191,13 @@ class SubscribableElectrumXClient {
|
||||||
final scripthash = params.first as String;
|
final scripthash = params.first as String;
|
||||||
final taskId = "blockchain.scripthash.subscribe:$scripthash";
|
final taskId = "blockchain.scripthash.subscribe:$scripthash";
|
||||||
|
|
||||||
_tasks[taskId]?.subscription?.response = params.last;
|
_tasks[taskId]?.subscription?.addToStream(params.last);
|
||||||
break;
|
break;
|
||||||
case "blockchain.headers.subscribe":
|
case "blockchain.headers.subscribe":
|
||||||
final params = response["params"];
|
final params = response["params"];
|
||||||
const taskId = "blockchain.headers.subscribe";
|
const taskId = "blockchain.headers.subscribe";
|
||||||
|
|
||||||
_tasks[taskId]?.subscription?.response = params.first;
|
_tasks[taskId]?.subscription?.addToStream(params.first);
|
||||||
break;
|
break;
|
||||||
default:
|
default:
|
||||||
break;
|
break;
|
||||||
|
@ -230,7 +229,7 @@ class SubscribableElectrumXClient {
|
||||||
if (!(_tasks[id]?.isSubscription ?? false)) {
|
if (!(_tasks[id]?.isSubscription ?? false)) {
|
||||||
_tasks.remove(id);
|
_tasks.remove(id);
|
||||||
} else {
|
} else {
|
||||||
_tasks[id]?.subscription?.response = data;
|
_tasks[id]?.subscription?.addToStream(data);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
@ -35,7 +35,8 @@ mixin ElectrumXInterface<T extends Bip39HDCurrency> on Bip39HDWallet<T> {
|
||||||
|
|
||||||
int? get maximumFeerate => null;
|
int? get maximumFeerate => null;
|
||||||
|
|
||||||
int? latestHeight;
|
int? _latestHeight;
|
||||||
|
StreamSubscription<dynamic>? _heightSubscription;
|
||||||
|
|
||||||
static const _kServerBatchCutoffVersion = [1, 6];
|
static const _kServerBatchCutoffVersion = [1, 6];
|
||||||
List<int>? _serverVersion;
|
List<int>? _serverVersion;
|
||||||
|
@ -800,14 +801,10 @@ mixin ElectrumXInterface<T extends Bip39HDCurrency> on Bip39HDWallet<T> {
|
||||||
final Completer<int> completer = Completer<int>();
|
final Completer<int> completer = Completer<int>();
|
||||||
|
|
||||||
try {
|
try {
|
||||||
// Subscribe to block headers.
|
// Don't set a stream subscription if one already exists.
|
||||||
final subscription =
|
if (_heightSubscription != null) {
|
||||||
subscribableElectrumXClient.subscribeToBlockHeaders();
|
if (_latestHeight != null) {
|
||||||
|
return _latestHeight!;
|
||||||
// Don't add a listener if one already exists.
|
|
||||||
if (subscription.hasListeners) {
|
|
||||||
if (latestHeight != null) {
|
|
||||||
return latestHeight!;
|
|
||||||
} else {
|
} else {
|
||||||
// Wait for first response.
|
// Wait for first response.
|
||||||
return completer.future;
|
return completer.future;
|
||||||
|
@ -817,15 +814,20 @@ mixin ElectrumXInterface<T extends Bip39HDCurrency> on Bip39HDWallet<T> {
|
||||||
// Make sure we only complete once.
|
// Make sure we only complete once.
|
||||||
bool isFirstResponse = true;
|
bool isFirstResponse = true;
|
||||||
|
|
||||||
// Add listener.
|
// Subscribe to block headers.
|
||||||
subscription.addListener(() {
|
final subscription =
|
||||||
final response = subscription.response;
|
subscribableElectrumXClient.subscribeToBlockHeaders();
|
||||||
|
|
||||||
|
// set stream subscription
|
||||||
|
_heightSubscription =
|
||||||
|
subscription.responseStream.asBroadcastStream().listen((event) {
|
||||||
|
final response = event;
|
||||||
if (response != null &&
|
if (response != null &&
|
||||||
response is Map &&
|
response is Map &&
|
||||||
response.containsKey('height')) {
|
response.containsKey('height')) {
|
||||||
final int chainHeight = response['height'] as int;
|
final int chainHeight = response['height'] as int;
|
||||||
// print("Current chain height: $chainHeight");
|
// print("Current chain height: $chainHeight");
|
||||||
latestHeight = chainHeight;
|
_latestHeight = chainHeight;
|
||||||
|
|
||||||
if (isFirstResponse) {
|
if (isFirstResponse) {
|
||||||
isFirstResponse = false;
|
isFirstResponse = false;
|
||||||
|
|
Loading…
Reference in a new issue