dev
dart 396 lines 13.8 KB
Raw
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 }