diff --git a/MeshService/mesh/transport.py b/MeshService/mesh/transport.py index a7ebdc5..4abab7a 100644 --- a/MeshService/mesh/transport.py +++ b/MeshService/mesh/transport.py @@ -20,6 +20,20 @@ * ``seen_msg_ids`` is paired with a FIFO ``_seen_order`` deque and bounded at 4096 entries so a long-running peer doesn't leak as chatter accumulates. All de-dup writes go through ``_remember_msg_id``. + +ACK protocol: + When a node receives a message, it sends back an ACK packet to the + original sender. The sender tracks pending ACKs and retries up to + MAX_RETRIES times with ACK_TIMEOUT_S between attempts. + + ACK packet format: + { + "type": "ack", + "ack_msg_id": "", + "from_id": "", + "from_name": "", + "mesh_id": "" + } """ import hashlib @@ -37,6 +51,9 @@ SEEN_MSG_ID_CAP = 4096 DEFAULT_MESH_ID = "beacon-default" +ACK_TIMEOUT_S = 2.0 # seconds to wait before retry +MAX_RETRIES = 3 # total send attempts (1 original + 2 retries) + def _autogen_port(node_id: str) -> int: """Deterministic per-node-id port in [MSG_PORT_BASE, MSG_PORT_BASE+1000). @@ -57,17 +74,19 @@ def __init__(self, node_id, name, port=None, *, mesh_id=DEFAULT_MESH_ID): self.mesh_id = mesh_id self.peers = {} # node_id -> {"name": str, "addr": (ip, port), "last_seen": float} self.on_message = None # callback(src_id, src_name, dst_id, text, hop_count, msg_id) + self.on_ack = None # callback(msg_id, from_id, from_name) self.running = False - # Bounded de-dup ring. The set is the membership index; the deque - # is the FIFO eviction order. Both are maintained in lockstep via - # ``_remember_msg_id`` under ``_seen_lock`` so the listener thread - # and the sender thread can't race them out of sync. + # Bounded de-dup ring. self.seen_msg_ids: set[str] = set() self._seen_order: deque[str] = deque() self._seen_lock = threading.Lock() - # message socket — receives direct messages + # ACK tracking: msg_id -> {"packet": bytes, "attempts": int, "time": float, "acked_by": set} + self._pending_acks: dict[str, dict] = {} + self._ack_lock = threading.Lock() + + # message socket — receives direct messages + ACKs self.msg_sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) self.msg_sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) @@ -78,13 +97,7 @@ def __init__(self, node_id, name, port=None, *, mesh_id=DEFAULT_MESH_ID): self.bc_sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEPORT, 1) def _remember_msg_id(self, mid: str) -> None: - """Record a msg_id in the bounded de-dup ring. - - Idempotent — re-adding an already-known id is a no-op (the entry - is not re-ordered to the back of the FIFO; we deliberately do - *not* implement LRU semantics, the goal here is leak prevention, - not optimal cache hit rate). - """ + """Record a msg_id in the bounded de-dup ring.""" with self._seen_lock: if mid in self.seen_msg_ids: return @@ -102,6 +115,7 @@ def start(self): threading.Thread(target=self._listen_messages, daemon=True).start() threading.Thread(target=self._listen_discovery, daemon=True).start() threading.Thread(target=self._send_heartbeats, daemon=True).start() + threading.Thread(target=self._retry_unacked, daemon=True).start() # announce ourselves immediately self._broadcast_discovery() @@ -128,10 +142,26 @@ def send(self, dst_id, text, msg_id=None, hop_count=0, src_id=None, max_hops=3): "mesh_id": self.mesh_id, }).encode() - # send to all known peers (flood-based mesh relay) + # Track for ACK (only for messages we originate, not relays) + if src_id == self.node_id and hop_count == 0: + with self._ack_lock: + self._pending_acks[msg_id] = { + "packet": packet, + "attempts": 1, + "time": time.time(), + "dst_id": dst_id, + "acked_by": set(), + } + + self._flood_packet(packet) + + def broadcast(self, text): + self.send("^all", text) + + def _flood_packet(self, packet: bytes): + """Send packet to all known peers.""" for peer_id, peer in self.peers.items(): addr = peer["addr"] - # try both the discovered IP and localhost (for same-machine nodes) targets = [addr] if addr[0] != "127.0.0.1": targets.append(("127.0.0.1", addr[1])) @@ -141,8 +171,87 @@ def send(self, dst_id, text, msg_id=None, hop_count=0, src_id=None, max_hops=3): except OSError: pass - def broadcast(self, text): - self.send("^all", text) + def _send_ack(self, msg_id: str, src_id: str): + """Send an ACK back to the message originator.""" + ack_packet = json.dumps({ + "type": "ack", + "ack_msg_id": msg_id, + "from_id": self.node_id, + "from_name": self.name, + "mesh_id": self.mesh_id, + }).encode() + + # Send to the source peer if we know them + if src_id in self.peers: + addr = self.peers[src_id]["addr"] + targets = [addr] + if addr[0] != "127.0.0.1": + targets.append(("127.0.0.1", addr[1])) + for target in targets: + try: + self.msg_sock.sendto(ack_packet, target) + except OSError: + pass + else: + # Don't know the peer directly — flood the ACK + self._flood_packet(ack_packet) + + def _handle_ack(self, packet: dict): + """Process an incoming ACK.""" + ack_msg_id = packet.get("ack_msg_id") + from_id = packet.get("from_id") + from_name = packet.get("from_name", from_id) + + if not ack_msg_id: + return + + with self._ack_lock: + if ack_msg_id in self._pending_acks: + entry = self._pending_acks[ack_msg_id] + entry["acked_by"].add(from_id) + dst_id = entry["dst_id"] + + # For direct messages, one ACK is enough + # For broadcasts, we keep collecting but don't retry once we have any + if dst_id != "^all" or len(entry["acked_by"]) > 0: + del self._pending_acks[ack_msg_id] + + print(f"\n [ACK] {from_name} ({from_id}) confirmed msg {ack_msg_id[:12]}...") + print("mesh> ", end="", flush=True) + + if self.on_ack: + self.on_ack(ack_msg_id, from_id, from_name) + + def _retry_unacked(self): + """Background thread: retry messages that haven't been ACKed.""" + while self.running: + time.sleep(1.0) + now = time.time() + to_retry = [] + to_drop = [] + + with self._ack_lock: + for msg_id, entry in self._pending_acks.items(): + elapsed = now - entry["time"] + if elapsed >= ACK_TIMEOUT_S: + if entry["attempts"] < MAX_RETRIES: + entry["attempts"] += 1 + entry["time"] = now + to_retry.append((msg_id, entry["packet"], entry["attempts"])) + else: + to_drop.append(msg_id) + + for msg_id in to_drop: + del self._pending_acks[msg_id] + + for msg_id, packet, attempt in to_retry: + print(f"\n [RETRY] msg {msg_id[:12]}... attempt {attempt}/{MAX_RETRIES}") + print("mesh> ", end="", flush=True) + self._flood_packet(packet) + + for msg_id in to_drop: + print(f"\n [FAIL] msg {msg_id[:12]}... no ACK after {MAX_RETRIES} attempts") + print("mesh> ", end="", flush=True) def _broadcast_discovery(self): packet = json.dumps({ @@ -179,9 +288,6 @@ def _listen_discovery(self): continue if packet.get("node_id") == self.node_id: continue - # Charlie filter: drop foreign meshes. A packet without the - # field is treated as "beacon-default" for backward-compat - # with pre-Wave-1 peers. if packet.get("mesh_id", DEFAULT_MESH_ID) != self.mesh_id: continue peer_id = packet["node_id"] @@ -202,11 +308,21 @@ def _listen_messages(self): try: data, addr = self.msg_sock.recvfrom(BUFFER_SIZE) packet = json.loads(data.decode()) - if packet.get("type") != "msg": + pkt_type = packet.get("type") + + # Handle ACKs + if pkt_type == "ack": + if packet.get("mesh_id", DEFAULT_MESH_ID) != self.mesh_id: + continue + if packet.get("from_id") == self.node_id: + continue # ignore our own ACKs + self._handle_ack(packet) continue - # Charlie filter on msg packets too — same backward-compat - # default as discovery. + if pkt_type != "msg": + continue + + # Charlie filter if packet.get("mesh_id", DEFAULT_MESH_ID) != self.mesh_id: continue @@ -216,17 +332,21 @@ def _listen_messages(self): self._remember_msg_id(msg_id) dst_id = packet["dst_id"] + src_id = packet["src_id"] + # deliver if for us or broadcast if dst_id == self.node_id or dst_id == "^all": if self.on_message: self.on_message( - packet["src_id"], + src_id, packet["src_name"], dst_id, packet["text"], packet["hop_count"], msg_id, ) + # Send ACK back to originator + self._send_ack(msg_id, src_id) # relay if hops remain if packet["hop_count"] < packet["max_hops"]: @@ -235,7 +355,7 @@ def _listen_messages(self): text=packet["text"], msg_id=msg_id, hop_count=packet["hop_count"] + 1, - src_id=packet["src_id"], + src_id=src_id, max_hops=packet["max_hops"], ) except OSError: diff --git a/MeshService/server.py b/MeshService/server.py index 2b1f597..8054934 100644 --- a/MeshService/server.py +++ b/MeshService/server.py @@ -9,16 +9,19 @@ GET /status — node info + peer list GET /peers — list discovered peers GET /messages — all received messages - POST /send — send message: {"to": "bravo", "text": "hello"} - POST /broadcast — broadcast: {"text": "hello everyone"} + POST /send — send envelope: {"to": "bravo", "payload": {...}} + POST /broadcast — broadcast envelope: {"payload": {...}} WS /ws — real-time message stream """ import argparse import asyncio +import json import random import string import time +import uuid +from datetime import datetime, timezone from contextlib import asynccontextmanager from fastapi import FastAPI, WebSocket, WebSocketDisconnect, Request @@ -28,15 +31,40 @@ from mesh.transport import UDPTransport +MESH_ID = "beacon-default" + # --- globals filled in at startup --- transport: UDPTransport = None message_log: list[dict] = [] ws_clients: list[WebSocket] = [] +def make_envelope(payload: dict) -> dict: + """Wrap a payload in a MeshEnvelope.""" + return { + "event_uuid": str(uuid.uuid4()), + "mesh_id": MESH_ID, + "origin_node_id": transport.node_id, + "timestamp": datetime.now(timezone.utc).isoformat(), + "confirm_status": "n/a", + "payload": payload, + } + + def on_message(src_id, src_name, dst_id, text, hop_count, msg_id): target = "broadcast" if dst_id == "^all" else f"to {dst_id}" - print(f"\n >> [{src_name}] ({target}): \"{text}\" (hops: {hop_count})") + + # try to parse envelope from text + try: + envelope = json.loads(text) + payload = envelope.get("payload", {}) + display = f"{payload.get('operator', src_name)} · {payload.get('text', text)}" + except (json.JSONDecodeError, AttributeError): + envelope = None + display = text + + print(f"\n >> [{src_name}] ({target}): \"{display}\" (hops: {hop_count})") + entry = { "msg_id": msg_id, "from_id": src_id, @@ -46,6 +74,9 @@ def on_message(src_id, src_name, dst_id, text, hop_count, msg_id): "hops": hop_count, "timestamp": time.time(), } + if envelope: + entry["envelope"] = envelope + message_log.append(entry) # push to all connected websocket clients for ws in list(ws_clients): @@ -55,6 +86,22 @@ def on_message(src_id, src_name, dst_id, text, hop_count, msg_id): pass +def on_ack(msg_id, from_id, from_name): + """Called when an ACK is received for a message we sent.""" + entry = { + "type": "ack", + "ack_msg_id": msg_id, + "from_id": from_id, + "from_name": from_name, + "timestamp": time.time(), + } + for ws in list(ws_clients): + try: + asyncio.run_coroutine_threadsafe(ws.send_json(entry), loop) + except Exception: + pass + + loop: asyncio.AbstractEventLoop = None @@ -63,6 +110,7 @@ async def lifespan(app: FastAPI): global loop loop = asyncio.get_event_loop() transport.on_message = on_message + transport.on_ack = on_ack transport.start() yield transport.stop() @@ -84,6 +132,7 @@ def get_status(): "node_id": transport.node_id, "name": transport.name, "port": transport.port, + "mesh_id": MESH_ID, "peer_count": len(transport.peers), "message_count": len(message_log), } @@ -106,23 +155,80 @@ def get_messages(): async def send_message(request: Request): body = await request.json() to = body.get("to") - text = body.get("text", "") - if not to or not text: - return {"error": "provide 'to' and 'text'"} + if not to: + return {"error": "provide 'to'"} if to not in transport.peers: return {"error": f"unknown peer '{to}'", "known_peers": list(transport.peers.keys())} - transport.send(to, text) - return {"status": "sent", "to": to, "text": text} + + # Build envelope from payload fields + payload = body.get("payload") or { + "kind": body.get("kind", "operator_query"), + "operator": body.get("operator", transport.name), + "text": body.get("text", ""), + } + envelope = make_envelope(payload) + envelope_json = json.dumps(envelope) + + transport.send(to, envelope_json) + + # Echo to local WS clients + entry = { + "msg_id": f"local-{int(time.time()*1000)}", + "from_id": transport.node_id, + "from_name": transport.name, + "to": to, + "text": envelope_json, + "hops": 0, + "timestamp": time.time(), + "envelope": envelope, + } + message_log.append(entry) + for ws in list(ws_clients): + try: + await ws.send_json(entry) + except Exception: + pass + + return {"status": "sent", "to": to, "envelope": envelope} @app.post("/broadcast") async def broadcast_message(request: Request): body = await request.json() - text = body.get("text", "") - if not text: - return {"error": "provide 'text'"} - transport.broadcast(text) - return {"status": "broadcast", "text": text} + + # Build envelope from payload fields + payload = body.get("payload") or { + "kind": body.get("kind", "operator_query"), + "operator": body.get("operator", transport.name), + "text": body.get("text", ""), + } + if not payload.get("text"): + return {"error": "provide 'text' or 'payload.text'"} + + envelope = make_envelope(payload) + envelope_json = json.dumps(envelope) + + transport.broadcast(envelope_json) + + # Echo to local WS clients + entry = { + "msg_id": f"local-{int(time.time()*1000)}", + "from_id": transport.node_id, + "from_name": transport.name, + "to": "^all", + "text": envelope_json, + "hops": 0, + "timestamp": time.time(), + "envelope": envelope, + } + message_log.append(entry) + for ws in list(ws_clients): + try: + await ws.send_json(entry) + except Exception: + pass + + return {"status": "broadcast", "envelope": envelope} @app.websocket("/ws") @@ -131,7 +237,6 @@ async def websocket_endpoint(ws: WebSocket): ws_clients.append(ws) try: while True: - # keep connection alive, ignore client messages await ws.receive_text() except WebSocketDisconnect: ws_clients.remove(ws) @@ -155,6 +260,7 @@ def main(): print("=== Mesh Node API Server ===") print(f" Node: {name} ({node_id})") + print(f" Mesh: {MESH_ID}") print(f" UDP: {port}") print(f" HTTP: http://localhost:{args.http_port}") print(f" WS: ws://localhost:{args.http_port}/ws") diff --git a/SoldierOS/SoldierOS.xcodeproj/project.pbxproj b/SoldierOS/SoldierOS.xcodeproj/project.pbxproj index dd038a5..0e5f36e 100644 --- a/SoldierOS/SoldierOS.xcodeproj/project.pbxproj +++ b/SoldierOS/SoldierOS.xcodeproj/project.pbxproj @@ -6,6 +6,10 @@ objectVersion = 77; objects = { +/* Begin PBXBuildFile section */ + 0874BD072FA715AD00D99098 /* WhisperKit in Frameworks */ = {isa = PBXBuildFile; productRef = 0874BD062FA715AD00D99098 /* WhisperKit */; }; +/* End PBXBuildFile section */ + /* Begin PBXFileReference section */ 0874BC802FA6BABC00D99098 /* SoldierOS.app */ = {isa = PBXFileReference; explicitFileType = wrapper.application; includeInIndex = 0; path = SoldierOS.app; sourceTree = BUILT_PRODUCTS_DIR; }; /* End PBXFileReference section */ @@ -23,6 +27,7 @@ isa = PBXFrameworksBuildPhase; buildActionMask = 2147483647; files = ( + 0874BD072FA715AD00D99098 /* WhisperKit in Frameworks */, ); runOnlyForDeploymentPostprocessing = 0; }; @@ -65,6 +70,7 @@ ); name = SoldierOS; packageProductDependencies = ( + 0874BD062FA715AD00D99098 /* WhisperKit */, ); productName = SoldierOS; productReference = 0874BC802FA6BABC00D99098 /* SoldierOS.app */; @@ -94,6 +100,9 @@ ); mainGroup = 0874BC772FA6BABC00D99098; minimizedProjectReferenceProxies = 1; + packageReferences = ( + 0874BD052FA715AD00D99098 /* XCRemoteSwiftPackageReference "WhisperKit" */, + ); preferredProjectObjectVersion = 77; productRefGroup = 0874BC812FA6BABC00D99098 /* Products */; projectDirPath = ""; @@ -251,10 +260,22 @@ DEVELOPMENT_TEAM = V2FT48C373; ENABLE_APP_SANDBOX = YES; ENABLE_HARDENED_RUNTIME = YES; + ENABLE_INCOMING_NETWORK_CONNECTIONS = NO; + ENABLE_OUTGOING_NETWORK_CONNECTIONS = YES; ENABLE_PREVIEWS = YES; + ENABLE_RESOURCE_ACCESS_AUDIO_INPUT = NO; + ENABLE_RESOURCE_ACCESS_BLUETOOTH = NO; + ENABLE_RESOURCE_ACCESS_CALENDARS = NO; + ENABLE_RESOURCE_ACCESS_CAMERA = NO; + ENABLE_RESOURCE_ACCESS_CONTACTS = NO; + ENABLE_RESOURCE_ACCESS_LOCATION = NO; + ENABLE_RESOURCE_ACCESS_PRINTING = NO; + ENABLE_RESOURCE_ACCESS_USB = NO; ENABLE_USER_SELECTED_FILES = readonly; GENERATE_INFOPLIST_FILE = YES; INFOPLIST_KEY_LSApplicationCategoryType = ""; + INFOPLIST_KEY_NSMicrophoneUsageDescription = "SoldierOS needs microphone access for voice input."; + INFOPLIST_KEY_NSSpeechRecognitionUsageDescription = "SoldierOS needs speech recognition to transcribe audio into text."; "INFOPLIST_KEY_UIApplicationSceneManifest_Generation[sdk=iphoneos*]" = YES; "INFOPLIST_KEY_UIApplicationSceneManifest_Generation[sdk=iphonesimulator*]" = YES; "INFOPLIST_KEY_UIApplicationSupportsIndirectInputEvents[sdk=iphoneos*]" = YES; @@ -296,10 +317,22 @@ DEVELOPMENT_TEAM = V2FT48C373; ENABLE_APP_SANDBOX = YES; ENABLE_HARDENED_RUNTIME = YES; + ENABLE_INCOMING_NETWORK_CONNECTIONS = NO; + ENABLE_OUTGOING_NETWORK_CONNECTIONS = YES; ENABLE_PREVIEWS = YES; + ENABLE_RESOURCE_ACCESS_AUDIO_INPUT = NO; + ENABLE_RESOURCE_ACCESS_BLUETOOTH = NO; + ENABLE_RESOURCE_ACCESS_CALENDARS = NO; + ENABLE_RESOURCE_ACCESS_CAMERA = NO; + ENABLE_RESOURCE_ACCESS_CONTACTS = NO; + ENABLE_RESOURCE_ACCESS_LOCATION = NO; + ENABLE_RESOURCE_ACCESS_PRINTING = NO; + ENABLE_RESOURCE_ACCESS_USB = NO; ENABLE_USER_SELECTED_FILES = readonly; GENERATE_INFOPLIST_FILE = YES; INFOPLIST_KEY_LSApplicationCategoryType = ""; + INFOPLIST_KEY_NSMicrophoneUsageDescription = "SoldierOS needs microphone access for voice input."; + INFOPLIST_KEY_NSSpeechRecognitionUsageDescription = "SoldierOS needs speech recognition to transcribe audio into text."; "INFOPLIST_KEY_UIApplicationSceneManifest_Generation[sdk=iphoneos*]" = YES; "INFOPLIST_KEY_UIApplicationSceneManifest_Generation[sdk=iphonesimulator*]" = YES; "INFOPLIST_KEY_UIApplicationSupportsIndirectInputEvents[sdk=iphoneos*]" = YES; @@ -353,6 +386,25 @@ defaultConfigurationName = Release; }; /* End XCConfigurationList section */ + +/* Begin XCRemoteSwiftPackageReference section */ + 0874BD052FA715AD00D99098 /* XCRemoteSwiftPackageReference "WhisperKit" */ = { + isa = XCRemoteSwiftPackageReference; + repositoryURL = "https://github.com/argmaxinc/WhisperKit.git"; + requirement = { + branch = main; + kind = branch; + }; + }; +/* End XCRemoteSwiftPackageReference section */ + +/* Begin XCSwiftPackageProductDependency section */ + 0874BD062FA715AD00D99098 /* WhisperKit */ = { + isa = XCSwiftPackageProductDependency; + package = 0874BD052FA715AD00D99098 /* XCRemoteSwiftPackageReference "WhisperKit" */; + productName = WhisperKit; + }; +/* End XCSwiftPackageProductDependency section */ }; rootObject = 0874BC782FA6BABC00D99098 /* Project object */; } diff --git a/SoldierOS/SoldierOS.xcodeproj/project.xcworkspace/xcshareddata/swiftpm/Package.resolved b/SoldierOS/SoldierOS.xcodeproj/project.xcworkspace/xcshareddata/swiftpm/Package.resolved new file mode 100644 index 0000000..381de62 --- /dev/null +++ b/SoldierOS/SoldierOS.xcodeproj/project.xcworkspace/xcshareddata/swiftpm/Package.resolved @@ -0,0 +1,24 @@ +{ + "originHash" : "be2b508cf91e00ccfdfd83e8177b94c9aee963c0dbf9b604f03fa19a3d85c0d4", + "pins" : [ + { + "identity" : "swift-argument-parser", + "kind" : "remoteSourceControl", + "location" : "https://github.com/apple/swift-argument-parser.git", + "state" : { + "revision" : "626b5b7b2f45e1b0b1c6f4a309296d1d21d7311b", + "version" : "1.7.1" + } + }, + { + "identity" : "whisperkit", + "kind" : "remoteSourceControl", + "location" : "https://github.com/argmaxinc/WhisperKit.git", + "state" : { + "branch" : "main", + "revision" : "c9cf2033012c900d1b15a2f3f423c2901cbbf40d" + } + } + ], + "version" : 3 +} diff --git a/SoldierOS/SoldierOS.xcodeproj/project.xcworkspace/xcuserdata/angelzzz23.xcuserdatad/UserInterfaceState.xcuserstate b/SoldierOS/SoldierOS.xcodeproj/project.xcworkspace/xcuserdata/angelzzz23.xcuserdatad/UserInterfaceState.xcuserstate index e4389e0..118ab02 100644 Binary files a/SoldierOS/SoldierOS.xcodeproj/project.xcworkspace/xcuserdata/angelzzz23.xcuserdatad/UserInterfaceState.xcuserstate and b/SoldierOS/SoldierOS.xcodeproj/project.xcworkspace/xcuserdata/angelzzz23.xcuserdatad/UserInterfaceState.xcuserstate differ diff --git a/SoldierOS/SoldierOS/ContentView.swift b/SoldierOS/SoldierOS/ContentView.swift index 4a0dcd5..b3972e5 100644 --- a/SoldierOS/SoldierOS/ContentView.swift +++ b/SoldierOS/SoldierOS/ContentView.swift @@ -1,9 +1,10 @@ import SwiftUI + struct ContentView: View { @AppStorage("meshHost") private var host = "localhost" @AppStorage("meshPort") private var port = 8009 - @StateObject private var mesh = MeshService() + @ObservedObject var mesh: MeshService @State private var messageText = "" @State private var showSettings = false @@ -16,6 +17,7 @@ struct ContentView: View { } .navigationTitle("SoldierOS") .toolbar { + #if os(iOS) ToolbarItem(placement: .topBarLeading) { Button { showSettings = true @@ -32,6 +34,24 @@ struct ContentView: View { } } } + #else + ToolbarItem(placement: .navigation) { + Button { + showSettings = true + } label: { + Image(systemName: "gear") + } + } + ToolbarItem(placement: .automatic) { + Button(mesh.isConnected ? "Disconnect" : "Connect") { + if mesh.isConnected { + mesh.disconnect() + } else { + mesh.reconnect(host: host, port: port) + } + } + } + #endif } .sheet(isPresented: $showSettings) { SettingsView() @@ -75,6 +95,28 @@ struct ContentView: View { .font(.body) HStack { Text("hops: \(msg.hops)") + if let status = msg.deliveryStatus { + HStack(spacing: 3) { + switch status { + case .sending: + Image(systemName: "clock") + .foregroundStyle(.secondary) + case .sent: + Image(systemName: "checkmark") + .foregroundStyle(.secondary) + case .delivered: + Image(systemName: "checkmark.circle.fill") + .foregroundStyle(.green) + case .failed: + Image(systemName: "exclamationmark.triangle.fill") + .foregroundStyle(.red) + } + Text(status.rawValue) + if let acker = msg.ackedBy { + Text("by \(acker)") + } + } + } Spacer() Text(Date(timeIntervalSince1970: msg.timestamp), style: .time) } @@ -119,5 +161,5 @@ struct ContentView: View { } #Preview { - ContentView() + ContentView(mesh: MeshService()) } diff --git a/SoldierOS/SoldierOS/MeshService.swift b/SoldierOS/SoldierOS/MeshService.swift index cdd8c8c..ad88ada 100644 --- a/SoldierOS/SoldierOS/MeshService.swift +++ b/SoldierOS/SoldierOS/MeshService.swift @@ -1,7 +1,14 @@ import Foundation import Combine -struct MeshMessage: Identifiable, Codable { +enum DeliveryStatus: String { + case sending = "Sending..." + case sent = "Sent" + case delivered = "Delivered" + case failed = "Failed" +} + +struct MeshMessage: Identifiable { let msgId: String let fromId: String let fromName: String @@ -9,15 +16,10 @@ struct MeshMessage: Identifiable, Codable { let text: String let hops: Int let timestamp: Double + var deliveryStatus: DeliveryStatus? + var ackedBy: String? var id: String { msgId } - - enum CodingKeys: String, CodingKey { - case msgId = "msg_id" - case fromId = "from_id" - case fromName = "from_name" - case to, text, hops, timestamp - } } class MeshService: NSObject, ObservableObject, URLSessionWebSocketDelegate { @@ -121,9 +123,73 @@ class MeshService: NSObject, ObservableObject, URLSessionWebSocketDelegate { } private func handleData(_ data: Data) { - let decoder = JSONDecoder() - if let msg = try? decoder.decode(MeshMessage.self, from: data) { - messages.append(msg) + guard let json = try? JSONSerialization.jsonObject(with: data) as? [String: Any] else { return } + + // Check if this is an ACK + if json["type"] as? String == "ack" { + handleAck(json) + return + } + + // Regular message + guard let msgId = json["msg_id"] as? String, + let fromId = json["from_id"] as? String, + let fromName = json["from_name"] as? String, + let to = json["to"] as? String, + let rawText = json["text"] as? String, + let hops = json["hops"] as? Int, + let timestamp = json["timestamp"] as? Double else { return } + + // Extract readable text from envelope JSON, or use raw text + let displayText = extractPayloadText(from: rawText) + + // Check if this is our own echo (from a pending send) + if let idx = messages.firstIndex(where: { $0.msgId == "pending" && $0.text == displayText }) { + messages[idx] = MeshMessage( + msgId: msgId, + fromId: fromId, + fromName: fromName, + to: to, + text: displayText, + hops: hops, + timestamp: timestamp, + deliveryStatus: .sent + ) + return + } + + let msg = MeshMessage( + msgId: msgId, + fromId: fromId, + fromName: fromName, + to: to, + text: displayText, + hops: hops, + timestamp: timestamp, + deliveryStatus: nil + ) + messages.append(msg) + } + + private func extractPayloadText(from raw: String) -> String { + // Try to parse as envelope JSON and extract payload.text + guard let data = raw.data(using: .utf8), + let envelope = try? JSONSerialization.jsonObject(with: data) as? [String: Any], + let payload = envelope["payload"] as? [String: Any], + let text = payload["text"] as? String else { + return raw + } + return text + } + + private func handleAck(_ json: [String: Any]) { + guard let ackMsgId = json["ack_msg_id"] as? String, + let fromName = json["from_name"] as? String else { return } + + // Find the message and update its delivery status + if let idx = messages.firstIndex(where: { $0.msgId == ackMsgId }) { + messages[idx].deliveryStatus = .delivered + messages[idx].ackedBy = fromName } } @@ -134,7 +200,14 @@ class MeshService: NSObject, ObservableObject, URLSessionWebSocketDelegate { var request = URLRequest(url: url) request.httpMethod = "POST" request.setValue("application/json", forHTTPHeaderField: "Content-Type") - let body: [String: String] = ["to": peerId, "text": text] + let body: [String: Any] = [ + "to": peerId, + "payload": [ + "kind": "operator_query", + "operator": "debug_1", + "text": text, + ] as [String: String], + ] request.httpBody = try? JSONSerialization.data(withJSONObject: body) session.dataTask(with: request) { _, _, error in if let error = error { @@ -144,15 +217,41 @@ class MeshService: NSObject, ObservableObject, URLSessionWebSocketDelegate { } func broadcast(text: String) { + // Add local message immediately with "sending" status + let pending = MeshMessage( + msgId: "pending", + fromId: "me", + fromName: "You", + to: "^all", + text: text, + hops: 0, + timestamp: Date().timeIntervalSince1970, + deliveryStatus: .sending + ) + messages.append(pending) + let url = baseURL.appendingPathComponent("broadcast") var request = URLRequest(url: url) request.httpMethod = "POST" request.setValue("application/json", forHTTPHeaderField: "Content-Type") - let body: [String: String] = ["text": text] + let body: [String: Any] = [ + "payload": [ + "kind": "operator_query", + "operator": "debug_1", + "text": text, + ] as [String: String], + ] request.httpBody = try? JSONSerialization.data(withJSONObject: body) + session.dataTask(with: request) { _, _, error in if let error = error { print("[MeshService] Broadcast error: \(error.localizedDescription)") + DispatchQueue.main.async { + // Mark as failed + if let idx = self.messages.lastIndex(where: { $0.msgId == "pending" && $0.text == text }) { + self.messages[idx].deliveryStatus = .failed + } + } } }.resume() } diff --git a/SoldierOS/SoldierOS/NORTH-1-impacts-inside-the-wire.mp3 b/SoldierOS/SoldierOS/NORTH-1-impacts-inside-the-wire.mp3 new file mode 100644 index 0000000..5413fa1 Binary files /dev/null and b/SoldierOS/SoldierOS/NORTH-1-impacts-inside-the-wire.mp3 differ diff --git a/SoldierOS/SoldierOS/NORTH-1-three-more-impacts.mp3 b/SoldierOS/SoldierOS/NORTH-1-three-more-impacts.mp3 new file mode 100644 index 0000000..3394a98 Binary files /dev/null and b/SoldierOS/SoldierOS/NORTH-1-three-more-impacts.mp3 differ diff --git a/SoldierOS/SoldierOS/NORTH-1-two-contacts.mp3 b/SoldierOS/SoldierOS/NORTH-1-two-contacts.mp3 new file mode 100644 index 0000000..87fdd0c Binary files /dev/null and b/SoldierOS/SoldierOS/NORTH-1-two-contacts.mp3 differ diff --git a/SoldierOS/SoldierOS/OP-6-call-for-fire.mp3 b/SoldierOS/SoldierOS/OP-6-call-for-fire.mp3 new file mode 100644 index 0000000..6b41127 Binary files /dev/null and b/SoldierOS/SoldierOS/OP-6-call-for-fire.mp3 differ diff --git a/SoldierOS/SoldierOS/OP-6-contact-broken.mp3 b/SoldierOS/SoldierOS/OP-6-contact-broken.mp3 new file mode 100644 index 0000000..009f251 Binary files /dev/null and b/SoldierOS/SoldierOS/OP-6-contact-broken.mp3 differ diff --git a/SoldierOS/SoldierOS/OP-6-sitrep.mp3 b/SoldierOS/SoldierOS/OP-6-sitrep.mp3 new file mode 100644 index 0000000..bb52155 Binary files /dev/null and b/SoldierOS/SoldierOS/OP-6-sitrep.mp3 differ diff --git a/SoldierOS/SoldierOS/SOUTH-MED-medevac-update.mp3 b/SoldierOS/SoldierOS/SOUTH-MED-medevac-update.mp3 new file mode 100644 index 0000000..a26df50 Binary files /dev/null and b/SoldierOS/SoldierOS/SOUTH-MED-medevac-update.mp3 differ diff --git a/SoldierOS/SoldierOS/SOUTH-MED-nine-line-how-copy.mp3 b/SoldierOS/SoldierOS/SOUTH-MED-nine-line-how-copy.mp3 new file mode 100644 index 0000000..eed5734 Binary files /dev/null and b/SoldierOS/SoldierOS/SOUTH-MED-nine-line-how-copy.mp3 differ diff --git a/SoldierOS/SoldierOS/SOUTH-RTO-fuel-point-hit.mp3 b/SoldierOS/SoldierOS/SOUTH-RTO-fuel-point-hit.mp3 new file mode 100644 index 0000000..0f3137b Binary files /dev/null and b/SoldierOS/SoldierOS/SOUTH-RTO-fuel-point-hit.mp3 differ diff --git a/SoldierOS/SoldierOS/SOUTH-RTO-jp-8-at-three-zero.mp3 b/SoldierOS/SoldierOS/SOUTH-RTO-jp-8-at-three-zero.mp3 new file mode 100644 index 0000000..431a88a Binary files /dev/null and b/SoldierOS/SoldierOS/SOUTH-RTO-jp-8-at-three-zero.mp3 differ diff --git a/SoldierOS/SoldierOS/SOUTH-RTO-uav-bearing-three-five-zero.mp3 b/SoldierOS/SoldierOS/SOUTH-RTO-uav-bearing-three-five-zero.mp3 new file mode 100644 index 0000000..6168968 Binary files /dev/null and b/SoldierOS/SoldierOS/SOUTH-RTO-uav-bearing-three-five-zero.mp3 differ diff --git a/SoldierOS/SoldierOS/SettingsView.swift b/SoldierOS/SoldierOS/SettingsView.swift index 167e119..00ef827 100644 --- a/SoldierOS/SoldierOS/SettingsView.swift +++ b/SoldierOS/SoldierOS/SettingsView.swift @@ -14,15 +14,16 @@ struct SettingsView: View { Spacer() TextField("localhost", text: $host) .multilineTextAlignment(.trailing) - .autocorrectionDisabled() - .autocapitalization(.none) + .disableAutocorrection(true) } HStack { Text("Port") Spacer() TextField("8009", value: $port, format: .number) .multilineTextAlignment(.trailing) + #if os(iOS) .keyboardType(.numberPad) + #endif } } @@ -36,9 +37,15 @@ struct SettingsView: View { } .navigationTitle("Settings") .toolbar { + #if os(iOS) ToolbarItem(placement: .topBarTrailing) { Button("Done") { dismiss() } } + #else + ToolbarItem(placement: .automatic) { + Button("Done") { dismiss() } + } + #endif } } } diff --git a/SoldierOS/SoldierOS/SoldierOSApp.swift b/SoldierOS/SoldierOS/SoldierOSApp.swift index 0612736..f7cd2b6 100644 --- a/SoldierOS/SoldierOS/SoldierOSApp.swift +++ b/SoldierOS/SoldierOS/SoldierOSApp.swift @@ -1,17 +1,23 @@ -// -// SoldierOSApp.swift -// SoldierOS -// -// Created by angel zambrano on 5/2/26. -// - import SwiftUI @main struct SoldierOSApp: App { + @StateObject private var mesh = MeshService() + var body: some Scene { WindowGroup { - ContentView() + TabView { + ContentView(mesh: mesh) + .tabItem { + Image(systemName: "antenna.radiowaves.left.and.right") + Text("Mesh") + } + TrascriptionTest(mesh: mesh) + .tabItem { + Image(systemName: "waveform") + Text("Transcription") + } + } } } } diff --git a/SoldierOS/SoldierOS/WEST-1-contact.mp3 b/SoldierOS/SoldierOS/WEST-1-contact.mp3 new file mode 100644 index 0000000..330bdb8 Binary files /dev/null and b/SoldierOS/SoldierOS/WEST-1-contact.mp3 differ diff --git a/SoldierOS/SoldierOS/WEST-1-uav-swarm-bearing.mp3 b/SoldierOS/SoldierOS/WEST-1-uav-swarm-bearing.mp3 new file mode 100644 index 0000000..9be84b0 Binary files /dev/null and b/SoldierOS/SoldierOS/WEST-1-uav-swarm-bearing.mp3 differ diff --git a/SoldierOS/SoldierOS/transcription/TrascriptionTest.swift b/SoldierOS/SoldierOS/transcription/TrascriptionTest.swift new file mode 100644 index 0000000..c586cda --- /dev/null +++ b/SoldierOS/SoldierOS/transcription/TrascriptionTest.swift @@ -0,0 +1,326 @@ +import SwiftUI +import AVFoundation +import WhisperKit + +struct AudioEntry: Identifiable { + let id = UUID() + let filename: String + var transcription: String = "" + var status: TranscriptionStatus = .pending + + enum TranscriptionStatus { + case pending, playing, transcribing, done, error(String) + } +} + +enum TranscribeMode: String, CaseIterable { + case withAudio = "Play & Transcribe" + case transcribeOnly = "Transcribe Only" +} + +struct TrascriptionTest: View { + @ObservedObject var mesh: MeshService + @State private var audioEntries: [AudioEntry] = [ + AudioEntry(filename: "NORTH-1-two-contacts"), + AudioEntry(filename: "WEST-1-uav-swarm-bearing"), + AudioEntry(filename: "SOUTH-RTO-uav-bearing-three-five-zero"), + AudioEntry(filename: "NORTH-1-impacts-inside-the-wire"), + AudioEntry(filename: "SOUTH-RTO-fuel-point-hit"), + AudioEntry(filename: "SOUTH-MED-nine-line-how-copy"), + AudioEntry(filename: "WEST-1-contact"), + AudioEntry(filename: "SOUTH-RTO-jp-8-at-three-zero"), + AudioEntry(filename: "OP-6-call-for-fire"), + AudioEntry(filename: "SOUTH-MED-medevac-update"), + AudioEntry(filename: "NORTH-1-three-more-impacts"), + AudioEntry(filename: "OP-6-sitrep"), + AudioEntry(filename: "OP-6-contact-broken"), + ] + @State private var isRunning = false + @State private var currentIndex = 0 + @State private var errorMessage: String? + @State private var audioPlayer: AVAudioPlayer? + @State private var whisperKit: WhisperKit? + @State private var modelLoaded = false + @State private var loadingModel = false + @State private var cancelled = false + @State private var mode: TranscribeMode = .withAudio + + var body: some View { + NavigationStack { + VStack(spacing: 0) { + if !modelLoaded { + Button(action: loadModel) { + HStack { + if loadingModel { + ProgressView() + .progressViewStyle(CircularProgressViewStyle()) + Text("Loading Whisper model...") + } else { + Image(systemName: "arrow.down.circle") + Text("Load Whisper Model") + } + } + .frame(maxWidth: .infinity) + .padding() + .background(Color.orange) + .foregroundColor(.white) + .cornerRadius(10) + } + .disabled(loadingModel) + .padding() + } + + // Mode picker + Picker("Mode", selection: $mode) { + ForEach(TranscribeMode.allCases, id: \.self) { m in + Text(m.rawValue).tag(m) + } + } + .pickerStyle(.segmented) + .disabled(isRunning) + .padding(.horizontal) + .padding(.top, 8) + + // Start / Stop buttons + HStack(spacing: 12) { + Button(action: startTranscription) { + HStack { + Image(systemName: "play.fill") + Text(isRunning ? "\(completedCount)/\(audioEntries.count)..." : "Start") + } + .frame(maxWidth: .infinity) + .padding() + .background(isRunning || !modelLoaded ? Color.gray : Color.blue) + .foregroundColor(.white) + .cornerRadius(10) + } + .disabled(isRunning || !modelLoaded) + + if isRunning { + Button(action: stopTranscription) { + HStack { + Image(systemName: "stop.fill") + Text("Stop") + } + .frame(maxWidth: .infinity) + .padding() + .background(Color.red) + .foregroundColor(.white) + .cornerRadius(10) + } + } + } + .padding() + + if let errorMessage { + Text(errorMessage) + .foregroundStyle(.red) + .font(.caption) + .padding(.horizontal) + } + + List { + ForEach(audioEntries) { entry in + VStack(alignment: .leading, spacing: 6) { + HStack { + statusIcon(for: entry.status) + Text(entry.filename) + .font(.headline) + .lineLimit(1) + } + if !entry.transcription.isEmpty { + Text(entry.transcription) + .font(.body) + .foregroundStyle(.primary) + } else if case .error(let msg) = entry.status { + Text(msg) + .font(.caption) + .foregroundStyle(.red) + } + } + .padding(.vertical, 4) + } + } + .listStyle(.plain) + } + .navigationTitle("Transcription") + } + } + + private var completedCount: Int { + audioEntries.filter { + if case .done = $0.status { return true } + if case .error = $0.status { return true } + return false + }.count + } + + @ViewBuilder + private func statusIcon(for status: AudioEntry.TranscriptionStatus) -> some View { + switch status { + case .pending: + Image(systemName: "circle") + .foregroundStyle(.secondary) + case .playing: + Image(systemName: "speaker.wave.2.fill") + .foregroundStyle(.orange) + case .transcribing: + Image(systemName: "waveform.circle.fill") + .foregroundStyle(.blue) + case .done: + Image(systemName: "checkmark.circle.fill") + .foregroundStyle(.green) + case .error: + Image(systemName: "xmark.circle.fill") + .foregroundStyle(.red) + } + } + + // MARK: - Load + + private func loadModel() { + loadingModel = true + errorMessage = nil + Task { + do { + let kit = try await WhisperKit() + await MainActor.run { + whisperKit = kit + modelLoaded = true + loadingModel = false + } + } catch { + await MainActor.run { + errorMessage = "Failed to load model: \(error.localizedDescription)" + loadingModel = false + } + } + } + } + + // MARK: - Start / Stop + + private func startTranscription() { + errorMessage = nil + cancelled = false + for i in audioEntries.indices { + audioEntries[i].transcription = "" + audioEntries[i].status = .pending + } + currentIndex = 0 + isRunning = true + + if mode == .withAudio { + playAndTranscribeNext() + } else { + transcribeOnlyNext() + } + } + + private func stopTranscription() { + cancelled = true + audioPlayer?.stop() + audioPlayer = nil + isRunning = false + } + + // MARK: - Sequential: play audio then transcribe one at a time + + private func playAndTranscribeNext() { + guard !cancelled, currentIndex < audioEntries.count else { + isRunning = false + audioPlayer = nil + return + } + + let filename = audioEntries[currentIndex].filename + + guard let audioURL = Bundle.main.url(forResource: filename, withExtension: "mp3") else { + audioEntries[currentIndex].status = .error("File not found") + currentIndex += 1 + playAndTranscribeNext() + return + } + + audioEntries[currentIndex].status = .playing + do { + audioPlayer = try AVAudioPlayer(contentsOf: audioURL) + audioPlayer?.play() + } catch { + print("[Audio] Playback error: \(error.localizedDescription)") + } + + let duration = audioPlayer?.duration ?? 1.0 + DispatchQueue.main.asyncAfter(deadline: .now() + duration + 0.3) { + guard !self.cancelled else { return } + self.transcribeOne(index: self.currentIndex, audioURL: audioURL) { + self.currentIndex += 1 + self.playAndTranscribeNext() + } + } + } + + // MARK: - Sequential: transcribe only (no audio playback) + + private func transcribeOnlyNext() { + guard !cancelled, currentIndex < audioEntries.count else { + isRunning = false + return + } + + let filename = audioEntries[currentIndex].filename + + guard let audioURL = Bundle.main.url(forResource: filename, withExtension: "mp3") else { + audioEntries[currentIndex].status = .error("File not found") + currentIndex += 1 + transcribeOnlyNext() + return + } + + transcribeOne(index: currentIndex, audioURL: audioURL) { + // Send transcription over mesh + let text = self.audioEntries[self.currentIndex].transcription + if !text.isEmpty && text != "(no speech detected)" { + self.mesh.broadcast(text: text) + } + self.currentIndex += 1 + self.transcribeOnlyNext() + } + } + + // MARK: - Transcribe single file + + private func transcribeOne(index idx: Int, audioURL: URL, completion: @escaping () -> Void) { + guard !cancelled else { return } + audioEntries[idx].status = .transcribing + + guard let whisperKit = whisperKit else { + audioEntries[idx].status = .error("Model not loaded") + completion() + return + } + + Task { + do { + let results = try await whisperKit.transcribe(audioPath: audioURL.path) + let text = results.map { $0.text }.joined().trimmingCharacters(in: .whitespacesAndNewlines) + await MainActor.run { + guard !cancelled else { return } + audioEntries[idx].transcription = text.isEmpty ? "(no speech detected)" : text + audioEntries[idx].status = .done + completion() + } + } catch { + await MainActor.run { + guard !cancelled else { return } + audioEntries[idx].status = .error(error.localizedDescription) + completion() + } + } + } + } +} + +#Preview { + TrascriptionTest(mesh: MeshService()) +}