refactor(data): page the log export, remove dedupe leftovers (#7449)

This commit is contained in:
James Rich authored and GitHub committed 2026-09-29 20:58:28 +00:00
1 parent 2bc627283b
commit 733647bcbe
10 files changed
+150 -73

No files matched your search

@@ -17,6 +17,7 @@
package org.meshtastic.core.data.repository
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.channelFlow
import kotlinx.coroutines.flow.conflate
import kotlinx.coroutines.flow.distinctUntilChanged
import kotlinx.coroutines.flow.firstOrNull
@@ -35,6 +36,7 @@ import org.meshtastic.core.database.entity.asExternalModel
import org.meshtastic.core.di.CoroutineDispatchers
import org.meshtastic.core.model.MeshLog
import org.meshtastic.core.model.util.TELEMETRY_CHANNEL_COUNT
import org.meshtastic.core.model.util.TimeConstants
import org.meshtastic.core.model.util.adcVoltage
import org.meshtastic.core.model.util.oneWireTemperature
import org.meshtastic.core.model.util.withAdcVoltage
@@ -71,10 +73,22 @@ open class MeshLogRepositoryImpl(
.map { list -> list.map { it.asExternalModel() } }
.flowOn(dispatchers.io)
/** Retrieves all [MeshLog]s in the database in the order they were received. */
override fun getAllLogsInReceiveOrder(maxItem: Int): Flow<List<MeshLog>> = dbManager
.observeCurrentDb { it.meshLogDao().getAllLogsInReceiveOrder(maxItem) }
.map { list -> list.map { it.asExternalModel() } }
/** Pages through one database under a single read lease, so the logs all come from the same device. */
override fun readAllLogsInReceiveOrder(): Flow<MeshLog> = channelFlow {
dbManager.withReadDb { db ->
val dao = db.meshLogDao()
var afterReceivedDate = Long.MIN_VALUE
var afterRowId = Long.MIN_VALUE
do {
val page = dao.getLogsInReceiveOrderAfter(afterReceivedDate, afterRowId, RECEIVE_ORDER_PAGE_SIZE)
page.forEach { send(it.log.asExternalModel()) }
page.lastOrNull()?.let { last ->
afterReceivedDate = last.log.received_date
afterRowId = last.rowId
}
} while (page.size == RECEIVE_ORDER_PAGE_SIZE)
}
}
.flowOn(dispatchers.io)
/** Retrieves all [MeshLog]s in the database without any limit. */
@@ -132,7 +146,7 @@ open class MeshLogRepositoryImpl(
telemetry
.newBuilder()
.also { wb ->
wb.time = (log.received_date / MILLIS_PER_SEC).toInt()
wb.time = (log.received_date / TimeConstants.MS_PER_SEC).toInt()
wb.environment_metrics = telemetry.environment_metrics?.withSentinelsForAbsentReadings()
}
.build()
@@ -241,7 +255,6 @@ open class MeshLogRepositoryImpl(
}
companion object {
private const val MILLIS_PER_SEC = 1000L
private const val TELEMETRY_SNAPSHOT_PAGE_SIZE = 512
/**
@@ -250,6 +263,8 @@ open class MeshLogRepositoryImpl(
* auto-checkpoint.
*/
internal const val RETENTION_DELETE_BATCH_SIZE = 500
internal const val RECEIVE_ORDER_PAGE_SIZE = 500
}
}
@@ -76,23 +76,29 @@ class PacketRepositoryImpl(private val dbManager: DatabaseProvider, private val
.flow
.map { pagingData -> pagingData.map { it.contact_key to it.data } }
override suspend fun getMessageCount(contact: String): Int =
dbManager.withReadDb { it.packetDao().getMessageCount(contact) }
override suspend fun getMessageCount(contact: String): Int = dbManager.withReadDb {
it.packetDao().getMessageCount(contact)
}
override suspend fun getUnreadCount(contact: String): Int =
dbManager.withReadDb { it.packetDao().getUnreadCount(contact) }
override suspend fun getUnreadCount(contact: String): Int = dbManager.withReadDb {
it.packetDao().getUnreadCount(contact)
}
override fun getUnreadCountFlow(contact: String): Flow<Int> =
dbManager.observeCurrentDb { db -> db.packetDao().getUnreadCountFlow(contact) }
override fun getUnreadCountFlow(contact: String): Flow<Int> = dbManager.observeCurrentDb { db ->
db.packetDao().getUnreadCountFlow(contact)
}
override fun getFirstUnreadMessageUuid(contact: String): Flow<Long?> =
dbManager.observeCurrentDb { db -> db.packetDao().getFirstUnreadMessageUuid(contact) }
override fun getFirstUnreadMessageUuid(contact: String): Flow<Long?> = dbManager.observeCurrentDb { db ->
db.packetDao().getFirstUnreadMessageUuid(contact)
}
override fun hasUnreadMessages(contact: String): Flow<Boolean> =
dbManager.observeCurrentDb { db -> db.packetDao().hasUnreadMessages(contact) }
override fun hasUnreadMessages(contact: String): Flow<Boolean> = dbManager.observeCurrentDb { db ->
db.packetDao().hasUnreadMessages(contact)
}
override fun getUnreadCountTotal(): Flow<Int> =
dbManager.observeCurrentDb { db -> db.packetDao().getUnreadCountTotal() }
override fun getUnreadCountTotal(): Flow<Int> = dbManager.observeCurrentDb { db ->
db.packetDao().getUnreadCountTotal()
}
// One-shot writes go through withDb so they register with the cross-transport merge drain barrier. The callback
// is never replayed after it starts; callers needing retries must make that policy explicit where idempotency is
@@ -390,23 +396,14 @@ class PacketRepositoryImpl(private val dbManager: DatabaseProvider, private val
read: Boolean,
filtered: Boolean,
) {
val packetToSave =
RoomPacket(
uuid = 0L,
myNodeNum = myNodeNum,
packetId = packet.id,
port_num = packet.dataType,
contact_key = contactKey,
received_time = receivedTime,
read = read,
data = packet,
snr = packet.snr,
rssi = packet.rssi,
hopsAway = packet.hopsAway,
filtered = filtered,
messageText = packet.text.orEmpty(),
)
insertRoomPacket(packetToSave)
savePacket(
myNodeNum = myNodeNum,
contactKey = contactKey,
packet = packet,
receivedTime = receivedTime,
read = read,
filtered = filtered,
)
}
override suspend fun update(packet: DataPacket, routingError: Int): Unit =
@@ -498,11 +495,13 @@ class PacketRepositoryImpl(private val dbManager: DatabaseProvider, private val
withContext(dispatchers.io) { dbManager.withDb { it.packetDao().update(reaction) } }
}
override fun getFilteredCountFlow(contactKey: String): Flow<Int> =
dbManager.observeCurrentDb { db -> db.packetDao().getFilteredCountFlow(contactKey) }
override fun getFilteredCountFlow(contactKey: String): Flow<Int> = dbManager.observeCurrentDb { db ->
db.packetDao().getFilteredCountFlow(contactKey)
}
override suspend fun getFilteredCount(contactKey: String): Int =
dbManager.withReadDb { it.packetDao().getFilteredCount(contactKey) }
override suspend fun getFilteredCount(contactKey: String): Int = dbManager.withReadDb {
it.packetDao().getFilteredCount(contactKey)
}
override suspend fun setContactFilteringDisabled(contactKey: String, disabled: Boolean) {
withContext(dispatchers.io) {
@@ -25,6 +25,7 @@ import kotlinx.coroutines.Job
import kotlinx.coroutines.NonCancellable
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.first
import kotlinx.coroutines.flow.toList
import kotlinx.coroutines.launch
import kotlinx.coroutines.test.UnconfinedTestDispatcher
import kotlinx.coroutines.test.runTest
@@ -338,6 +339,28 @@ abstract class CommonMeshLogRepositoryTest {
private fun retentionLog(uuid: String, receivedDate: Long) =
MeshLog(uuid = uuid, message_type = "TEXT", received_date = receivedDate, raw_message = "")
@Test
fun `readAllLogsInReceiveOrder returns every log across pages oldest first and ties in insertion order`() =
runTest(testDispatcher) {
val count = MeshLogRepositoryImpl.RECEIVE_ORDER_PAGE_SIZE * 5 / 2
// Seven logs share most received dates, so tie groups straddle both page boundaries, and uuids run
// against insertion order so ordering ties by uuid would reverse every group.
val logs =
List(count) { i ->
MeshLog(
uuid = (count - i).toString().padStart(5, '0'),
message_type = "TEXT",
received_date = (i % 179).toLong(),
raw_message = "",
)
}
dbProvider.currentDb.value.meshLogDao().insertIgnore(logs.map { it.asEntity() })
val read = repository.readAllLogsInReceiveOrder().toList()
assertEquals(logs.sortedBy { it.received_date }.map { it.uuid }, read.map { it.uuid })
}
@Test
fun `parseTelemetryLog lifts legacy one-wire list onto per-channel fields`() = runTest(testDispatcher) {
// Firmware before 2.8 emitted the repeated field; stored logs must still chart after the repoint.
@@ -24,6 +24,7 @@ import androidx.room3.Transaction
import kotlinx.coroutines.flow.Flow
import org.meshtastic.core.database.DatabaseConstants.SQLITE_MAX_BIND_PARAMETERS
import org.meshtastic.core.database.entity.MeshLog
import org.meshtastic.core.database.entity.MeshLogRow
@Dao
@Suppress("TooManyFunctions")
@@ -42,8 +43,20 @@ interface MeshLogDao {
@Query("SELECT * FROM log ORDER BY received_date DESC LIMIT :maxItem")
fun getAllLogs(maxItem: Int): Flow<List<MeshLog>>
@Query("SELECT * FROM log ORDER BY received_date ASC LIMIT :maxItem")
fun getAllLogsInReceiveOrder(maxItem: Int): Flow<List<MeshLog>>
/**
* Returns up to [pageSize] logs after ([afterReceivedDate], [afterRowId]), oldest first. Equal received dates stay
* in rowid order, which is the order a scan of the received_date index returns. Start from [Long.MIN_VALUE] for
* both and continue from the last row returned.
*/
@Query(
"""
SELECT rowid AS log_rowid, * FROM log
WHERE (received_date, rowid) > (:afterReceivedDate, :afterRowId)
ORDER BY received_date ASC, rowid ASC
LIMIT :pageSize
""",
)
suspend fun getLogsInReceiveOrderAfter(afterReceivedDate: Long, afterRowId: Long, pageSize: Int): List<MeshLogRow>
/**
* Retrieves [MeshLog]s matching 'from_num' (nodeNum) and 'port_num' (PortNum).
@@ -17,6 +17,7 @@
package org.meshtastic.core.database.entity
import androidx.room3.ColumnInfo
import androidx.room3.Embedded
import androidx.room3.Entity
import androidx.room3.Index
import androidx.room3.PrimaryKey
@@ -88,6 +89,9 @@ data class MeshLog(
}
}
/** A [MeshLog] with its rowid, the cursor that continues a receive-order page read. */
data class MeshLogRow(@ColumnInfo(name = "log_rowid") val rowId: Long, @Embedded val log: MeshLog)
fun MeshLog.asExternalModel() = ExternalMeshLog(
uuid = uuid,
message_type = message_type,
@@ -144,31 +144,7 @@ data class Packet(
@ColumnInfo(name = "message_text", defaultValue = "") val messageText: String = "",
@ColumnInfo(name = "translated_text") val translatedText: String? = null,
@ColumnInfo(name = "show_translated", defaultValue = "0") val showTranslated: Boolean = false,
) {
companion object {
const val RELAY_NODE_SUFFIX_MASK = 0xFF
fun getRelayNode(relayNodeId: Int, nodes: List<Node>, ourNodeNum: Int?): Node? {
val relayNodeIdSuffix = relayNodeId and RELAY_NODE_SUFFIX_MASK
val candidateRelayNodes =
nodes.filter {
it.num != ourNodeNum &&
it.lastHeard != 0 &&
(it.num and RELAY_NODE_SUFFIX_MASK) == relayNodeIdSuffix
}
val closestRelayNode =
if (candidateRelayNodes.size == 1) {
candidateRelayNodes.first()
} else {
candidateRelayNodes.minByOrNull { it.hopsAway }
}
return closestRelayNode
}
}
}
)
@Suppress("ConstructorParameterNaming")
@Entity(tableName = "contact_settings")
@@ -16,7 +16,6 @@
*/
package org.meshtastic.core.domain.usecase.settings
import kotlinx.coroutines.flow.first
import kotlinx.datetime.TimeZone
import kotlinx.datetime.toLocalDateTime
import okio.BufferedSink
@@ -59,7 +58,7 @@ constructor(
"\"date\",\"time\",\"from\",\"sender name\",\"sender lat\",\"sender long\",\"rx lat\",\"rx long\",\"rx elevation\",\"rx snr\",\"distance(m)\",\"hop limit\",\"hop start\",\"relay node\",\"payload\"\n",
)
meshLogRepository.getAllLogsInReceiveOrder(Int.MAX_VALUE).first().forEach { packet ->
meshLogRepository.readAllLogsInReceiveOrder().collect { packet ->
packet.nodeInfo?.let { nodeInfo ->
positionToPos.invoke(nodeInfo.position)?.let { nodePositions[nodeInfo.num] = nodeInfo.position }
}
@@ -28,6 +28,7 @@ import org.meshtastic.proto.MeshPacket
import org.meshtastic.proto.PortNum
import kotlin.test.BeforeTest
import kotlin.test.Test
import kotlin.test.assertEquals
import kotlin.test.assertTrue
class ExportDataUseCaseTest {
@@ -170,4 +171,44 @@ class ExportDataUseCaseTest {
// relay_node defaults to 0 (unset) -> blank field before the payload
assertTrue(output.contains("\"0\",\"\",\"Hello\""))
}
@Test
fun `invoke writes one row per log in the order the repository reads them`() = runTest {
val senders = listOf(30, 10, 20, 50, 40)
meshLogRepository.setLogs(
senders.mapIndexed { index, from ->
MeshLog(
uuid = "$index",
message_type = "TEXT",
received_date = 1000000000L + index,
raw_message = "",
fromRadio =
FromRadio.Builder()
.also { wb ->
wb.packet =
MeshPacket.Builder()
.also { wb ->
wb.from = from
wb.rx_snr = 5.0f
wb.decoded =
Data.Builder()
.also { wb ->
wb.portnum = PortNum.TEXT_MESSAGE_APP
wb.payload = "Hi".encodeUtf8()
}
.build()
}
.build()
}
.build(),
)
},
)
val buffer = Buffer()
useCase(buffer, 1)
val rows = buffer.readUtf8().lines().drop(1).filter { it.isNotEmpty() }
assertEquals(senders.map { it.toString() }, rows.map { it.split("\",\"")[2] })
}
}
@@ -36,8 +36,11 @@ interface MeshLogRepository {
/** Retrieves all [MeshLog]s in the database, up to [maxItem]. */
fun getAllLogs(maxItem: Int = DEFAULT_MAX_LOGS): Flow<List<MeshLog>>
/** Retrieves all [MeshLog]s in the database in the order they were received. */
fun getAllLogsInReceiveOrder(maxItem: Int = DEFAULT_MAX_LOGS): Flow<List<MeshLog>>
/**
* Emits every [MeshLog] once, oldest first, then completes. The logs are read a page at a time rather than held in
* memory together.
*/
fun readAllLogsInReceiveOrder(): Flow<MeshLog>
/** Retrieves all [MeshLog]s in the database without any limit. */
fun getAllLogsUnbounded(): Flow<List<MeshLog>>
@@ -18,6 +18,7 @@ package org.meshtastic.core.testing
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.flow
import kotlinx.coroutines.flow.map
import org.meshtastic.core.common.util.nowMillis
import org.meshtastic.core.model.MeshLog
@@ -63,12 +64,15 @@ class FakeMeshLogRepository :
override fun getAllLogs(maxItem: Int): Flow<List<MeshLog>> = logsFlow.map { it.take(maxItem) }
override fun getAllLogsInReceiveOrder(maxItem: Int): Flow<List<MeshLog>> = logsFlow.map { it.take(maxItem) }
override fun readAllLogsInReceiveOrder(): Flow<MeshLog> = flow {
logsFlow.value.sortedBy { it.received_date }.forEach { emit(it) }
}
override fun getAllLogsUnbounded(): Flow<List<MeshLog>> = logsFlow
override fun getLogsFrom(nodeNum: Int, portNum: Int): Flow<List<MeshLog>> =
logsFlow.map { it.filter { log -> log.fromNum == nodeNum && log.portNum == portNum } }
override fun getLogsFrom(nodeNum: Int, portNum: Int): Flow<List<MeshLog>> = logsFlow.map {
it.filter { log -> log.fromNum == nodeNum && log.portNum == portNum }
}
override fun getMeshPacketsFrom(nodeNum: Int, portNum: Int): Flow<List<MeshPacket>> = MutableStateFlow(emptyList())