feat(ble): CommandQueue · 응답 대기 후 다음 명령 전송 (FW 팀 지시)
배경 (2026-08-04):
FW 팀 지시 — 이전 명령의 응답 도착 전에 다음 명령 보내지 말 것.
FW GATT queue 가 얕아 백투백 write 시 silent drop / freeze 유발.
기존 isMtbBusy gate 는 일부만 방어 (battery/IMU polling ↔ mtb 중간)
→ 모든 명령에 통일된 sequential 큐 도입.
CommandQueue.kt (신규 · 155 line):
- QueuedCommand(data, label, isDone(tag3), timeoutMs=3000)
- enqueue → pending==null 이면 즉시 dispatch, 아니면 대기
- 응답 tag 마다 onResponseTag() → pending.isDone 체크
- 공통 에러 응답 (rxx/rxd/rxn/rxc/rxs) 은 무조건 완료 처리
- 3초 timeout → onDrop 콜백 (UI observer 갱신 포함)
- clear(reason) — disconnect 경로에서 pending + 큐 전체 drop
- Thread-safe (synchronized on internal lock)
BleManager.kt:
- sendRaw() → sendRawWrite() private (실제 GATT write)
- 새 sendRaw() = 단일 응답 (m→r 자동 치환) 명령용 wrapper · enqueue 위임
- 스트리밍 명령은 자체 send fn 에서 explicit enqueue:
· maa/mbb → 종료 = raa && piezoCollector.isComplete
· mtb → 종료 = (rim|ric|raa) && piezo done && imu samples 채워짐
· mim → 종료 = (rim|ric) && imu samples 채워짐
· mec → 종료 = raa
- processReceivedData 말미에 commandQueue.onResponseTag(tag3) hook
- 공통 에러 응답 5종 (rxx/rxd/rxn/rxc/rxs) case 신설 · 로그
- disconnect 5 경로에 commandQueue.clear("disconnect") 추가
- onDrop 콜백에서 mtb timeout 시 UI observer (lastMtbTimeoutAt +
consecutiveMtbTimeouts) 갱신 · 기존 UX 유지
BLE 스펙 참조: ppt/ble.xlsx (25 명령 · 응답 태그 표).
This commit is contained in:
@@ -370,6 +370,7 @@ class BleManager private constructor(private val context: Context) {
|
|||||||
stopBatteryPolling()
|
stopBatteryPolling()
|
||||||
stopRssiPolling()
|
stopRssiPolling()
|
||||||
rssi.value = null
|
rssi.value = null
|
||||||
|
commandQueue.clear("disconnect")
|
||||||
com.medithings.vesiscan.services.BleForegroundService.stop(context)
|
com.medithings.vesiscan.services.BleForegroundService.stop(context)
|
||||||
connectionTimer?.let { handler.removeCallbacks(it) }
|
connectionTimer?.let { handler.removeCallbacks(it) }
|
||||||
connectionTimer = null
|
connectionTimer = null
|
||||||
@@ -462,15 +463,40 @@ class BleManager private constructor(private val context: Context) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Send
|
// Send
|
||||||
fun sendRaw(data: ByteArray) {
|
/**
|
||||||
|
* 명령 큐 (2026-08-04 · FW 팀 지시로 도입). 모든 명령은 depth=1 sequential.
|
||||||
|
* 이전 응답 도착 · 또는 3초 timeout 전까지 다음 명령은 대기.
|
||||||
|
*
|
||||||
|
* 사용:
|
||||||
|
* - `sendRaw(data)` 는 이제 단일 응답 (m→r) 명령용 wrapper (자동 enqueue).
|
||||||
|
* - 스트리밍 응답 (mtb/maa/mbb/mim/mec) 은 각 send fn 에서 `commandQueue.enqueue()`
|
||||||
|
* 로 직접 완료 조건 (piezoCollector.isComplete + imuCollector.samples 등) 지정.
|
||||||
|
*/
|
||||||
|
val commandQueue = CommandQueue(
|
||||||
|
handler = handler,
|
||||||
|
writeRaw = { data -> sendRawWrite(data) },
|
||||||
|
onDrop = { label, reason ->
|
||||||
|
debugLogger.warn("CMDQ drop label=$label · $reason")
|
||||||
|
// mtb timeout 이면 기존 UI observers 유지 (연속 timeout 카운터 등)
|
||||||
|
if (label.startsWith("mtb") && reason == "timeout") {
|
||||||
|
lastMtbTimeoutAt.value = System.currentTimeMillis()
|
||||||
|
consecutiveMtbTimeouts.value = consecutiveMtbTimeouts.value + 1
|
||||||
|
piezoCollector.reset()
|
||||||
|
imuCollector.reset()
|
||||||
|
}
|
||||||
|
},
|
||||||
|
)
|
||||||
|
|
||||||
|
/** 실제 GATT write — CommandQueue 내부에서만 호출. 외부는 sendRaw() 사용. */
|
||||||
|
private fun sendRawWrite(data: ByteArray) {
|
||||||
val characteristic = txCharacteristic
|
val characteristic = txCharacteristic
|
||||||
val gatt = bluetoothGatt
|
val gatt = bluetoothGatt
|
||||||
if (characteristic == null) {
|
if (characteristic == null) {
|
||||||
loge { "sendRaw: txCharacteristic is null!" }
|
loge { "sendRawWrite: txCharacteristic is null!" }
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
if (gatt == null) {
|
if (gatt == null) {
|
||||||
loge { "sendRaw: bluetoothGatt is null!" }
|
loge { "sendRawWrite: bluetoothGatt is null!" }
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
val cmdPreview = if (data.size >= 4) String(data, 0, 3, Charsets.US_ASCII) else "?"
|
val cmdPreview = if (data.size >= 4) String(data, 0, 3, Charsets.US_ASCII) else "?"
|
||||||
@@ -495,21 +521,42 @@ class BleManager private constructor(private val context: Context) {
|
|||||||
"mbb" -> "full measurement (battery+IMU+temp+6ch)"
|
"mbb" -> "full measurement (battery+IMU+temp+6ch)"
|
||||||
else -> ""
|
else -> ""
|
||||||
}
|
}
|
||||||
logd { "sendRaw: $cmdPreview (${data.size} bytes)" }
|
logd { "sendRawWrite: $cmdPreview (${data.size} bytes)" }
|
||||||
debugLogger.tx("$cmdPreview $cmdDetail", data.size)
|
debugLogger.tx("$cmdPreview $cmdDetail", data.size)
|
||||||
if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.TIRAMISU) {
|
if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.TIRAMISU) {
|
||||||
val result = gatt.writeCharacteristic(characteristic, data, BluetoothGattCharacteristic.WRITE_TYPE_NO_RESPONSE)
|
val result = gatt.writeCharacteristic(characteristic, data, BluetoothGattCharacteristic.WRITE_TYPE_NO_RESPONSE)
|
||||||
logd { "sendRaw result: $result" }
|
logd { "sendRawWrite result: $result" }
|
||||||
} else {
|
} else {
|
||||||
@Suppress("DEPRECATION")
|
@Suppress("DEPRECATION")
|
||||||
characteristic.value = data
|
characteristic.value = data
|
||||||
characteristic.writeType = BluetoothGattCharacteristic.WRITE_TYPE_NO_RESPONSE
|
characteristic.writeType = BluetoothGattCharacteristic.WRITE_TYPE_NO_RESPONSE
|
||||||
@Suppress("DEPRECATION")
|
@Suppress("DEPRECATION")
|
||||||
val ok = gatt.writeCharacteristic(characteristic)
|
val ok = gatt.writeCharacteristic(characteristic)
|
||||||
logd { "sendRaw success: $ok" }
|
logd { "sendRawWrite success: $ok" }
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 단일 응답 명령 (m→r 대응) 을 큐에 삽입.
|
||||||
|
*
|
||||||
|
* 자동 완료 조건: 3-byte 접두 (예: msn) → r + 뒤 2byte (rsn) 도착 시 완료.
|
||||||
|
* 스트리밍 명령 (maa/mtb/mbb/mim/mec) 은 자체 send fn 에서 `commandQueue.enqueue()` 로
|
||||||
|
* 직접 완료 조건 지정. 이 함수로 넣으면 잘못된 tag ("rmm" 등) 를 기다리다 3초 timeout.
|
||||||
|
*/
|
||||||
|
fun sendRaw(data: ByteArray) {
|
||||||
|
if (data.size < 4) {
|
||||||
|
loge { "sendRaw: too short (${data.size}B)" }
|
||||||
|
return
|
||||||
|
}
|
||||||
|
val cmd3 = String(data, 0, 3, Charsets.US_ASCII)
|
||||||
|
val expected = "r" + cmd3.substring(1) // msn → rsn, mid → rid, mls → rls
|
||||||
|
commandQueue.enqueue(QueuedCommand(
|
||||||
|
data = data,
|
||||||
|
label = cmd3,
|
||||||
|
isDone = { it == expected },
|
||||||
|
))
|
||||||
|
}
|
||||||
|
|
||||||
// Piezo Commands
|
// Piezo Commands
|
||||||
val piezoCollector = PiezoPacketCollector()
|
val piezoCollector = PiezoPacketCollector()
|
||||||
// IMU rim: 응답 파서. mtb? 명령 시 채워짐. PlacementGuide/Clinical에서 콜백 등록.
|
// IMU rim: 응답 파서. mtb? 명령 시 채워짐. PlacementGuide/Clinical에서 콜백 등록.
|
||||||
@@ -526,14 +573,24 @@ class BleManager private constructor(private val context: Context) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
fun sendBurst(freqOption: Int, cycles: Int, delayUs: Int, piezoCh: Int = 0) {
|
fun sendBurst(freqOption: Int, cycles: Int, delayUs: Int, piezoCh: Int = 0) {
|
||||||
sendRaw(CRC16.buildCommandBE("mec", intArrayOf(freqOption, delayUs, 140, cycles, 1, piezoCh)))
|
// mec: reb×1 → raa. 종료 = raa.
|
||||||
|
commandQueue.enqueue(QueuedCommand(
|
||||||
|
data = CRC16.buildCommandBE("mec", intArrayOf(freqOption, delayUs, 140, cycles, 1, piezoCh)),
|
||||||
|
label = "mec burst",
|
||||||
|
isDone = { it == "raa" },
|
||||||
|
))
|
||||||
}
|
}
|
||||||
|
|
||||||
fun sendAllChannels(mode: Int = 0): Boolean {
|
fun sendAllChannels(mode: Int = 0): Boolean {
|
||||||
if (!canSendMaa("sendAllChannels")) return false
|
if (!canSendMaa("sendAllChannels")) return false
|
||||||
lastMaaSentMs = System.currentTimeMillis()
|
lastMaaSentMs = System.currentTimeMillis()
|
||||||
piezoCollector.startMultiChannel(6)
|
piezoCollector.startMultiChannel(6)
|
||||||
sendRaw(CRC16.buildCommandASCII("maa", " "))
|
// maa: reb×6 → raa. 종료 = raa.
|
||||||
|
commandQueue.enqueue(QueuedCommand(
|
||||||
|
data = CRC16.buildCommandASCII("maa", " "),
|
||||||
|
label = "maa 6ch",
|
||||||
|
isDone = { it == "raa" },
|
||||||
|
))
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -548,22 +605,18 @@ class BleManager private constructor(private val context: Context) {
|
|||||||
lastMaaSentMs = System.currentTimeMillis()
|
lastMaaSentMs = System.currentTimeMillis()
|
||||||
piezoCollector.startMultiChannel(6)
|
piezoCollector.startMultiChannel(6)
|
||||||
imuCollector.reset()
|
imuCollector.reset()
|
||||||
sendRaw(CRC16.buildCommandASCII("mtb", " "))
|
// mtb: reb×6 → raa → rim (또는 ric 분할). 종료 = piezo raa 완료 AND imu 샘플 도착.
|
||||||
// 2026-07-08: 3초 timeout. raa 응답이 오면 취소 (아래 processReceivedData 참조).
|
// parseRic 는 마지막 chunk 도착 시에만 samples 를 채움 → 중간 ric 은 자연스레 skip.
|
||||||
// Timeout 시 lastMtbTimeoutAt 갱신 + consecutiveMtbTimeouts 증가 → 상위 UI 안내.
|
commandQueue.enqueue(QueuedCommand(
|
||||||
mtbTimeoutRunnable?.let { handler.removeCallbacks(it) }
|
data = CRC16.buildCommandASCII("mtb", " "),
|
||||||
val to = Runnable {
|
label = "mtb 6ch+imu",
|
||||||
if (isMtbBusy) {
|
isDone = { tag ->
|
||||||
lastMtbTimeoutAt.value = System.currentTimeMillis()
|
(tag == "rim" || tag == "ric" || tag == "raa") &&
|
||||||
consecutiveMtbTimeouts.value = consecutiveMtbTimeouts.value + 1
|
piezoCollector.isComplete &&
|
||||||
debugLogger.warn("MTB_TIMEOUT after ${MTB_TIMEOUT_MS}ms (consecutive=${consecutiveMtbTimeouts.value})")
|
imuCollector.samples.isNotEmpty()
|
||||||
// collector 상태 해제 — 다음 mtb 는 정상적으로 시작 가능
|
},
|
||||||
piezoCollector.reset()
|
timeoutMs = MTB_TIMEOUT_MS,
|
||||||
imuCollector.reset()
|
))
|
||||||
}
|
|
||||||
}
|
|
||||||
mtbTimeoutRunnable = to
|
|
||||||
handler.postDelayed(to, MTB_TIMEOUT_MS)
|
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -642,13 +695,25 @@ class BleManager private constructor(private val context: Context) {
|
|||||||
*/
|
*/
|
||||||
fun sendImuFifoQuery() {
|
fun sendImuFifoQuery() {
|
||||||
imuCollector.reset()
|
imuCollector.reset()
|
||||||
sendRaw(CRC16.buildCommandASCII("mim", " "))
|
// mim: rim: (또는 ric 분할). 종료 = imu 샘플 채워짐 (마지막 rim 또는 마지막 ric).
|
||||||
|
commandQueue.enqueue(QueuedCommand(
|
||||||
|
data = CRC16.buildCommandASCII("mim", " "),
|
||||||
|
label = "mim IMU",
|
||||||
|
isDone = { tag ->
|
||||||
|
(tag == "rim" || tag == "ric") && imuCollector.samples.isNotEmpty()
|
||||||
|
},
|
||||||
|
))
|
||||||
}
|
}
|
||||||
|
|
||||||
/** 전체 측정: 배터리+IMU+온도+6ch piezo (rbb→reb×6→raa) */
|
/** 전체 측정: 배터리+IMU+온도+6ch piezo (rbb→reb×6→raa) */
|
||||||
fun sendFullMeasurement() {
|
fun sendFullMeasurement() {
|
||||||
piezoCollector.startMultiChannel(6)
|
piezoCollector.startMultiChannel(6)
|
||||||
sendRaw(CRC16.buildCommandASCII("mbb", " "))
|
// mbb: rbb → reb×6 → raa. 종료 = raa (piezo 완료).
|
||||||
|
commandQueue.enqueue(QueuedCommand(
|
||||||
|
data = CRC16.buildCommandASCII("mbb", " "),
|
||||||
|
label = "mbb full",
|
||||||
|
isDone = { it == "raa" && piezoCollector.isComplete },
|
||||||
|
))
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -683,7 +748,12 @@ class BleManager private constructor(private val context: Context) {
|
|||||||
if (!canSendMaa("sendChannelsOnly")) return false
|
if (!canSendMaa("sendChannelsOnly")) return false
|
||||||
lastMaaSentMs = System.currentTimeMillis()
|
lastMaaSentMs = System.currentTimeMillis()
|
||||||
piezoCollector.startMultiChannel(6)
|
piezoCollector.startMultiChannel(6)
|
||||||
sendRaw(CRC16.buildCommandASCII("maa", " "))
|
// maa: reb×6 → raa. 종료 = raa (piezo 완료).
|
||||||
|
commandQueue.enqueue(QueuedCommand(
|
||||||
|
data = CRC16.buildCommandASCII("maa", " "),
|
||||||
|
label = "maa 6ch",
|
||||||
|
isDone = { it == "raa" && piezoCollector.isComplete },
|
||||||
|
))
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -833,6 +903,7 @@ class BleManager private constructor(private val context: Context) {
|
|||||||
stopBatteryPolling()
|
stopBatteryPolling()
|
||||||
stopRssiPolling()
|
stopRssiPolling()
|
||||||
rssi.value = null
|
rssi.value = null
|
||||||
|
commandQueue.clear("disconnect")
|
||||||
stopWatchdog()
|
stopWatchdog()
|
||||||
isConnected.value = false
|
isConnected.value = false
|
||||||
isServiceReady.value = false
|
isServiceReady.value = false
|
||||||
@@ -1310,10 +1381,10 @@ class BleManager private constructor(private val context: Context) {
|
|||||||
val tag = prefix.take(3)
|
val tag = prefix.take(3)
|
||||||
debugLogger.rx(tag, data.size, if (tag == "raa") "all-ch complete" else "single-ch end")
|
debugLogger.rx(tag, data.size, if (tag == "raa") "all-ch complete" else "single-ch end")
|
||||||
piezoCollector.addPacket(data)
|
piezoCollector.addPacket(data)
|
||||||
// 2026-07-08: mtb 응답 완료 → timeout runnable 취소 + 연속 실패 카운터 리셋.
|
// raa 성공 → 연속 timeout 카운터 리셋 (mtb 완료 UI 관찰).
|
||||||
if (tag == "raa") {
|
// 3초 timeout 은 CommandQueue 가 소유 (기존 mtbTimeoutRunnable 은 dead code · 유지만).
|
||||||
mtbTimeoutRunnable?.let { handler.removeCallbacks(it); mtbTimeoutRunnable = null }
|
if (tag == "raa" && consecutiveMtbTimeouts.value != 0) {
|
||||||
if (consecutiveMtbTimeouts.value != 0) consecutiveMtbTimeouts.value = 0
|
consecutiveMtbTimeouts.value = 0
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
"rsn:" -> {
|
"rsn:" -> {
|
||||||
@@ -1380,10 +1451,20 @@ class BleManager private constructor(private val context: Context) {
|
|||||||
imuCollector.parseRsp(data)
|
imuCollector.parseRsp(data)
|
||||||
onImuReceived?.invoke(data)
|
onImuReceived?.invoke(data)
|
||||||
}
|
}
|
||||||
"rxs:" -> {
|
// ── 공통 에러 응답 (xlsx §공통 프레임) — CommandQueue 가 pending 완료 처리 ──
|
||||||
|
// rxx: Unknown command · rxd: Disabled · rxn: NULL handler · rxc: CRC fail · rxs: Too short
|
||||||
|
"rxx:", "rxd:", "rxn:", "rxc:", "rxs:" -> {
|
||||||
|
val tag = prefix.take(3)
|
||||||
val echoCmd = if (data.size > 4) String(data, 4, minOf(data.size - 4, 10), Charsets.US_ASCII).trim() else "?"
|
val echoCmd = if (data.size > 4) String(data, 4, minOf(data.size - 4, 10), Charsets.US_ASCII).trim() else "?"
|
||||||
android.util.Log.w("ImuDbg", "rxs cmd_not_supported: $echoCmd")
|
val reason = when (tag) {
|
||||||
debugLogger.rx("rxs", data.size, "cmd_not_supported: $echoCmd")
|
"rxx" -> "unknown_cmd"
|
||||||
|
"rxd" -> "disabled_cmd"
|
||||||
|
"rxn" -> "null_handler"
|
||||||
|
"rxc" -> "crc_fail"
|
||||||
|
"rxs" -> "too_short"
|
||||||
|
else -> "err"
|
||||||
|
}
|
||||||
|
debugLogger.rx(tag, data.size, "$reason: $echoCmd")
|
||||||
}
|
}
|
||||||
"rim:" -> {
|
"rim:" -> {
|
||||||
val num = if (data.size >= 6) ((data[4].toInt() and 0xFF) shl 8) or (data[5].toInt() and 0xFF) else 0
|
val num = if (data.size >= 6) ((data[4].toInt() and 0xFF) shl 8) or (data[5].toInt() and 0xFF) else 0
|
||||||
@@ -1404,5 +1485,12 @@ class BleManager private constructor(private val context: Context) {
|
|||||||
logw { "Unknown prefix: '$prefix' hex=$hex" }
|
logw { "Unknown prefix: '$prefix' hex=$hex" }
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
// ── CommandQueue hook — 파서가 collector 상태 갱신한 뒤 완료 판정 ──
|
||||||
|
// 3-char 태그 (예: "raa", "rim", "rxc") 로 현재 pending 명령의 isDone 실행.
|
||||||
|
// True 면 pending 해제 + 큐에 대기 중이던 다음 명령 dispatch.
|
||||||
|
if (data.size >= 4) {
|
||||||
|
val tag3 = prefix.take(3)
|
||||||
|
commandQueue.onResponseTag(tag3)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,164 @@
|
|||||||
|
/*
|
||||||
|
* CommandQueue — BLE 명령 순차 전송 큐 (depth = 1).
|
||||||
|
*
|
||||||
|
* 설계 배경 (2026-08-04 · FW 팀 지시):
|
||||||
|
* - 이전 명령의 응답이 도착하기 전에 다음 명령을 보내지 말 것.
|
||||||
|
* - FW GATT queue 가 얕아 백투백 write 시 silent drop 가능.
|
||||||
|
* - 특히 스트리밍 응답 (mtb: reb×6 → raa → rim) 중간에 다른 명령이 끼면
|
||||||
|
* 파싱이 꼬여 freeze 유발 (기존 isMtbBusy gate 로 부분 방어 중이었음).
|
||||||
|
*
|
||||||
|
* 동작:
|
||||||
|
* - enqueue() 로 명령 삽입. pending==null 이면 즉시 dispatch, 아니면 큐 대기.
|
||||||
|
* - 각 응답 tag 도착 시 onResponseTag() → 현재 pending 의 isDone 체크.
|
||||||
|
* true 면 pending 해제 + 다음 dispatch.
|
||||||
|
* - 3초 timeout 이면 강제 완료 (다음 명령 진행 · deadlock 방지).
|
||||||
|
* - 공통 에러 응답 (rxx/rxd/rxc/rxs/rxn) 은 항상 pending 완료 처리
|
||||||
|
* (응답은 도착한 것이므로 대기 종료 · 다음 명령 dispatch).
|
||||||
|
*
|
||||||
|
* BLE 스펙 참조: docs/../ppt/ble.xlsx (25 명령 · 응답 태그 표).
|
||||||
|
*/
|
||||||
|
package com.medithings.vesiscan.ble
|
||||||
|
|
||||||
|
import android.os.Handler
|
||||||
|
import android.util.Log
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 큐에 들어가는 단일 명령 단위.
|
||||||
|
*
|
||||||
|
* @param data 실제 GATT write payload (CRC 포함).
|
||||||
|
* @param label 로그·debugLogger 표시명 (예: "mtb 6ch+imu").
|
||||||
|
* @param isDone 응답 tag 3자 (예: "raa", "rim") 를 받아 완료 여부 반환.
|
||||||
|
* 스트리밍 응답은 collector 상태를 함께 참조 가능.
|
||||||
|
* @param timeoutMs 응답 대기 최대 시간. 초과 시 강제 완료.
|
||||||
|
*/
|
||||||
|
data class QueuedCommand(
|
||||||
|
val data: ByteArray,
|
||||||
|
val label: String,
|
||||||
|
val isDone: (tag3: String) -> Boolean,
|
||||||
|
val timeoutMs: Long = 3000L,
|
||||||
|
) {
|
||||||
|
override fun equals(other: Any?): Boolean =
|
||||||
|
other is QueuedCommand && data.contentEquals(other.data) && label == other.label
|
||||||
|
override fun hashCode(): Int = 31 * data.contentHashCode() + label.hashCode()
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* BleManager 가 소유하는 명령 큐.
|
||||||
|
*
|
||||||
|
* @param handler 메인 루퍼 Handler (timeout · dispatch 는 메인 스레드에서 실행).
|
||||||
|
* @param writeRaw 실제 GATT write. commandQueue 는 이걸 호출만 함 (구현 몰라도 됨).
|
||||||
|
* @param onDrop (선택) 큐가 clear 되거나 timeout 시 호출. UI/로그용.
|
||||||
|
*/
|
||||||
|
class CommandQueue(
|
||||||
|
private val handler: Handler,
|
||||||
|
private val writeRaw: (ByteArray) -> Unit,
|
||||||
|
private val onDrop: ((label: String, reason: String) -> Unit)? = null,
|
||||||
|
) {
|
||||||
|
|
||||||
|
companion object {
|
||||||
|
private const val TAG = "CommandQueue"
|
||||||
|
/** 큐 soft cap — 이 개수 이상 쌓이면 신규 enqueue 를 debug log 로 경고. */
|
||||||
|
private const val WARN_QUEUE_DEPTH = 8
|
||||||
|
/** 공통 에러 응답 prefix (xlsx §공통 프레임): rxx/rxd/rxn/rxc/rxs. */
|
||||||
|
private val ERROR_TAGS = setOf("rxx", "rxd", "rxn", "rxc", "rxs")
|
||||||
|
}
|
||||||
|
|
||||||
|
private val lock = Any()
|
||||||
|
private val queue = ArrayDeque<QueuedCommand>()
|
||||||
|
@Volatile private var pending: QueuedCommand? = null
|
||||||
|
@Volatile private var pendingStartMs: Long = 0L
|
||||||
|
private var timeoutRunnable: Runnable? = null
|
||||||
|
|
||||||
|
/** 현재 처리 중인 명령이 있는지 (기존 isMtbBusy 대체). */
|
||||||
|
val isBusy: Boolean get() = pending != null
|
||||||
|
|
||||||
|
/** 현재 pending 명령의 label (디버그 UI 용). */
|
||||||
|
val pendingLabel: String? get() = pending?.label
|
||||||
|
|
||||||
|
/** 대기 중인 큐 길이 (디버그). */
|
||||||
|
val queuedSize: Int get() = synchronized(lock) { queue.size }
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 명령을 큐에 삽입. 현재 처리 중이면 대기, 아니면 즉시 dispatch.
|
||||||
|
* 스레드 안전 (여러 스레드에서 호출 가능).
|
||||||
|
*/
|
||||||
|
fun enqueue(cmd: QueuedCommand) {
|
||||||
|
synchronized(lock) {
|
||||||
|
queue.addLast(cmd)
|
||||||
|
if (queue.size >= WARN_QUEUE_DEPTH) {
|
||||||
|
Log.w(TAG, "queue depth ${queue.size} — 상위 caller throttle 확인 필요")
|
||||||
|
}
|
||||||
|
if (pending == null) dispatchNextLocked()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 응답 태그 도착 시 호출 (processReceivedData 말미에서 hook).
|
||||||
|
* 공통 에러 응답이면 무조건 pending 완료, 아니면 pending.isDone 체크.
|
||||||
|
*/
|
||||||
|
fun onResponseTag(tag3: String) {
|
||||||
|
val cur = pending ?: return
|
||||||
|
val isError = tag3 in ERROR_TAGS
|
||||||
|
val done = isError || run {
|
||||||
|
try { cur.isDone(tag3) } catch (e: Exception) {
|
||||||
|
Log.w(TAG, "isDone threw for ${cur.label}: ${e.message}")
|
||||||
|
false
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if (done) {
|
||||||
|
val reason = if (isError) "error:$tag3" else "response:$tag3"
|
||||||
|
completePendingLocked(reason)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 큐 전체 제거 + pending drop. 연결 해제 등에서 호출.
|
||||||
|
*/
|
||||||
|
fun clear(reason: String) {
|
||||||
|
synchronized(lock) {
|
||||||
|
timeoutRunnable?.let { handler.removeCallbacks(it) }
|
||||||
|
timeoutRunnable = null
|
||||||
|
val dropped = pending
|
||||||
|
if (dropped != null) {
|
||||||
|
Log.w(TAG, "clear: dropped pending=${dropped.label} + queued=${queue.size} · $reason")
|
||||||
|
onDrop?.invoke(dropped.label, "cleared:$reason")
|
||||||
|
}
|
||||||
|
queue.forEach { onDrop?.invoke(it.label, "cleared:$reason") }
|
||||||
|
pending = null
|
||||||
|
queue.clear()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// ── internal ──
|
||||||
|
|
||||||
|
private fun dispatchNextLocked() {
|
||||||
|
val next = queue.removeFirstOrNull() ?: return
|
||||||
|
pending = next
|
||||||
|
pendingStartMs = System.currentTimeMillis()
|
||||||
|
try {
|
||||||
|
writeRaw(next.data)
|
||||||
|
} catch (e: Exception) {
|
||||||
|
Log.e(TAG, "writeRaw threw for ${next.label}: ${e.message}")
|
||||||
|
// write 자체 실패 → 응답도 없을 것 → 즉시 timeout 처리로 이어감
|
||||||
|
}
|
||||||
|
val to = Runnable {
|
||||||
|
val cur = pending
|
||||||
|
if (cur != null && cur === next) {
|
||||||
|
Log.w(TAG, "TIMEOUT: ${next.label} after ${next.timeoutMs}ms")
|
||||||
|
onDrop?.invoke(next.label, "timeout")
|
||||||
|
synchronized(lock) { completePendingLocked("timeout") }
|
||||||
|
}
|
||||||
|
}
|
||||||
|
timeoutRunnable = to
|
||||||
|
handler.postDelayed(to, next.timeoutMs)
|
||||||
|
}
|
||||||
|
|
||||||
|
private fun completePendingLocked(reason: String) {
|
||||||
|
synchronized(lock) {
|
||||||
|
timeoutRunnable?.let { handler.removeCallbacks(it) }
|
||||||
|
timeoutRunnable = null
|
||||||
|
pending = null
|
||||||
|
if (queue.isNotEmpty()) dispatchNextLocked()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user