Skip to content
Open
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
@@ -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
Comment thread
0xbrayo marked this conversation as resolved.
}
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)
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -31,9 +31,7 @@ class MediaWatcher : NotificationListenerService() {
private const val TAG = "MediaWatcher"
private const val BUCKET_ID = "aw-watcher-android-media"
private const val BUCKET_TYPE = "media.playback"
// Heartbeat pulsetime: merge events within 60s (same track playing continuously)
private const val PULSETIME = 60.0
// How often to poll active media sessions to send heartbeats
// How often to poll active media sessions (also bounds how late a segment is written)
private const val POLL_INTERVAL_MS = 15000L

fun isNotificationAccessGranted(context: android.content.Context): Boolean {
Expand All @@ -56,7 +54,15 @@ class MediaWatcher : NotificationListenerService() {
// access these maps concurrently; CHM prevents ConcurrentModificationException.
private val activeControllers = ConcurrentHashMap<MediaSession.Token, MediaController>()
private val activeCallbacks = ConcurrentHashMap<MediaSession.Token, MediaController.Callback>()
private val lastEventKeys = ConcurrentHashMap<String, String>()
// Only touched on handlerThread: callbacks, polling, the session scan and the final
// writes all run there, so writes never block the service's main thread.
private val segments = MediaPlaybackSegments<MediaSession.Token> { start, durationSeconds, data ->
ri?.insertEvent(BUCKET_ID, start, durationSeconds, data)
}

// Handler thread only. False between teardown and the next session scan, so a callback
// or poll already queued when the listener disconnected can't reopen a segment.
private var connected = false

// Polling mechanism to prevent 60-second cutoffs
private var handler: android.os.Handler? = null
Expand Down Expand Up @@ -96,22 +102,21 @@ class MediaWatcher : NotificationListenerService() {
override fun onListenerConnected() {
super.onListenerConnected()
Log.i(TAG, "MediaWatcher listener connected")
registerActiveSessionListener()
handler?.post { registerActiveSessionListener() }
Comment thread
0xbrayo marked this conversation as resolved.
pollingRunnable?.let { handler?.postDelayed(it, POLL_INTERVAL_MS) }
}

override fun onListenerDisconnected() {
super.onListenerDisconnected()
Log.i(TAG, "MediaWatcher listener disconnected")
unregisterAllCallbacks()
handler?.removeCallbacksAndMessages(null)
disconnect()
}

override fun onDestroy() {
super.onDestroy()
Log.i(TAG, "MediaWatcher destroyed")
unregisterAllCallbacks()
handler?.removeCallbacksAndMessages(null)
disconnect()
// quitSafely still runs the teardown disconnect() just posted.
handlerThread.quitSafely()
}

Expand All @@ -129,6 +134,7 @@ class MediaWatcher : NotificationListenerService() {
* This is called once on service creation and handles all session lifecycle.
*/
private fun registerActiveSessionListener() {
connected = true
val componentName = ComponentName(this, MediaWatcher::class.java)
try {
val listener = MediaSessionManager.OnActiveSessionsChangedListener { controllers ->
Expand All @@ -152,13 +158,16 @@ class MediaWatcher : NotificationListenerService() {
* Registers callbacks for new sessions and cleans up stale ones.
*/
private fun onActiveSessionsChanged(controllers: List<MediaController>?) {
if (controllers == null) return
if (controllers == null || !connected) return

val currentTokens = controllers.map { it.sessionToken }.toSet()

// Remove callbacks for sessions that are no longer active
val staleTokens = activeControllers.keys - currentTokens
val now = Instant.now()
for (token in staleTokens) {
// An inactive session no longer reports state, so close its playback now.
segments.end(token, now)
val controller = activeControllers.remove(token)
val callback = activeCallbacks.remove(token)
if (controller != null && callback != null) {
Expand Down Expand Up @@ -194,10 +203,17 @@ class MediaWatcher : NotificationListenerService() {
private fun createMediaCallback(controller: MediaController): MediaController.Callback {
return object : MediaController.Callback() {
override fun onPlaybackStateChanged(state: PlaybackState?) {
val metadata = controller.metadata ?: return
if (state != null) {
handlePlaybackChange(controller, state, metadata)
if (state == null) return
val metadata = controller.metadata
if (metadata == null) {
// Some players clear metadata when they pause or stop. Nothing can be
// logged without it, but playback has still ended.
if (state.state != PlaybackState.STATE_PLAYING) {
segments.end(controller.sessionToken, Instant.now())
}
return
}
handlePlaybackChange(controller, state, metadata)
}

override fun onMetadataChanged(metadata: MediaMetadata?) {
Expand All @@ -211,7 +227,7 @@ class MediaWatcher : NotificationListenerService() {
val token = controller.sessionToken
activeControllers.remove(token)
activeCallbacks.remove(token)?.let { controller.unregisterCallback(it) }
controller.packageName?.let { lastEventKeys.remove(it) }
segments.end(token, Instant.now())
Log.d(TAG, "Session destroyed for ${controller.packageName}")
}
}
Expand All @@ -225,7 +241,9 @@ class MediaWatcher : NotificationListenerService() {
state: PlaybackState,
metadata: MediaMetadata
) {
if (!connected) return
val packageName = controller.packageName ?: return
val token = controller.sessionToken

val title = metadata.getString(MediaMetadata.METADATA_KEY_TITLE) ?: ""
val artist = metadata.getString(MediaMetadata.METADATA_KEY_ARTIST)
Expand All @@ -237,11 +255,19 @@ class MediaWatcher : NotificationListenerService() {
PlaybackState.STATE_PAUSED -> "paused"
PlaybackState.STATE_STOPPED -> "stopped"
PlaybackState.STATE_BUFFERING -> "buffering"
else -> return // Ignore transitional states (none, connecting, etc.)
else -> {
// Transitional states (none, connecting, etc.) aren't logged, but they
// aren't playback either.
segments.end(token, Instant.now())
return
}
}

// Skip events with no useful metadata
if (title.isEmpty() && artist.isEmpty()) return
// Skip events with no useful metadata, ending any playback they replace.
if (title.isEmpty() && artist.isEmpty()) {
segments.end(token, Instant.now())
return
}

// Resolve app name from package
val appName = try {
Expand All @@ -264,18 +290,12 @@ class MediaWatcher : NotificationListenerService() {
put("state", playbackState)
}

// Deduplicate: don't send identical heartbeats
val eventKey = "$packageName|$title|$artist|$playbackState"
val lastKey = lastEventKeys[packageName]
if (eventKey == lastKey && playbackState == "playing") {
// Same track still playing — let heartbeat merging handle it
ri?.heartbeatHelper(BUCKET_ID, Instant.now(), 0.0, data, PULSETIME)
return
// Everything written to the event, so a segment never outlives the data it records.
val eventKey = "$packageName|$title|$artist|$album|$playbackState"
val changed = segments.observe(token, eventKey, data, playbackState == "playing", Instant.now())
if (changed) {
Log.i(TAG, "Media event: $playbackState — $artist - $title ($appName)")
}
lastEventKeys[packageName] = eventKey

Log.i(TAG, "Media event: $playbackState — $artist - $title ($appName)")
ri?.heartbeatHelper(BUCKET_ID, Instant.now(), 0.0, data, PULSETIME)
}

private fun pollActiveSessions() {
Expand All @@ -288,7 +308,18 @@ class MediaWatcher : NotificationListenerService() {
}
}

private fun unregisterAllCallbacks() {
// Called on the main thread. Teardown runs on the handler thread, after any session scan
// already running there, so a scan can never register callbacks after it.
private fun disconnect() {
// Drops queued polls, callbacks and a session scan that hasn't started yet.
handler?.removeCallbacksAndMessages(null)
val now = Instant.now()
handler?.post { unregisterAllCallbacks(now) }
}

// Handler thread only. Ends at [disconnectedAt] whatever is still playing.
private fun unregisterAllCallbacks(disconnectedAt: Instant) {
connected = false
// Remove the active sessions listener to prevent leaks
activeSessionsListener?.let { listener ->
sessionManager?.removeOnActiveSessionsChangedListener(listener)
Expand All @@ -301,6 +332,6 @@ class MediaWatcher : NotificationListenerService() {
}
activeControllers.clear()
activeCallbacks.clear()
lastEventKeys.clear()
segments.endAll(disconnectedAt)
}
}
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)
}
}
Loading