| 1 | import 'dart:async'; |
| 2 | import 'dart:io'; |
| 3 | import 'dart:isolate'; |
| 4 | import 'dart:math'; |
| 5 | import 'dart:typed_data'; |
| 6 | |
| 7 | import 'package:bitcoin_base/bitcoin_base.dart'; |
| 8 | import 'package:cw_bitcoin/bitcoin_wallet.dart'; |
| 9 | import 'package:cw_bitcoin/bitcoin_wallet_addresses.dart'; |
| 10 | import 'package:cw_bitcoin/payjoin/payjoin_persister.dart'; |
| 11 | import 'package:cw_bitcoin/payjoin/payjoin_receive_worker.dart'; |
| 12 | import 'package:cw_bitcoin/payjoin/payjoin_send_worker.dart'; |
| 13 | import 'package:cw_bitcoin/payjoin/payjoin_session_errors.dart'; |
| 14 | import 'package:cw_bitcoin/payjoin/storage.dart'; |
| 15 | import 'package:cw_bitcoin/psbt/signer.dart'; |
| 16 | import 'package:cw_bitcoin/psbt/utils.dart'; |
| 17 | import 'package:cw_core/pathForWallet.dart'; |
| 18 | import 'package:cw_core/utils/print_verbose.dart'; |
| 19 | import 'package:payjoin_flutter/common.dart'; |
| 20 | import 'package:payjoin_flutter/receive.dart'; |
| 21 | import 'package:payjoin_flutter/send.dart'; |
| 22 | import 'package:payjoin_flutter/src/config.dart' as pj_config; |
| 23 | import 'package:payjoin_flutter/src/generated/api.dart' as pj_api; |
| 24 | import 'package:payjoin_flutter/uri.dart' as PayjoinUri; |
| 25 | |
| 26 | class PayjoinManager { |
| 27 | PayjoinManager(this._payjoinStorage, this._wallet); |
| 28 | |
| 29 | final PayjoinStorage _payjoinStorage; |
| 30 | final BitcoinWalletBase _wallet; |
| 31 | final Map<String, PayjoinPollerSession> _activePollers = {}; |
| 32 | |
| 33 | static const List<String> ohttpRelayUrls = [ |
| 34 | 'https://pj.bobspacebkk.com', |
| 35 | 'https://ohttp.achow101.com', |
| 36 | 'https://ohttp.cakewallet.com', |
| 37 | ]; |
| 38 | |
| 39 | static String randomOhttpRelayUrl() => |
| 40 | ohttpRelayUrls[Random.secure().nextInt(ohttpRelayUrls.length)]; |
| 41 | |
| 42 | static const payjoinDirectoryUrl = 'https://payjo.in'; |
| 43 | |
| 44 | var _logStreamController = StreamController<String>.broadcast(); |
| 45 | Stream<String> get logStream => _logStreamController.stream; |
| 46 | StreamSubscription<String>? _logSubscription; |
| 47 | |
| 48 | Future<void> initPayjoin() async { |
| 49 | await initLogging(); |
| 50 | |
| 51 | await pj_config.PConfig.initializeApp(); |
| 52 | } |
| 53 | |
| 54 | Future<void> initLogging() async { |
| 55 | _logSubscription?.cancel(); |
| 56 | _logStreamController = StreamController<String>.broadcast(); |
| 57 | |
| 58 | try { |
| 59 | final path = await pathForWalletDir(name: _wallet.name, type: _wallet.type); |
| 60 | File("$path/payjoin.log") |
| 61 | .create() |
| 62 | .then(_subscribeToLogStream) |
| 63 | .onError((e, s) async => printV("$e\n$s")); |
| 64 | } catch (e) { |
| 65 | printV(e); |
| 66 | } |
| 67 | } |
| 68 | |
| 69 | void _subscribeToLogStream(File logFile) { |
| 70 | _logSubscription = logStream.listen((logEntry) { |
| 71 | logFile.writeAsString("$logEntry\n", mode: FileMode.writeOnlyAppend); |
| 72 | }); |
| 73 | } |
| 74 | |
| 75 | void writePayjoinLog(String message) { |
| 76 | if (_logStreamController.isClosed) return; |
| 77 | |
| 78 | try { |
| 79 | _logStreamController.add(message); |
| 80 | } catch (_) {} |
| 81 | } |
| 82 | |
| 83 | Future<void> resumeSessions() async { |
| 84 | await initLogging(); |
| 85 | final allSessions = _payjoinStorage.readAllOpenSessions(_wallet.id); |
| 86 | |
| 87 | final spawnedSessions = allSessions.map((session) { |
| 88 | try { |
| 89 | if (session.isSenderSession) { |
| 90 | printV("Resuming Payjoin Sender Session ${session.pjUri!}"); |
| 91 | return _spawnSender( |
| 92 | sender: Sender.fromJson(json: session.sender!), |
| 93 | pjUri: session.pjUri!, |
| 94 | ); |
| 95 | } |
| 96 | final receiver = Receiver.fromJson(json: session.receiver!); |
| 97 | printV("Resuming Payjoin Receiver Session ${receiver.id()}"); |
| 98 | return spawnReceiver(receiver: receiver); |
| 99 | } on pj_api.FfiSerdeJsonError catch (_) { |
| 100 | _payjoinStorage.markSenderSessionUnrecoverable(session.pjUri!, "Outdated Session"); |
| 101 | } |
| 102 | }).nonNulls; |
| 103 | |
| 104 | printV("Resumed ${spawnedSessions.length} Payjoin Sessions"); |
| 105 | await Future.wait(spawnedSessions); |
| 106 | } |
| 107 | |
| 108 | Future<Sender> initSender( |
| 109 | String pjUriString, String originalPsbt, int networkFeesSatPerVb) async { |
| 110 | try { |
| 111 | final pjUri = (await PayjoinUri.Uri.fromStr(pjUriString)).checkPjSupported(); |
| 112 | final minFeeRateSatPerKwu = BigInt.from(networkFeesSatPerVb * 250); |
| 113 | final senderBuilder = await SenderBuilder.fromPsbtAndUri( |
| 114 | psbtBase64: originalPsbt, |
| 115 | pjUri: pjUri, |
| 116 | ); |
| 117 | final persister = PayjoinSenderPersister.impl(); |
| 118 | final newSender = await senderBuilder.buildRecommended(minFeeRate: minFeeRateSatPerKwu); |
| 119 | final senderToken = await newSender.persist(persister: persister); |
| 120 | |
| 121 | return Sender.load(token: senderToken, persister: persister); |
| 122 | } catch (e) { |
| 123 | throw Exception('Error initializing Payjoin Sender: $e'); |
| 124 | } |
| 125 | } |
| 126 | |
| 127 | Future<void> spawnNewSender({ |
| 128 | required Sender sender, |
| 129 | required String pjUrl, |
| 130 | required BigInt amount, |
| 131 | bool isTestnet = false, |
| 132 | }) async { |
| 133 | final pjUri = Uri.parse(pjUrl).queryParameters['pj']!; |
| 134 | await _payjoinStorage.insertSenderSession(sender, pjUri, _wallet.id, amount); |
| 135 | |
| 136 | return _spawnSender(isTestnet: isTestnet, sender: sender, pjUri: pjUri); |
| 137 | } |
| 138 | |
| 139 | Future<void> _spawnSender({ |
| 140 | required Sender sender, |
| 141 | required String pjUri, |
| 142 | bool isTestnet = false, |
| 143 | }) async { |
| 144 | final completer = Completer(); |
| 145 | final receivePort = ReceivePort(); |
| 146 | |
| 147 | receivePort.listen((message) async { |
| 148 | if (message is Map<String, dynamic>) { |
| 149 | try { |
| 150 | switch (message['type'] as PayjoinSenderRequestTypes) { |
| 151 | case PayjoinSenderRequestTypes.requestPosted: |
| 152 | writePayjoinLog("Sender($pjUri) PayjoinSenderRequestTypes.requestPosted"); |
| 153 | return; |
| 154 | case PayjoinSenderRequestTypes.psbtToSign: |
| 155 | writePayjoinLog("Sender($pjUri) PayjoinSenderRequestTypes.psbtToSign"); |
| 156 | |
| 157 | final proposalPsbt = message['psbt'] as String; |
| 158 | writePayjoinLog("Sender($pjUri) proposedPSBT: $proposalPsbt"); |
| 159 | |
| 160 | final utxos = _wallet.getUtxoWithPrivateKeys(); |
| 161 | final finalizedPsbt = await _wallet.signPsbt(proposalPsbt, utxos); |
| 162 | writePayjoinLog("Sender($pjUri) finalizedPsbt: $finalizedPsbt"); |
| 163 | |
| 164 | final txId = getTxIdFromPsbtV0(finalizedPsbt); |
| 165 | writePayjoinLog("Sender($pjUri) expected: $txId"); |
| 166 | |
| 167 | _wallet.commitPsbt(finalizedPsbt); |
| 168 | |
| 169 | _cleanupSession(pjUri); |
| 170 | await _payjoinStorage.markSenderSessionComplete(pjUri, txId); |
| 171 | completer.complete(); |
| 172 | } |
| 173 | } catch (e) { |
| 174 | writePayjoinLog("[ERROR] Sender($pjUri) $e"); |
| 175 | _cleanupSession(pjUri); |
| 176 | await _payjoinStorage.markSenderSessionUnrecoverable(pjUri, e.toString()); |
| 177 | completer.complete(); |
| 178 | } |
| 179 | } else if (message is PayjoinSessionError) { |
| 180 | _cleanupSession(pjUri); |
| 181 | if (message is UnrecoverableError) { |
| 182 | writePayjoinLog("[ERROR] Sender($pjUri) UnrecoverableError(${message.message})"); |
| 183 | await _payjoinStorage.markSenderSessionUnrecoverable(pjUri, message.message); |
| 184 | completer.complete(); |
| 185 | } else if (message is RecoverableError) { |
| 186 | writePayjoinLog("[ERROR] Sender($pjUri) RecoverableError(${message.message})"); |
| 187 | |
| 188 | completer.complete(); |
| 189 | } else { |
| 190 | writePayjoinLog("[ERROR] Sender($pjUri) GenericError(${message.message})"); |
| 191 | |
| 192 | completer.completeError(message); |
| 193 | } |
| 194 | } |
| 195 | }); |
| 196 | |
| 197 | final isolate = await Isolate.spawn( |
| 198 | PayjoinSenderWorker.run, |
| 199 | [receivePort.sendPort, sender.toJson(), pjUri], |
| 200 | ); |
| 201 | |
| 202 | _activePollers[pjUri] = PayjoinPollerSession(isolate, receivePort); |
| 203 | |
| 204 | return completer.future; |
| 205 | } |
| 206 | |
| 207 | Future<Receiver> getUnusedReceiver(String address, [bool isTestnet = false]) async { |
| 208 | final session = _payjoinStorage.getUnusedActiveReceiverSession(_wallet.id); |
| 209 | |
| 210 | if (session != null) { |
| 211 | await PayjoinUri.Url.fromStr(payjoinDirectoryUrl); |
| 212 | |
| 213 | return Receiver.fromJson(json: session.receiver!); |
| 214 | } |
| 215 | |
| 216 | return initReceiver(address); |
| 217 | } |
| 218 | |
| 219 | Future<Receiver> initReceiver(String address, |
| 220 | [bool isTestnet = false, int retryCount = 0]) async { |
| 221 | if (retryCount > 0) writePayjoinLog("Retrying initReceiver ${retryCount + 1} attempt"); |
| 222 | |
| 223 | try { |
| 224 | final ohttpKeys = await PayjoinUri.fetchOhttpKeys( |
| 225 | ohttpRelay: await randomOhttpRelayUrl(), |
| 226 | payjoinDirectory: payjoinDirectoryUrl, |
| 227 | ); |
| 228 | |
| 229 | final newReceiver = await NewReceiver.create( |
| 230 | address: address, |
| 231 | network: isTestnet ? Network.testnet : Network.bitcoin, |
| 232 | directory: payjoinDirectoryUrl, |
| 233 | ohttpKeys: ohttpKeys, |
| 234 | ); |
| 235 | final persister = PayjoinReceiverPersister.impl(); |
| 236 | final receiverToken = await newReceiver.persist(persister: persister); |
| 237 | final receiver = await Receiver.load(persister: persister, token: receiverToken); |
| 238 | |
| 239 | await _payjoinStorage.insertReceiverSession(receiver, _wallet.id); |
| 240 | |
| 241 | return receiver; |
| 242 | } catch (e) { |
| 243 | writePayjoinLog(e.toString()); |
| 244 | if (e.toString().contains("error sending request for url") && retryCount < 5) { |
| 245 | return initReceiver(address, isTestnet, ++retryCount); |
| 246 | } else { |
| 247 | rethrow; |
| 248 | } |
| 249 | } |
| 250 | } |
| 251 | |
| 252 | Future<void> spawnReceiver({ |
| 253 | required Receiver receiver, |
| 254 | bool isTestnet = false, |
| 255 | }) async { |
| 256 | final completer = Completer(); |
| 257 | final receivePort = ReceivePort(); |
| 258 | |
| 259 | SendPort? mainToIsolateSendPort; |
| 260 | List<UtxoWithPrivateKey> utxos = []; |
| 261 | String rawAmount = '0'; |
| 262 | |
| 263 | receivePort.listen((message) async { |
| 264 | if (message is Map<String, dynamic>) { |
| 265 | try { |
| 266 | switch (message['type'] as PayjoinReceiverRequestTypes) { |
| 267 | case PayjoinReceiverRequestTypes.processOriginalTx: |
| 268 | final tx = message['tx'] as String; |
| 269 | rawAmount = getOutputAmountFromTx(tx, _wallet); |
| 270 | break; |
| 271 | case PayjoinReceiverRequestTypes.checkIsOwned: |
| 272 | (_wallet.walletAddresses as BitcoinWalletAddresses).newPayjoinReceiver(); |
| 273 | _payjoinStorage.markReceiverSessionInProgress(receiver.id()); |
| 274 | |
| 275 | final inputScript = message['input_script'] as Uint8List; |
| 276 | final isOwned = _wallet.isMine(Script.fromRaw(byteData: inputScript)); |
| 277 | mainToIsolateSendPort?.send({ |
| 278 | 'requestId': message['requestId'], |
| 279 | 'result': isOwned, |
| 280 | }); |
| 281 | break; |
| 282 | |
| 283 | case PayjoinReceiverRequestTypes.checkIsReceiverOutput: |
| 284 | final outputScript = message['output_script'] as Uint8List; |
| 285 | final isReceiverOutput = _wallet.isMine(Script.fromRaw(byteData: outputScript)); |
| 286 | mainToIsolateSendPort?.send({ |
| 287 | 'requestId': message['requestId'], |
| 288 | 'result': isReceiverOutput, |
| 289 | }); |
| 290 | break; |
| 291 | |
| 292 | case PayjoinReceiverRequestTypes.getCandidateInputs: |
| 293 | utxos = _wallet.getUtxoWithPrivateKeys(confirmedOnly: true); |
| 294 | if (utxos.isEmpty) { |
| 295 | await _wallet.updateAllUnspents(); |
| 296 | utxos = _wallet.getUtxoWithPrivateKeys(confirmedOnly: true); |
| 297 | } |
| 298 | // Candidates arrive in wallet scan order (address, then age), which is |
| 299 | // predictable; shuffle so the receiver's input choice can't mirror it. |
| 300 | utxos.shuffle(Random.secure()); |
| 301 | mainToIsolateSendPort?.send({ |
| 302 | 'requestId': message['requestId'], |
| 303 | 'result': utxos, |
| 304 | }); |
| 305 | break; |
| 306 | |
| 307 | case PayjoinReceiverRequestTypes.processPsbt: |
| 308 | final psbt = message['psbt'] as String; |
| 309 | writePayjoinLog( |
| 310 | "Receiver(${receiver.id()}) PayjoinReceiverRequestTypes.processPsbt: $psbt"); |
| 311 | |
| 312 | final signedPsbt = await _wallet.signPsbt(psbt, utxos); |
| 313 | mainToIsolateSendPort?.send({ |
| 314 | 'requestId': message['requestId'], |
| 315 | 'result': signedPsbt, |
| 316 | }); |
| 317 | break; |
| 318 | |
| 319 | case PayjoinReceiverRequestTypes.proposalSent: |
| 320 | _cleanupSession(receiver.id()); |
| 321 | final psbt = message['psbt'] as String; |
| 322 | writePayjoinLog( |
| 323 | "Receiver(${receiver.id()}) PayjoinReceiverRequestTypes.proposalSent: $psbt"); |
| 324 | |
| 325 | await _payjoinStorage.markReceiverSessionComplete( |
| 326 | receiver.id(), getTxIdFromPsbtV0(psbt), rawAmount); |
| 327 | completer.complete(); |
| 328 | } |
| 329 | } catch (e) { |
| 330 | _cleanupSession(receiver.id()); |
| 331 | writePayjoinLog("[ERROR] Receiver(${receiver.id()}) $e"); |
| 332 | |
| 333 | await _payjoinStorage.markReceiverSessionUnrecoverable(receiver.id(), e.toString()); |
| 334 | completer.completeError(e); |
| 335 | } |
| 336 | } else if (message is PayjoinSessionError) { |
| 337 | _cleanupSession(receiver.id()); |
| 338 | if (message is UnrecoverableError) { |
| 339 | writePayjoinLog( |
| 340 | "[ERROR] Receiver(${receiver.id()}) UnrecoverableError(${message.message})"); |
| 341 | |
| 342 | await _payjoinStorage.markReceiverSessionUnrecoverable(receiver.id(), message.message); |
| 343 | completer.complete(); |
| 344 | } else if (message is RecoverableError) { |
| 345 | writePayjoinLog( |
| 346 | "[ERROR] Receiver(${receiver.id()}) RecoverableError(${message.message})"); |
| 347 | |
| 348 | completer.complete(); |
| 349 | } else { |
| 350 | writePayjoinLog("[ERROR] Receiver(${receiver.id()}) GenericError(${message.message})"); |
| 351 | |
| 352 | completer.completeError(message); |
| 353 | } |
| 354 | } else if (message is SendPort) { |
| 355 | mainToIsolateSendPort = message; |
| 356 | } |
| 357 | }); |
| 358 | |
| 359 | final isolate = await Isolate.spawn( |
| 360 | PayjoinReceiverWorker.run, |
| 361 | [receivePort.sendPort, receiver.toJson()], |
| 362 | ); |
| 363 | |
| 364 | _activePollers[receiver.id()] = PayjoinPollerSession(isolate, receivePort); |
| 365 | |
| 366 | return completer.future; |
| 367 | } |
| 368 | |
| 369 | void cleanupSessions() { |
| 370 | final sessionIds = _activePollers.keys.toList(); |
| 371 | for (final sessionId in sessionIds) { |
| 372 | _cleanupSession(sessionId); |
| 373 | } |
| 374 | |
| 375 | _logSubscription?.cancel(); |
| 376 | _logStreamController.close(); |
| 377 | } |
| 378 | |
| 379 | void _cleanupSession(String sessionId) { |
| 380 | writePayjoinLog("Cleaning $sessionId"); |
| 381 | _activePollers[sessionId]?.close(); |
| 382 | _activePollers.remove(sessionId); |
| 383 | } |
| 384 | } |
| 385 | |
| 386 | class PayjoinPollerSession { |
| 387 | final Isolate isolate; |
| 388 | final ReceivePort port; |
| 389 | |
| 390 | PayjoinPollerSession(this.isolate, this.port); |
| 391 | |
| 392 | void close() { |
| 393 | isolate.kill(); |
| 394 | port.close(); |
| 395 | } |
| 396 | } |