Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -765,27 +765,39 @@ class PaykitPaymentRequestRepo @Inject constructor(
return PaykitPaymentRequestCreation(request, creatorIdentity, wasPublishedToActiveState)
}

suspend fun accept(request: PaykitPaymentRequest): Result<Unit> {
if (!request.requiresAcceptance) {
return withContext(ioDispatcher) {
runSuspendCatching {
operationMutex.withLock {
if (_pendingRequests.value.none { it.id == request.id }) {
throw PaykitPaymentRequestError.RequestUnavailable
}
}
suspend fun ensurePaymentAllowed(request: PaykitPaymentRequest): Result<Unit> = withContext(ioDispatcher) {
runSuspendCatching {
if (
paykitSdkService.linkedPeers().any {
it.state == LinkedPeerState.BLOCKED &&
PubkyPublicKeyFormat.matches(it.counterparty, request.counterparty) &&
it.counterpartyReceiverPath == request.counterpartyReceiverPath
}
) {
throw PaykitPaymentRequestError.RequestUnavailable
}
}
return updateRequest(
request = request,
resultingState = PaymentRequestLifecycleState.ACCEPTED,
) {
paykitSdkService.acceptPaymentRequest(
counterparty = it.counterparty,
counterpartyReceiverPath = it.counterpartyReceiverPath,
paymentRequestId = it.paymentRequestId,
)
}

suspend fun accept(request: PaykitPaymentRequest): Result<Unit> = withContext(ioDispatcher) {
runSuspendCatching {
if (!request.requiresAcceptance) {
operationMutex.withLock {
ensurePaymentAllowed(request).getOrThrow()
if (_pendingRequests.value.none { it.id == request.id }) {
throw PaykitPaymentRequestError.RequestUnavailable
}
}
} else {
updateRequest(request, PaymentRequestLifecycleState.ACCEPTED) {
ensurePaymentAllowed(it).getOrThrow()
paykitSdkService.acceptPaymentRequest(
counterparty = it.counterparty,
counterpartyReceiverPath = it.counterpartyReceiverPath,
paymentRequestId = it.paymentRequestId,
)
}.getOrThrow()
}
}.onFailure {
Logger.warn("Failed to accept incoming Paykit payment request", it, context = TAG)
}
Expand Down Expand Up @@ -899,16 +911,39 @@ class PaykitPaymentRequestRepo @Inject constructor(
paykitSdkService.receivePrivateMessagesFromLinkedPeers().also(::logIntakeFailures)
val now = clock.now()
val records = paykitSdkService.paymentRequests()
val blockedPeers = paykitSdkService.linkedPeers().filter { it.state == LinkedPeerState.BLOCKED }
val availableRecords = records.filterNot { record ->
blockedPeers.any {
PubkyPublicKeyFormat.matches(it.counterparty, record.counterparty) &&
it.counterpartyReceiverPath == record.counterpartyReceiverPath
}
}
val locallyCompletedProofKinds = expectedIdentity
?.let(paymentProofStore::completedRequestProofKindsAwaitingSubmission)
.orEmpty()
val locallyCompletedRequestIds = locallyCompletedProofKinds.keys
val locallyInFlightRequestIds = expectedIdentity
?.let(paymentProofStore::inFlightRequestIds)
.orEmpty()
val subscriptions = records.mapNotNull(PaymentRequestRecord::toPaykitSubscription)
val allSubscriptions = records.mapNotNull(PaymentRequestRecord::toPaykitSubscription)
.map { it.withExpiredLifecycle(now) }
val restoredAcceptances = subscriptions
val blockedSubscriptionIds = allSubscriptions.mapNotNull { subscription ->
subscription.id.takeIf {
blockedPeers.any {
PubkyPublicKeyFormat.matches(it.counterparty, subscription.counterparty) &&
it.counterpartyReceiverPath == subscription.counterpartyReceiverPath
}
}
}.toSet()
val actionableSubscriptions = allSubscriptions.filterNot { it.id in blockedSubscriptionIds }
val subscriptions = allSubscriptions.filterNot { subscription ->
subscription.id in blockedSubscriptionIds && (
subscription.paidPeriods.isEmpty() ||
subscription.isProposalVisible(now) ||
subscription.isActive(now)
)
}
val restoredAcceptances = allSubscriptions
.filter {
it.isPayer &&
(it.lifecycleState == PaymentRequestLifecycleState.ACTIVE_RECURRING || it.paidPeriods.isNotEmpty())
Expand All @@ -921,11 +956,12 @@ class PaykitPaymentRequestRepo @Inject constructor(
subscription.id to acceptedAt
}
val updatedSubscriptionAcceptedAt = subscriptionAcceptedAt + restoredAcceptances
val recurringRequestsBySubscription = subscriptions.filter { it.isPayer }.associateWith { subscription ->
updatedSubscriptionAcceptedAt[subscription.id]
?.let { subscription.requestsThrough(now, it) }
.orEmpty()
}
val recurringRequestsBySubscription = actionableSubscriptions.filter { it.isPayer }
.associateWith { subscription ->
updatedSubscriptionAcceptedAt[subscription.id]
?.let { subscription.requestsThrough(now, it) }
.orEmpty()
}
val activeRecurringRequestIds = recurringRequestsBySubscription
.filterKeys { it.lifecycleState == PaymentRequestLifecycleState.ACTIVE_RECURRING }
.values
Expand All @@ -942,7 +978,11 @@ class PaykitPaymentRequestRepo @Inject constructor(
it.id !in locallyInFlightRequestIds &&
it.id !in updatedDismissedPaymentIds
}
val recurringHistory = recurringRequestsBySubscription.values.flatten().mapNotNull { request ->
val recurringHistory = allSubscriptions.filter { it.isPayer }.flatMap { subscription ->
updatedSubscriptionAcceptedAt[subscription.id]
?.let { subscription.requestsThrough(now, it) }
.orEmpty()
}.mapNotNull { request ->
when {
request.lifecycleState == PaymentRequestLifecycleState.PROOF_SUBMITTED -> request
request.id in locallyCompletedRequestIds -> request.copy(
Expand All @@ -952,7 +992,7 @@ class PaykitPaymentRequestRepo @Inject constructor(
else -> null
}
}
val oneTimeIncoming = records.mapNotNull { record ->
val oneTimeIncoming = availableRecords.mapNotNull { record ->
when (val result = record.parseIncomingPaykitPaymentRequest(now)) {
is PaykitPaymentRequestParseResult.Parsed -> result.request
is PaykitPaymentRequestParseResult.Rejected -> {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1430,11 +1430,13 @@ class PrivatePaykitRepo @Inject constructor(
).forEach { receiverPath ->
runSuspendCatching {
val report = paykitSdkService.clearPrivatePaymentList(publicKey, receiverPath)
?: return@runSuspendCatching false
if (report.failedToQueue.isNotEmpty() || report.failedToDeliver.isNotEmpty()) {
throw PrivatePaykitError.PrivateUnavailable
}
true
}.onSuccess {
clearedRetryKeys += PrivateMessageDrainRetryKey(publicKey, receiverPath)
if (it) clearedRetryKeys += PrivateMessageDrainRetryKey(publicKey, receiverPath)
}.onFailure {
failedPublicKeys += publicKey
firstError = firstError ?: it
Expand Down
117 changes: 77 additions & 40 deletions app/src/main/java/to/bitkit/repositories/PubkyRepo.kt
Original file line number Diff line number Diff line change
Expand Up @@ -65,8 +65,12 @@ sealed class PubkyContactError(message: String) : AppError(message) {
data object AlreadyExists : PubkyContactError("Contact already exists")
data object CannotAddSelf : PubkyContactError("Cannot add your own pubky as a contact")
data object InvalidFormat : PubkyContactError("Invalid pubky key format")
data object ActiveSubscription : PubkyContactError("Contact has an active subscription")
}

private fun Throwable.containsActiveSubscriptionError(): Boolean =
generateSequence(this) { it.cause }.any { it is PubkyContactError.ActiveSubscription }

data object PubkyAlreadySignedInError : AppError("Already signed in")

@Suppress("TooManyFunctions", "LargeClass", "LongParameterList")
Expand Down Expand Up @@ -94,6 +98,8 @@ class PubkyRepo @Inject constructor(
private val initializeMutex = Mutex()
private val loadProfileMutex = Mutex()
private val loadContactsMutex = Mutex()
private val contactsLock = Any()
private var contactsRevision = 0L
private val adoptedSourceCheckMutex = Mutex()
private var isServiceInitialized = false

Expand Down Expand Up @@ -717,7 +723,7 @@ class PubkyRepo @Inject constructor(
}
pubkyStore.update { it.copy(contactProfileOverrides = emptyMap()) }
notifyBackupStateChanged()
_contacts.update { emptyList() }
updateContacts { emptyList() }
markContactsLoaded()
Logger.info("Deleted all contacts", context = TAG)
}
Expand Down Expand Up @@ -776,41 +782,52 @@ class PubkyRepo @Inject constructor(
_isLoadingContacts.update { true }
var shouldMarkLoadCompleted = false
try {
runSuspendCatching {
withContext(ioDispatcher) {
val records = pubkyService.contactRecords()
val overrides = pubkyStore.data.first().contactProfileOverrides

coroutineScope {
records.map { record ->
async {
runSuspendCatching {
contactProfile(record.publicKey, record.label, record.profile, overrides)
}.onFailure {
Logger.warn(
"Failed to load contact '${redacted(record.publicKey)}'",
it,
context = TAG,
)
}.getOrElse {
PubkyProfile.placeholder(record.publicKey.ensurePubkyPrefix())
var reload: Boolean
do {
val revision = synchronized(contactsLock) { contactsRevision }
reload = false
runSuspendCatching {
withContext(ioDispatcher) {
val records = pubkyService.contactRecords()
val overrides = pubkyStore.data.first().contactProfileOverrides

coroutineScope {
records.map { record ->
async {
runSuspendCatching {
contactProfile(record.publicKey, record.label, record.profile, overrides)
}.onFailure {
Logger.warn(
"Failed to load contact '${redacted(record.publicKey)}'",
it,
context = TAG,
)
}.getOrElse {
PubkyProfile.placeholder(record.publicKey.ensurePubkyPrefix())
}
}
}
}.awaitAll().sortedBy { it.name.lowercase() }
}.awaitAll().sortedBy { it.name.lowercase() }
}
}
}.onSuccess { loadedContacts ->
if (_publicKey.value != pk) {
Logger.debug("Skipped stale contacts load for '${redacted(pk)}'", context = TAG)
return@onSuccess
}
synchronized(contactsLock) {
if (contactsRevision != revision) {
reload = true
return@onSuccess
}
_contacts.update { loadedContacts }
}
markContactsLoaded()
shouldMarkLoadCompleted = true
}.onFailure {
shouldMarkLoadCompleted = _publicKey.value == pk
Logger.error("Failed to load contacts", it, context = TAG)
}
}.onSuccess { loadedContacts ->
if (_publicKey.value != pk) {
Logger.debug("Skipped stale contacts load for '${redacted(pk)}'", context = TAG)
return@onSuccess
}
_contacts.update { loadedContacts }
markContactsLoaded()
shouldMarkLoadCompleted = true
}.onFailure {
shouldMarkLoadCompleted = _publicKey.value == pk
Logger.error("Failed to load contacts", it, context = TAG)
}
} while (reload && _publicKey.value == pk)
} finally {
_isLoadingContacts.update { false }
loadContactsMutex.unlock()
Expand Down Expand Up @@ -844,8 +861,13 @@ class PubkyRepo @Inject constructor(
val profile = existingProfile?.copy(publicKey = prefixedKey)
?: resolveContactProfile(prefixedKey).getOrThrow()
?: PubkyProfile.placeholder(prefixedKey)
pubkyService.saveContact(prefixedKey, profile.name, relevantReceiverPaths(prefixedKey))
_contacts.update { current ->
pubkyService.saveContact(
prefixedKey,
profile.name,
relevantReceiverPaths(prefixedKey),
restorePrivateConnection = true,
)
updateContacts { current ->
(current.filter { it.publicKey != prefixedKey } + profile)
.sortedBy { it.name.lowercase() }
}
Expand Down Expand Up @@ -886,7 +908,7 @@ class PubkyRepo @Inject constructor(
)
pubkyService.saveContact(prefixedKey, name)
upsertContactProfileOverride(updatedProfile)
_contacts.update { current ->
updateContacts { current ->
current.map { if (it.publicKey == prefixedKey) updatedProfile else it }
.sortedBy { it.name.lowercase() }
}
Expand All @@ -900,10 +922,13 @@ class PubkyRepo @Inject constructor(
val prefixedKey = publicKey.ensurePubkyPrefix()
pubkyService.removeContact(prefixedKey)
removeContactProfileOverride(prefixedKey)
_contacts.update { current -> current.filter { it.publicKey != prefixedKey } }
updateContacts { current -> current.filter { it.publicKey != prefixedKey } }
markContactsLoaded()
Logger.info("Removed contact '${redacted(prefixedKey)}'", context = TAG)
}
}.recoverCatching {
if (it.containsActiveSubscriptionError()) throw PubkyContactError.ActiveSubscription
throw it
}

suspend fun importContacts(publicKeys: List<String>): Result<Unit> = runSuspendCatching {
Expand All @@ -915,15 +940,20 @@ class PubkyRepo @Inject constructor(
runSuspendCatching {
val profile = resolveContactProfile(prefixedKey).getOrThrow()
?: PubkyProfile.placeholder(prefixedKey)
pubkyService.saveContact(prefixedKey, profile.name, relevantReceiverPaths(prefixedKey))
pubkyService.saveContact(
prefixedKey,
profile.name,
relevantReceiverPaths(prefixedKey),
restorePrivateConnection = true,
)
profile
}.onFailure {
Logger.warn("Failed to import contact '${redacted(prefixedKey)}'", it, context = TAG)
}.getOrNull()
}
}.awaitAll().filterNotNull()
}
_contacts.update { current ->
updateContacts { current ->
val existing = current.map { it.publicKey }.toSet()
(current + imported.filter { it.publicKey !in existing })
.sortedBy { it.name.lowercase() }
Expand Down Expand Up @@ -1407,13 +1437,20 @@ class PubkyRepo @Inject constructor(
}
_publicKey.update { null }
_profile.update { null }
_contacts.update { emptyList() }
updateContacts { emptyList() }
_contactsLoadVersion.update { 0L }
_contactsLoadCompletionVersion.update { 0L }
clearPendingImport()
if (clearRestorationFailure) _sessionRestorationFailed.update { false }
}

private fun updateContacts(transform: (List<PubkyProfile>) -> List<PubkyProfile>) {
synchronized(contactsLock) {
contactsRevision++
_contacts.update(transform)
}
}

private fun markContactsLoaded() {
_contactsLoadVersion.update { it + 1 }
}
Expand Down
Loading
Loading