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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -122,7 +122,9 @@ fun HomeScreen(nav: NavHostController) {
if (recentList.isNotEmpty()) {
val ids = recentList.map { it.track.id }.distinct()
item(key = "h.continue", contentType = "header") { SectionHeader(stringResource(R.string.home_continue), Modifier.animateItem()) }
items(recentList.take(5), key = { "r" + it.playedAt + it.track.id }, contentType = { "track" }) { entry ->
// One row per play (start + track is its identity, and a lazy list throws on a
// repeated key): older databases can hold the same play twice.
items(recentList.distinctBy { it.playedAt to it.track.id }.take(5), key = { "r" + it.playedAt + it.track.id }, contentType = { "track" }) { entry ->
TrackRow(entry.track, onClick = { client.dispatch(Commands.playTracks(serverId, ids, ids.indexOf(entry.track.id), continueLabel)) },
modifier = Modifier.animateItem(),
trailing = { Text(entryAgo(entry), style = MaterialTheme.typography.labelSmall, color = MaterialTheme.colorScheme.onSurfaceVariant, modifier = Modifier.padding(start = 8.dp)) })
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,12 @@ class ConnectRouteProvider : MediaRoute2ProviderService() {
private var state = ConnectRoutes.State()
/** The device a switcher pick asked for, and when, until the core's state shows it playing. */
private var pending: Pair<String, Long>? = null
/**
* When "Stop casting" asked for playback back here, until the core's state shows it here. The
* session stays released meanwhile: recreating it while the other device still reports playing
* would bring the chip back mid-handoff, and each further tap would start another handoff.
*/
private var returningSince: Long? = null
private var sessionRoute: String? = null

override fun onCreate() {
Expand All @@ -67,6 +73,10 @@ class ConnectRouteProvider : MediaRoute2ProviderService() {

/** The device the session should show as selected: the core's, or a pick still in flight. */
private fun target(): DeviceInfo? {
returningSince?.let { since ->
if (state.remote == null || SystemClock.elapsedRealtime() - since > PENDING_MS) returningSince = null
else return null
}
state.remote?.let { remote -> if (pending?.first == remote.id) pending = null; return remote }
val (id, at) = pending ?: return null
if (SystemClock.elapsedRealtime() - at > PENDING_MS) { pending = null; return null }
Expand Down Expand Up @@ -115,6 +125,7 @@ class ConnectRouteProvider : MediaRoute2ProviderService() {
}

private fun handoffTo(deviceId: String) {
returningSince = null
pending = deviceId to SystemClock.elapsedRealtime()
dispatch(Commands.handoffTo(deviceId))
}
Expand All @@ -138,7 +149,10 @@ class ConnectRouteProvider : MediaRoute2ProviderService() {
sessionRoute = null
notifySessionReleased(sessionId)
val self = state.selfId
if (state.remote != null && self != null) dispatch(Commands.handoffTo(self))
if (state.remote != null && self != null) {
returningSince = SystemClock.elapsedRealtime()
dispatch(Commands.handoffTo(self))
}
}

override fun onTransferToRoute(requestId: Long, sessionId: String, routeId: String) {
Expand Down
41 changes: 41 additions & 0 deletions crates/hocket-core/src/db/migrations/0004_play_history_unique.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,41 @@
-- 0004: a play is recorded once. A play's identity is its track and its
-- original start (`played_at`, the session clock's `startedAt`, carried
-- unchanged across handoffs). A play that moved back to a device that had
-- already recorded it (quick handoffs back and forth) reached its threshold
-- there again and was recorded twice, and the Android Home screen, keyed by
-- start and track, crashed on the pair.
--
-- Keep the first row of each play (scrobbled if any copy was), take the
-- extra local play count bumps back, drop the copies, and let the index
-- refuse any new one.
UPDATE play_history SET scrobbled = 1
WHERE scrobbled = 0 AND EXISTS (
SELECT 1 FROM play_history o
WHERE o.server_id = play_history.server_id AND o.track_id = play_history.track_id
AND o.played_at = play_history.played_at AND o.scrobbled = 1
);

UPDATE tracks SET local_play_count = MAX(0, local_play_count - (
SELECT count(*) FROM play_history d
WHERE d.server_id = tracks.server_id AND d.track_id = tracks.id
AND EXISTS (
SELECT 1 FROM play_history k
WHERE k.server_id = d.server_id AND k.track_id = d.track_id
AND k.played_at = d.played_at AND k.id < d.id
)
))
WHERE EXISTS (
SELECT 1 FROM play_history d JOIN play_history k
ON k.server_id = d.server_id AND k.track_id = d.track_id
AND k.played_at = d.played_at AND k.id < d.id
WHERE d.server_id = tracks.server_id AND d.track_id = tracks.id
);

DELETE FROM play_history
WHERE EXISTS (
SELECT 1 FROM play_history k
WHERE k.server_id = play_history.server_id AND k.track_id = play_history.track_id
AND k.played_at = play_history.played_at AND k.id < play_history.id
);

CREATE UNIQUE INDEX IF NOT EXISTS idx_play_history_play ON play_history(server_id, track_id, played_at);
68 changes: 67 additions & 1 deletion crates/hocket-core/src/db/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -40,10 +40,15 @@ pub const MIGRATIONS: &[(u32, &str, &str)] = &[
"0003_stream_cache_spans",
include_str!("migrations/0003_stream_cache_spans.sql"),
),
(
4,
"0004_play_history_unique",
include_str!("migrations/0004_play_history_unique.sql"),
),
];

/// Current schema version (the last migration number).
pub const SCHEMA_VERSION: u32 = 3;
pub const SCHEMA_VERSION: u32 = 4;

/// Tables that are a cache of the server and may be dropped and re-synced.
pub const MIRROR_TABLES: &[&str] = &[
Expand Down Expand Up @@ -507,6 +512,67 @@ mod tests {
assert_eq!(v as u32, SCHEMA_VERSION);
}

#[test]
fn migration_4_drops_duplicate_plays() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("hocket.db");
{
// A version 3 database holding one play recorded twice.
let db = Db::open(&path, &dir.path().join("backups")).unwrap();
db.with_conn(|c| {
c.execute_batch(
"DROP INDEX idx_play_history_play;
DELETE FROM schema_version WHERE version = 4;
INSERT INTO tracks(id, server_id, title, local_play_count) VALUES ('t', 's', 'T', 3), ('u', 's', 'U', 1);
INSERT INTO play_history(server_id, track_id, played_at, played_ms, scrobbled, device_id) VALUES
('s', 't', 1790294961606.5, 126260, 0, 'd'),
('s', 't', 1790294961606.5, 132388, 1, 'd'),
('s', 't', 1790294000000.0, 126260, 0, 'd'),
('s', 'u', 1790294961606.5, 126260, 0, 'd');",
)?;
Ok(())
})
.unwrap();
}
let db = Db::open(&path, &dir.path().join("backups")).unwrap();
assert_eq!(db.schema_version().unwrap(), SCHEMA_VERSION);
let rows: Vec<(String, f64, i64, i64)> = db
.with_conn(|c| {
let mut st = c.prepare(
"SELECT track_id, played_at, played_ms, scrobbled FROM play_history ORDER BY id",
)?;
let rows = st.query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?)))?;
Ok(rows.collect::<Result<Vec<_>, _>>()?)
})
.unwrap();
assert_eq!(
rows,
vec![
("t".into(), 1790294961606.5, 126260, 1),
("t".into(), 1790294000000.0, 126260, 0),
("u".into(), 1790294961606.5, 126260, 0),
]
);
let counts: Vec<i64> = db
.with_conn(|c| {
let mut st = c.prepare("SELECT local_play_count FROM tracks ORDER BY id")?;
let rows = st.query_map([], |r| r.get(0))?;
Ok(rows.collect::<Result<Vec<_>, _>>()?)
})
.unwrap();
assert_eq!(counts, vec![2, 1]);
// The index refuses a copy from now on.
assert!(db
.with_conn(|c| {
c.execute(
"INSERT INTO play_history(server_id, track_id, played_at, played_ms, scrobbled, device_id) VALUES ('s', 'u', 1790294961606.5, 1, 0, 'd')",
[],
)?;
Ok(())
})
.is_err());
}

#[test]
fn newer_database_is_refused() {
let dir = tempfile::tempdir().unwrap();
Expand Down
28 changes: 28 additions & 0 deletions crates/hocket-core/src/db/queries.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1274,9 +1274,31 @@ impl From<Vec<Value>> for WhereClause {
#[allow(dead_code)]
fn _assert_tosql(_: &dyn ToSql) {}

/// The `play_history` row already recorded for the play of `track_id`
/// that started at `played_at` (a play's identity), if any.
pub fn recorded_play_in(
tx: &Connection,
server_id: &str,
track_id: &str,
played_at: f64,
) -> DbResult<Option<i64>> {
Ok(tx
.query_row(
"SELECT id FROM play_history WHERE server_id = ?1 AND track_id = ?2 AND played_at = ?3",
params![server_id, track_id, played_at],
|r| r.get(0),
)
.optional()?)
}

/// [`Db::record_play`] inside the caller's transaction: a `play_history`
/// row (returning its id) and the track's `local_play_count` /
/// `local_last_played` bump. The one place these rows are written.
///
/// A play is recorded once: when this play (track and start) already has a
/// row (it came back to this device after a handoff and reached its
/// threshold here again), that row's id is returned and nothing is counted
/// twice; it is only marked `scrobbled` if this verdict says so.
pub fn record_play_in(
tx: &Connection,
server_id: &str,
Expand All @@ -1286,6 +1308,12 @@ pub fn record_play_in(
scrobbled: bool,
device_id: &str,
) -> DbResult<i64> {
if let Some(id) = recorded_play_in(tx, server_id, track_id, played_at)? {
if scrobbled {
tx.execute("UPDATE play_history SET scrobbled = 1 WHERE id = ?1", [id])?;
}
return Ok(id);
}
tx.execute(
"INSERT INTO play_history(server_id, track_id, played_at, played_ms, scrobbled, device_id) VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
params![server_id, track_id, played_at, played_ms, scrobbled as i64, device_id],
Expand Down
81 changes: 80 additions & 1 deletion crates/hocket-core/src/outbox/scrobbler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -336,6 +336,10 @@ impl ScrobbleRecorder {
also: impl FnOnce(&rusqlite::Transaction) -> DbResult<()>,
) -> DbResult<()> {
self.db.with_tx(|tx| {
// Asked again about a play already recorded here (it came back
// after a handoff): one row, and at most one submission.
let recorded =
crate::db::queries::recorded_play_in(tx, server_id, track_id, played_at)?.is_some();
let history_id = crate::db::queries::record_play_in(
tx,
server_id,
Expand All @@ -345,7 +349,7 @@ impl ScrobbleRecorder {
verdict == Verdict::ScrobbledElsewhere,
&self.device_id,
)?;
if verdict == Verdict::Submit {
if verdict == Verdict::Submit && !recorded {
self.outbox.enqueue_in(
tx,
server_id,
Expand Down Expand Up @@ -713,6 +717,81 @@ mod tests {
assert_eq!((lpc, llp), (1, Some(4_000.0)));
}

#[test]
fn a_play_asked_about_again_is_recorded_and_submitted_once() {
// The play came back to this device after a handoff and reached its
// threshold here again: same track, same start, more time played.
let db = Db::open_in_memory().unwrap();
db.upsert_tracks(
&[crate::api::Track {
id: "t".into(),
server_id: "srv".into(),
title: "T".into(),
..Default::default()
}],
&[],
1,
)
.unwrap();
let clock = Arc::new(TestClock(Mutex::new(5_000.0)));
let outbox = Outbox::new(db.clone(), clock.clone());
let rec = ScrobbleRecorder::new(db.clone(), outbox.clone(), "dev");
let started_at = 1_790_294_961_606.5;
rec.record_verdict("srv", "t", started_at, 126_260, Verdict::Submit, |_| Ok(()))
.unwrap();
rec.record_verdict("srv", "t", started_at, 132_388, Verdict::Submit, |_| Ok(()))
.unwrap();
rec.record_verdict(
"srv",
"t",
started_at,
140_000,
Verdict::ScrobbledElsewhere,
|_| Ok(()),
)
.unwrap();
let h = db.recently_played(5).unwrap();
assert_eq!(h.len(), 1);
assert_eq!(h[0].played_ms, 126_260);
assert!(h[0].scrobbled);
let submissions = outbox
.pending()
.unwrap()
.into_iter()
.filter(|p| {
matches!(
p.mutation,
Mutation::Scrobble {
submission: true,
..
}
)
})
.count();
assert_eq!(submissions, 1);
let lpc: i64 = db
.with_conn(|c| {
Ok(c.query_row(
"SELECT local_play_count FROM tracks WHERE id='t'",
[],
|r| r.get(0),
)?)
})
.unwrap();
assert_eq!(lpc, 1);
// A later play of the same track is a play of its own.
rec.record_verdict(
"srv",
"t",
started_at + 300_000.0,
126_260,
Verdict::Submit,
|_| Ok(()),
)
.unwrap();
assert_eq!(db.recently_played(5).unwrap().len(), 2);
}

mod props {
use super::*;
use proptest::prelude::*;
Expand Down
Loading