fix(pow): resolve audit findings — durability, coverage, perf and l10n

Durability & correctness:
- Keep the job entry + disk checkpoint alive until the publish continuation
  completes (was: dropped at mining-complete); failed publishes keep their
  checkpoint for restart retry and surface through a failures SharedFlow
- Move draft deletion into the publish continuations in every composer so a
  cancelled or process-killed mining job can't destroy the only copy of a post
- Persist gift-wrap mining: split NIP17Factory into createSeals (signer
  interaction, runs inline) + wrapSeal (pure CPU, runs on the queue), new
  REPLAY_WRAPS records restore pending DM wraps after process death
- Clamp synced NIP-78 difficulty (PoWPolicy.MAX_DIFFICULTY), require a sane
  target in PoWMiner, bound PoWRankEvaluator against short ids
- Reaction double-tap while mining now toggles (dedupeKey + cancelByKey)
  instead of publishing duplicates
- Restore checkpoints for every loaded account, not just the active one
- logOff purges the account's checkpoints and cancels its queued jobs
- Private notes: composer chip now gates on the gift-wrap kind and the
  per-post override reaches sendPrivateNote

Coverage:
- Public/live chat (kinds 42 + 1311) and voice replies (1244 + kind-1 audio
  replies) now route through the mining gate

Perf:
- FGS start() dedupes with a running flag; PendingIntents built once
- Banner 1 Hz clock only ticks while a job shows elapsed time
- PoWEstimator benchmark is single-flight behind a Mutex

UX / l10n:
- Post-mining failures toast with retry information
- Count strings converted to <plurals>; elapsed time via DateUtils; settings
  estimate uses localized units; shared powKindLabelRes replaces three
  duplicated kind→label maps; dead pow_chip_* strings removed
- CLI: pow mine lowercases the pubkey before mining (uppercase hex mined an
  id that never matches the signed event) and validates via quartz Hex

New queue tests: checkpoint lifetime, failure reporting, cancelByKey toggle,
per-owner cancellation.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01ADb3dez9jPk6QqyQ1rTx4V
This commit is contained in:
Claude
2026-07-11 00:31:05 +00:00
parent 81dc4449bd
commit 6dd1e4d31d
27 changed files with 985 additions and 353 deletions
@@ -888,13 +888,13 @@ class AppModules(
}
}
// Resume PoW mining jobs that were checkpointed before a process death.
// Resume PoW mining jobs that were checkpointed before a process death,
// for EVERY loaded account (the always-on service preloads non-active
// accounts, whose pending posts must not stay stranded on disk).
// Idempotent (the queue dedupes by job id), so re-emissions are safe.
applicationIOScope.launch {
sessionManager.accountContent.collectLatest { state ->
if (state is AccountState.LoggedIn) {
powJobRestorer.restore(state.account)
}
accountsCache.accounts.collect { loaded ->
loaded.values.forEach { powJobRestorer.restore(it) }
}
}
@@ -52,6 +52,7 @@ import com.vitorpamplona.amethyst.commons.onchain.OnchainZapSendStage
import com.vitorpamplona.amethyst.commons.onchain.OnchainZapSender
import com.vitorpamplona.amethyst.commons.onchain.OnchainZapShare
import com.vitorpamplona.amethyst.commons.richtext.RichTextParser
import com.vitorpamplona.amethyst.commons.service.pow.PersistedPoWJob
import com.vitorpamplona.amethyst.commons.service.pow.PoWCategory
import com.vitorpamplona.amethyst.commons.service.pow.PoWPolicy
import com.vitorpamplona.amethyst.commons.service.pow.PoWPublishQueue
@@ -831,7 +832,7 @@ class Account(
work: suspend (isActive: () -> Boolean) -> Unit,
): Boolean {
val queue = powQueue() ?: return false
queue.enqueueWork(kind, difficulty, work)
queue.enqueueWork(kind, difficulty, owner = signer.pubKey, work = work)
return true
}
@@ -871,6 +872,68 @@ class Account(
return true
}
/**
* The one-liner for template send paths: when [template]'s kind should be
* mined (per settings and the optional composer [overrideDifficulty]),
* enqueue it and run [send] with the mined template once the nonce is
* found; otherwise run [send] with [template] right now.
*/
suspend fun <T : Event> sendMined(
template: EventTemplate<T>,
replay: PoWReplay?,
overrideDifficulty: Int? = null,
send: suspend (EventTemplate<T>) -> Unit,
) {
val difficulty = powDifficultyFor(template.kind, overrideDifficulty)
if (difficulty == null || !mineTemplateInBackground(template, difficulty, replay, send)) {
send(template)
}
}
/**
* Queues wrap mining for pre-signed [seals] (see NIP17Factory.createSeals):
* each seal gets its ephemeral-key envelope mined at [difficulty] on the
* worker pool, then the wraps broadcast. Checkpointed under
* [PersistedPoWJob.REPLAY_WRAPS] (the seals are already-signed ciphertext,
* safe to persist) unless an [existingRecord] from the restorer is passed.
* Returns false when no queue is wired.
*/
fun mineWrapsInBackground(
seals: List<NIP17Factory.AddressedSeal>,
expirationDelta: Long?,
difficulty: Int,
existingRecord: PersistedPoWJob? = null,
): Boolean {
val queue = powQueue() ?: return false
if (seals.isEmpty()) return true
val record =
existingRecord
?: PersistedPoWJob(
id = RandomInstance.randomChars(16),
accountPubkey = signer.pubKey,
kind = GiftWrapEvent.KIND,
difficulty = difficulty,
templateJson = "",
replayType = PersistedPoWJob.REPLAY_WRAPS,
extraEventsJson = seals.map { it.seal.toJson() },
recipientPubkeys = seals.map { it.recipient },
wrapExpirationDelta = expirationDelta,
createdAtSec = TimeUtils.now(),
)
queue.enqueueStaged(
kind = GiftWrapEvent.KIND,
difficulty = difficulty,
persistAs = record,
mine = { isActive ->
seals.map { NIP17Factory().wrapSeal(it, expirationDelta, powDifficulty = difficulty, powIsActive = isActive) }
},
publish = { wraps -> broadcastPrivately(wraps) },
)
return true
}
private fun <T : Event> withFinalSignerTags(template: EventTemplate<T>): EventTemplate<T> {
val currentSigner = signer
if (currentSigner !is NostrSignerWithClientTag) return template
@@ -913,19 +976,28 @@ class Account(
val isPrivateTarget = note.event is NIP17Group || note.isPrivateRumor()
val powDifficulty = if (isPrivateTarget) null else powDifficultyFor(ReactionEvent.KIND)
if (powDifficulty != null &&
mineInBackground(ReactionEvent.KIND, powDifficulty) { isActive ->
ReactionAction.reactTo(
note = note,
reaction = reaction,
by = userProfile(),
signer = miningSigner(powDifficulty, setOf(ReactionEvent.KIND), isActive),
onPublic = ::sendAutomatic,
onPrivate = ::broadcastPrivately,
)
if (powDifficulty != null) {
val queue = powQueue()
if (queue != null) {
// toggle semantics while mining: a second tap on the same
// reaction un-likes by cancelling the pending job instead of
// publishing a duplicate (the mined event doesn't exist yet,
// so hasReacted can't dedupe).
val dedupeKey = "reaction:${note.idHex}:$reaction"
if (queue.cancelByKey(dedupeKey)) return
queue.enqueueWork(ReactionEvent.KIND, powDifficulty, dedupeKey, owner = signer.pubKey) { isActive ->
ReactionAction.reactTo(
note = note,
reaction = reaction,
by = userProfile(),
signer = miningSigner(powDifficulty, setOf(ReactionEvent.KIND), isActive),
onPublic = ::sendAutomatic,
onPrivate = ::broadcastPrivately,
)
}
return
}
) {
return
}
ReactionAction.reactTo(
@@ -2732,12 +2804,15 @@ class Account(
if (!isWriteable()) return
val powDifficulty = powDifficultyFor(GiftWrapEvent.KIND)
if (powDifficulty != null &&
mineInBackground(GiftWrapEvent.KIND, powDifficulty) { isActive ->
broadcastPrivately(NIP17Factory().createEncryptedFileNIP17(template, signer, wrapPowDifficulty = powDifficulty, wrapPowIsActive = isActive))
}
) {
return
if (powDifficulty != null) {
// Sign the inner event and every seal NOW, in the caller's
// interaction context — an external signer (Amber/bunker) cannot
// prompt from a background mining worker. Only the local-CPU
// ephemeral-key wrap mining goes to the queue, checkpointed so a
// process death mid-mine cannot lose the file announcement.
val senderMessage = signer.sign(template)
val seals = NIP17Factory().createSeals(senderMessage, senderMessage.groupMembers(), signer)
if (mineWrapsInBackground(seals.seals, seals.expirationDelta, powDifficulty)) return
}
broadcastPrivately(NIP17Factory().createEncryptedFileNIP17(template, signer))
@@ -2745,12 +2820,11 @@ class Account(
override suspend fun sendNip17PrivateMessage(template: EventTemplate<ChatMessageEvent>) {
val powDifficulty = powDifficultyFor(GiftWrapEvent.KIND)
if (powDifficulty != null &&
mineInBackground(GiftWrapEvent.KIND, powDifficulty) { isActive ->
broadcastPrivately(NIP17Factory().createMessageNIP17(template, signer, wrapPowDifficulty = powDifficulty, wrapPowIsActive = isActive))
}
) {
return
if (powDifficulty != null) {
// See sendNip17EncryptedFile: sign inline, queue only wrap mining.
val senderMessage = signer.sign(template)
val seals = NIP17Factory().createSeals(senderMessage, senderMessage.groupMembers(), signer)
if (mineWrapsInBackground(seals.seals, seals.expirationDelta, powDifficulty)) return
}
broadcastPrivately(NIP17Factory().createMessageNIP17(template, signer))
@@ -2762,17 +2836,23 @@ class Account(
* to the recipient's DM relays. Used for private replies (the parent's
* author and participants are already p-tagged) and for private posts
* (the Notify list is the audience). Nothing reaches public relays.
*
* [powOverrideDifficulty] is the composer chip's per-post override:
* null follows the account's gift-wrap setting, 0 disables mining.
*/
suspend fun sendPrivateNote(template: EventTemplate<TextNoteEvent>) {
suspend fun sendPrivateNote(
template: EventTemplate<TextNoteEvent>,
powOverrideDifficulty: Int? = null,
) {
if (!isWriteable()) return
val powDifficulty = powDifficultyFor(GiftWrapEvent.KIND)
if (powDifficulty != null &&
mineInBackground(GiftWrapEvent.KIND, powDifficulty) { isActive ->
broadcastPrivately(NIP17Factory().createNoteNIP17(template, signer, wrapPowDifficulty = powDifficulty, wrapPowIsActive = isActive))
}
) {
return
val powDifficulty = powDifficultyFor(GiftWrapEvent.KIND, powOverrideDifficulty)
if (powDifficulty != null) {
// See sendNip17EncryptedFile: sign inline, queue only wrap mining.
val senderNote = signer.sign(template)
val recipients = senderNote.taggedUserIds().plus(signer.pubKey).toSet()
val seals = NIP17Factory().createSeals(senderNote, recipients, signer)
if (mineWrapsInBackground(seals.seals, seals.expirationDelta, powDifficulty)) return
}
broadcastPrivately(NIP17Factory().createNoteNIP17(template, signer))
@@ -2785,8 +2865,10 @@ class Account(
}
}
suspend fun broadcastPrivately(signedEvents: NIP17Factory.Result) {
val mine = signedEvents.wraps.filter { (it.recipientPubKey() == signer.pubKey) }
suspend fun broadcastPrivately(signedEvents: NIP17Factory.Result) = broadcastPrivately(signedEvents.wraps)
suspend fun broadcastPrivately(wraps: List<GiftWrapEvent>) {
val mine = wraps.filter { (it.recipientPubKey() == signer.pubKey) }
mine.forEach { giftWrap ->
cache.justConsumeMyOwnEvent(giftWrap)
@@ -2795,7 +2877,7 @@ class Account(
val id = mine.firstOrNull()?.id
val mineNote = if (id == null) null else cache.getNoteIfExists(id)
signedEvents.wraps.forEach { wrap ->
wraps.forEach { wrap ->
// Creates an alias
if (mineNote != null && wrap.recipientPubKey() != signer.pubKey) {
cache.getOrAddAliasNote(wrap.id, mineNote)
@@ -23,6 +23,7 @@ package com.vitorpamplona.amethyst.model
import androidx.compose.runtime.Stable
import com.vitorpamplona.amethyst.commons.audio.VisualizerStyle
import com.vitorpamplona.amethyst.commons.service.pow.PoWCategory
import com.vitorpamplona.amethyst.commons.service.pow.PoWPolicy
import com.vitorpamplona.amethyst.ui.screen.loggedIn.notifications.equalImmutableLists
import com.vitorpamplona.quartz.nip17Dm.base.ChatroomKey
import com.vitorpamplona.quartz.nip57Zaps.LnZapEvent
@@ -197,8 +198,12 @@ class AccountSyncedSettings(
chats.pinnedChatrooms.tryEmit(newPinnedChatrooms)
}
if (proofOfWork.difficulty.value != syncedSettingsInternal.proofOfWork.difficulty) {
proofOfWork.difficulty.tryEmit(syncedSettingsInternal.proofOfWork.difficulty)
// clamp like the local setter: a synced NIP-78 event from another
// client could carry an out-of-range value that would crash the miner
// (>256) or mine forever (41+).
val newDifficulty = syncedSettingsInternal.proofOfWork.difficulty.coerceIn(0, PoWPolicy.MAX_DIFFICULTY)
if (proofOfWork.difficulty.value != newDifficulty) {
proofOfWork.difficulty.tryEmit(newDifficulty)
}
val newPoWCategories = PoWCategory.fromIds(syncedSettingsInternal.proofOfWork.enabledCategories)
@@ -313,13 +318,18 @@ class AccountPoWPreferences(
val difficulty: MutableStateFlow<Int> = MutableStateFlow(0),
val enabledCategories: MutableStateFlow<Set<PoWCategory>> = MutableStateFlow(PoWCategory.DEFAULT_ENABLED),
) {
fun updateDifficulty(newDifficulty: Int): Boolean =
if (difficulty.value != newDifficulty) {
difficulty.tryEmit(newDifficulty.coerceIn(0, MAX_POW_DIFFICULTY))
fun updateDifficulty(newDifficulty: Int): Boolean {
// compare the coerced value: reporting a change for an out-of-range
// input that clamps to the current value would republish identical
// settings to relays.
val coerced = newDifficulty.coerceIn(0, MAX_POW_DIFFICULTY)
return if (difficulty.value != coerced) {
difficulty.tryEmit(coerced)
true
} else {
false
}
}
fun updateCategory(
category: PoWCategory,
@@ -336,8 +346,7 @@ class AccountPoWPreferences(
}
companion object {
// above ~40 bits a phone would mine for days; treat it as a config error.
const val MAX_POW_DIFFICULTY = 40
const val MAX_POW_DIFFICULTY = PoWPolicy.MAX_DIFFICULTY
}
}
@@ -28,6 +28,7 @@ import com.vitorpamplona.amethyst.service.scheduledposts.ScheduledPostStore
import com.vitorpamplona.quartz.nip01Core.core.Event
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.NormalizedRelayUrl
import com.vitorpamplona.quartz.nip01Core.signers.EventTemplate
import com.vitorpamplona.quartz.nip17Dm.NIP17Factory
import com.vitorpamplona.quartz.utils.Log
import java.util.UUID
@@ -49,6 +50,16 @@ class PowJobRestorer(
Log.d(TAG) { "Restoring ${records.size} pending PoW job(s) for ${account.signer.pubKey.take(8)}…" }
records.forEach { record ->
if (record.difficulty <= 0) {
store.remove(record.id)
return@forEach
}
if (record.replayType == PersistedPoWJob.REPLAY_WRAPS) {
restoreWraps(account, record)
return@forEach
}
val template =
try {
EventTemplate.fromJson(record.templateJson)
@@ -58,11 +69,6 @@ class PowJobRestorer(
return@forEach
}
if (record.difficulty <= 0) {
store.remove(record.id)
return@forEach
}
queue.enqueue(
template = template,
pubKey = account.signer.pubKey,
@@ -78,6 +84,46 @@ class PowJobRestorer(
}
}
/**
* Wrap jobs carry no template — the pre-signed seals live in
* [PersistedPoWJob.extraEventsJson], one recipient per seal in
* [PersistedPoWJob.recipientPubkeys]. Re-enqueueing goes through the same
* [Account.mineWrapsInBackground] path the send used, with the existing
* record so the checkpoint id (and restore idempotence) is preserved.
*/
private fun restoreWraps(
account: Account,
record: PersistedPoWJob,
) {
if (record.extraEventsJson.size != record.recipientPubkeys.size) {
Log.w(TAG) { "Dropping malformed wrap PoW job ${record.id}: ${record.extraEventsJson.size} seal(s) vs ${record.recipientPubkeys.size} recipient(s)" }
store.remove(record.id)
return
}
val seals =
record.extraEventsJson.zip(record.recipientPubkeys).mapNotNull { (sealJson, recipient) ->
try {
NIP17Factory.AddressedSeal(recipient = recipient, seal = Event.fromJson(sealJson))
} catch (e: Exception) {
Log.w(TAG, "Dropping unreadable seal of wrap PoW job ${record.id}", e)
null
}
}
if (seals.isEmpty()) {
store.remove(record.id)
return
}
account.mineWrapsInBackground(
seals = seals,
expirationDelta = record.wrapExpirationDelta,
difficulty = record.difficulty,
existingRecord = record,
)
}
private suspend fun replay(
account: Account,
record: PersistedPoWJob,
@@ -0,0 +1,51 @@
/*
* Copyright (c) 2025 Vitor Pamplona
*
* Permission is hereby granted, free of charge, to any person obtaining a copy of
* this software and associated documentation files (the "Software"), to deal in
* the Software without restriction, including without limitation the rights to use,
* copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the
* Software, and to permit persons to whom the Software is furnished to do so,
* subject to the following conditions:
*
* The above copyright notice and this permission notice shall be included in all
* copies or substantial portions of the Software.
*
* THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
* IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS
* FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR
* COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN
* AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION
* WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
*/
package com.vitorpamplona.amethyst.service.pow
import androidx.annotation.StringRes
import com.vitorpamplona.amethyst.R
import com.vitorpamplona.quartz.nip18Reposts.GenericRepostEvent
import com.vitorpamplona.quartz.nip18Reposts.RepostEvent
import com.vitorpamplona.quartz.nip25Reactions.ReactionEvent
import com.vitorpamplona.quartz.nip28PublicChat.message.ChannelMessageEvent
import com.vitorpamplona.quartz.nip53LiveActivities.chat.LiveActivitiesChatMessageEvent
import com.vitorpamplona.quartz.nip56Reports.ReportEvent
import com.vitorpamplona.quartz.nip59Giftwrap.wraps.GiftWrapEvent
import com.vitorpamplona.quartz.nipA0VoiceMessages.VoiceEvent
import com.vitorpamplona.quartz.nipA0VoiceMessages.VoiceReplyEvent
/**
* The one user-facing label for "what is being mined": shared by the
* broadcast banner, the mining foreground notification, and failure toasts
* so a job is described the same way everywhere it appears.
*/
@StringRes
fun powKindLabelRes(kind: Int): Int =
when (kind) {
ReactionEvent.KIND -> R.string.reaction
RepostEvent.KIND, GenericRepostEvent.KIND -> R.string.boost
VoiceEvent.KIND -> R.string.voice_post
VoiceReplyEvent.KIND -> R.string.voice_reply
ReportEvent.KIND -> R.string.pow_kind_report
GiftWrapEvent.KIND -> R.string.private_message
ChannelMessageEvent.KIND, LiveActivitiesChatMessageEvent.KIND -> R.string.pow_kind_chat_message
else -> R.string.post
}
@@ -20,6 +20,7 @@
*/
package com.vitorpamplona.amethyst.service.pow
import android.app.Notification
import android.app.NotificationChannel
import android.app.NotificationManager
import android.app.PendingIntent
@@ -35,12 +36,8 @@ import com.vitorpamplona.amethyst.Amethyst
import com.vitorpamplona.amethyst.R
import com.vitorpamplona.amethyst.commons.service.pow.PoWJobState
import com.vitorpamplona.amethyst.ui.MainActivity
import com.vitorpamplona.amethyst.ui.pluralStringRes
import com.vitorpamplona.amethyst.ui.stringRes
import com.vitorpamplona.quartz.nip18Reposts.GenericRepostEvent
import com.vitorpamplona.quartz.nip18Reposts.RepostEvent
import com.vitorpamplona.quartz.nip25Reactions.ReactionEvent
import com.vitorpamplona.quartz.nipA0VoiceMessages.VoiceEvent
import com.vitorpamplona.quartz.nipA0VoiceMessages.VoiceReplyEvent
import com.vitorpamplona.quartz.utils.Log
import kotlinx.collections.immutable.ImmutableList
import kotlinx.coroutines.CoroutineScope
@@ -75,6 +72,33 @@ class PowMiningForegroundService : Service() {
private var sessionTotal = 0
private var lastQueueSize = 0
// Built once per service instance: the intents never change, and
// buildNotification runs on every queue update.
private val tapIntent: PendingIntent by lazy {
PendingIntent.getActivity(
this,
0,
Intent(this, MainActivity::class.java).apply {
addFlags(Intent.FLAG_ACTIVITY_NEW_TASK or Intent.FLAG_ACTIVITY_CLEAR_TOP)
},
PendingIntent.FLAG_IMMUTABLE or PendingIntent.FLAG_UPDATE_CURRENT,
)
}
private val cancelIntent: PendingIntent by lazy {
PendingIntent.getService(
this,
1,
Intent(this, PowMiningForegroundService::class.java).setAction(ACTION_CANCEL_ALL),
PendingIntent.FLAG_IMMUTABLE or PendingIntent.FLAG_UPDATE_CURRENT,
)
}
override fun onCreate() {
super.onCreate()
running = true
}
override fun onBind(intent: Intent?): IBinder? = null
override fun onStartCommand(
@@ -114,6 +138,7 @@ class PowMiningForegroundService : Service() {
}
override fun onDestroy() {
running = false
scope.cancel()
super.onDestroy()
}
@@ -158,14 +183,14 @@ class PowMiningForegroundService : Service() {
}
}
private fun buildNotification(jobs: ImmutableList<PoWJobState>): android.app.Notification {
private fun buildNotification(jobs: ImmutableList<PoWJobState>): Notification {
val done = (sessionTotal - jobs.size).coerceAtLeast(0)
val total = (done + jobs.size).coerceAtLeast(1)
val current = jobs.firstOrNull { it.isMining } ?: jobs.firstOrNull()
val text =
current?.let {
stringRes(this, R.string.pow_mining_job, kindLabel(this, it.kind), it.difficulty.toString())
pluralStringRes(this, R.plurals.pow_mining_job, it.difficulty, stringRes(this, powKindLabelRes(it.kind)), it.difficulty)
} ?: stringRes(this, R.string.pow_mining_title)
val progressStyle: NotificationCompat.ProgressStyle =
@@ -178,30 +203,12 @@ class PowMiningForegroundService : Service() {
.setProgress(done)
}
val tapIntent =
PendingIntent.getActivity(
this,
0,
Intent(this, MainActivity::class.java).apply {
addFlags(Intent.FLAG_ACTIVITY_NEW_TASK or Intent.FLAG_ACTIVITY_CLEAR_TOP)
},
PendingIntent.FLAG_IMMUTABLE or PendingIntent.FLAG_UPDATE_CURRENT,
)
val cancelIntent =
PendingIntent.getService(
this,
1,
Intent(this, PowMiningForegroundService::class.java).setAction(ACTION_CANCEL_ALL),
PendingIntent.FLAG_IMMUTABLE or PendingIntent.FLAG_UPDATE_CURRENT,
)
return NotificationCompat
.Builder(this, CHANNEL_ID)
.setSmallIcon(R.drawable.amethyst)
.setContentTitle(
if (jobs.size > 1) {
stringRes(this, R.string.pow_mining_progress, jobs.size.toString())
pluralStringRes(this, R.plurals.pow_mining_progress, jobs.size, jobs.size)
} else {
stringRes(this, R.string.pow_mining_title)
},
@@ -223,6 +230,13 @@ class PowMiningForegroundService : Service() {
private const val NOTIFICATION_ID = 0x504F57 // "POW"
private const val ACTION_CANCEL_ALL = "com.vitorpamplona.amethyst.pow.CANCEL_ALL"
// Best-effort de-dup for start(): the queue calls it on EVERY enqueue,
// and each call otherwise round-trips through system_server. A stale
// false only costs one redundant startForegroundService (which Android
// routes to the existing instance's onStartCommand anyway).
@Volatile
private var running = false
/**
* Best-effort start: enqueue happens while the user is interacting
* with the app, so the foreground-start allowance normally holds. A
@@ -230,6 +244,7 @@ class PowMiningForegroundService : Service() {
* proceeds unprotected and the service starts on the next enqueue.
*/
fun start(context: Context) {
if (running) return
try {
context.startForegroundService(Intent(context, PowMiningForegroundService::class.java))
} catch (e: Exception) {
@@ -251,17 +266,5 @@ class PowMiningForegroundService : Service() {
},
)
}
private fun kindLabel(
context: Context,
kind: Int,
): String =
when (kind) {
ReactionEvent.KIND -> stringRes(context, R.string.reaction)
RepostEvent.KIND, GenericRepostEvent.KIND -> stringRes(context, R.string.boost)
VoiceEvent.KIND -> stringRes(context, R.string.voice_post)
VoiceReplyEvent.KIND -> stringRes(context, R.string.voice_reply)
else -> stringRes(context, R.string.post)
}
}
}
@@ -20,6 +20,7 @@
*/
package com.vitorpamplona.amethyst.ui.broadcast
import android.text.format.DateUtils
import androidx.compose.animation.AnimatedVisibility
import androidx.compose.animation.animateContentSize
import androidx.compose.animation.core.RepeatMode
@@ -57,6 +58,7 @@ import androidx.compose.runtime.setValue
import androidx.compose.ui.Alignment
import androidx.compose.ui.Modifier
import androidx.compose.ui.graphics.Color
import androidx.compose.ui.res.pluralStringResource
import androidx.compose.ui.text.style.TextOverflow
import androidx.compose.ui.tooling.preview.Preview
import androidx.compose.ui.unit.dp
@@ -68,6 +70,7 @@ import com.vitorpamplona.amethyst.commons.service.broadcast.BroadcastEvent
import com.vitorpamplona.amethyst.commons.service.broadcast.BroadcastStatus
import com.vitorpamplona.amethyst.commons.service.broadcast.RelayResult
import com.vitorpamplona.amethyst.commons.service.pow.PoWJobState
import com.vitorpamplona.amethyst.service.pow.powKindLabelRes
import com.vitorpamplona.amethyst.ui.stringRes
import com.vitorpamplona.amethyst.ui.theme.ThemeComparisonColumn
import com.vitorpamplona.quartz.nip01Core.core.Event
@@ -172,12 +175,15 @@ private fun MiningContent(
label = "boltAlpha",
)
// 1 Hz clock driving the per-job elapsed labels
// 1 Hz clock driving the per-job elapsed labels; only ticks while some
// job actually shows an elapsed time (queued-only banners don't need it).
var nowSec by remember { mutableLongStateOf(TimeUtils.now()) }
LaunchedEffect(Unit) {
while (true) {
delay(1_000)
nowSec = TimeUtils.now()
if (miningJobs.any { it.miningStartedAt != null }) {
LaunchedEffect(Unit) {
while (true) {
nowSec = TimeUtils.now()
delay(1_000)
}
}
}
@@ -195,7 +201,7 @@ private fun MiningContent(
)
Text(
text = stringRes(R.string.pow_mining_progress, miningJobs.size),
text = pluralStringResource(R.plurals.pow_mining_progress, miningJobs.size, miningJobs.size),
style = MaterialTheme.typography.bodySmall,
color = MaterialTheme.colorScheme.onSurface,
maxLines = 1,
@@ -213,12 +219,13 @@ private fun MiningContent(
Spacer(Modifier.width(26.dp))
val base =
stringRes(
if (job.isMining) R.string.pow_mining_job else R.string.pow_queued_job,
pluralStringResource(
if (job.isMining) R.plurals.pow_mining_job else R.plurals.pow_queued_job,
job.difficulty,
kindToName(job.kind),
job.difficulty.toString(),
job.difficulty,
)
val elapsed = job.miningStartedAt?.let { formatElapsed((nowSec - it).coerceAtLeast(0)) }
val elapsed = job.miningStartedAt?.let { DateUtils.formatElapsedTime((nowSec - it).coerceAtLeast(0)) }
Text(
text = if (elapsed != null) "$base • $elapsed" else base,
@@ -229,16 +236,20 @@ private fun MiningContent(
modifier = Modifier.weight(1f),
)
IconButton(
onClick = { onCancelJob(job.id) },
modifier = Modifier.size(22.dp),
) {
Icon(
symbol = MaterialSymbols.Close,
contentDescription = stringRes(R.string.pow_notification_cancel_all),
tint = MaterialTheme.colorScheme.onSurfaceVariant,
modifier = Modifier.size(14.dp),
)
// once the nonce is found the job is signing/broadcasting —
// there is nothing safe to abort anymore.
if (job.isCancellable) {
IconButton(
onClick = { onCancelJob(job.id) },
modifier = Modifier.size(22.dp),
) {
Icon(
symbol = MaterialSymbols.Close,
contentDescription = stringRes(R.string.pow_notification_cancel_all),
tint = MaterialTheme.colorScheme.onSurfaceVariant,
modifier = Modifier.size(14.dp),
)
}
}
}
}
@@ -253,13 +264,6 @@ private fun MiningContent(
}
}
private fun formatElapsed(seconds: Long): String =
if (seconds < 60) {
"${seconds}s"
} else {
"${seconds / 60}m ${seconds % 60}s"
}
@Composable
private fun SingleBroadcastContent(broadcast: BroadcastEvent) {
Row(
@@ -576,14 +580,7 @@ fun Event.toKindName(): String =
}
@Composable
fun kindToName(kind: Int): String =
when (kind) {
ReactionEvent.KIND -> stringRes(R.string.reaction)
RepostEvent.KIND, GenericRepostEvent.KIND -> stringRes(R.string.boost)
VoiceEvent.KIND -> stringRes(R.string.voice_post)
VoiceReplyEvent.KIND -> stringRes(R.string.voice_reply)
else -> stringRes(R.string.post)
}
fun kindToName(kind: Int): String = stringRes(powKindLabelRes(kind))
@Preview
@Composable
@@ -37,6 +37,7 @@ import androidx.compose.runtime.remember
import androidx.compose.runtime.setValue
import androidx.compose.ui.Alignment
import androidx.compose.ui.Modifier
import androidx.compose.ui.res.pluralStringResource
import androidx.compose.ui.text.font.FontWeight
import androidx.compose.ui.tooling.preview.Preview
import androidx.compose.ui.unit.dp
@@ -110,7 +111,7 @@ fun PowOverrideButton(
text = {
Text(
if (defaultDifficulty != null && defaultDifficulty > 0) {
stringRes(R.string.pow_option_default_on, defaultDifficulty)
pluralStringResource(R.plurals.pow_option_default_on, defaultDifficulty, defaultDifficulty)
} else {
stringRes(R.string.pow_option_default_off)
},
@@ -138,7 +139,7 @@ fun PowOverrideButton(
DropdownMenuItem(
text = {
Text(
stringRes(R.string.pow_option_bits, preset),
pluralStringResource(R.plurals.pow_option_bits, preset, preset),
fontWeight = if (isOverridden && effectiveDifficulty == preset) FontWeight.Bold else null,
)
},
@@ -537,9 +537,15 @@ open class CommentPostViewModel :
val draftToDelete = draftNote
val anonymous = wantsAnonymousPost
val powDifficulty = accountViewModel.account.powDifficultyFor(template.kind, powOverride)
// captured before cancel() resets the chip
val chosenPow = powOverride
cancel()
// Draft deletion lives INSIDE each publish continuation: when the post
// is mined first, the draft must survive until the mined event is
// actually signed and dispatched — a cancelled or process-killed
// mining job would otherwise have destroyed the only copy of the text.
// A reply within a NIP-29 group is group content: pin it to the group's
// host relay (the relay the thread was seen on) instead of the author's
// outbox, so it reaches the group and — for a private/closed group — is
@@ -559,8 +565,10 @@ open class CommentPostViewModel :
if (anonymous) {
// The anonymous key signs without a client tag, so the template is
// mined as-is against the throwaway pubkey.
// mined as-is against the throwaway pubkey — and never checkpointed
// to disk, so the key and content can't outlive the process.
val anonSigner = anonymousSigner()
val powDifficulty = accountViewModel.account.powDifficultyFor(template.kind, chosenPow)
val enqueued =
powDifficulty != null &&
accountViewModel.account.mineInBackground(template.kind, powDifficulty) { isActive ->
@@ -569,37 +577,27 @@ open class CommentPostViewModel :
val fresh = EventTemplate<Event>(TimeUtils.now(), template.kind, template.tags, template.content)
val mined = PoWMiner.run(fresh, anonSigner.pubKey, powDifficulty, isActive)
accountViewModel.account.signAnonymouslyAndBroadcast(mined, extraNotesToBroadcast, anonSigner)
accountViewModel.account.deleteDraftIgnoreErrors(draftToDelete)
}
if (!enqueued) {
accountViewModel.account.signAnonymouslyAndBroadcast(template, extraNotesToBroadcast, anonSigner)
accountViewModel.account.deleteDraftIgnoreErrors(draftToDelete)
}
} else if (replyGroupId != null) {
// Group content: route to the resolved host. If it couldn't be resolved, publish to the
// parent's relays (possibly empty) rather than broadcasting to the outbox — better to
// under-deliver a group reply than to leak group participation to unrelated relays.
val relays = groupHostRelays ?: replyingTo?.relays.orEmpty()
val enqueued =
powDifficulty != null &&
accountViewModel.account.mineTemplateInBackground(template, powDifficulty, PoWReplay.ToRelays(relays)) { mined ->
accountViewModel.account.signAndSendPrivatelyOrBroadcast(mined) { relays }
}
if (!enqueued) {
accountViewModel.account.signAndSendPrivatelyOrBroadcast(template) { relays }
accountViewModel.account.sendMined(template, PoWReplay.ToRelays(relays), chosenPow) { readyTemplate ->
accountViewModel.account.signAndSendPrivatelyOrBroadcast(readyTemplate) { relays }
accountViewModel.account.deleteDraftIgnoreErrors(draftToDelete)
}
} else {
val enqueued =
powDifficulty != null &&
accountViewModel.account.mineTemplateInBackground(template, powDifficulty, PoWReplay.Broadcast(extraNotesToBroadcast)) { mined ->
accountViewModel.account.signAndComputeBroadcast(mined, extraNotesToBroadcast)
}
if (!enqueued) {
accountViewModel.account.signAndComputeBroadcast(template, extraNotesToBroadcast)
accountViewModel.account.sendMined(template, PoWReplay.Broadcast(extraNotesToBroadcast), chosenPow) { readyTemplate ->
accountViewModel.account.signAndComputeBroadcast(readyTemplate, extraNotesToBroadcast)
accountViewModel.account.deleteDraftIgnoreErrors(draftToDelete)
}
}
accountViewModel.viewModelScope.launch(Dispatchers.IO) {
accountViewModel.account.deleteDraftIgnoreErrors(draftToDelete)
}
}
suspend fun sendDraftSync() {
@@ -389,15 +389,26 @@ class AccountSessionManager(
localPreferences.deleteAccount(accountInfo)
accountsCache.removeAccount(hex)
accountsCache.deleteAccountFiles(hex)
Amethyst.instance.scheduledPostStore.removeForAccount(hex)
purgePendingPosts(hex)
loginWithDefaultAccount()
} else {
// delete without switching logins
localPreferences.deleteAccount(accountInfo)
accountsCache.removeAccount(hex)
accountsCache.deleteAccountFiles(hex)
Amethyst.instance.scheduledPostStore.removeForAccount(hex)
purgePendingPosts(hex)
}
}
}
/**
* Drops everything a deleted account left in the publish pipelines: parked
* scheduled posts, checkpointed PoW mining jobs, and any of its jobs still
* queued or mining (a post must not publish after its account is gone).
*/
private suspend fun purgePendingPosts(hex: String) {
Amethyst.instance.scheduledPostStore.removeForAccount(hex)
Amethyst.instance.powPublishQueue.cancelForOwner(hex)
Amethyst.instance.powJobStore.removeForAccount(hex)
}
}
@@ -78,6 +78,7 @@ import com.vitorpamplona.amethyst.service.checkNotInMainThread
import com.vitorpamplona.amethyst.service.lnurl.LightningAddressResolver
import com.vitorpamplona.amethyst.service.location.LocationState
import com.vitorpamplona.amethyst.service.notifications.NotificationUtils.dismissNotificationForEvent
import com.vitorpamplona.amethyst.service.pow.powKindLabelRes
import com.vitorpamplona.amethyst.service.relayClient.reqCommand.RelaySubscriptionsCoordinator
import com.vitorpamplona.amethyst.service.relayClient.reqCommand.nwc.NWCPaymentFilterAssembler
import com.vitorpamplona.amethyst.ui.actions.Dao
@@ -278,6 +279,19 @@ class AccountViewModel(
// receivers can reach callManager + accountViewModel.
com.vitorpamplona.amethyst.service.call.CallSessionBridge
.set(callManager, this)
// A mined post that fails to sign or broadcast would otherwise die
// silently — the composer already returned when it was enqueued.
viewModelScope.launch {
Amethyst.instance.powPublishQueue.failures.collect { failure ->
val kindLabel = stringRes(Amethyst.instance.appContext, powKindLabelRes(failure.kind))
if (failure.willRetryOnRestart) {
toastManager.toast(R.string.pow_settings_title, R.string.pow_publish_failed_retry, kindLabel)
} else {
toastManager.toast(R.string.pow_settings_title, R.string.pow_publish_failed, kindLabel, failure.message.orEmpty())
}
}
}
}
/**
@@ -40,6 +40,7 @@ import com.vitorpamplona.amethyst.commons.model.nip29RelayGroups.RelayGroupChann
import com.vitorpamplona.amethyst.commons.model.nip30CustomEmojis.EmojiPackState
import com.vitorpamplona.amethyst.commons.model.nip53LiveActivities.LiveActivitiesChannel
import com.vitorpamplona.amethyst.commons.richtext.UrlParser
import com.vitorpamplona.amethyst.commons.service.pow.PoWReplay
import com.vitorpamplona.amethyst.commons.ui.text.currentWord
import com.vitorpamplona.amethyst.commons.ui.text.insertUrlAtCursor
import com.vitorpamplona.amethyst.commons.ui.text.replaceCurrentWord
@@ -316,10 +317,14 @@ open class ChannelNewMessageViewModel :
val draftToDelete = draftNote
cancel()
accountViewModel.account.signAndSendPrivatelyOrBroadcast(template) {
channelRelays.toList()
}
accountViewModel.viewModelScope.launch(Dispatchers.IO) {
// Kinds 42/1311 are the PUBLIC_CHAT PoW category: route through the
// mining gate (a no-op when the category is off). Draft deletion runs
// inside the publish continuation so a cancelled mining job can't
// destroy the only copy of the text.
accountViewModel.account.sendMined(template, PoWReplay.ToRelays(channelRelays.toList())) { readyTemplate ->
accountViewModel.account.signAndSendPrivatelyOrBroadcast(readyTemplate) {
channelRelays.toList()
}
accountViewModel.account.deleteDraftIgnoreErrors(draftToDelete)
}
}
@@ -352,49 +352,37 @@ class LongFormPostViewModel :
val template = createTemplate() ?: return
val draftToDelete = draftNote
val powDifficulty = accountViewModel.account.powDifficultyFor(template.kind)
cancel()
val enqueued =
powDifficulty != null &&
accountViewModel.account.mineTemplateInBackground(template, powDifficulty, PoWReplay.Broadcast()) { mined ->
broadcastArticle(mined)
}
if (!enqueued) {
if (accountViewModel.settings.useTrackedBroadcasts()) {
val (event, relays, extras) = accountViewModel.account.createPostEvent(template, emptyList())
accountViewModel.viewModelScope.launch(Dispatchers.IO) {
accountViewModel.broadcastTracker.trackBroadcast(
event = event,
relays = relays,
client = accountViewModel.account.client,
)
accountViewModel.account.consumePostEvent(event, relays, extras)
}
} else {
accountViewModel.account.signAndComputeBroadcast(template, emptyList())
}
}
accountViewModel.launchSigner {
// Draft deletion runs INSIDE the publish continuation: when the
// article is mined first, the draft must survive until the mined event
// is actually signed and dispatched — a cancelled or process-killed
// mining job would otherwise have destroyed the only copy of the text.
accountViewModel.account.sendMined(template, PoWReplay.Broadcast()) { readyTemplate ->
broadcastArticle(readyTemplate)
accountViewModel.account.deleteDraftIgnoreErrors(draftToDelete)
}
}
/**
* The post-mining continuation: same tracked/untracked split as the direct
* path, but running on the mining queue's scope — the composer's
* viewModelScope may already be gone by the time the nonce is found.
* The publish step shared by the direct and post-mining paths. Tracked
* broadcasting is launched fire-and-forget on the account scope — the
* composer must not wait for relay acks before navigating away, and a
* mined job must not hold its queue entry while acks trickle in. Runs on
* the mining queue's scope when PoW is on, so it must not touch
* viewModelScope.
*/
private suspend fun broadcastArticle(template: EventTemplate<out Event>) {
if (accountViewModel.settings.useTrackedBroadcasts()) {
val (event, relays, extras) = accountViewModel.account.createPostEvent(template, emptyList())
accountViewModel.broadcastTracker.trackBroadcast(
event = event,
relays = relays,
client = accountViewModel.account.client,
)
accountViewModel.account.consumePostEvent(event, relays, extras)
accountViewModel.account.scope.launch {
accountViewModel.broadcastTracker.trackBroadcast(
event = event,
relays = relays,
client = accountViewModel.account.client,
)
accountViewModel.account.consumePostEvent(event, relays, extras)
}
} else {
accountViewModel.account.signAndComputeBroadcast(template, emptyList())
}
@@ -139,6 +139,7 @@ import com.vitorpamplona.quartz.nip57Zaps.splits.zapSplitSetup
import com.vitorpamplona.quartz.nip57Zaps.splits.zapSplits
import com.vitorpamplona.quartz.nip57Zaps.zapraiser.zapraiser
import com.vitorpamplona.quartz.nip57Zaps.zapraiser.zapraiserAmount
import com.vitorpamplona.quartz.nip59Giftwrap.wraps.GiftWrapEvent
import com.vitorpamplona.quartz.nip72ModCommunities.definition.CommunityDefinitionEvent
import com.vitorpamplona.quartz.nip7DThreads.ThreadEvent
import com.vitorpamplona.quartz.nip88Polls.poll.PollEvent
@@ -175,6 +176,7 @@ import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.collectLatest
import kotlinx.coroutines.launch
import java.util.UUID
enum class UserSuggestionAnchor {
MAIN_MESSAGE,
@@ -388,8 +390,11 @@ open class ShortNotePostViewModel :
var powOverride by mutableStateOf<Int?>(null)
// Best guess of the kind createTemplate() will produce, for the PoW chip.
// A private note is gift-wrapped, so what gets mined (and what the
// settings gate on) is the kind-1059 wrap, not the inner kind-1.
private fun anticipatedPowKind(): Int =
when {
wantsPrivateNote -> GiftWrapEvent.KIND
wantsPoll -> PollEvent.KIND
wantsZapPoll -> ZapPollEvent.KIND
voiceRecording != null -> VoiceEvent.KIND
@@ -981,21 +986,20 @@ open class ShortNotePostViewModel :
val scheduledFor = scheduledForSec
val privately = wantsPrivateNote
val threadTarget = groupThreadTarget
val powDifficulty = accountViewModel.account.powDifficultyFor(template.kind, powOverride)
// captured before cancel() resets the chip
val chosenPow = powOverride
cancel()
// Draft deletion lives INSIDE each publish continuation: when the post
// is mined first, the draft must survive until the mined event is
// actually signed and dispatched — a cancelled or process-killed
// mining job would otherwise have destroyed the only copy of the text.
if (threadTarget != null) {
// NIP-29 group thread: publish only to the group's host relay, never the account's
// outbox — bypass the private/scheduled/anonymous paths entirely.
val enqueued =
powDifficulty != null &&
accountViewModel.account.mineTemplateInBackground(template, powDifficulty, PoWReplay.ToRelays(threadTarget.relays)) { mined ->
accountViewModel.account.signAndSendPrivatelyOrBroadcast(mined) { threadTarget.relays }
}
if (!enqueued) {
accountViewModel.account.signAndSendPrivatelyOrBroadcast(template) { threadTarget.relays }
}
accountViewModel.launchSigner {
accountViewModel.account.sendMined(template, PoWReplay.ToRelays(threadTarget.relays), chosenPow) { readyTemplate ->
accountViewModel.account.signAndSendPrivatelyOrBroadcast(readyTemplate) { threadTarget.relays }
accountViewModel.account.deleteDraftIgnoreErrors(draftToDelete)
}
return
@@ -1005,19 +1009,22 @@ open class ShortNotePostViewModel :
// Gift-wrap to the p-tagged users instead of publishing. Private
// wins over the anonymous and scheduled modes: a locked private
// reply must never fall through to a public publish path (the UI
// hides those toggles while private mode is on).
// hides those toggles while private mode is on). The inner note and
// seals are signed inline; only wrap mining is queued — the content
// is committed (and checkpointed) by the time this returns.
@Suppress("UNCHECKED_CAST")
accountViewModel.account.sendPrivateNote(template as EventTemplate<TextNoteEvent>)
accountViewModel.launchSigner {
accountViewModel.account.deleteDraftIgnoreErrors(draftToDelete)
}
accountViewModel.account.sendPrivateNote(template as EventTemplate<TextNoteEvent>, chosenPow)
accountViewModel.account.deleteDraftIgnoreErrors(draftToDelete)
return
}
if (scheduledFor != null && !anonymous) {
// Re-stamp the template with created_at = scheduled time so the post,
// when published later, shows up at its scheduled moment in feeds
// rather than as N minutes/hours old (= compose time).
// rather than as N minutes/hours old (= compose time). Mining
// commits the future created_at into the hashed id, and the worker
// publishes the stored signed JSON verbatim, so the nonce is still
// valid at publish time.
val rescheduledTemplate =
EventTemplate<Event>(
createdAt = scheduledFor,
@@ -1026,23 +1033,12 @@ open class ShortNotePostViewModel :
content = template.content,
)
// Mining commits the future created_at into the hashed id, and the
// worker publishes the stored signed JSON verbatim, so the nonce is
// still valid at publish time.
val enqueued =
powDifficulty != null &&
accountViewModel.account.mineTemplateInBackground(
rescheduledTemplate,
powDifficulty,
PoWReplay.Schedule(scheduledFor, extraNotesToBroadcast),
) { mined ->
storeScheduledPost(mined, extraNotesToBroadcast, scheduledFor)
}
if (!enqueued) {
storeScheduledPost(rescheduledTemplate, extraNotesToBroadcast, scheduledFor)
}
accountViewModel.launchSigner {
accountViewModel.account.sendMined(
rescheduledTemplate,
PoWReplay.Schedule(scheduledFor, extraNotesToBroadcast),
chosenPow,
) { readyTemplate ->
storeScheduledPost(readyTemplate, extraNotesToBroadcast, scheduledFor)
accountViewModel.account.deleteDraftIgnoreErrors(draftToDelete)
}
return
@@ -1050,8 +1046,10 @@ open class ShortNotePostViewModel :
if (anonymous) {
// The anonymous key signs without a client tag, so the template is
// mined as-is against the throwaway pubkey.
// mined as-is against the throwaway pubkey — and never checkpointed
// to disk, so the key and content can't outlive the process.
val anonSigner = anonymousSigner()
val powDifficulty = accountViewModel.account.powDifficultyFor(template.kind, chosenPow)
val enqueued =
powDifficulty != null &&
accountViewModel.account.mineInBackground(template.kind, powDifficulty) { isActive ->
@@ -1060,38 +1058,17 @@ open class ShortNotePostViewModel :
val fresh = EventTemplate<Event>(TimeUtils.now(), template.kind, template.tags, template.content)
val mined = PoWMiner.run(fresh, anonSigner.pubKey, powDifficulty, isActive)
accountViewModel.account.signAnonymouslyAndBroadcast(mined, extraNotesToBroadcast, anonSigner)
accountViewModel.account.deleteDraftIgnoreErrors(draftToDelete)
}
if (!enqueued) {
accountViewModel.account.signAnonymouslyAndBroadcast(template, extraNotesToBroadcast, anonSigner)
accountViewModel.account.deleteDraftIgnoreErrors(draftToDelete)
}
} else {
val enqueued =
powDifficulty != null &&
accountViewModel.account.mineTemplateInBackground(template, powDifficulty, PoWReplay.Broadcast(extraNotesToBroadcast)) { mined ->
broadcastPublicPost(mined, extraNotesToBroadcast)
}
if (!enqueued) {
if (accountViewModel.settings.useTrackedBroadcasts()) {
// Tracked broadcasting with progress feedback (non-blocking)
val (event, relays, extras) = accountViewModel.account.createPostEvent(template, extraNotesToBroadcast)
// Launch broadcast in background - don't wait for completion
accountViewModel.viewModelScope.launch(Dispatchers.IO) {
accountViewModel.broadcastTracker.trackBroadcast(
event = event,
relays = relays,
client = accountViewModel.account.client,
)
accountViewModel.account.consumePostEvent(event, relays, extras)
}
} else {
// Fire-and-forget (original behavior)
accountViewModel.account.signAndComputeBroadcast(template, extraNotesToBroadcast)
}
}
return
}
accountViewModel.launchSigner {
accountViewModel.account.sendMined(template, PoWReplay.Broadcast(extraNotesToBroadcast), chosenPow) { readyTemplate ->
broadcastPublicPost(readyTemplate, extraNotesToBroadcast)
accountViewModel.account.deleteDraftIgnoreErrors(draftToDelete)
}
}
@@ -1110,10 +1087,7 @@ open class ShortNotePostViewModel :
val (event, relays, extras) = accountViewModel.account.createPostEvent(template, extraNotesToBroadcast)
Amethyst.instance.scheduledPostStore.add(
ScheduledPost(
id =
java.util.UUID
.randomUUID()
.toString(),
id = UUID.randomUUID().toString(),
accountPubkey = event.pubKey,
signedEventJson = event.toJson(),
relayUrls = relays.map { it.url },
@@ -1125,9 +1099,13 @@ open class ShortNotePostViewModel :
}
/**
* The post-mining continuation: same tracked/untracked split as the direct
* path, but running on the mining queue's scope — the composer's
* viewModelScope may already be gone by the time the nonce is found.
* The publish step shared by the direct and post-mining paths: signs the
* template and dispatches the broadcast. Tracked broadcasting is launched
* fire-and-forget on the account scope — the composer must not wait for
* relay acks before navigating away, and a mined job must not hold its
* queue entry while acks trickle in (the event is signed and dispatched;
* only progress reporting remains). Runs on the mining queue's scope when
* PoW is on, so it must not touch viewModelScope.
*/
private suspend fun broadcastPublicPost(
template: EventTemplate<out Event>,
@@ -1135,12 +1113,14 @@ open class ShortNotePostViewModel :
) {
if (accountViewModel.settings.useTrackedBroadcasts()) {
val (event, relays, extras) = accountViewModel.account.createPostEvent(template, extraNotesToBroadcast)
accountViewModel.broadcastTracker.trackBroadcast(
event = event,
relays = relays,
client = accountViewModel.account.client,
)
accountViewModel.account.consumePostEvent(event, relays, extras)
accountViewModel.account.scope.launch {
accountViewModel.broadcastTracker.trackBroadcast(
event = event,
relays = relays,
client = accountViewModel.account.client,
)
accountViewModel.account.consumePostEvent(event, relays, extras)
}
} else {
accountViewModel.account.signAndComputeBroadcast(template, extraNotesToBroadcast)
}
@@ -28,6 +28,7 @@ import androidx.lifecycle.ViewModel
import androidx.lifecycle.viewModelScope
import com.vitorpamplona.amethyst.Amethyst
import com.vitorpamplona.amethyst.R
import com.vitorpamplona.amethyst.commons.service.pow.PoWReplay
import com.vitorpamplona.amethyst.commons.util.deleteOrWarn
import com.vitorpamplona.amethyst.model.Note
import com.vitorpamplona.amethyst.service.uploads.CompressorQuality
@@ -251,8 +252,11 @@ class VoiceReplyViewModel : ViewModel() {
// Check if replying to a voice event
val voiceHint = note.toEventHint<BaseVoiceEvent>()
if (voiceHint != null) {
// Create VoiceReplyEvent (KIND 1244) for voice-to-voice replies
accountViewModel.account.signAndComputeBroadcast(VoiceReplyEvent.build(audioMeta, voiceHint))
// Create VoiceReplyEvent (KIND 1244) for voice-to-voice replies,
// routed through the PoW gate (VOICE category; no-op when off)
accountViewModel.account.sendMined(VoiceReplyEvent.build(audioMeta, voiceHint), PoWReplay.Broadcast()) {
accountViewModel.account.signAndComputeBroadcast(it)
}
} else {
// Create TextNoteEvent (KIND 1) with audio IMeta for voice replies to regular notes
val textHint = note.toEventHint<TextNoteEvent>()
@@ -270,7 +274,9 @@ class VoiceReplyViewModel : ViewModel() {
// Add audio as IMeta attachment
add(audioMeta.toIMetaArray())
}
accountViewModel.account.signAndComputeBroadcast(template)
accountViewModel.account.sendMined(template, PoWReplay.Broadcast()) {
accountViewModel.account.signAndComputeBroadcast(it)
}
}
accountViewModel.account.settings.changeDefaultFileServer(server)
@@ -20,6 +20,7 @@
*/
package com.vitorpamplona.amethyst.ui.screen.loggedIn.settings
import android.content.Context
import androidx.annotation.StringRes
import androidx.compose.foundation.layout.Arrangement
import androidx.compose.foundation.layout.Column
@@ -40,6 +41,7 @@ import androidx.compose.runtime.collectAsState
import androidx.compose.runtime.getValue
import androidx.compose.runtime.produceState
import androidx.compose.ui.Modifier
import androidx.compose.ui.platform.LocalContext
import androidx.compose.ui.tooling.preview.Preview
import androidx.compose.ui.unit.dp
import androidx.lifecycle.compose.collectAsStateWithLifecycle
@@ -54,6 +56,7 @@ import com.vitorpamplona.amethyst.model.UiSettingsFlow
import com.vitorpamplona.amethyst.ui.navigation.navs.INav
import com.vitorpamplona.amethyst.ui.navigation.topbars.TopBarWithBackButton
import com.vitorpamplona.amethyst.ui.note.creators.pow.POW_PRESETS
import com.vitorpamplona.amethyst.ui.pluralStringRes
import com.vitorpamplona.amethyst.ui.screen.loggedIn.AccountViewModel
import com.vitorpamplona.amethyst.ui.screen.loggedIn.mockAccountViewModel
import com.vitorpamplona.amethyst.ui.stringRes
@@ -191,10 +194,11 @@ private fun PowDifficultyTile(accountViewModel: AccountViewModel) {
private fun PowTimeEstimate(difficulty: Int) {
if (difficulty <= 0) return
val context = LocalContext.current
val estimate by
produceState<String?>(initialValue = null, difficulty) {
val rate = PoWEstimator.hashesPerSecond()
value = formatEstimate(PoWEstimator.estimateSeconds(difficulty, rate))
value = formatEstimate(context, PoWEstimator.estimateSeconds(difficulty, rate))
}
estimate?.let {
@@ -206,14 +210,23 @@ private fun PowTimeEstimate(difficulty: Int) {
}
}
private fun formatEstimate(seconds: Double): String =
when {
seconds < 1.0 -> "<1s"
seconds < 90.0 -> "${seconds.roundToLong()}s"
seconds < 90.0 * 60.0 -> "${(seconds / 60.0).roundToLong()}m"
seconds < 48.0 * 3600.0 -> "${(seconds / 3600.0).roundToLong()}h"
else -> "${(seconds / 86400.0).roundToLong()}d"
private fun formatEstimate(
context: Context,
seconds: Double,
): String {
fun quantity(
id: Int,
count: Long,
) = pluralStringRes(context, id, count.toInt(), count.toInt())
return when {
seconds < 1.0 -> stringRes(context, R.string.pow_estimate_instant)
seconds < 90.0 -> quantity(R.plurals.pow_estimate_seconds, seconds.roundToLong())
seconds < 90.0 * 60.0 -> quantity(R.plurals.pow_estimate_minutes, (seconds / 60.0).roundToLong())
seconds < 48.0 * 3600.0 -> quantity(R.plurals.pow_estimate_hours, (seconds / 3600.0).roundToLong())
else -> quantity(R.plurals.pow_estimate_days, (seconds / 86400.0).roundToLong())
}
}
@Composable
private fun PowCategoryChecklist(accountViewModel: AccountViewModel) {
+41 -7
View File
@@ -3764,19 +3764,53 @@
<string name="pow_category_other_public">Other public content</string>
<string name="pow_category_other_public_explainer">Polls, live statuses, classifieds and everything else public</string>
<string name="pow_mining_title">Mining proof of work</string>
<string name="pow_mining_progress">Mining proof of work… (%1$s in queue)</string>
<string name="pow_mining_job">%1$s • mining at %2$s bits</string>
<string name="pow_queued_job">%1$s • waiting to mine at %2$s bits</string>
<string name="pow_chip_active">PoW %1$d</string>
<string name="pow_chip_off">PoW off</string>
<string name="pow_option_default_on">Default (%1$d bits)</string>
<plurals name="pow_mining_progress">
<item quantity="one">Mining proof of work… (%1$d post in queue)</item>
<item quantity="other">Mining proof of work… (%1$d posts in queue)</item>
</plurals>
<plurals name="pow_mining_job">
<item quantity="one">%1$s • mining at %2$d bit</item>
<item quantity="other">%1$s • mining at %2$d bits</item>
</plurals>
<plurals name="pow_queued_job">
<item quantity="one">%1$s • waiting to mine at %2$d bit</item>
<item quantity="other">%1$s • waiting to mine at %2$d bits</item>
</plurals>
<plurals name="pow_option_default_on">
<item quantity="one">Default (%1$d bit)</item>
<item quantity="other">Default (%1$d bits)</item>
</plurals>
<string name="pow_option_default_off">Default (off)</string>
<string name="pow_option_off">Off for this post</string>
<string name="pow_option_bits">%1$d bits</string>
<plurals name="pow_option_bits">
<item quantity="one">%1$d bit</item>
<item quantity="other">%1$d bits</item>
</plurals>
<string name="pow_notification_channel_name">Proof of work mining</string>
<string name="pow_notification_channel_description">Shows posts that are still mining their NIP-13 proof of work so they can finish after you leave the app.</string>
<string name="pow_notification_cancel_all">Cancel</string>
<string name="pow_difficulty_estimate">≈ %1$s per post on this device</string>
<string name="pow_estimate_instant">&lt;1 s</string>
<plurals name="pow_estimate_seconds">
<item quantity="one">%1$d second</item>
<item quantity="other">%1$d seconds</item>
</plurals>
<plurals name="pow_estimate_minutes">
<item quantity="one">%1$d minute</item>
<item quantity="other">%1$d minutes</item>
</plurals>
<plurals name="pow_estimate_hours">
<item quantity="one">%1$d hour</item>
<item quantity="other">%1$d hours</item>
</plurals>
<plurals name="pow_estimate_days">
<item quantity="one">%1$d day</item>
<item quantity="other">%1$d days</item>
</plurals>
<string name="pow_publish_failed">Your %1$s finished mining but could not be published: %2$s</string>
<string name="pow_publish_failed_retry">Your %1$s finished mining but could not be published. It will be retried the next time the app starts.</string>
<string name="pow_kind_chat_message">Chat message</string>
<string name="pow_kind_report">Report</string>
<string name="ai_writing_use_this">Use This</string>
<string name="ai_writing_dismiss">Dismiss</string>
<string name="ai_tone_correct">Correct</string>
@@ -24,6 +24,7 @@ import com.vitorpamplona.amethyst.cli.Args
import com.vitorpamplona.amethyst.cli.Context
import com.vitorpamplona.amethyst.cli.DataDir
import com.vitorpamplona.amethyst.cli.Output
import com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
import com.vitorpamplona.quartz.nip10Notes.TextNoteEvent
import com.vitorpamplona.quartz.nip13Pow.miner.PoWMiner
import com.vitorpamplona.quartz.nip13Pow.pow
@@ -69,11 +70,7 @@ object PostCommand {
Context.open(dataDir).use { ctx ->
ctx.prepare()
val outbox = ctx.outboxRelays()
val extraNormalized =
extraRelays.mapNotNull {
com.vitorpamplona.quartz.nip01Core.relay.normalizer.RelayUrlNormalizer
.normalizeOrNull(it)
}
val extraNormalized = extraRelays.mapNotNull { RelayUrlNormalizer.normalizeOrNull(it) }
val targets = (outbox + extraNormalized).toSet()
if (targets.isEmpty()) {
return Output.error("no_relays", "no outbox relays configured; pass --relay or run `amy relay add`")
@@ -34,6 +34,7 @@ import com.vitorpamplona.quartz.nip13Pow.commitedPoW
import com.vitorpamplona.quartz.nip13Pow.miner.PoWMiner
import com.vitorpamplona.quartz.nip13Pow.miner.PoWRankEvaluator
import com.vitorpamplona.quartz.nip13Pow.pow
import com.vitorpamplona.quartz.utils.Hex
import com.vitorpamplona.quartz.utils.sha256.sha256
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.withContext
@@ -49,6 +50,9 @@ import kotlin.math.roundToLong
* pass `--pubkey`, or omit it to mine for the active account.
*/
object PowCommands {
// Deliberately above PoWPolicy.MAX_DIFFICULTY (the app's UI ceiling, 40):
// amy is a power tool that may mine on beefy hardware or run deliberate
// long jobs. Still bounded well under the miner's hard 256-bit limit.
private const val MAX_DIFFICULTY = 64
suspend fun dispatch(
@@ -119,14 +123,19 @@ object PowCommands {
return Output.error("bad_template", e.message)
}
// normalized to lowercase: the pubkey is serialized verbatim into the
// id preimage, and NIP-01 ids/keys are lowercase hex — an uppercase
// pubkey would mine an id that never matches the signed event.
val pubKey =
args.flags["pubkey"]
?: try {
Context.open(dataDir).use { it.signer.pubKey }
} catch (e: Exception) {
return Output.error("bad_args", "no account available; pass --pubkey (${e.message})")
}
if (pubKey.length != 64 || pubKey.any { it !in "0123456789abcdefABCDEF" }) {
(
args.flags["pubkey"]
?: try {
Context.open(dataDir).use { it.signer.pubKey }
} catch (e: Exception) {
return Output.error("bad_args", "no account available; pass --pubkey (${e.message})")
}
).lowercase()
if (pubKey.length != 64 || !Hex.isHex(pubKey)) {
return Output.error("bad_args", "--pubkey must be 64 hex characters")
}
@@ -23,6 +23,8 @@ package com.vitorpamplona.amethyst.commons.service.pow
import com.vitorpamplona.quartz.utils.sha256.sha256
import kotlinx.coroutines.CoroutineDispatcher
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.sync.Mutex
import kotlinx.coroutines.sync.withLock
import kotlinx.coroutines.withContext
import kotlin.math.pow
import kotlin.time.Duration.Companion.milliseconds
@@ -47,10 +49,16 @@ object PoWEstimator {
private val BENCH_DURATION = 250.milliseconds
private var cachedRate: Double? = null
private val benchLock = Mutex()
suspend fun hashesPerSecond(dispatcher: CoroutineDispatcher = Dispatchers.Default): Double =
cachedRate ?: withContext(dispatcher) {
cachedRate ?: benchmark().also { cachedRate = it }
// single-flight: concurrent first callers (e.g. the settings screen
// recomposing while the composer chip opens) share one ~250 ms
// benchmark instead of each burning a core.
benchLock.withLock {
cachedRate ?: benchmark().also { cachedRate = it }
}
}
fun estimateSeconds(
@@ -36,17 +36,26 @@ data class PersistedPoWJob(
val accountPubkey: String,
val kind: Int,
val difficulty: Int,
/** Unsigned template for the template replay types; empty for [REPLAY_WRAPS]. */
val templateJson: String,
val replayType: String,
val relayUrls: List<String> = emptyList(),
/** Extra pre-signed events to broadcast; for [REPLAY_WRAPS], the signed seals. */
val extraEventsJson: List<String> = emptyList(),
val publishAtSec: Long? = null,
/** [REPLAY_WRAPS]: recipient of each seal in [extraEventsJson], same order. */
val recipientPubkeys: List<String> = emptyList(),
/** [REPLAY_WRAPS]: expiration delta to stamp on each wrap. */
val wrapExpirationDelta: Long? = null,
val createdAtSec: Long = 0,
) {
companion object {
const val REPLAY_BROADCAST = "broadcast"
const val REPLAY_RELAYS = "relays"
const val REPLAY_SCHEDULE = "schedule"
/** Mine one gift wrap per pre-signed seal, then broadcast the wraps. */
const val REPLAY_WRAPS = "wraps"
}
}
@@ -20,6 +20,11 @@
*/
package com.vitorpamplona.amethyst.commons.service.pow
import com.vitorpamplona.quartz.nip01Core.core.Kind
import com.vitorpamplona.quartz.nip01Core.core.isEphemeral
import com.vitorpamplona.quartz.nip01Core.core.isReplaceable
import com.vitorpamplona.quartz.nip01Core.metadata.MetadataEvent
import com.vitorpamplona.quartz.nip02FollowList.ContactListEvent
import com.vitorpamplona.quartz.nip03Timestamp.OtsEvent
import com.vitorpamplona.quartz.nip10Notes.TextNoteEvent
import com.vitorpamplona.quartz.nip18Reposts.GenericRepostEvent
@@ -28,6 +33,7 @@ import com.vitorpamplona.quartz.nip22Comments.CommentEvent
import com.vitorpamplona.quartz.nip23LongContent.LongTextNoteEvent
import com.vitorpamplona.quartz.nip25Reactions.ReactionEvent
import com.vitorpamplona.quartz.nip28PublicChat.message.ChannelMessageEvent
import com.vitorpamplona.quartz.nip37Drafts.DraftWrapEvent
import com.vitorpamplona.quartz.nip53LiveActivities.chat.LiveActivitiesChatMessageEvent
import com.vitorpamplona.quartz.nip56Reports.ReportEvent
import com.vitorpamplona.quartz.nip57Zaps.LnZapRequestEvent
@@ -73,7 +79,13 @@ enum class PoWCategory(
* keystroke debounce.
*/
object PoWPolicy {
private const val DRAFT_WRAP_KIND = 31234 // quartz's DraftWrapEvent (NIP-37)
/**
* Practical UI ceiling for the difficulty setting: above ~40 bits a phone
* would mine for days; treat anything larger (including values arriving
* from a synced NIP-78 settings event) as a config error and clamp.
*/
const val MAX_DIFFICULTY = 40
private const val LONG_FORM_DRAFT_KIND = 30024
/** NIP-51 sets and other settings-like addressable kinds. */
@@ -95,26 +107,26 @@ object PoWPolicy {
private val NEVER_EXPLICIT =
setOf(
0, // metadata
3, // contact list
MetadataEvent.KIND,
ContactListEvent.KIND,
LnZapRequestEvent.KIND, // blocks the invoice fetch
OtsEvent.KIND, // machine-generated companion events
DRAFT_WRAP_KIND, // re-signed on a 1s debounce while typing
DraftWrapEvent.KIND, // re-signed on a 1s debounce while typing
)
/**
* Kinds that must never be mined regardless of user settings.
*
* The replaceable range (10000..19999) covers relay lists, NIP-51 standard
* lists, NWC info and other settings sync; the ephemeral range
* (20000..29999) covers relay AUTH (22242), NWC RPC (23194..23196), NIP-46
* bunker messages (24133), Blossom auth (24242) and HTTP auth (27235) —
* all time-critical request/response events where mining only adds latency.
* The replaceable range covers relay lists, NIP-51 standard lists, NWC
* info and other settings sync; the ephemeral range covers relay AUTH
* (22242), NWC RPC (23194..23196), NIP-46 bunker messages (24133),
* Blossom auth (24242) and HTTP auth (27235) — all time-critical
* request/response events where mining only adds latency.
*/
fun neverMine(kind: Int): Boolean =
fun neverMine(kind: Kind): Boolean =
kind in NEVER_EXPLICIT ||
kind in 10000..19999 ||
kind in 20000..29999 ||
kind.isReplaceable() ||
kind.isEphemeral() ||
kind in NEVER_ADDRESSABLE
fun categoryOf(kind: Int): PoWCategory =
@@ -133,7 +145,10 @@ object PoWPolicy {
/**
* Returns the difficulty to mine [kind] at, or null when the event should
* be published without proof of work.
* be published without proof of work. The difficulty is clamped to
* [MAX_DIFFICULTY] as defense in depth against out-of-range values synced
* from other clients — an unclamped 300 would crash the miner and 41+
* would mine effectively forever.
*/
fun shouldMine(
kind: Int,
@@ -143,6 +158,6 @@ object PoWPolicy {
if (difficulty <= 0) return null
if (neverMine(kind)) return null
if (categoryOf(kind) !in enabledCategories) return null
return difficulty
return difficulty.coerceAtMost(MAX_DIFFICULTY)
}
}
@@ -39,8 +39,11 @@ import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.channels.Channel
import kotlinx.coroutines.channels.Channel.Factory.UNLIMITED
import kotlinx.coroutines.currentCoroutineContext
import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.SharedFlow
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asSharedFlow
import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.flow.update
import kotlinx.coroutines.isActive
@@ -49,18 +52,41 @@ import kotlinx.coroutines.launch
import kotlin.concurrent.Volatile
import kotlin.coroutines.cancellation.CancellationException
enum class PoWJobPhase {
/** Waiting for a mining worker. Cancellable. */
QUEUED,
/** A worker is searching for the nonce. Cancellable. */
MINING,
/** Nonce found; signing and broadcasting. No longer cancellable. */
PUBLISHING,
}
/**
* Snapshot of one queued/mining publish job, for display in the broadcast banner.
* [miningStartedAt] (epoch seconds) is set when a worker picks the job up, so
* the UI can show how long the current nonce search has been running.
* Snapshot of one queued/mining/publishing job, for display in the broadcast
* banner. [miningStartedAt] (epoch seconds) is set when a worker picks the
* job up, so the UI can show how long the nonce search has been running.
*/
@Immutable
data class PoWJobState(
val id: String,
val kind: Int,
val difficulty: Int,
val isMining: Boolean,
val phase: PoWJobPhase = PoWJobPhase.QUEUED,
val miningStartedAt: Long? = null,
) {
val isMining: Boolean get() = phase == PoWJobPhase.MINING
val isCancellable: Boolean get() = phase != PoWJobPhase.PUBLISHING
}
/** Emitted when a job dies after mining or while publishing, so the UI can toast. */
@Immutable
data class PoWJobFailure(
val kind: Int,
/** true when a checkpoint survived and the restorer will retry it on next login. */
val willRetryOnRestart: Boolean,
val message: String?,
)
/**
@@ -70,14 +96,17 @@ data class PoWJobState(
* to the normal sign+broadcast continuation captured at enqueue time.
*
* FIFO with at most [maxConcurrent] concurrent miners so a burst of posts
* queues up instead of spawning unbounded CPU work. Each job is cancellable
* while queued or mining.
* queues up instead of spawning unbounded CPU work. Jobs are cancellable
* while QUEUED or MINING; once PUBLISHING starts, cancel is a no-op.
*
* Template jobs enqueued with a [PersistedPoWJob] record are checkpointed to
* [persistence] and removed when they finish or are cancelled, so the platform
* layer can re-enqueue them after process death. Opaque [enqueueWork] jobs
* (reactions, reposts, gift wraps — and anonymous posts, deliberately, so a
* throwaway key and its content never touch disk) stay in-memory only.
* [persistence]. The checkpoint is deleted only after the [onMined]
* continuation COMPLETES (publish included) or the user cancels — never at
* mining-complete — so a process death mid-sign or mid-broadcast is replayed
* by the restorer on the next login. A continuation that throws keeps its
* checkpoint (retry on restart) and reports through [failures]. Opaque
* [enqueueWork] jobs (reactions, reposts — and anonymous posts, deliberately,
* so a throwaway key and its content never touch disk) stay in-memory only.
*
* [onQueueActive] fires on every enqueue; the Android layer uses it to start
* the mining foreground service so backgrounding the app doesn't freeze the
@@ -94,19 +123,36 @@ class PoWPublishQueue(
val id: String,
val kind: Int,
val difficulty: Int,
val persisted: Boolean,
val dedupeKey: String?,
val owner: HexKey?,
val work: suspend (isActive: () -> Boolean) -> Unit,
) {
@Volatile
var cancelled = false
@Volatile
var publishing = false
}
private val queue = Channel<MiningJob>(UNLIMITED)
private val _jobs = MutableStateFlow<ImmutableList<PoWJobState>>(persistentListOf())
/** Queued + currently-mining jobs, in enqueue order. */
/** Queued + mining + publishing jobs, in enqueue order. */
val jobs: StateFlow<ImmutableList<PoWJobState>> = _jobs.asStateFlow()
private val _failures = MutableSharedFlow<PoWJobFailure>(extraBufferCapacity = 16)
/** Post-mining failures (signer rejected, broadcast threw). */
val failures: SharedFlow<PoWJobFailure> = _failures.asSharedFlow()
// Jobs the workers haven't finished yet, so cancel() can reach the flag of
// a job that is still sitting in the channel. StateFlow.update gives us
// atomic CAS updates across the UI thread and the mining workers. This map
// and _jobs are only ever mutated together inside addJob/setPhase/removeEntry.
private val pending = MutableStateFlow<PersistentMap<String, MiningJob>>(persistentMapOf())
init {
repeat(maxConcurrent.coerceAtLeast(1)) {
scope.launch(miningDispatcher) {
@@ -120,12 +166,12 @@ class PoWPublishQueue(
/**
* Mines [template] at [difficulty] and hands the mined template to
* [onMined] on the queue's scope. [onMined] should run the exact
* sign+broadcast path the caller would have used without PoW.
* sign+broadcast path the caller would have used without PoW; the job
* (and its checkpoint) lives until that continuation finishes.
*
* When [persistAs] is given, the job is checkpointed (under the record's
* id) until it finishes or is cancelled, so it can be restored after
* process death. Re-enqueueing an id already in the queue is a no-op —
* that makes restore-on-login idempotent.
* When [persistAs] is given, the job is checkpointed under the record's
* id. Re-enqueueing an id already in the queue is a no-op — that makes
* restore-on-login idempotent.
*
* [refreshCreatedAtOnStart] re-stamps the template's created_at to "now"
* when a worker picks the job up — NIP-13 recommends updating created_at
@@ -140,22 +186,55 @@ class PoWPublishQueue(
persistAs: PersistedPoWJob? = null,
refreshCreatedAtOnStart: Boolean = false,
onMined: suspend (EventTemplate<T>) -> Unit,
) = addJob(
id = persistAs?.id ?: RandomInstance.randomChars(16),
) = enqueueStaged(
kind = template.kind,
difficulty = difficulty,
persistAs = persistAs,
) { isActive ->
val toMine =
if (refreshCreatedAtOnStart) {
EventTemplate<T>(TimeUtils.now(), template.kind, template.tags, template.content)
} else {
template
}
val mined = PoWMiner.run(toMine, pubKey, difficulty, isActive)
// frees the mining worker: signing may wait on an external signer
// (Amber/bunker) and broadcasting is IO, neither belongs on the pool.
scope.launch { onMined(mined) }
owner = pubKey,
mine = { isActive ->
val toMine =
if (refreshCreatedAtOnStart) {
EventTemplate<T>(TimeUtils.now(), template.kind, template.tags, template.content)
} else {
template
}
PoWMiner.run(toMine, pubKey, difficulty, isActive)
},
publish = onMined,
)
/**
* The staged primitive behind [enqueue]: [mine] runs on the capped worker
* pool (CPU only — it must not touch the user's signer or the network);
* [publish] runs detached on the queue's scope so a slow external signer
* or broadcast never holds a mining slot. The job entry and its optional
* checkpoint live until [publish] completes.
*/
fun <R> enqueueStaged(
kind: Int,
difficulty: Int,
persistAs: PersistedPoWJob? = null,
dedupeKey: String? = null,
owner: HexKey? = null,
mine: suspend (isActive: () -> Boolean) -> R,
publish: suspend (R) -> Unit,
) {
val id = persistAs?.id ?: RandomInstance.randomChars(16)
addJob(
id = id,
kind = kind,
difficulty = difficulty,
persistAs = persistAs,
dedupeKey = dedupeKey,
owner = owner ?: persistAs?.accountPubkey,
) { isActive ->
val mined = mine(isActive)
// frees the mining worker: signing may wait on an external signer
// (Amber/bunker) and broadcasting is IO, neither belongs on the
// pool. The job entry + checkpoint survive until publish finishes.
setPhase(id, PoWJobPhase.PUBLISHING)
scope.launch { finishDetached(id, kind, persisted = persistAs != null) { publish(mined) } }
}
}
/**
@@ -164,29 +243,41 @@ class PoWPublishQueue(
* recipient's ephemeral-key wrap is mined right before its local
* signature) and by anonymous posts, whose throwaway key must not be
* written to disk. Lost on process death.
*
* A non-null [dedupeKey] makes the enqueue idempotent against pending
* jobs with the same key (see [cancelByKey] for toggle semantics).
*/
fun enqueueWork(
kind: Int,
difficulty: Int,
dedupeKey: String? = null,
owner: HexKey? = null,
work: suspend (isActive: () -> Boolean) -> Unit,
) = addJob(RandomInstance.randomChars(16), kind, difficulty, persistAs = null, work = work)
) = addJob(RandomInstance.randomChars(16), kind, difficulty, persistAs = null, dedupeKey = dedupeKey, owner = owner, work = work)
private fun addJob(
id: String,
kind: Int,
difficulty: Int,
persistAs: PersistedPoWJob?,
dedupeKey: String?,
owner: HexKey?,
work: suspend (isActive: () -> Boolean) -> Unit,
) {
if (_jobs.value.any { it.id == id }) {
val current = pending.value
if (current.containsKey(id)) {
Log.d(TAG) { "PoW job $id already queued; skipping duplicate enqueue" }
return
}
if (dedupeKey != null && current.values.any { it.dedupeKey == dedupeKey && !it.cancelled }) {
Log.d(TAG) { "PoW job with key $dedupeKey already queued; skipping duplicate enqueue" }
return
}
val job = MiningJob(id, kind, difficulty, work)
val job = MiningJob(id, kind, difficulty, persisted = persistAs != null, dedupeKey = dedupeKey, owner = owner, work = work)
persistAs?.let { persistence?.save(it) }
pending.update { it.put(job.id, job) }
_jobs.update { (it + PoWJobState(job.id, job.kind, job.difficulty, isMining = false)).toImmutableList() }
_jobs.update { (it + PoWJobState(job.id, job.kind, job.difficulty)).toImmutableList() }
Log.d(TAG) {
val durability = if (persistAs != null) "persisted" else "in-memory only, lost on process death"
"Enqueued PoW job ${job.id} kind=${job.kind} difficulty=${job.difficulty} ($durability)"
@@ -195,58 +286,128 @@ class PoWPublishQueue(
onQueueActive()
}
/** Cancels a queued or mining job. No-op if the job already finished. */
/**
* Cancels a queued or mining job. No-op if the job already finished or is
* already publishing (the event may be half-signed/half-broadcast — there
* is nothing safe to abort).
*/
fun cancel(jobId: String) {
val job = pending.value[jobId] ?: return
if (job.publishing) return
job.cancelled = true
remove(jobId)
removeEntry(jobId, dropCheckpoint = true, persisted = job.persisted)
Log.d(TAG) { "Cancelled PoW job $jobId" }
}
/**
* Cancels the pending job carrying [dedupeKey], if any. Returns true when
* a job was cancelled — callers use this for toggle semantics (tapping
* "like" again while the first like is still mining un-likes it).
*/
fun cancelByKey(dedupeKey: String): Boolean {
val job = pending.value.values.firstOrNull { it.dedupeKey == dedupeKey && !it.publishing } ?: return false
cancel(job.id)
return true
}
/** Cancels everything still queued or mining. */
fun cancelAll() {
pending.value.keys.forEach { cancel(it) }
}
// Jobs the workers haven't finished yet, so cancel() can reach the flag of
// a job that is still sitting in the channel. StateFlow.update gives us
// atomic CAS updates across the UI thread and the mining workers.
private val pending = MutableStateFlow<PersistentMap<String, MiningJob>>(persistentMapOf())
/**
* Cancels every queued/mining job enqueued for [owner]'s account — used at
* log-off so a deleted account's posts can't publish after the fact.
* Publishing jobs are left to finish (cancel is unsafe mid-broadcast).
*/
fun cancelForOwner(owner: HexKey) {
pending.value.values
.filter { it.owner == owner }
.forEach { cancel(it.id) }
}
private suspend fun process(job: MiningJob) {
if (job.cancelled) {
remove(job.id)
removeEntry(job.id, dropCheckpoint = true, persisted = job.persisted)
return
}
markMining(job.id)
setPhase(job.id, PoWJobPhase.MINING)
val workerJob = currentCoroutineContext().job
var detached = false
try {
job.work { !job.cancelled && workerJob.isActive }
Log.d(TAG) { "Finished PoW job ${job.id} kind=${job.kind} difficulty=${job.difficulty}" }
detached = job.publishing
Log.d(TAG) { "Finished mining PoW job ${job.id} kind=${job.kind} difficulty=${job.difficulty}" }
} catch (e: CancellationException) {
// the worker itself was cancelled (scope teardown): propagate.
if (!currentCoroutineContext().isActive) throw e
Log.d(TAG) { "PoW job ${job.id} cancelled while mining" }
} catch (e: Exception) {
Log.w(TAG, "PoW job ${job.id} kind=${job.kind} failed", e)
// mining-stage failures are deterministic (bad difficulty, broken
// template): drop the checkpoint too, or every restore re-crashes.
Log.w(TAG, "PoW job ${job.id} kind=${job.kind} failed while mining", e)
_failures.tryEmit(PoWJobFailure(job.kind, willRetryOnRestart = false, message = e.message))
} finally {
remove(job.id)
// detached template jobs remove themselves in finishDetached once
// the sign+broadcast continuation completes.
if (!detached) {
removeEntry(job.id, dropCheckpoint = true, persisted = job.persisted)
}
}
}
private fun markMining(jobId: String) {
/**
* Runs the post-mining continuation off the worker pool. Success drops the
* checkpoint; failure keeps it (the restorer replays it headlessly on the
* next login) and surfaces through [failures].
*/
private suspend fun finishDetached(
jobId: String,
kind: Int,
persisted: Boolean,
onMined: suspend () -> Unit,
) {
try {
onMined()
Log.d(TAG) { "Published PoW job $jobId kind=$kind" }
removeEntry(jobId, dropCheckpoint = true, persisted = persisted)
} catch (e: CancellationException) {
throw e
} catch (e: Exception) {
Log.w(TAG, "PoW job $jobId kind=$kind failed after mining", e)
removeEntry(jobId, dropCheckpoint = false, persisted = persisted)
_failures.tryEmit(PoWJobFailure(kind, willRetryOnRestart = persisted, message = e.message))
}
}
private fun setPhase(
jobId: String,
phase: PoWJobPhase,
) {
if (phase == PoWJobPhase.PUBLISHING) pending.value[jobId]?.publishing = true
_jobs.update { list ->
list.map { if (it.id == jobId) it.copy(isMining = true, miningStartedAt = TimeUtils.now()) else it }.toImmutableList()
list
.map {
when {
it.id != jobId -> it
phase == PoWJobPhase.MINING -> it.copy(phase = phase, miningStartedAt = TimeUtils.now())
else -> it.copy(phase = phase)
}
}.toImmutableList()
}
}
private fun remove(jobId: String) {
private fun removeEntry(
jobId: String,
dropCheckpoint: Boolean,
persisted: Boolean,
) {
pending.update { it.remove(jobId) }
_jobs.update { list -> list.filter { it.id != jobId }.toImmutableList() }
persistence?.remove(jobId)
if (persisted && dropCheckpoint) persistence?.remove(jobId)
}
companion object {
@@ -28,7 +28,9 @@ import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.cancel
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.first
import kotlinx.coroutines.launch
import kotlinx.coroutines.test.runTest
import kotlinx.coroutines.withContext
import kotlinx.coroutines.withTimeout
@@ -213,4 +215,116 @@ class PoWPublishQueueTest {
withContext(Dispatchers.Default) { withTimeout(60_000) { queue.jobs.first { it.isEmpty() } } }
scope.cancel()
}
@Test
fun failedPublishKeepsCheckpointAndReportsFailure() =
runTest {
val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default)
val persistence = FakePersistence()
val queue = PoWPublishQueue(scope, maxConcurrent = 1, persistence = persistence)
val failure = CompletableDeferred<PoWJobFailure>()
scope.launch { queue.failures.collect { failure.complete(it) } }
// give the collector a beat to subscribe (failures has no replay)
withContext(Dispatchers.Default) { delay(50) }
queue.enqueueStaged(
kind = TextNoteEvent.KIND,
difficulty = 10,
persistAs = recordFor("job-d"),
mine = { "nonce" },
publish = { throw IllegalStateException("signer rejected") },
)
val reported = withContext(Dispatchers.Default) { withTimeout(10_000) { failure.await() } }
assertEquals(TextNoteEvent.KIND, reported.kind)
assertTrue(reported.willRetryOnRestart, "persisted job must be retried by the restorer")
withContext(Dispatchers.Default) { withTimeout(10_000) { queue.jobs.first { it.isEmpty() } } }
assertFalse("job-d" in persistence.removed, "checkpoint must survive a failed publish for restart retry")
scope.cancel()
}
@Test
fun checkpointSurvivesUntilPublishCompletes() =
runTest {
val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default)
val persistence = FakePersistence()
val queue = PoWPublishQueue(scope, maxConcurrent = 1, persistence = persistence)
val publishing = CompletableDeferred<Unit>()
val publishGate = CompletableDeferred<Unit>()
queue.enqueueStaged(
kind = TextNoteEvent.KIND,
difficulty = 10,
persistAs = recordFor("job-e"),
mine = { "nonce" },
publish = {
publishing.complete(Unit)
publishGate.await()
},
)
withContext(Dispatchers.Default) { withTimeout(10_000) { publishing.await() } }
assertFalse("job-e" in persistence.removed, "checkpoint must not be dropped at mining-complete")
val job = queue.jobs.value.first()
assertEquals(PoWJobPhase.PUBLISHING, job.phase)
assertFalse(job.isCancellable)
// cancel is a no-op mid-publish: the event may be half-broadcast
queue.cancel(job.id)
assertEquals(1, queue.jobs.value.size)
publishGate.complete(Unit)
withContext(Dispatchers.Default) { withTimeout(10_000) { queue.jobs.first { it.isEmpty() } } }
assertTrue("job-e" in persistence.removed, "checkpoint drops once publish completes")
scope.cancel()
}
@Test
fun cancelByKeyTogglesAPendingJob() =
runTest {
val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default)
val queue = PoWPublishQueue(scope, maxConcurrent = 1)
val gate = CompletableDeferred<Unit>()
var reactionRan = false
queue.enqueueWork(kind = 1, difficulty = 10) { gate.await() }
queue.enqueueWork(kind = 7, difficulty = 10, dedupeKey = "reaction:abc:+") { reactionRan = true }
// same key while pending → deduped
queue.enqueueWork(kind = 7, difficulty = 10, dedupeKey = "reaction:abc:+") { reactionRan = true }
assertEquals(2, queue.jobs.value.size)
assertTrue(queue.cancelByKey("reaction:abc:+"), "first toggle cancels the pending like")
assertFalse(queue.cancelByKey("reaction:abc:+"), "nothing left to cancel")
gate.complete(Unit)
withContext(Dispatchers.Default) { withTimeout(10_000) { queue.jobs.first { it.isEmpty() } } }
assertFalse(reactionRan, "cancelled reaction must never publish")
scope.cancel()
}
@Test
fun cancelForOwnerOnlyDropsThatAccountsJobs() =
runTest {
val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default)
val queue = PoWPublishQueue(scope, maxConcurrent = 1)
val gate = CompletableDeferred<Unit>()
var otherRan = false
queue.enqueueWork(kind = 1, difficulty = 10) { gate.await() }
queue.enqueueWork(kind = 1, difficulty = 10, owner = "account-a") {}
queue.enqueueWork(kind = 1, difficulty = 10, owner = "account-b") { otherRan = true }
queue.cancelForOwner("account-a")
assertEquals(2, queue.jobs.value.size, "only account-a's job leaves the queue")
gate.complete(Unit)
withContext(Dispatchers.Default) { withTimeout(10_000) { queue.jobs.first { it.isEmpty() } } }
assertTrue(otherRan, "account-b's job still runs")
scope.cancel()
}
}
@@ -83,6 +83,10 @@ class PoWMiner(
desiredPoW: Int,
isActive: () -> Boolean = { true },
): EventTemplate<T> {
// sha256 ids have 256 bits; anything outside would index past the
// hash (or never terminate) deep inside the hot loop.
require(desiredPoW in 1..256) { "desiredPoW must be in 1..256, was $desiredPoW" }
var nextSize = STARTING_NONCE_SIZE
do {
@@ -106,7 +106,9 @@ class PoWRankEvaluator {
minPoW: Int,
emptyBytes: Int,
): Boolean {
for (index in 0 until emptyBytes) {
// emptyBytes is minPoW/8; clamp so an oversized target can never
// index past the 32-byte hash.
for (index in 0 until emptyBytes.coerceAtMost(id.size)) {
if (id[index] != R8) return false
}
@@ -47,6 +47,81 @@ class NIP17Factory {
val wraps: List<GiftWrapEvent>,
)
/** A seal encrypted to one recipient, waiting to be gift-wrapped. */
data class AddressedSeal(
val recipient: HexKey,
val seal: Event,
)
data class SealsForWrapping(
val seals: List<AddressedSeal>,
val expirationDelta: Long?,
)
/**
* Phase one of a split wrap build: signs one seal per recipient using
* [signer]. This is the only step that talks to the user's signer
* (external signers must prompt in the user's interaction context);
* the remaining wrap step ([wrapSeal]) is pure local CPU work that can
* run on a background mining worker.
*/
suspend fun createSeals(
event: Event,
to: Set<HexKey>,
signer: NostrSigner,
): SealsForWrapping {
val innerExpDelta =
event.expiration()?.let {
if (it > event.createdAt) {
it - event.createdAt
} else {
null
}
}
val bunkerLimiter = if (signer is NostrSignerRemote) Semaphore(BUNKER_PARALLELISM) else null
val seals =
mapNotNullAsync(to.toList()) { next ->
val build: suspend () -> AddressedSeal = {
AddressedSeal(
recipient = next,
seal =
SealedRumorEvent.create(
event = event,
encryptTo = next,
expirationDelta = innerExpDelta,
signer = signer,
),
)
}
bunkerLimiter?.withPermit { build() } ?: build()
}
return SealsForWrapping(seals, innerExpDelta)
}
/**
* Phase two: wraps one pre-signed seal in its ephemeral-key envelope,
* optionally mining a NIP-13 proof of work into the wrap. Local-only —
* no user signer involved.
*/
fun wrapSeal(
addressed: AddressedSeal,
expirationDelta: Long?,
recipientRelayHint: NormalizedRelayUrl? = null,
powDifficulty: Int? = null,
powIsActive: () -> Boolean = { true },
): GiftWrapEvent =
GiftWrapEvent.create(
event = addressed.seal,
recipientPubKey = addressed.recipient,
expirationDelta = expirationDelta,
recipientRelayHint = recipientRelayHint,
powDifficulty = powDifficulty,
powIsActive = powIsActive,
)
/**
* Build one NIP-59 gift wrap per recipient.
*