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
4 changes: 4 additions & 0 deletions dsm_client/android/app/src/main/AndroidManifest.xml
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,10 @@
<uses-permission android:name="android.permission.FOREGROUND_SERVICE"/>
<uses-permission android:name="android.permission.FOREGROUND_SERVICE_CONNECTED_DEVICE"/>
<uses-permission android:name="android.permission.FOREGROUND_SERVICE_DATA_SYNC"/>
<!-- The one-time "let DSM run in the background" dialog (BatteryExemption):
the wallet keeps its own connection to the storage nodes while the
phone sleeps, with no push service in between. -->
<uses-permission android:name="android.permission.REQUEST_IGNORE_BATTERY_OPTIMIZATIONS"/>

<!-- Notifications (Android 13+) for FGS persistent notif -->
<uses-permission android:name="android.permission.POST_NOTIFICATIONS"/>
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,65 @@
// SPDX-License-Identifier: MIT OR Apache-2.0

package com.dsm.wallet.ui

import android.app.Activity
import android.content.ActivityNotFoundException
import android.content.Context
import android.content.Intent
import android.os.PowerManager
import android.provider.Settings
import android.util.Log
import androidx.core.content.pm.PackageInfoCompat
import androidx.core.net.toUri

/**
* Asks once, per install, for DSM to be left out of Android's battery
* optimisation, so the wallet keeps its own connection to the storage nodes
* while the phone sleeps. A transfer still settling when the app is left is
* finished by the background service; with the app optimised, Android cuts its
* network in Doze and the counterparty waits until the phone wakes.
*
* Nothing goes through Google: this is the wallet's own connection, kept open
* by the system's own setting. The user answers in the system's dialog; a
* refusal is kept and not asked again.
*/
internal object BatteryExemption {
private const val TAG = "BatteryExemption"
private const val PREFS = "dsm_battery_exemption"

/** The app version that asked, once asked. */
private const val KEY_ASKED_BY_VERSION = "asked_by_version"

/**
* Opens the system's "let this app run in the background" dialog, once,
* when the app is still optimised. Opens the system's list of optimised
* apps instead where the dialog does not exist.
*/
fun askOnce(activity: Activity) {
val power = activity.getSystemService(Context.POWER_SERVICE) as? PowerManager
if (power == null) {
Log.w(TAG, "no power service: the battery exemption cannot be asked for")
return
}
if (power.isIgnoringBatteryOptimizations(activity.packageName)) return
val prefs = activity.getSharedPreferences(PREFS, Context.MODE_PRIVATE)
if (prefs.contains(KEY_ASKED_BY_VERSION)) return

val info = activity.packageManager.getPackageInfo(activity.packageName, 0)
prefs.edit()
.putLong(KEY_ASKED_BY_VERSION, PackageInfoCompat.getLongVersionCode(info))
.apply()
val ask = Intent(Settings.ACTION_REQUEST_IGNORE_BATTERY_OPTIMIZATIONS)
.setData("package:${activity.packageName}".toUri())
try {
activity.startActivity(ask)
} catch (e: ActivityNotFoundException) {
Log.w(TAG, "no battery exemption dialog on this phone; opening the list", e)
try {
activity.startActivity(Intent(Settings.ACTION_IGNORE_BATTERY_OPTIMIZATION_SETTINGS))
} catch (none: ActivityNotFoundException) {
Log.w(TAG, "no battery optimisation settings on this phone", none)
}
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -1340,6 +1340,7 @@ class MainActivity : AppCompatActivity(), NfcAdapter.ReaderCallback {
} catch (t: Throwable) {
Log.w(tag, "onStart: startForegroundService failed", t)
}
BatteryExemption.askOnce(this)
}
val intent = Intent(this, BleBackgroundService::class.java)
try {
Expand Down
125 changes: 125 additions & 0 deletions dsm_client/deterministic_state_machine/dsm_sdk/src/sdk/b0x_sdk.rs
Original file line number Diff line number Diff line change
Expand Up @@ -530,6 +530,126 @@ pub(crate) fn certresync_message_id(method: &str, recipient_tip: &[u8], body: &[
h.finalize().as_bytes()[..16].to_vec()
}

/// Where `storage.sync` last read each member's spool at each address to its
/// end, in this process: `(address, endpoint)` to the position after the last
/// entry seen. A wait on that spool ([`wait_on_member`]) wakes for what lands
/// from there. Kept in memory only: a process that has not yet read a spool
/// to its end waits from its read position, and its first sync records the
/// end.
static SPOOL_ENDS: once_cell::sync::Lazy<std::sync::Mutex<HashMap<(String, String), u64>>> =
once_cell::sync::Lazy::new(|| std::sync::Mutex::new(HashMap::new()));

/// The spool ends. A writer that panicked left at worst one stale end, and a
/// stale end only wakes a wait early.
fn spool_ends() -> std::sync::MutexGuard<'static, HashMap<(String, String), u64>> {
match SPOOL_ENDS.lock() {
Ok(ends) => ends,
Err(poisoned) => poisoned.into_inner(),
}
}

fn record_spool_end(address: &str, endpoint: &str, end: u64) {
spool_ends().insert((address.to_string(), endpoint.to_string()), end);
}

/// Where a sync last read the spool at `address` on `endpoint` to its end,
/// in this process.
pub(crate) fn spool_end(address: &str, endpoint: &str) -> Option<u64> {
spool_ends()
.get(&(address.to_string(), endpoint.to_string()))
.copied()
}

/// How long a member holds a wait before answering that nothing landed: the
/// deployed node's bound (`dsm_storage_node::api::transport::b0x::MAX_WAIT`).
pub(crate) const MEMBER_WAIT_BOUND: std::time::Duration = std::time::Duration::from_secs(25);

/// How much longer than the member's bound a wait request may take: the
/// round trip, and a member slow to answer at the end of its bound. A member
/// that has not answered by then did not answer.
const WAIT_TRANSPORT_SLACK: std::time::Duration = std::time::Duration::from_secs(10);

/// What a member answered a wait on this device's spools.
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum WaitAnswer {
/// These spools hold an entry at or after their mark.
Ready(Vec<String>),
/// Nothing landed while the member held the wait (`204`).
Quiet,
/// The member does not hold waits: a node that does not serve them yet
/// (`404`), or one holding as many as it can (`503`).
NotHeld(reqwest::StatusCode),
}

/// Hold a wait on `marks` (spool address, position) at the member at
/// `endpoint` until it answers (storage spec §8, long-poll): at once when an
/// entry is past its mark, when one lands, or `204` after its bound. Reads
/// nothing and moves nothing: what landed is read by the next sync.
pub(crate) async fn wait_on_member(
client: &reqwest::Client,
endpoint: &str,
marks: &[(String, u64)],
) -> Result<WaitAnswer, DsmError> {
let request = dsm::types::proto::B0xWaitRequest {
marks: marks
.iter()
.map(|(address, from_seq)| dsm::types::proto::B0xWaitMark {
address: address.clone(),
from_seq: *from_seq,
})
.collect(),
};
let url = format!("{}/api/v2/b0x/wait", endpoint.trim_end_matches('/'));
let resp = client
.post(&url)
.header("Content-Type", "application/protobuf")
.header("Accept", "application/protobuf")
.timeout(MEMBER_WAIT_BOUND + WAIT_TRANSPORT_SLACK)
.body(request.encode_to_vec())
.send()
.await
.map_err(|e| {
DsmError::network(
format!("b0x wait at {endpoint}: {e}"),
None::<std::io::Error>,
)
})?;
let status = resp.status();
if status == reqwest::StatusCode::NO_CONTENT {
return Ok(WaitAnswer::Quiet);
}
if status == reqwest::StatusCode::NOT_FOUND
|| status == reqwest::StatusCode::SERVICE_UNAVAILABLE
{
return Ok(WaitAnswer::NotHeld(status));
}
if status != reqwest::StatusCode::OK {
return Err(DsmError::network(
format!("b0x wait at {endpoint} answered HTTP {status}"),
None::<std::io::Error>,
));
}
let bytes = resp.bytes().await.map_err(|e| {
DsmError::network(
format!("b0x wait answer from {endpoint}: {e}"),
None::<std::io::Error>,
)
})?;
let answer = dsm::types::proto::B0xWaitResponse::decode(bytes.as_ref()).map_err(|e| {
DsmError::network(
format!("b0x wait answer from {endpoint} does not decode: {e}"),
None::<std::io::Error>,
)
})?;
if answer.ready.is_empty() {
return Err(DsmError::network(
format!("b0x wait at {endpoint} answered 200 naming no spool"),
None::<std::io::Error>,
));
}
Ok(WaitAnswer::Ready(answer.ready))
}

impl B0xSDK {
fn hash_b0x_component(
domain_tag: dsm::crypto::domain::TaggedHashDomain<'_>,
Expand Down Expand Up @@ -2380,6 +2500,11 @@ impl B0xSDK {
// not reach is read from there again.
let read_answered = !failed && malformed.is_none();
let capped = read_answered && stop == ReadStop::PageCap;
// A sync that read this member's spool to its end has seen all of
// it: a wait on this spool wakes for what lands after.
if from == ReadFrom::Resume && read_answered && stop == ReadStop::End {
record_spool_end(b0x_address, &epc, cursor);
}
if from == ReadFrom::Resume && read_answered {
b0x_consumed::set_scan_cursor(b0x_address, &epc, capped.then_some(cursor))
.map_err(local)?;
Expand Down
108 changes: 102 additions & 6 deletions dsm_client/deterministic_state_machine/dsm_sdk/src/sdk/inbox_poller.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@
//! by making the frontend the authority over inbox discovery timing.

use prost::Message;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::Arc;
use tokio::sync::Notify;

Expand Down Expand Up @@ -41,6 +41,21 @@ const FOREGROUND_POLL_INTERVAL_MS: u64 = 5_000;
/// Number of consecutive eager-interval polls before reverting to default.
const EAGER_POLL_CYCLES: u32 = 5;

/// Poll interval while settlement work is outstanding and the inbox waiter's
/// waits cover the fleet ([`crate::sdk::inbox_waiter::covers_fleet`]): a
/// reply or a certificate wakes a sync the moment it lands, and this poll
/// only retries what an arrival does not drive — a release object not yet
/// fetchable, a delivery that failed.
const COVERED_SETTLEMENT_POLL_INTERVAL_MS: u64 = 5_000;

/// Poll interval while the app is on screen, nothing is settling, and the
/// waits cover the fleet: a safety net only, an arrival wakes a sync at once.
const COVERED_FOREGROUND_POLL_INTERVAL_MS: u64 = 30_000;

/// Poller cycles completed in this process, and the signal each one gives.
static CYCLES_COMPLETED: AtomicU64 = AtomicU64::new(0);
static CYCLE_DONE: once_cell::sync::Lazy<Notify> = once_cell::sync::Lazy::new(Notify::new);

/// Global poller state.
static POLLER_RUNNING: AtomicBool = AtomicBool::new(false);
static POLLER_STOP: AtomicBool = AtomicBool::new(false);
Expand Down Expand Up @@ -161,6 +176,8 @@ pub fn start_poller() {
}

let (processed, pulled, more_pending) = run_inbox_sync_cycle_counted("poll").await;
CYCLES_COMPLETED.fetch_add(1, Ordering::SeqCst);
CYCLE_DONE.notify_waiters();
#[cfg(test)]
POLLER_CYCLE_DONE.notify_waiters();
// Settlement-urgent covers BOTH directions: the sender awaiting an
Expand Down Expand Up @@ -189,15 +206,21 @@ pub fn start_poller() {
eager_remaining = eager_remaining.saturating_sub(1);
}

let interval_ms = if pending_gate_active {
PENDING_GATE_POLL_INTERVAL_MS
let activity = if pending_gate_active {
Activity::Settling
} else if crate::sdk::session_manager::app_in_foreground() {
FOREGROUND_POLL_INTERVAL_MS
Activity::OnScreen
} else if eager_remaining > 0 {
EAGER_POLL_INTERVAL_MS
Activity::Eager
} else {
DEFAULT_POLL_INTERVAL_MS
Activity::Idle
};
let waits = if crate::sdk::inbox_waiter::covers_fleet() {
Waits::CoverFleet
} else {
Waits::DoNotCover
};
let interval_ms = poll_interval_ms(activity, waits);

// Wait for either the poll interval or a wake-up signal.
tokio::select! {
Expand All @@ -210,6 +233,59 @@ pub fn start_poller() {

log::info!("[inbox_poller] Background poller stopped");
});
crate::sdk::inbox_waiter::start();
}

/// What the device is doing, as the poller's cadence reads it.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum Activity {
/// Settlement work is outstanding ([`has_pending_settlement_work`]).
Settling,
/// The app is on screen and nothing is settling.
OnScreen,
/// A recent cycle took something: follow-ups come soon.
Eager,
/// None of these.
Idle,
}

/// Whether the inbox waiter's waits cover the fleet
/// ([`crate::sdk::inbox_waiter::covers_fleet`]).
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum Waits {
CoverFleet,
DoNotCover,
}

/// The time to the next poll: the settlement cadence while a transfer is
/// settling, the on-screen cadence while the app is shown, the eager one
/// after an exchange, else the idle minute. While the waits cover the fleet
/// an arrival wakes a sync at once, so the settlement and on-screen cadences
/// slow to a safety net.
pub(crate) fn poll_interval_ms(activity: Activity, waits: Waits) -> u64 {
match (activity, waits) {
(Activity::Settling, Waits::CoverFleet) => COVERED_SETTLEMENT_POLL_INTERVAL_MS,
(Activity::Settling, Waits::DoNotCover) => PENDING_GATE_POLL_INTERVAL_MS,
(Activity::OnScreen, Waits::CoverFleet) => COVERED_FOREGROUND_POLL_INTERVAL_MS,
(Activity::OnScreen, Waits::DoNotCover) => FOREGROUND_POLL_INTERVAL_MS,
(Activity::Eager, _) => EAGER_POLL_INTERVAL_MS,
(Activity::Idle, _) => DEFAULT_POLL_INTERVAL_MS,
}
}

/// Poller cycles completed in this process.
pub(crate) fn cycles_completed() -> u64 {
CYCLES_COMPLETED.load(Ordering::SeqCst)
}

/// Resolves when the next poller cycle completes.
pub(crate) fn cycle_done() -> tokio::sync::futures::Notified<'static> {
CYCLE_DONE.notified()
}

/// Whether the poller has been told to stop.
pub(crate) fn poller_stopping() -> bool {
POLLER_STOP.load(Ordering::SeqCst)
}

/// True while this device owes the network a settlement step that only polling
Expand Down Expand Up @@ -271,12 +347,14 @@ pub fn stop_poller_for_lifecycle() -> anyhow::Result<bool> {
pub fn stop_poller() {
POLLER_STOP.store(true, Ordering::SeqCst);
POLLER_WAKE.notify_one();
crate::sdk::inbox_waiter::wake_to_stop();
}

/// Wake the poller immediately (e.g. app foreground, bilateral commit).
pub fn resume_poller() {
if POLLER_RUNNING.load(Ordering::SeqCst) {
POLLER_WAKE.notify_one();
crate::sdk::inbox_waiter::start();
} else {
// If poller isn't running, start it.
start_poller();
Expand Down Expand Up @@ -663,6 +741,24 @@ mod tests {

// ── push_inbox_event_to_webview is no-op on non-android ──

/// While a transfer settles the poller checks every 2 s, and every 5 s
/// while the app is on screen; with the waits covering the fleet an
/// arrival syncs at once, so those slow to 5 s and 30 s. An eager burst
/// and the idle minute are the same either way.
#[test]
fn the_waits_covering_the_fleet_slow_the_settling_and_on_screen_polls() {
use Activity::{Eager, Idle, OnScreen, Settling};
use Waits::{CoverFleet, DoNotCover};
assert_eq!(poll_interval_ms(Settling, DoNotCover), 2_000);
assert_eq!(poll_interval_ms(Settling, CoverFleet), 5_000);
assert_eq!(poll_interval_ms(OnScreen, DoNotCover), 5_000);
assert_eq!(poll_interval_ms(OnScreen, CoverFleet), 30_000);
assert_eq!(poll_interval_ms(Eager, DoNotCover), 8_000);
assert_eq!(poll_interval_ms(Eager, CoverFleet), 8_000);
assert_eq!(poll_interval_ms(Idle, DoNotCover), 60_000);
assert_eq!(poll_interval_ms(Idle, CoverFleet), 60_000);
}

#[test]
fn a_pending_route_hurries_the_poller_only_while_entries_are_taken() {
assert!(enters_eager_mode(1, 0, 0), "something processed");
Expand Down
Loading
Loading