Repository navigation
fix(media-watcher): write playback as discrete per-player segments #343
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Open
0xbrayo
wants to merge
3
commits into
ActivityWatch:master
Choose a base branch
from
0xbrayo:fix/media-watcher-concurrent-players
base: master
Could not load branches
Branch not found: {{ refName }}
Loading
Could not load tags
Nothing to show
Loading
Are you sure you want to change the base?
Some commits from the old base branch may be removed from the timeline,
and old review comments may become outdated.
+217
−30
Open
Changes from all commits
Commits
Show all changes
3 commits
Select commit
Hold shift + click to select a range
File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
77 changes: 77 additions & 0 deletions
77
mobile/src/main/java/net/activitywatch/android/watcher/MediaPlaybackSegments.kt
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,77 @@ | ||
| package net.activitywatch.android.watcher | ||
|
|
||
| import org.json.JSONObject | ||
| import org.threeten.bp.Duration | ||
| import org.threeten.bp.Instant | ||
|
|
||
| // Longest stretch of playback held in memory before it is written. Bounds how much is | ||
| // lost if the process dies mid-track, and how stale the bucket is during a long episode. | ||
| internal val MEDIA_SEGMENT_MAX_LENGTH: Duration = Duration.ofMinutes(5) | ||
|
|
||
| /** | ||
| * Turns per-player playback observations into discrete events. | ||
| * | ||
| * Playback used to be sent as heartbeats, but the server only merges a heartbeat into the | ||
| * newest event in the bucket. With two players active, their heartbeats alternated, never | ||
| * merged, and every poll inserted new zero-duration events. Instead each player's playing | ||
| * time is tracked as an open segment with a known start and written once, with its real | ||
| * duration, when the track or state changes, the session ends, or the segment reaches | ||
| * [MEDIA_SEGMENT_MAX_LENGTH]. | ||
| * | ||
| * Not thread-safe; callers serialize access. | ||
| */ | ||
| internal class MediaPlaybackSegments<P>( | ||
| private val emit: (start: Instant, durationSeconds: Double, data: JSONObject) -> Unit, | ||
| ) { | ||
| private class Segment(val key: String, val data: JSONObject, val start: Instant) | ||
|
|
||
| // Keyed by player (one media session), not by app: an app can run several sessions. | ||
| private val open = HashMap<P, Segment>() | ||
| private val lastStateKeys = HashMap<P, String>() | ||
|
|
||
| /** Returns true when this observation is a change worth logging. */ | ||
| fun observe(player: P, key: String, data: JSONObject, playing: Boolean, now: Instant): Boolean { | ||
| val segment = open[player] | ||
| if (playing) { | ||
| if (segment != null && segment.key == key) { | ||
| if (Duration.between(segment.start, now) >= MEDIA_SEGMENT_MAX_LENGTH) { | ||
| flush(segment, now) | ||
| open[player] = Segment(key, data, now) | ||
| } | ||
| return false | ||
| } | ||
| segment?.let { flush(it, now) } | ||
| open[player] = Segment(key, data, now) | ||
| lastStateKeys[player] = key | ||
| return true | ||
| } | ||
|
|
||
| segment?.let { flush(it, now) } | ||
| open.remove(player) | ||
| if (lastStateKeys[player] == key) return false | ||
| lastStateKeys[player] = key | ||
| // Pauses and stops have no duration; record the transition itself. | ||
| emit(now, 0.0, data) | ||
| return true | ||
| } | ||
|
|
||
| /** | ||
| * Writes and forgets the player's open segment: its session went away, or it reported | ||
| * something that isn't recognisable playback (no metadata, an unknown state). | ||
| */ | ||
| fun end(player: P, now: Instant) { | ||
| open.remove(player)?.let { flush(it, now) } | ||
| lastStateKeys.remove(player) | ||
| } | ||
|
|
||
| fun endAll(now: Instant) { | ||
| open.values.forEach { flush(it, now) } | ||
| open.clear() | ||
| lastStateKeys.clear() | ||
| } | ||
|
|
||
| private fun flush(segment: Segment, end: Instant) { | ||
| val millis = Duration.between(segment.start, end).toMillis() | ||
| if (millis > 0) emit(segment.start, millis / 1000.0, segment.data) | ||
| } | ||
| } | ||
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
79 changes: 79 additions & 0 deletions
79
mobile/src/test/java/net/activitywatch/android/watcher/MediaPlaybackSegmentsTest.kt
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,79 @@ | ||
| package net.activitywatch.android.watcher | ||
|
|
||
| import org.json.JSONObject | ||
| import org.junit.Assert.assertEquals | ||
| import org.junit.Test | ||
| import org.threeten.bp.Instant | ||
|
|
||
| class MediaPlaybackSegmentsTest { | ||
| private data class Emitted(val start: Long, val seconds: Double, val title: String) | ||
|
|
||
| private val emitted = mutableListOf<Emitted>() | ||
| private val segments = MediaPlaybackSegments<String> { start, seconds, data -> | ||
| emitted += Emitted(start.epochSecond, seconds, data.getString("title")) | ||
| } | ||
|
|
||
| private fun at(seconds: Long) = Instant.ofEpochSecond(seconds) | ||
|
|
||
| private fun observe(player: String, title: String, playing: Boolean, seconds: Long) { | ||
| val state = if (playing) "playing" else "paused" | ||
| segments.observe( | ||
| player, "$player|$title|$state", JSONObject().put("title", title), playing, at(seconds) | ||
| ) | ||
| } | ||
|
|
||
| @Test | ||
| fun twoPlayersPollingTogetherWriteOneEventEachWhenTheyStop() { | ||
| // Both players are polled every 15s. The old heartbeats alternated between the two | ||
| // and inserted a new zero-duration event on every poll. | ||
| for (t in 0L..120L step 15) { | ||
| observe("music", "song", playing = true, seconds = t) | ||
| observe("podcast", "episode", playing = true, seconds = t) | ||
| } | ||
| assertEquals(emptyList<Emitted>(), emitted) | ||
|
|
||
| observe("music", "song", playing = false, seconds = 130) | ||
| segments.end("podcast", at(140)) | ||
|
|
||
| assertEquals( | ||
| listOf( | ||
| Emitted(0, 130.0, "song"), | ||
| Emitted(130, 0.0, "song"), // the pause itself | ||
| Emitted(0, 140.0, "episode"), | ||
| ), | ||
| emitted, | ||
| ) | ||
| } | ||
|
|
||
| @Test | ||
| fun trackChangeEndsThePreviousTrack() { | ||
| observe("music", "first", playing = true, seconds = 0) | ||
| observe("music", "second", playing = true, seconds = 200) | ||
| segments.endAll(at(260)) | ||
|
|
||
| assertEquals( | ||
| listOf(Emitted(0, 200.0, "first"), Emitted(200, 60.0, "second")), | ||
| emitted, | ||
| ) | ||
| } | ||
|
|
||
| @Test | ||
| fun longPlaybackIsWrittenInBoundedChunks() { | ||
| val chunk = MEDIA_SEGMENT_MAX_LENGTH.seconds | ||
| for (t in 0L..chunk + 30 step 15) { | ||
| observe("podcast", "episode", playing = true, seconds = t) | ||
| } | ||
| // The first chunk is written as soon as a poll sees it reach the limit, without | ||
| // waiting for the episode to end. | ||
| assertEquals(listOf(Emitted(0, chunk.toDouble(), "episode")), emitted) | ||
| } | ||
|
|
||
| @Test | ||
| fun repeatedPauseIsRecordedOnce() { | ||
| observe("music", "song", playing = true, seconds = 0) | ||
| observe("music", "song", playing = false, seconds = 10) | ||
| observe("music", "song", playing = false, seconds = 25) | ||
|
|
||
| assertEquals(listOf(Emitted(0, 10.0, "song"), Emitted(10, 0.0, "song")), emitted) | ||
| } | ||
| } |
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.