AeroToss v0.4.0: полная реализация M1-M4
- KMP проект: Kotlin 2.3.21, Compose 1.11.1, Gradle 8.13, AGP 8.13.2 - Обнаружение устройств: mDNS (JMDNS/Android NSD) + Wi-Fi Direct P2P - Передача файлов: Java Socket, стриминг чанками, SHA-256 checksum - UI: HomeScreen, SendScreen, ReceiveScreen, навигация - Защита: санитизация имён файлов (path traversal) - 65+ тестов (WireProtocol, Transfer, FileUtils, Wi-Fi Direct)
This commit is contained in:
@@ -0,0 +1,30 @@
|
||||
package com.aerotoss
|
||||
|
||||
import androidx.compose.runtime.*
|
||||
import androidx.compose.ui.window.Window
|
||||
import androidx.compose.ui.window.application
|
||||
import com.aerotoss.core.AeroTossManager
|
||||
import com.aerotoss.discovery.JmdnsDiscovery
|
||||
import com.aerotoss.transfer.DesktopFileTransferManager
|
||||
import com.aerotoss.ui.App
|
||||
|
||||
fun main() = application {
|
||||
val transferManager = remember { DesktopFileTransferManager() }
|
||||
val discoveryManager = remember { JmdnsDiscovery() }
|
||||
val manager = remember { AeroTossManager(discoveryManager, transferManager) }
|
||||
|
||||
LaunchedEffect(Unit) {
|
||||
val port = transferManager.startServer(0)
|
||||
manager.start(port)
|
||||
}
|
||||
|
||||
Window(
|
||||
onCloseRequest = {
|
||||
manager.dispose()
|
||||
exitApplication()
|
||||
},
|
||||
title = "AeroToss"
|
||||
) {
|
||||
App(manager)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,148 @@
|
||||
package com.aerotoss.discovery
|
||||
|
||||
import com.aerotoss.model.Device
|
||||
import com.aerotoss.model.DeviceType
|
||||
import com.aerotoss.util.generateDeviceId
|
||||
import com.aerotoss.util.getDeviceName
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.MutableStateFlow
|
||||
import kotlinx.coroutines.flow.asStateFlow
|
||||
import java.net.InetAddress
|
||||
import java.net.NetworkInterface
|
||||
import java.util.concurrent.ConcurrentHashMap
|
||||
import javax.jmdns.JmDNS
|
||||
import javax.jmdns.ServiceEvent
|
||||
import javax.jmdns.ServiceListener
|
||||
import javax.jmdns.ServiceInfo
|
||||
|
||||
class JmdnsDiscovery : DiscoveryManager {
|
||||
private val _devices = MutableStateFlow<List<Device>>(emptyList())
|
||||
override val devices: Flow<List<Device>> = _devices.asStateFlow()
|
||||
|
||||
private var jmdns: JmDNS? = null
|
||||
private var serviceListener: ServiceListener? = null
|
||||
private var discoveryThread: Thread? = null
|
||||
private val deviceId = generateDeviceId()
|
||||
private val deviceName = getDeviceName()
|
||||
private val discoveredDevices = ConcurrentHashMap<String, Device>()
|
||||
|
||||
override fun startDiscovery(servicePort: Int) {
|
||||
if (discoveryThread != null) return
|
||||
|
||||
discoveryThread = Thread {
|
||||
try {
|
||||
val addr = findLocalAddress() ?: return@Thread
|
||||
jmdns = JmDNS.create(addr, "aerotoss-$deviceId")
|
||||
|
||||
val serviceInfo = ServiceInfo.create(
|
||||
SERVICE_TYPE,
|
||||
SERVICE_NAME,
|
||||
servicePort,
|
||||
"path=/ aerotoss=1 id=$deviceId name=$deviceName"
|
||||
)
|
||||
jmdns?.registerService(serviceInfo)
|
||||
|
||||
serviceListener = object : ServiceListener {
|
||||
override fun serviceAdded(event: ServiceEvent) {
|
||||
jmdns?.requestServiceInfo(event.type, event.name, true)
|
||||
}
|
||||
|
||||
override fun serviceRemoved(event: ServiceEvent) {
|
||||
discoveredDevices.remove(event.name)
|
||||
_devices.value = discoveredDevices.values.toList()
|
||||
}
|
||||
|
||||
override fun serviceResolved(event: ServiceEvent) {
|
||||
val info = event.info
|
||||
val hostAddresses = info.hostAddresses
|
||||
if (hostAddresses.isNotEmpty()) {
|
||||
val attributes = info.textBytes?.let { parseAttributes(it) } ?: emptyMap()
|
||||
val id = attributes["id"] ?: event.name
|
||||
val name = attributes["name"] ?: event.name
|
||||
val device = Device(
|
||||
id = id,
|
||||
name = name,
|
||||
type = DeviceType.DESKTOP,
|
||||
hostAddress = hostAddresses.first(),
|
||||
port = info.port
|
||||
)
|
||||
if (id != deviceId) {
|
||||
discoveredDevices[event.name] = device
|
||||
_devices.value = discoveredDevices.values.toList()
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
jmdns?.addServiceListener(SERVICE_TYPE, serviceListener)
|
||||
Thread.currentThread().join()
|
||||
} catch (e: Exception) {
|
||||
e.printStackTrace()
|
||||
}
|
||||
}.apply {
|
||||
isDaemon = true
|
||||
name = "aerotoss-mdns"
|
||||
start()
|
||||
}
|
||||
}
|
||||
|
||||
override fun stopDiscovery() {
|
||||
discoveryThread?.interrupt()
|
||||
discoveryThread = null
|
||||
try {
|
||||
serviceListener?.let { jmdns?.removeServiceListener(SERVICE_TYPE, it) }
|
||||
jmdns?.unregisterAllServices()
|
||||
jmdns?.close()
|
||||
} catch (_: Exception) {
|
||||
}
|
||||
jmdns = null
|
||||
serviceListener = null
|
||||
discoveredDevices.clear()
|
||||
_devices.value = emptyList()
|
||||
}
|
||||
|
||||
override fun dispose() {
|
||||
stopDiscovery()
|
||||
}
|
||||
|
||||
private fun findLocalAddress(): InetAddress? {
|
||||
return try {
|
||||
NetworkInterface.getNetworkInterfaces()?.toList()
|
||||
?.filter { it.isUp && !it.isLoopback && !isVirtual(it) }
|
||||
?.flatMap { it.inetAddresses.toList() }
|
||||
?.firstOrNull { it is java.net.Inet4Address }
|
||||
?: InetAddress.getLocalHost()
|
||||
} catch (_: Exception) {
|
||||
try { InetAddress.getLocalHost() } catch (_: Exception) { null }
|
||||
}
|
||||
}
|
||||
|
||||
private fun isVirtual(iface: NetworkInterface): Boolean {
|
||||
return iface.isVirtual || iface.name.startsWith("vmnet") || iface.name.startsWith("veth")
|
||||
}
|
||||
|
||||
private fun parseAttributes(raw: ByteArray): Map<String, String> {
|
||||
val result = mutableMapOf<String, String>()
|
||||
var i = 0
|
||||
while (i < raw.size) {
|
||||
val keyLen = raw[i].toInt() and 0xFF
|
||||
i++
|
||||
if (i + keyLen > raw.size) break
|
||||
val key = String(raw, i, keyLen, Charsets.UTF_8)
|
||||
i += keyLen
|
||||
if (i >= raw.size) break
|
||||
val valueLen = raw[i].toInt() and 0xFF
|
||||
i++
|
||||
if (i + valueLen > raw.size) break
|
||||
val value = String(raw, i, valueLen, Charsets.UTF_8)
|
||||
i += valueLen
|
||||
result[key] = value
|
||||
}
|
||||
return result
|
||||
}
|
||||
|
||||
companion object {
|
||||
const val SERVICE_TYPE = "_aerotoss._tcp.local."
|
||||
const val SERVICE_NAME = "AeroToss"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,302 @@
|
||||
package com.aerotoss.transfer
|
||||
|
||||
import com.aerotoss.model.TransferProgress
|
||||
import com.aerotoss.model.TransferRequest
|
||||
import com.aerotoss.model.TransferState
|
||||
import com.aerotoss.util.FileUtils
|
||||
import com.aerotoss.util.getDeviceName
|
||||
import kotlinx.coroutines.*
|
||||
import kotlinx.coroutines.flow.*
|
||||
import kotlinx.serialization.json.Json
|
||||
import java.io.*
|
||||
import java.net.ServerSocket
|
||||
import java.net.Socket
|
||||
import java.security.MessageDigest
|
||||
import java.util.UUID
|
||||
import java.util.concurrent.ConcurrentHashMap
|
||||
import java.util.concurrent.atomic.AtomicBoolean
|
||||
|
||||
@OptIn(ExperimentalCoroutinesApi::class)
|
||||
class DesktopFileTransferManager : TransferManager {
|
||||
private val _incomingTransfers = MutableStateFlow<List<TransferProgress>>(emptyList())
|
||||
override val incomingTransfers: Flow<TransferProgress> = _incomingTransfers.flatMapLatest { list ->
|
||||
flow { list.forEach { emit(it) } }
|
||||
}
|
||||
|
||||
private val _outgoingTransfers = MutableStateFlow<List<TransferProgress>>(emptyList())
|
||||
override val outgoingTransfers: Flow<TransferProgress> = _outgoingTransfers.flatMapLatest { list ->
|
||||
flow { list.forEach { emit(it) } }
|
||||
}
|
||||
|
||||
private val transferStates = ConcurrentHashMap<String, MutableStateFlow<TransferProgress>>()
|
||||
|
||||
private var serverSocket: ServerSocket? = null
|
||||
private var serverThread: Thread? = null
|
||||
private val scope = CoroutineScope(Dispatchers.IO + SupervisorJob())
|
||||
private val activeJobs = ConcurrentHashMap<String, Job>()
|
||||
private val running = AtomicBoolean(false)
|
||||
|
||||
private val downloadsDir: File = File(System.getProperty("user.home"), "AeroTossDownloads").apply {
|
||||
mkdirs()
|
||||
}
|
||||
|
||||
override fun getServerPort(): Int = serverSocket?.localPort ?: 0
|
||||
|
||||
fun startServer(port: Int = 0): Int {
|
||||
if (running.get()) return serverSocket?.localPort ?: 0
|
||||
|
||||
try {
|
||||
val socket = ServerSocket(port)
|
||||
serverSocket = socket
|
||||
running.set(true)
|
||||
|
||||
serverThread = Thread {
|
||||
while (running.get() && !socket.isClosed) {
|
||||
try {
|
||||
val clientSocket = socket.accept()
|
||||
scope.launch { handleIncomingConnection(clientSocket) }
|
||||
} catch (_: Exception) {
|
||||
if (running.get()) break
|
||||
}
|
||||
}
|
||||
}.apply {
|
||||
isDaemon = true
|
||||
name = "aerotoss-server"
|
||||
start()
|
||||
}
|
||||
|
||||
return socket.localPort
|
||||
} catch (e: Exception) {
|
||||
e.printStackTrace()
|
||||
return 0
|
||||
}
|
||||
}
|
||||
|
||||
private suspend fun handleIncomingConnection(socket: Socket) {
|
||||
withContext(Dispatchers.IO) {
|
||||
val progressId = UUID.randomUUID().toString()
|
||||
try {
|
||||
socket.use { sock ->
|
||||
sock.soTimeout = 30_000
|
||||
val input = DataInputStream(sock.getInputStream())
|
||||
val output = DataOutputStream(sock.getOutputStream())
|
||||
|
||||
val requestJson = input.readUTF()
|
||||
val request = Json.decodeFromString<TransferRequest>(requestJson)
|
||||
|
||||
val progress = TransferProgress(
|
||||
id = progressId,
|
||||
request = request,
|
||||
bytesTransferred = 0,
|
||||
totalBytes = request.fileSize,
|
||||
state = TransferState.PENDING
|
||||
)
|
||||
updateIncoming(progress)
|
||||
|
||||
output.writeUTF(Json.encodeToString(TransferRequest.serializer(), request))
|
||||
|
||||
val accepted = input.readBoolean()
|
||||
if (!accepted) {
|
||||
updateIncoming(progress.copy(state = TransferState.CANCELLED))
|
||||
return@withContext
|
||||
}
|
||||
|
||||
val file = FileUtils.resolveUniqueFile(downloadsDir, request.fileName)
|
||||
val sha256 = MessageDigest.getInstance("SHA-256")
|
||||
var bytesWritten = 0L
|
||||
|
||||
updateIncoming(progress.copy(state = TransferState.TRANSFERRING))
|
||||
|
||||
file.outputStream().use { fos ->
|
||||
val buffer = ByteArray(65536)
|
||||
var remaining = request.fileSize
|
||||
while (remaining > 0) {
|
||||
val toRead = minOf(buffer.size.toLong(), remaining).toInt()
|
||||
val read = input.read(buffer, 0, toRead)
|
||||
if (read == -1) break
|
||||
fos.write(buffer, 0, read)
|
||||
sha256.update(buffer, 0, read)
|
||||
bytesWritten += read
|
||||
remaining -= read
|
||||
updateIncomingById(progressId, TransferProgress(
|
||||
id = progressId,
|
||||
request = request,
|
||||
bytesTransferred = bytesWritten,
|
||||
totalBytes = request.fileSize,
|
||||
state = TransferState.TRANSFERRING
|
||||
))
|
||||
}
|
||||
}
|
||||
|
||||
val checksum = "sha256:${sha256.digest().joinToString("") { "%02x".format(it) }}"
|
||||
output.writeUTF(checksum)
|
||||
output.writeLong(bytesWritten)
|
||||
|
||||
if (bytesWritten == request.fileSize) {
|
||||
updateIncomingById(progressId, TransferProgress(
|
||||
id = progressId,
|
||||
request = request,
|
||||
bytesTransferred = bytesWritten,
|
||||
totalBytes = request.fileSize,
|
||||
state = TransferState.COMPLETED
|
||||
))
|
||||
} else {
|
||||
FileUtils.deleteIfExists(file)
|
||||
updateIncomingById(progressId, TransferProgress(
|
||||
id = progressId,
|
||||
request = request,
|
||||
bytesTransferred = bytesWritten,
|
||||
totalBytes = request.fileSize,
|
||||
state = TransferState.FAILED,
|
||||
error = "Incomplete transfer: expected ${request.fileSize}, got $bytesWritten"
|
||||
))
|
||||
}
|
||||
}
|
||||
} catch (e: Exception) {
|
||||
e.printStackTrace()
|
||||
val current = _incomingTransfers.value.find { it.id == progressId }
|
||||
if (current != null && current.state != TransferState.COMPLETED &&
|
||||
current.state != TransferState.FAILED && current.state != TransferState.CANCELLED
|
||||
) {
|
||||
FileUtils.deleteIfExists(File(downloadsDir, FileUtils.sanitizeFileName(current.request.fileName)))
|
||||
updateIncomingById(progressId, current.copy(
|
||||
state = TransferState.FAILED,
|
||||
error = e.message ?: "Unknown error"
|
||||
))
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
override suspend fun sendFile(file: File, targetHost: String, targetPort: Int): Flow<TransferProgress> {
|
||||
val requestId = UUID.randomUUID().toString()
|
||||
val request = TransferRequest(
|
||||
fileName = file.name,
|
||||
fileSize = file.length(),
|
||||
mimeType = "application/octet-stream",
|
||||
senderId = "",
|
||||
senderName = getDeviceName()
|
||||
)
|
||||
|
||||
val stateFlow = MutableStateFlow(TransferProgress(
|
||||
id = requestId,
|
||||
request = request,
|
||||
bytesTransferred = 0,
|
||||
totalBytes = file.length(),
|
||||
state = TransferState.PENDING
|
||||
))
|
||||
transferStates[requestId] = stateFlow
|
||||
updateOutgoing(stateFlow.value)
|
||||
|
||||
val job = scope.launch {
|
||||
try {
|
||||
val socket = Socket(targetHost, targetPort)
|
||||
socket.use { sock ->
|
||||
sock.soTimeout = 30_000
|
||||
val input = DataInputStream(sock.getInputStream())
|
||||
val output = DataOutputStream(sock.getOutputStream())
|
||||
|
||||
output.writeUTF(Json.encodeToString(TransferRequest.serializer(), request))
|
||||
|
||||
val serverResponseJson = input.readUTF()
|
||||
try {
|
||||
Json.decodeFromString<TransferRequest>(serverResponseJson)
|
||||
} catch (_: Exception) {
|
||||
}
|
||||
|
||||
output.writeBoolean(true)
|
||||
|
||||
stateFlow.value = stateFlow.value.copy(state = TransferState.TRANSFERRING)
|
||||
updateOutgoing(stateFlow.value)
|
||||
|
||||
val sha256 = MessageDigest.getInstance("SHA-256")
|
||||
var bytesSent = 0L
|
||||
val buffer = ByteArray(65536)
|
||||
file.inputStream().use { fis ->
|
||||
while (true) {
|
||||
val read = fis.read(buffer)
|
||||
if (read == -1) break
|
||||
output.write(buffer, 0, read)
|
||||
sha256.update(buffer, 0, read)
|
||||
bytesSent += read
|
||||
stateFlow.value = stateFlow.value.copy(
|
||||
bytesTransferred = bytesSent,
|
||||
state = TransferState.TRANSFERRING
|
||||
)
|
||||
updateOutgoing(stateFlow.value)
|
||||
}
|
||||
}
|
||||
output.flush()
|
||||
|
||||
val serverChecksum = input.readUTF()
|
||||
val bytesReceived = input.readLong()
|
||||
|
||||
val localChecksum = "sha256:${sha256.digest().joinToString("") { "%02x".format(it) }}"
|
||||
|
||||
if (bytesReceived == file.length() && serverChecksum == localChecksum) {
|
||||
stateFlow.value = stateFlow.value.copy(
|
||||
bytesTransferred = file.length(),
|
||||
state = TransferState.COMPLETED
|
||||
)
|
||||
} else {
|
||||
stateFlow.value = stateFlow.value.copy(
|
||||
bytesTransferred = bytesSent,
|
||||
state = TransferState.FAILED,
|
||||
error = "Checksum mismatch or incomplete transfer"
|
||||
)
|
||||
}
|
||||
updateOutgoing(stateFlow.value)
|
||||
}
|
||||
} catch (e: Exception) {
|
||||
e.printStackTrace()
|
||||
stateFlow.value = stateFlow.value.copy(
|
||||
state = TransferState.FAILED,
|
||||
error = e.message ?: "Unknown error"
|
||||
)
|
||||
updateOutgoing(stateFlow.value)
|
||||
} finally {
|
||||
activeJobs.remove(requestId)
|
||||
}
|
||||
}
|
||||
|
||||
activeJobs[requestId] = job
|
||||
return stateFlow
|
||||
}
|
||||
|
||||
override fun cancelTransfer(requestId: String) {
|
||||
activeJobs[requestId]?.cancel()
|
||||
activeJobs.remove(requestId)
|
||||
transferStates[requestId]?.let { flow ->
|
||||
flow.value = flow.value.copy(state = TransferState.CANCELLED)
|
||||
updateOutgoing(flow.value)
|
||||
}
|
||||
}
|
||||
|
||||
override fun dispose() {
|
||||
running.set(false)
|
||||
activeJobs.values.forEach { it.cancel() }
|
||||
activeJobs.clear()
|
||||
transferStates.clear()
|
||||
scope.cancel()
|
||||
try { serverSocket?.close() } catch (_: Exception) {}
|
||||
serverThread?.interrupt()
|
||||
}
|
||||
|
||||
private fun updateIncoming(progress: TransferProgress) {
|
||||
_incomingTransfers.update { list ->
|
||||
list.filter { it.id != progress.id }.plus(progress)
|
||||
}
|
||||
}
|
||||
|
||||
private fun updateIncomingById(id: String, progress: TransferProgress) {
|
||||
_incomingTransfers.update { list ->
|
||||
list.filter { it.id != id }.plus(progress)
|
||||
}
|
||||
}
|
||||
|
||||
private fun updateOutgoing(progress: TransferProgress) {
|
||||
_outgoingTransfers.update { list ->
|
||||
list.filter { it.id != progress.id }.plus(progress)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,16 @@
|
||||
package com.aerotoss.util
|
||||
|
||||
import java.util.UUID
|
||||
|
||||
actual fun generateDeviceId(): String = UUID.randomUUID().toString()
|
||||
|
||||
actual fun getDeviceName(): String {
|
||||
val user = System.getProperty("user.name") ?: "Desktop"
|
||||
val os = System.getProperty("os.name") ?: ""
|
||||
val host = try {
|
||||
java.net.InetAddress.getLocalHost().hostName
|
||||
} catch (_: Exception) {
|
||||
"unknown"
|
||||
}
|
||||
return "$user@$host"
|
||||
}
|
||||
Reference in New Issue
Block a user