mirror of
https://github.com/BrendanGreenlee/openclaw-android-heartbeat.git
synced 2026-08-17 16:49:14 +00:00
Import official OpenClaw Android app (apps/android @ 71a59512ba476df3328cf485d84748121cf341f2)
Pristine fork source for PROJ-0088. v1 will turn this into a WebView web shell; the WebSocket node infrastructure stays intact for v2 heartbeat.
This commit is contained in:
commit
7a380c40ed
656 changed files with 209982 additions and 0 deletions
1
wear-shared/src/main/AndroidManifest.xml
Normal file
1
wear-shared/src/main/AndroidManifest.xml
Normal file
|
|
@ -0,0 +1 @@
|
|||
<manifest />
|
||||
|
|
@ -0,0 +1,363 @@
|
|||
package ai.openclaw.wear.shared
|
||||
|
||||
import kotlinx.serialization.SerialName
|
||||
import kotlinx.serialization.Serializable
|
||||
import kotlinx.serialization.SerializationException
|
||||
import kotlinx.serialization.json.Json
|
||||
import kotlinx.serialization.json.JsonArray
|
||||
import kotlinx.serialization.json.JsonElement
|
||||
import kotlinx.serialization.json.JsonObject
|
||||
import kotlinx.serialization.json.JsonPrimitive
|
||||
import kotlinx.serialization.json.buildJsonObject
|
||||
import kotlinx.serialization.json.intOrNull
|
||||
import kotlinx.serialization.json.jsonObject
|
||||
import java.nio.charset.CharacterCodingException
|
||||
import java.security.MessageDigest
|
||||
|
||||
object WearProtocol {
|
||||
const val VERSION = 1
|
||||
const val REQUEST_PATH = "/openclaw/wear/v1/request"
|
||||
const val RESPONSE_PATH = "/openclaw/wear/v1/response"
|
||||
const val EVENT_PATH = "/openclaw/wear/v1/event"
|
||||
const val LEGACY_REALTIME_AUDIO_CHANNEL_PATH = "/openclaw/wear/v1/realtime/audio"
|
||||
const val REALTIME_AUDIO_CHANNEL_PATH_PREFIX = "/openclaw/wear/v1/realtime/audio/"
|
||||
const val PHONE_CAPABILITY = "openclaw_phone_proxy_v1"
|
||||
const val WATCH_CAPABILITY = "openclaw_wear_companion_v1"
|
||||
|
||||
// MessageClient has a 100 KiB ceiling. Keep headroom for transport metadata and
|
||||
// force transcript pagination instead of depending on an edge-sized message.
|
||||
const val MAX_MESSAGE_BYTES = 64 * 1024
|
||||
const val MAX_REALTIME_AUDIO_FRAME_BYTES = 8 * 1024
|
||||
const val REALTIME_AUDIO_SAMPLE_RATE_HZ = 24_000
|
||||
const val REALTIME_AUDIO_FRAME_MILLIS = 20
|
||||
const val RPC_REQUEST_TIMEOUT_MILLIS = 10_000L
|
||||
|
||||
// The Watch opens the audio channel before sending talk.start. Keep the
|
||||
// pending phone-side channel through that RPC deadline plus setup margin.
|
||||
const val REALTIME_AUDIO_PENDING_CHANNEL_TIMEOUT_MILLIS = RPC_REQUEST_TIMEOUT_MILLIS + 5_000L
|
||||
|
||||
// Bound recursive JSON parsing at the untrusted Data Layer boundary.
|
||||
const val MAX_JSON_DEPTH = 32
|
||||
|
||||
fun realtimeAudioChannelPath(attemptId: String): String {
|
||||
require(attemptId.isNotBlank())
|
||||
val digest = MessageDigest.getInstance("SHA-256").digest(attemptId.encodeToByteArray())
|
||||
return buildString(REALTIME_AUDIO_CHANNEL_PATH_PREFIX.length + digest.size * 2) {
|
||||
append(REALTIME_AUDIO_CHANNEL_PATH_PREFIX)
|
||||
digest.forEach { byte ->
|
||||
val value = byte.toInt() and 0xff
|
||||
append(LOWER_HEX[value ushr 4])
|
||||
append(LOWER_HEX[value and 0x0f])
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fun isRealtimeAudioChannelPath(path: String): Boolean {
|
||||
if (path == LEGACY_REALTIME_AUDIO_CHANNEL_PATH) return true
|
||||
return isAttemptScopedRealtimeAudioChannelPath(path)
|
||||
}
|
||||
|
||||
fun isAttemptScopedRealtimeAudioChannelPath(path: String): Boolean {
|
||||
if (!path.startsWith(REALTIME_AUDIO_CHANNEL_PATH_PREFIX)) return false
|
||||
val token = path.substring(REALTIME_AUDIO_CHANNEL_PATH_PREFIX.length)
|
||||
return token.length == REALTIME_AUDIO_ATTEMPT_TOKEN_CHARS &&
|
||||
token.all { char -> char in '0'..'9' || char in 'a'..'f' }
|
||||
}
|
||||
|
||||
private const val REALTIME_AUDIO_ATTEMPT_TOKEN_CHARS = 64
|
||||
private const val LOWER_HEX = "0123456789abcdef"
|
||||
}
|
||||
|
||||
enum class WearProxyCapability(
|
||||
val wireValue: String,
|
||||
) {
|
||||
AgentControls(wireValue = "agent-controls"),
|
||||
GatewayControls(wireValue = "gateway-controls"),
|
||||
ModelControls(wireValue = "model-controls"),
|
||||
SessionSelectionLookup(wireValue = "session-selection-lookup"),
|
||||
AttemptScopedRealtimeAudio(wireValue = "attempt-scoped-realtime-audio"),
|
||||
;
|
||||
|
||||
companion object {
|
||||
fun fromWireValue(value: String): WearProxyCapability? = entries.firstOrNull { capability -> capability.wireValue == value }
|
||||
}
|
||||
}
|
||||
|
||||
enum class WearConnectionFailure(
|
||||
val wireValue: String,
|
||||
) {
|
||||
GatewayOffline(wireValue = "gateway_offline"),
|
||||
Incompatible(wireValue = "incompatible"),
|
||||
;
|
||||
|
||||
companion object {
|
||||
fun fromWireValue(value: String?): WearConnectionFailure? = entries.firstOrNull { failure -> failure.wireValue == value }
|
||||
}
|
||||
}
|
||||
|
||||
@Serializable
|
||||
enum class WearRpcMethod {
|
||||
@SerialName("proxy.status")
|
||||
ProxyStatus,
|
||||
|
||||
@SerialName("sessions.list")
|
||||
SessionsList,
|
||||
|
||||
@SerialName("agents.list")
|
||||
AgentsList,
|
||||
|
||||
@SerialName("agents.select")
|
||||
AgentsSelect,
|
||||
|
||||
@SerialName("models.list")
|
||||
ModelsList,
|
||||
|
||||
@SerialName("models.select")
|
||||
ModelsSelect,
|
||||
|
||||
@SerialName("gateway.connect")
|
||||
GatewayConnect,
|
||||
|
||||
@SerialName("gateway.disconnect")
|
||||
GatewayDisconnect,
|
||||
|
||||
@SerialName("chat.history")
|
||||
ChatHistory,
|
||||
|
||||
@SerialName("chat.send")
|
||||
ChatSend,
|
||||
|
||||
@SerialName("chat.abort")
|
||||
ChatAbort,
|
||||
|
||||
@SerialName("talk.start")
|
||||
TalkStart,
|
||||
|
||||
@SerialName("talk.stop")
|
||||
TalkStop,
|
||||
}
|
||||
|
||||
@Serializable
|
||||
enum class WearEventType {
|
||||
@SerialName("chat")
|
||||
Chat,
|
||||
|
||||
@SerialName("connection")
|
||||
Connection,
|
||||
|
||||
@SerialName("resync")
|
||||
Resync,
|
||||
|
||||
@SerialName("talk")
|
||||
Talk,
|
||||
}
|
||||
|
||||
@Serializable
|
||||
sealed interface WearMessage {
|
||||
val version: Int
|
||||
|
||||
@Serializable
|
||||
@SerialName("request")
|
||||
data class Request(
|
||||
override val version: Int = WearProtocol.VERSION,
|
||||
val requestId: String,
|
||||
val method: WearRpcMethod,
|
||||
val params: JsonObject = buildJsonObject {},
|
||||
) : WearMessage
|
||||
|
||||
@Serializable
|
||||
@SerialName("response")
|
||||
data class Response(
|
||||
override val version: Int = WearProtocol.VERSION,
|
||||
val requestId: String,
|
||||
val ok: Boolean,
|
||||
val result: JsonElement? = null,
|
||||
val error: WearRpcError? = null,
|
||||
val eventStreamId: String? = null,
|
||||
val eventSequence: Long? = null,
|
||||
) : WearMessage
|
||||
|
||||
@Serializable
|
||||
@SerialName("event")
|
||||
data class Event(
|
||||
override val version: Int = WearProtocol.VERSION,
|
||||
val streamId: String? = null,
|
||||
val sequence: Long,
|
||||
val event: WearEventType,
|
||||
val payload: JsonElement? = null,
|
||||
) : WearMessage
|
||||
}
|
||||
|
||||
@Serializable
|
||||
data class WearRpcError(
|
||||
val code: String,
|
||||
val message: String,
|
||||
)
|
||||
|
||||
enum class WearDecodeFailureReason {
|
||||
Empty,
|
||||
TooLarge,
|
||||
TooDeep,
|
||||
Malformed,
|
||||
UnsupportedVersion,
|
||||
InvalidEnvelope,
|
||||
}
|
||||
|
||||
sealed interface WearDecodeResult {
|
||||
data class Success(
|
||||
val message: WearMessage,
|
||||
) : WearDecodeResult
|
||||
|
||||
data class Failure(
|
||||
val reason: WearDecodeFailureReason,
|
||||
) : WearDecodeResult
|
||||
}
|
||||
|
||||
object WearProtocolCodec {
|
||||
private val json =
|
||||
Json {
|
||||
classDiscriminator = "type"
|
||||
encodeDefaults = true
|
||||
explicitNulls = false
|
||||
ignoreUnknownKeys = true
|
||||
}
|
||||
|
||||
fun encode(message: WearMessage): ByteArray {
|
||||
requireValid(message)
|
||||
require(hasValidPayloadDepth(message)) {
|
||||
"Wear message exceeds JSON depth ${WearProtocol.MAX_JSON_DEPTH}"
|
||||
}
|
||||
val encoded = json.encodeToString(WearMessage.serializer(), message)
|
||||
require(!exceedsJsonDepth(encoded)) {
|
||||
"Wear message exceeds JSON depth ${WearProtocol.MAX_JSON_DEPTH}"
|
||||
}
|
||||
val bytes =
|
||||
encoded.encodeToByteArray(throwOnInvalidSequence = true)
|
||||
require(bytes.size <= WearProtocol.MAX_MESSAGE_BYTES) {
|
||||
"Wear message exceeds ${WearProtocol.MAX_MESSAGE_BYTES} bytes"
|
||||
}
|
||||
return bytes
|
||||
}
|
||||
|
||||
private fun hasValidPayloadDepth(message: WearMessage): Boolean {
|
||||
val payloads =
|
||||
when (message) {
|
||||
is WearMessage.Request -> listOf(message.params)
|
||||
is WearMessage.Response -> listOfNotNull(message.result)
|
||||
is WearMessage.Event -> listOfNotNull(message.payload)
|
||||
}
|
||||
return payloads.all { element -> hasValidElementDepth(element, parentDepth = 1) }
|
||||
}
|
||||
|
||||
private fun hasValidElementDepth(
|
||||
element: JsonElement,
|
||||
parentDepth: Int,
|
||||
): Boolean {
|
||||
val pending = ArrayDeque<Pair<JsonElement, Int>>()
|
||||
pending.addLast(element to parentDepth)
|
||||
while (pending.isNotEmpty()) {
|
||||
val (current, parent) = pending.removeLast()
|
||||
val children =
|
||||
when (current) {
|
||||
is JsonArray -> current
|
||||
is JsonObject -> current.values
|
||||
else -> continue
|
||||
}
|
||||
val depth = parent + 1
|
||||
if (depth > WearProtocol.MAX_JSON_DEPTH) return false
|
||||
children.forEach { child -> pending.addLast(child to depth) }
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
fun decode(bytes: ByteArray): WearDecodeResult {
|
||||
if (bytes.isEmpty()) return WearDecodeResult.Failure(WearDecodeFailureReason.Empty)
|
||||
if (bytes.size > WearProtocol.MAX_MESSAGE_BYTES) {
|
||||
return WearDecodeResult.Failure(WearDecodeFailureReason.TooLarge)
|
||||
}
|
||||
|
||||
val text =
|
||||
try {
|
||||
bytes.decodeToString(throwOnInvalidSequence = true)
|
||||
} catch (_: CharacterCodingException) {
|
||||
return WearDecodeResult.Failure(WearDecodeFailureReason.Malformed)
|
||||
}
|
||||
if (exceedsJsonDepth(text)) {
|
||||
return WearDecodeResult.Failure(WearDecodeFailureReason.TooDeep)
|
||||
}
|
||||
val root =
|
||||
try {
|
||||
json.parseToJsonElement(text).jsonObject
|
||||
} catch (_: SerializationException) {
|
||||
return WearDecodeResult.Failure(WearDecodeFailureReason.Malformed)
|
||||
} catch (_: IllegalArgumentException) {
|
||||
return WearDecodeResult.Failure(WearDecodeFailureReason.Malformed)
|
||||
}
|
||||
val version =
|
||||
(root["version"] as? JsonPrimitive)?.intOrNull
|
||||
?: return WearDecodeResult.Failure(WearDecodeFailureReason.Malformed)
|
||||
if (version != WearProtocol.VERSION) {
|
||||
return WearDecodeResult.Failure(WearDecodeFailureReason.UnsupportedVersion)
|
||||
}
|
||||
val message =
|
||||
try {
|
||||
json.decodeFromJsonElement(WearMessage.serializer(), root)
|
||||
} catch (_: SerializationException) {
|
||||
return WearDecodeResult.Failure(WearDecodeFailureReason.Malformed)
|
||||
} catch (_: IllegalArgumentException) {
|
||||
return WearDecodeResult.Failure(WearDecodeFailureReason.Malformed)
|
||||
}
|
||||
if (!isValid(message)) {
|
||||
return WearDecodeResult.Failure(WearDecodeFailureReason.InvalidEnvelope)
|
||||
}
|
||||
return WearDecodeResult.Success(message)
|
||||
}
|
||||
|
||||
private fun exceedsJsonDepth(text: String): Boolean {
|
||||
var depth = 0
|
||||
var inString = false
|
||||
var escaped = false
|
||||
for (character in text) {
|
||||
if (inString) {
|
||||
when {
|
||||
escaped -> escaped = false
|
||||
character == '\\' -> escaped = true
|
||||
character == '"' -> inString = false
|
||||
}
|
||||
continue
|
||||
}
|
||||
|
||||
when (character) {
|
||||
'"' -> inString = true
|
||||
'{', '[' -> {
|
||||
depth += 1
|
||||
if (depth > WearProtocol.MAX_JSON_DEPTH) return true
|
||||
}
|
||||
'}', ']' -> depth -= 1
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
private fun requireValid(message: WearMessage) {
|
||||
require(message.version == WearProtocol.VERSION) { "Unsupported Wear protocol version: ${message.version}" }
|
||||
require(isValid(message)) { "Invalid Wear protocol envelope" }
|
||||
}
|
||||
|
||||
private fun isValid(message: WearMessage): Boolean =
|
||||
when (message) {
|
||||
is WearMessage.Request -> message.requestId.isNotBlank()
|
||||
is WearMessage.Response ->
|
||||
message.requestId.isNotBlank() &&
|
||||
(message.eventStreamId == null || message.eventStreamId.isNotBlank()) &&
|
||||
(message.eventSequence == null || message.eventSequence >= 0) &&
|
||||
if (message.ok) {
|
||||
message.error == null
|
||||
} else {
|
||||
message.error != null && message.result == null && message.error.code.isNotBlank()
|
||||
}
|
||||
is WearMessage.Event ->
|
||||
(message.streamId == null || message.streamId.isNotBlank()) &&
|
||||
message.sequence >= 0
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,121 @@
|
|||
package ai.openclaw.wear.shared
|
||||
|
||||
import kotlinx.serialization.Serializable
|
||||
import kotlinx.serialization.json.Json
|
||||
import kotlinx.serialization.json.JsonElement
|
||||
import java.io.DataInputStream
|
||||
import java.io.DataOutputStream
|
||||
import java.io.EOFException
|
||||
import java.io.InputStream
|
||||
import java.io.OutputStream
|
||||
|
||||
@Serializable
|
||||
data class WearRealtimeTalkSnapshot(
|
||||
val attemptId: String? = null,
|
||||
val active: Boolean = false,
|
||||
val listening: Boolean = false,
|
||||
val speaking: Boolean = false,
|
||||
val status: WearRealtimeTalkStatus = WearRealtimeTalkStatus.OFF,
|
||||
val statusText: String = "Off",
|
||||
val conversation: List<WearRealtimeTalkEntry> = emptyList(),
|
||||
)
|
||||
|
||||
@Serializable
|
||||
data class WearRealtimeTalkEntry(
|
||||
val id: String,
|
||||
val role: WearRealtimeTalkRole,
|
||||
val text: String,
|
||||
val streaming: Boolean = false,
|
||||
)
|
||||
|
||||
@Serializable
|
||||
enum class WearRealtimeTalkRole { USER, ASSISTANT }
|
||||
|
||||
@Serializable
|
||||
enum class WearRealtimeTalkStatus { OFF, CONNECTING, LISTENING, THINKING, SPEAKING, ERROR }
|
||||
|
||||
object WearRealtimeTalkCodec {
|
||||
private val json =
|
||||
Json {
|
||||
encodeDefaults = true
|
||||
explicitNulls = false
|
||||
ignoreUnknownKeys = true
|
||||
}
|
||||
|
||||
fun encode(snapshot: WearRealtimeTalkSnapshot): JsonElement = json.encodeToJsonElement(WearRealtimeTalkSnapshot.serializer(), snapshot)
|
||||
|
||||
fun decode(payload: JsonElement): WearRealtimeTalkSnapshot = json.decodeFromJsonElement(WearRealtimeTalkSnapshot.serializer(), payload)
|
||||
}
|
||||
|
||||
enum class WearRealtimeAudioFrameType(
|
||||
val wireValue: Int,
|
||||
) {
|
||||
INPUT_PCM(1),
|
||||
OUTPUT_PCM(2),
|
||||
CLEAR_OUTPUT(3),
|
||||
;
|
||||
|
||||
companion object {
|
||||
fun fromWireValue(value: Int): WearRealtimeAudioFrameType? = entries.firstOrNull { it.wireValue == value }
|
||||
}
|
||||
}
|
||||
|
||||
data class WearRealtimeAudioFrame(
|
||||
val type: WearRealtimeAudioFrameType,
|
||||
val payload: ByteArray,
|
||||
)
|
||||
|
||||
object WearRealtimeAudioFraming {
|
||||
fun write(
|
||||
output: OutputStream,
|
||||
type: WearRealtimeAudioFrameType,
|
||||
payload: ByteArray,
|
||||
) {
|
||||
requireValid(type, payload)
|
||||
DataOutputStream(output).apply {
|
||||
writeByte(type.wireValue)
|
||||
writeInt(payload.size)
|
||||
write(payload)
|
||||
flush()
|
||||
}
|
||||
}
|
||||
|
||||
fun read(input: InputStream): WearRealtimeAudioFrame? {
|
||||
val stream = DataInputStream(input)
|
||||
val typeValue = stream.read()
|
||||
if (typeValue < 0) return null
|
||||
val type =
|
||||
WearRealtimeAudioFrameType.fromWireValue(typeValue)
|
||||
?: throw IllegalArgumentException("Unknown Wear realtime audio frame type")
|
||||
val size =
|
||||
try {
|
||||
stream.readInt()
|
||||
} catch (err: EOFException) {
|
||||
throw IllegalArgumentException("Truncated Wear realtime audio frame", err)
|
||||
}
|
||||
if (size < 0 || size > WearProtocol.MAX_REALTIME_AUDIO_FRAME_BYTES) {
|
||||
throw IllegalArgumentException("Invalid Wear realtime audio frame size")
|
||||
}
|
||||
val payload = ByteArray(size)
|
||||
try {
|
||||
stream.readFully(payload)
|
||||
} catch (err: EOFException) {
|
||||
throw IllegalArgumentException("Truncated Wear realtime audio payload", err)
|
||||
}
|
||||
requireValid(type, payload)
|
||||
return WearRealtimeAudioFrame(type, payload)
|
||||
}
|
||||
|
||||
private fun requireValid(
|
||||
type: WearRealtimeAudioFrameType,
|
||||
payload: ByteArray,
|
||||
) {
|
||||
require(payload.size <= WearProtocol.MAX_REALTIME_AUDIO_FRAME_BYTES)
|
||||
when (type) {
|
||||
WearRealtimeAudioFrameType.INPUT_PCM,
|
||||
WearRealtimeAudioFrameType.OUTPUT_PCM,
|
||||
-> require(payload.isNotEmpty() && payload.size % 2 == 0)
|
||||
WearRealtimeAudioFrameType.CLEAR_OUTPUT -> require(payload.isEmpty())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,276 @@
|
|||
package ai.openclaw.wear.shared
|
||||
|
||||
import kotlinx.serialization.json.Json
|
||||
import kotlinx.serialization.json.JsonArray
|
||||
import kotlinx.serialization.json.JsonElement
|
||||
import kotlinx.serialization.json.JsonPrimitive
|
||||
import kotlinx.serialization.json.buildJsonObject
|
||||
import kotlinx.serialization.json.jsonObject
|
||||
import kotlinx.serialization.json.jsonPrimitive
|
||||
import kotlinx.serialization.json.put
|
||||
import org.junit.Assert.assertArrayEquals
|
||||
import org.junit.Assert.assertEquals
|
||||
import org.junit.Assert.assertFalse
|
||||
import org.junit.Assert.assertThrows
|
||||
import org.junit.Assert.assertTrue
|
||||
import org.junit.Test
|
||||
|
||||
class WearProtocolTest {
|
||||
@Test
|
||||
fun realtimeTalkSnapshotCarriesAttemptCorrelation() {
|
||||
val snapshot =
|
||||
WearRealtimeTalkSnapshot(
|
||||
attemptId = "attempt-7",
|
||||
active = true,
|
||||
listening = true,
|
||||
status = WearRealtimeTalkStatus.LISTENING,
|
||||
)
|
||||
|
||||
assertEquals(snapshot, WearRealtimeTalkCodec.decode(WearRealtimeTalkCodec.encode(snapshot)))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun roundTripsEveryEnvelopeKind() {
|
||||
val messages =
|
||||
listOf(
|
||||
WearMessage.Request(
|
||||
requestId = "req-1",
|
||||
method = WearRpcMethod.ChatHistory,
|
||||
params = buildJsonObject { put("sessionKey", "main") },
|
||||
),
|
||||
WearMessage.Response(
|
||||
requestId = "req-1",
|
||||
ok = true,
|
||||
result = buildJsonObject { put("count", 2) },
|
||||
eventStreamId = "phone-process-1",
|
||||
eventSequence = 7,
|
||||
),
|
||||
WearMessage.Response(
|
||||
requestId = "req-2",
|
||||
ok = false,
|
||||
error = WearRpcError(code = "unavailable", message = "Phone offline"),
|
||||
),
|
||||
WearMessage.Event(
|
||||
streamId = "phone-process-1",
|
||||
sequence = 7,
|
||||
event = WearEventType.Chat,
|
||||
payload = buildJsonObject { put("state", "delta") },
|
||||
),
|
||||
)
|
||||
|
||||
messages.forEach { message ->
|
||||
assertEquals(WearDecodeResult.Success(message), WearProtocolCodec.decode(WearProtocolCodec.encode(message)))
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun usesStableWireNamesAndPaths() {
|
||||
val methodNames =
|
||||
mapOf(
|
||||
WearRpcMethod.ProxyStatus to "proxy.status",
|
||||
WearRpcMethod.SessionsList to "sessions.list",
|
||||
WearRpcMethod.AgentsList to "agents.list",
|
||||
WearRpcMethod.AgentsSelect to "agents.select",
|
||||
WearRpcMethod.ModelsList to "models.list",
|
||||
WearRpcMethod.ModelsSelect to "models.select",
|
||||
WearRpcMethod.GatewayConnect to "gateway.connect",
|
||||
WearRpcMethod.GatewayDisconnect to "gateway.disconnect",
|
||||
WearRpcMethod.ChatHistory to "chat.history",
|
||||
WearRpcMethod.ChatSend to "chat.send",
|
||||
WearRpcMethod.ChatAbort to "chat.abort",
|
||||
WearRpcMethod.TalkStart to "talk.start",
|
||||
WearRpcMethod.TalkStop to "talk.stop",
|
||||
)
|
||||
methodNames.forEach { (method, wireName) ->
|
||||
val request = WearMessage.Request(requestId = "req-1", method = method)
|
||||
val root = Json.parseToJsonElement(WearProtocolCodec.encode(request).decodeToString()).jsonObject
|
||||
assertEquals("request", root.getValue("type").jsonPrimitive.content)
|
||||
assertEquals(wireName, root.getValue("method").jsonPrimitive.content)
|
||||
}
|
||||
|
||||
val eventNames =
|
||||
mapOf(
|
||||
WearEventType.Chat to "chat",
|
||||
WearEventType.Connection to "connection",
|
||||
WearEventType.Resync to "resync",
|
||||
WearEventType.Talk to "talk",
|
||||
)
|
||||
eventNames.forEach { (event, wireName) ->
|
||||
val message = WearMessage.Event(sequence = 1, event = event)
|
||||
val root = Json.parseToJsonElement(WearProtocolCodec.encode(message).decodeToString()).jsonObject
|
||||
assertEquals("event", root.getValue("type").jsonPrimitive.content)
|
||||
assertEquals(wireName, root.getValue("event").jsonPrimitive.content)
|
||||
}
|
||||
|
||||
assertEquals("/openclaw/wear/v1/request", WearProtocol.REQUEST_PATH)
|
||||
assertEquals("/openclaw/wear/v1/response", WearProtocol.RESPONSE_PATH)
|
||||
assertEquals("/openclaw/wear/v1/event", WearProtocol.EVENT_PATH)
|
||||
assertEquals(10_000L, WearProtocol.RPC_REQUEST_TIMEOUT_MILLIS)
|
||||
assertEquals(15_000L, WearProtocol.REALTIME_AUDIO_PENDING_CHANNEL_TIMEOUT_MILLIS)
|
||||
assertTrue(
|
||||
WearProtocol.REALTIME_AUDIO_PENDING_CHANNEL_TIMEOUT_MILLIS >
|
||||
WearProtocol.RPC_REQUEST_TIMEOUT_MILLIS,
|
||||
)
|
||||
val realtimePath = WearProtocol.realtimeAudioChannelPath("attempt-7")
|
||||
assertEquals(
|
||||
"/openclaw/wear/v1/realtime/audio",
|
||||
WearProtocol.LEGACY_REALTIME_AUDIO_CHANNEL_PATH,
|
||||
)
|
||||
assertEquals(
|
||||
"/openclaw/wear/v1/realtime/audio/9804dc90c374fd8e83c9b95a75611f9bec6e0c6ecdcbed5319d6491208417521",
|
||||
realtimePath,
|
||||
)
|
||||
assertEquals(realtimePath, WearProtocol.realtimeAudioChannelPath("attempt-7"))
|
||||
assertTrue(WearProtocol.isRealtimeAudioChannelPath(realtimePath))
|
||||
assertTrue(WearProtocol.isAttemptScopedRealtimeAudioChannelPath(realtimePath))
|
||||
assertTrue(WearProtocol.isRealtimeAudioChannelPath(WearProtocol.LEGACY_REALTIME_AUDIO_CHANNEL_PATH))
|
||||
assertFalse(
|
||||
WearProtocol.isAttemptScopedRealtimeAudioChannelPath(
|
||||
WearProtocol.LEGACY_REALTIME_AUDIO_CHANNEL_PATH,
|
||||
),
|
||||
)
|
||||
assertFalse(WearProtocol.isRealtimeAudioChannelPath("$realtimePath/extra"))
|
||||
assertEquals("openclaw_phone_proxy_v1", WearProtocol.PHONE_CAPABILITY)
|
||||
assertEquals("openclaw_wear_companion_v1", WearProtocol.WATCH_CAPABILITY)
|
||||
assertEquals("gateway_offline", WearConnectionFailure.GatewayOffline.wireValue)
|
||||
assertEquals("incompatible", WearConnectionFailure.Incompatible.wireValue)
|
||||
assertEquals(
|
||||
WearConnectionFailure.Incompatible,
|
||||
WearConnectionFailure.fromWireValue("incompatible"),
|
||||
)
|
||||
assertEquals(null, WearConnectionFailure.fromWireValue("future-failure"))
|
||||
assertEquals("agent-controls", WearProxyCapability.AgentControls.wireValue)
|
||||
assertEquals("gateway-controls", WearProxyCapability.GatewayControls.wireValue)
|
||||
assertEquals("model-controls", WearProxyCapability.ModelControls.wireValue)
|
||||
assertEquals("session-selection-lookup", WearProxyCapability.SessionSelectionLookup.wireValue)
|
||||
assertEquals(
|
||||
"attempt-scoped-realtime-audio",
|
||||
WearProxyCapability.AttemptScopedRealtimeAudio.wireValue,
|
||||
)
|
||||
assertEquals(WearProxyCapability.AgentControls, WearProxyCapability.fromWireValue("agent-controls"))
|
||||
assertEquals(null, WearProxyCapability.fromWireValue("future-capability"))
|
||||
}
|
||||
|
||||
@Test
|
||||
fun ignoresUnknownFieldsWithinCurrentVersion() {
|
||||
val bytes =
|
||||
"""{"type":"request","version":1,"requestId":"req-1","method":"proxy.status","params":{},"future":true}"""
|
||||
.encodeToByteArray()
|
||||
|
||||
assertEquals(
|
||||
WearDecodeResult.Success(
|
||||
WearMessage.Request(requestId = "req-1", method = WearRpcMethod.ProxyStatus),
|
||||
),
|
||||
WearProtocolCodec.decode(bytes),
|
||||
)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun rejectsMalformedUnsupportedAndInvalidMessages() {
|
||||
assertEquals(
|
||||
WearDecodeResult.Failure(WearDecodeFailureReason.Empty),
|
||||
WearProtocolCodec.decode(byteArrayOf()),
|
||||
)
|
||||
assertEquals(
|
||||
WearDecodeResult.Failure(WearDecodeFailureReason.Malformed),
|
||||
WearProtocolCodec.decode("not-json".encodeToByteArray()),
|
||||
)
|
||||
val invalidUtf8 =
|
||||
"""{"type":"request","version":1,"requestId":"""".encodeToByteArray() +
|
||||
byteArrayOf(0xc3.toByte(), 0x28) +
|
||||
"""","method":"proxy.status","params":{}}""".encodeToByteArray()
|
||||
assertEquals(
|
||||
WearDecodeResult.Failure(WearDecodeFailureReason.Malformed),
|
||||
WearProtocolCodec.decode(invalidUtf8),
|
||||
)
|
||||
assertEquals(
|
||||
WearDecodeResult.Failure(WearDecodeFailureReason.UnsupportedVersion),
|
||||
WearProtocolCodec.decode(
|
||||
"""{"type":"future-message","version":2,"futureRequiredField":true}"""
|
||||
.encodeToByteArray(),
|
||||
),
|
||||
)
|
||||
assertEquals(
|
||||
WearDecodeResult.Failure(WearDecodeFailureReason.InvalidEnvelope),
|
||||
WearProtocolCodec.decode(
|
||||
"""{"type":"response","version":1,"requestId":"req-1","ok":false}""".encodeToByteArray(),
|
||||
),
|
||||
)
|
||||
assertEquals(
|
||||
WearDecodeResult.Failure(WearDecodeFailureReason.InvalidEnvelope),
|
||||
WearProtocolCodec.decode(
|
||||
"""{"type":"event","version":1,"streamId":"","sequence":1,"event":"connection"}"""
|
||||
.encodeToByteArray(),
|
||||
),
|
||||
)
|
||||
assertEquals(
|
||||
WearDecodeResult.Failure(WearDecodeFailureReason.InvalidEnvelope),
|
||||
WearProtocolCodec.decode(
|
||||
"""{"type":"response","version":1,"requestId":"req-1","ok":true,"eventSequence":-1}"""
|
||||
.encodeToByteArray(),
|
||||
),
|
||||
)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun rejectsOversizedMessagesOnEncodeAndDecode() {
|
||||
val oversizedBytes = ByteArray(WearProtocol.MAX_MESSAGE_BYTES + 1)
|
||||
assertEquals(
|
||||
WearDecodeResult.Failure(WearDecodeFailureReason.TooLarge),
|
||||
WearProtocolCodec.decode(oversizedBytes),
|
||||
)
|
||||
|
||||
val oversizedMessage =
|
||||
WearMessage.Request(
|
||||
requestId = "req-1",
|
||||
method = WearRpcMethod.ChatSend,
|
||||
params = buildJsonObject { put("message", "x".repeat(WearProtocol.MAX_MESSAGE_BYTES)) },
|
||||
)
|
||||
assertThrows(IllegalArgumentException::class.java) {
|
||||
WearProtocolCodec.encode(oversizedMessage)
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun rejectsExcessiveJsonDepthBeforeParsing() {
|
||||
val nesting = WearProtocol.MAX_JSON_DEPTH + 1
|
||||
val deeplyNested =
|
||||
"""{"type":"request","version":1,"requestId":"req-1","method":"chat.send","params":{"payload":${"[".repeat(nesting)}0${"]".repeat(nesting)}}}"""
|
||||
|
||||
assertEquals(
|
||||
WearDecodeResult.Failure(WearDecodeFailureReason.TooDeep),
|
||||
WearProtocolCodec.decode(deeplyNested.encodeToByteArray()),
|
||||
)
|
||||
|
||||
val bracketsInString =
|
||||
WearMessage.Request(
|
||||
requestId = "req-2",
|
||||
method = WearRpcMethod.ChatSend,
|
||||
params = buildJsonObject { put("message", "[".repeat(WearProtocol.MAX_JSON_DEPTH + 1)) },
|
||||
)
|
||||
assertEquals(
|
||||
WearDecodeResult.Success(bracketsInString),
|
||||
WearProtocolCodec.decode(WearProtocolCodec.encode(bracketsInString)),
|
||||
)
|
||||
|
||||
var nestedPayload: JsonElement = JsonPrimitive(0)
|
||||
repeat(4_096) {
|
||||
nestedPayload = JsonArray(listOf(nestedPayload))
|
||||
}
|
||||
val deeplyNestedMessage =
|
||||
WearMessage.Request(
|
||||
requestId = "req-3",
|
||||
method = WearRpcMethod.ChatSend,
|
||||
params = buildJsonObject { put("payload", nestedPayload) },
|
||||
)
|
||||
assertThrows(IllegalArgumentException::class.java) {
|
||||
WearProtocolCodec.encode(deeplyNestedMessage)
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun encodingIsDeterministic() {
|
||||
val message = WearMessage.Request(requestId = "req-1", method = WearRpcMethod.SessionsList)
|
||||
assertArrayEquals(WearProtocolCodec.encode(message), WearProtocolCodec.encode(message))
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,35 @@
|
|||
package ai.openclaw.wear.shared
|
||||
|
||||
import org.junit.Assert.assertArrayEquals
|
||||
import org.junit.Assert.assertEquals
|
||||
import org.junit.Assert.assertNull
|
||||
import org.junit.Test
|
||||
import java.io.ByteArrayInputStream
|
||||
import java.io.ByteArrayOutputStream
|
||||
import java.io.DataOutputStream
|
||||
|
||||
class WearRealtimeTalkTest {
|
||||
@Test
|
||||
fun audioFramesRoundTripAndPreserveBoundaries() {
|
||||
val output = ByteArrayOutputStream()
|
||||
WearRealtimeAudioFraming.write(output, WearRealtimeAudioFrameType.INPUT_PCM, byteArrayOf(1, 2, 3, 4))
|
||||
WearRealtimeAudioFraming.write(output, WearRealtimeAudioFrameType.CLEAR_OUTPUT, byteArrayOf())
|
||||
val input = ByteArrayInputStream(output.toByteArray())
|
||||
|
||||
val audio = WearRealtimeAudioFraming.read(input)!!
|
||||
assertEquals(WearRealtimeAudioFrameType.INPUT_PCM, audio.type)
|
||||
assertArrayEquals(byteArrayOf(1, 2, 3, 4), audio.payload)
|
||||
assertEquals(WearRealtimeAudioFrameType.CLEAR_OUTPUT, WearRealtimeAudioFraming.read(input)!!.type)
|
||||
assertNull(WearRealtimeAudioFraming.read(input))
|
||||
}
|
||||
|
||||
@Test(expected = IllegalArgumentException::class)
|
||||
fun rejectsOversizedFrameBeforeAllocatingPayload() {
|
||||
val output = ByteArrayOutputStream()
|
||||
DataOutputStream(output).apply {
|
||||
writeByte(WearRealtimeAudioFrameType.INPUT_PCM.wireValue)
|
||||
writeInt(WearProtocol.MAX_REALTIME_AUDIO_FRAME_BYTES + 1)
|
||||
}
|
||||
WearRealtimeAudioFraming.read(ByteArrayInputStream(output.toByteArray()))
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue