Files
naowalk/naowalk.py
T

1321 lines
62 KiB
Python

#!/usr/bin/env python
# -*- encoding: UTF-8 -*-
import qi
import argparse
import sys
import time
import threading
import queue
import cv2
import numpy as np
import pygame
import subprocess
import os
import json
import struct
import socket
import importlib.util
import shutil
import tempfile
import inspect
# Import the new music module
from naomusic import NaoMusicPlayer
# --- Audio for mic ---
try:
import pyaudio
HAS_AUDIO = True
except ImportError:
HAS_AUDIO = False
# --- SSH for live mic relay (push-to-talk) ---
try:
import paramiko
HAS_SSH = True
except ImportError:
HAS_SSH = False
# ==========================================
# MEDIA RELAY WIRE PROTOCOL (see nao_video_server.py)
# One video connection (server->client only) and one audio connection
# (bidirectional - server->client AUDIO is the robot's own mic,
# client->server TALK is push-to-talk audio to play on the robot's
# speaker), both length-framed: [1 byte type][4 byte length][payload].
# Replaces both the old HTTP MJPEG pull AND the old SSH `nc | aplay`
# push-to-talk relay - everything audio/video now goes through this
# one pair of connections.
# ==========================================
TYPE_HELLO = 0x00
TYPE_VIDEO = 0x01
TYPE_AUDIO = 0x02
TYPE_TALK = 0x03
def decode_ulaw(q):
"""uint8 mu-law samples -> int16 PCM. Must match encode_ulaw() in
nao_video_server.py exactly, or this just produces noise."""
mu = 255.0
y = (q.astype(np.float64) / 255.0) * 2.0 - 1.0
x = np.sign(y) * (1.0 / mu) * (np.power(1.0 + mu, np.abs(y)) - 1.0)
return np.clip(x * 32768.0, -32768, 32767).astype(np.int16)
def encode_ulaw(pcm16):
"""int16 numpy array -> uint8 numpy array, mu-law companded. Mirror
of nao_video_server.py's encode_ulaw() - used here now to compress
push-to-talk audio before it crosses the (possibly tunneled) link,
the same way the robot's own mic audio already was. Must match
nao_video_server.py's decode_ulaw() exactly."""
mu = 255.0
x = np.clip(pcm16.astype(np.float64) / 32768.0, -1.0, 1.0)
y = np.sign(x) * np.log1p(mu * np.abs(x)) / np.log1p(mu)
return (((y + 1.0) / 2.0) * 255.0).round().astype(np.uint8)
def _recv_exact(sock, n):
"""Read exactly n bytes from a TCP socket (recv() can return short
reads/chunks at any point - never assume one call gives you n)."""
buf = b""
while len(buf) < n:
chunk = sock.recv(n - len(buf))
if not chunk:
raise IOError("media relay connection closed")
buf += chunk
return buf
def _send_framed(sock, msg_type, payload):
"""Mirror of nao_video_server.py's _send() - used to push TYPE_TALK
(push-to-talk audio) the other way, client -> robot."""
sock.sendall(struct.pack(">BI", msg_type, len(payload)))
sock.sendall(payload)
# ==========================================
# NAO MIC STREAM
# ==========================================
class SoundReceiver:
PREBUFFER_CHUNKS = 1 # was 3 - less silence before playback starts
TARGET_QUEUE_DEPTH = 4 # once a stall lets the backlog build past
# QUEUE_MAXSIZE and we start dropping, trim
# it back down to this instead of just
# shaving one chunk at a time - a played-back
# backlog sounds like the mic sped up, and
# that's more noticeable than a couple of
# dropped chunks
QUEUE_MAXSIZE = 12 # was 40 - bounds worst-case playback lag
SILENCE_CHUNK = b"\x00" * 4096
def __init__(self, output_device_index=None, mic_gain=8.0, pyaudio_instance=None):
self.running = True
self.underrun_count = 0
self.chunks_received = 0
self._write_errors = 0
self._peak_since_print = 0
self._last_status_print = 0
self.mic_gain = float(mic_gain)
self.muted = False
self._owns_pyaudio = False
if HAS_AUDIO:
self.p = pyaudio_instance if pyaudio_instance is not None else pyaudio.PyAudio()
self._owns_pyaudio = pyaudio_instance is None
print("🔈 Available output devices:")
for i in range(self.p.get_device_count()):
info = self.p.get_device_info_by_index(i)
if info.get("maxOutputChannels", 0) > 0:
print(f" [{i}] {info['name']}")
try:
chosen = (self.p.get_device_info_by_index(output_device_index) if output_device_index is not None
else self.p.get_default_output_device_info())
print(f"🔈 Using: [{chosen['index']}] {chosen['name']}")
except: pass
self.stream = self.p.open(format=pyaudio.paInt16, channels=1, rate=16000,
output=True, output_device_index=output_device_index,
frames_per_buffer=1024)
self._play_test_tone()
self.queue = queue.Queue(maxsize=self.QUEUE_MAXSIZE)
self.playback_thread = threading.Thread(target=self._playback_loop, daemon=True)
self.playback_thread.start()
def _play_test_tone(self):
try:
duration, freq, rate = 0.3, 440.0, 16000
t = np.linspace(0, duration, int(rate * duration), False)
tone = (np.sin(freq * t * 2 * np.pi) * 12000).astype(np.int16)
self.stream.write(tone.tobytes())
except: pass
def processRemote(self, nbOfChannels, nbOfSamplesByChannel, timeStamp, buffer):
self.feed(bytes(buffer))
def feed(self, raw_pcm_bytes):
"""Push one chunk of 16-bit PCM into the playback queue,
dropping the oldest chunk on overflow instead of blocking.
Used both by the direct qi audio callback (processRemote,
above) and by the TCP media relay path (which mu-law-decodes
back to PCM before calling this)."""
if not HAS_AUDIO: return
self.chunks_received += 1
try:
self.queue.put_nowait(raw_pcm_bytes)
except queue.Full:
try:
self.queue.get_nowait()
self.queue.put_nowait(raw_pcm_bytes)
except: pass
# A stall (network hiccup, etc.) can let the backlog build up
# even before the queue is literally Full - trim it back down
# toward real-time here rather than only reacting once it hits
# the ceiling above, so a brief stall degrades to a couple of
# dropped chunks instead of several seconds of played-back lag.
while self.queue.qsize() > self.TARGET_QUEUE_DEPTH:
try:
self.queue.get_nowait()
except queue.Empty:
break
def _playback_loop(self):
primed = []
while self.running and len(primed) < self.PREBUFFER_CHUNKS:
try:
primed.append(self.queue.get(timeout=1.0))
except queue.Empty:
break
for chunk in primed:
self._write_chunk(chunk)
while self.running:
try:
chunk = self.queue.get(timeout=0.2)
self._write_chunk(chunk)
except queue.Empty:
self.underrun_count += 1
try:
self.stream.write(self.SILENCE_CHUNK)
except: pass
now = time.time()
if now - self._last_status_print > 10:
self._last_status_print = now
print(f"🎤 mic: {self.chunks_received} chunks, {self.underrun_count} underruns, peak {self._peak_since_print}")
self._peak_since_print = 0
def _write_chunk(self, raw_bytes):
try:
if self.muted:
self.stream.write(b"\x00" * len(raw_bytes))
return
samples = np.frombuffer(raw_bytes, dtype=np.int16).astype(np.float32)
if samples.size:
amplified = np.clip(samples * self.mic_gain, -32768, 32767)
peak = int(np.abs(amplified).max())
if peak > self._peak_since_print:
self._peak_since_print = peak
self.stream.write(amplified.astype(np.int16).tobytes())
else:
self.stream.write(raw_bytes)
except: pass
def close(self):
self.running = False
if HAS_AUDIO:
try: self.playback_thread.join(timeout=1.0)
except: pass
try:
self.stream.stop_stream()
self.stream.close()
if self._owns_pyaudio:
self.p.terminate()
except: pass
# ==========================================
# PUSH-TO-TALK (over the media relay's audio connection)
# ==========================================
class PushToTalk:
"""
Push-to-talk: streams the local mic to NAO's speaker over the same
TCP audio connection the media relay already uses for the robot's
own mic (see NaoTeleop._audio_relay_loop / _send_talk_chunk) -
mu-law encoded, framed as TYPE_TALK, decoded and played by
nao_video_server.py's TalkPlayer.
This replaces the old SSH-based relay (`nc -l | aplay` over a
dedicated SSH channel): no SSH connection of its own, no per-press
socket handshake - start()/stop() just toggle a background thread
that reads the mic and hands chunks to send_fn, which writes
straight into the already-open, already-connected audio socket.
"""
def __init__(self, send_fn, is_ready_fn, rate=16000, chunk_frames=1024,
input_device_index=None, pyaudio_instance=None):
self.send_fn = send_fn
self.is_ready_fn = is_ready_fn
self.rate = rate
self.chunk_frames = chunk_frames
self.input_device_index = input_device_index
# Reuse the app's single PyAudio context - a second independent one
# opened from a background thread crashes PortAudio's PulseAudio
# backend (pa_atomic_load assertion).
self._p = pyaudio_instance if pyaudio_instance is not None else (pyaudio.PyAudio() if HAS_AUDIO else None)
self.active = False
# "idle" | "live" - read by the HUD to show status
self.state = "idle"
self._stream = None
self._thread = None
def start(self):
if self.active:
return
if not HAS_AUDIO:
print("🎙️ Live mic needs pyaudio installed")
return
if not self.is_ready_fn():
print("🎙️ Talk relay: audio connection not ready yet")
return
# If a previous relay thread is still tearing itself down, wait for
# it to finish before reusing self._stream for a new one.
if self._thread and self._thread.is_alive():
self._thread.join(timeout=2.0)
self.active = True
self.state = "live"
self._thread = threading.Thread(target=self._run, daemon=True)
self._thread.start()
def stop(self):
"""Signal the relay thread to stop and wait for it to exit.
Deliberately does NOT touch self._stream here: this is called
from the main/render thread on key-up, and closing a PortAudio
stream from a different thread while the relay thread is mid
blocking-read() on it is what caused the ALSA segfault on
release. Only _run() closes it, from its own thread.
"""
self.active = False
self.state = "idle"
if self._thread and self._thread.is_alive() and threading.current_thread() is not self._thread:
self._thread.join(timeout=2.0)
def _run(self):
try:
self._stream = self._p.open(format=pyaudio.paInt16, channels=1,
rate=self.rate, input=True,
input_device_index=self.input_device_index,
frames_per_buffer=self.chunk_frames)
print("🎙️ Talk relay: live")
while self.active:
try:
chunk = self._stream.read(self.chunk_frames, exception_on_overflow=False)
pcm = np.frombuffer(chunk, dtype=np.int16)
self.send_fn(encode_ulaw(pcm).tobytes())
except Exception:
break
except Exception as e:
print(f"🎙️ Talk relay error: {e}")
finally:
self.active = False
self.state = "idle"
if self._stream:
try:
self._stream.stop_stream()
self._stream.close()
except: pass
self._stream = None
print("🎙️ Talk relay: stopped")
def close(self):
"""No persistent connection of its own to tear down anymore -
kept as an alias for stop() so shutdown code doesn't need to
know which push-to-talk implementation it's holding."""
self.stop()
# ==========================================
# ON-ROBOT VIDEO/AUDIO RELAY SERVER (auto-deploy)
# ==========================================
class VideoServerLauncher:
"""
SCPs nao_video_server.py to the robot's home directory over SFTP and
launches it via a persistent, pty-backed SSH session - so the user
never has to manually copy the file over or SSH in to start/stop it.
Uses the same trick as LiveMicRelay: a pty ties the remote process
group to the SSH channel, so close() (killing the channel) reliably
kills the video server too instead of leaving it orphaned on the
robot, still holding VIDEO_PORT/AUDIO_PORT, for next time.
"""
def __init__(self, nao_ip, ssh_user="nao", ssh_password="nao", ssh_port=22,
local_path=None, remote_filename="nao_video_server.py"):
self.nao_ip = nao_ip
self.ssh_user = ssh_user
self.ssh_password = ssh_password
self.ssh_port = ssh_port
self.local_path = local_path or os.path.join(
os.path.dirname(os.path.abspath(__file__)), remote_filename)
self.remote_filename = remote_filename
self.ready = False
self._ssh_client = None
self._channel = None
self._drain_thread = None
def deploy_and_start(self):
if not HAS_SSH:
print("🎥 Video server auto-deploy needs paramiko: pip install paramiko")
return False
if not os.path.isfile(self.local_path):
print(f"🎥 Video server auto-deploy: can't find {self.local_path}")
return False
try:
self._ssh_client = paramiko.SSHClient()
self._ssh_client.set_missing_host_key_policy(paramiko.AutoAddPolicy())
self._ssh_client.connect(self.nao_ip, port=self.ssh_port, username=self.ssh_user,
password=self.ssh_password, timeout=10)
# Best-effort: a previous session that didn't get to call
# close() cleanly (laptop lost network, crashed, etc.) can
# leave a stale instance squatting on the ports - clear it
# before starting a fresh one.
try:
_, out, _ = self._ssh_client.exec_command(
f"pkill -f {self.remote_filename}", timeout=5)
out.channel.recv_exit_status()
time.sleep(0.3)
except Exception:
pass
sftp = self._ssh_client.open_sftp()
try:
remote_home = sftp.normalize(".") # robot's home dir
remote_path = remote_home + "/" + self.remote_filename
print(f"🎥 Copying {self.remote_filename} -> {self.nao_ip}:{remote_path}")
sftp.put(self.local_path, remote_path)
finally:
sftp.close()
self._channel = self._ssh_client.get_transport().open_session()
self._channel.get_pty()
remote_cmd = (
"export PYTHONPATH=/opt/aldebaran/lib/python2.7/site-packages && "
"export LD_LIBRARY_PATH=/opt/aldebaran/lib && "
f"python2 {remote_path}"
)
self._channel.exec_command(remote_cmd)
self._drain_thread = threading.Thread(target=self._drain_output, daemon=True)
self._drain_thread.start()
time.sleep(1.0) # give it a moment to bind before the client tries to connect
if self._channel.exit_status_ready():
print("🎥 Video server exited immediately - check the log lines above")
self.ready = False
return False
self.ready = True
print("🎥 Video server deployed and running on robot")
return True
except Exception as e:
print(f"🎥 Video server auto-deploy failed: {e}")
self.ready = False
return False
def _drain_output(self):
"""Keeps reading the remote process's stdout/stderr so its pty
buffer never fills up and stalls the server, and so its own
status prints (quality tier changes, client connects) surface
locally too."""
chan = self._channel
while chan and not chan.closed:
try:
got = False
if chan.recv_ready():
data = chan.recv(4096)
if data:
print(f"🎥 [nao_video_server] {data.decode('utf-8', 'replace').rstrip()}")
got = True
if chan.recv_stderr_ready():
data = chan.recv_stderr(4096)
if data:
print(f"🎥 [nao_video_server] {data.decode('utf-8', 'replace').rstrip()}")
got = True
if not got:
time.sleep(0.1)
except Exception:
break
def close(self):
"""Kill the pty-backed channel (which kills the remote python
process with it) and tear down the SSH connection. Call once at
app shutdown."""
if self._channel:
try: self._channel.close()
except: pass
self._channel = None
if self._ssh_client:
try: self._ssh_client.close()
except: pass
self._ssh_client = None
self.ready = False
# ==========================================
# CHOREGRAPHE BEHAVIOR DEPLOY + RUN (for gestures)
# ==========================================
class ChoregrapheBehaviorRunner:
"""Deploys a Choregraphe-exported behavior folder (a manifest.xml
plus one or more '<BehaviorName>/behavior.xar' dirs, i.e. exactly
what "Export as package" produces before you zip it) to the robot
over SFTP, installs it via the PackageManager service, and runs it
via ALBehaviorManager - the same three steps Choregraphe's own
"Connect to robot" + Play button does, done here in code so a
gesture script can just point at a local project folder.
Caches per local folder for the lifetime of this object: once a
package has been deployed+installed this session, repeat runs
just call runBehavior() again instead of re-uploading.
NOTE: PackageManager.install() is documented as "install a package
from a path" without saying whether that path can be a raw folder
or must be a zipped .pkg - this tries the raw folder first (since
that's genuinely what Choregraphe's live "current project" push
to a connected robot does) and falls back to zipping it into a
.pkg if the robot's PackageManager rejects the folder.
"""
def __init__(self, session, nao_ip, ssh_user, ssh_password, ssh_port):
self.nao_ip = nao_ip
self.ssh_user = ssh_user
self.ssh_password = ssh_password
self.ssh_port = ssh_port
try:
self.package_mgr = session.service("PackageManager")
except Exception:
self.package_mgr = session.service("ALPackageManager")
self.behavior_mgr = session.service("ALBehaviorManager")
self._installed = {} # local_dir -> resolved "<package>/<behavior>" name
def _sftp_put_dir(self, sftp, local_dir, remote_dir):
"""Recursively mirror local_dir onto remote_dir over an
already-open SFTP session - paramiko has no built-in recursive
put, only single-file put()."""
try:
sftp.mkdir(remote_dir)
except IOError:
pass # already exists, fine
for entry in sorted(os.listdir(local_dir)):
local_path = os.path.join(local_dir, entry)
remote_path = remote_dir + "/" + entry
if os.path.isdir(local_path):
self._sftp_put_dir(sftp, local_path, remote_path)
else:
sftp.put(local_path, remote_path)
def ensure_running(self, local_dir, behavior_subdir):
"""Deploy+install local_dir (only if not already done this
session), then block until '<package>/<behavior_subdir>' has
finished running on the robot. local_dir's own folder name
becomes the package id, matching how Choregraphe names an
exported project after its folder."""
if not HAS_SSH:
raise RuntimeError("Choregraphe gesture deploy needs paramiko: pip install paramiko")
if not os.path.isdir(local_dir):
raise RuntimeError(f"No Choregraphe project folder at {local_dir}")
package_name = os.path.basename(os.path.normpath(local_dir))
behavior_name = self._installed.get(local_dir)
if behavior_name is None:
ssh = paramiko.SSHClient()
ssh.set_missing_host_key_policy(paramiko.AutoAddPolicy())
ssh.connect(self.nao_ip, port=self.ssh_port, username=self.ssh_user,
password=self.ssh_password, timeout=10)
try:
sftp = ssh.open_sftp()
try:
remote_home = sftp.normalize(".")
remote_dir = f"{remote_home}/naowalk_gestures/{package_name}"
print(f"🕺 Copying {package_name} -> {self.nao_ip}:{remote_dir}")
self._sftp_put_dir(sftp, local_dir, remote_dir)
print(f"🕺 Installing package {package_name}...")
try:
self.package_mgr.install(remote_dir)
except Exception as e:
print(f"🕺 install() on the raw folder failed ({e}) - "
f"retrying as a zipped .pkg...")
local_zip = shutil.make_archive(
os.path.join(tempfile.gettempdir(), package_name), "zip", local_dir)
remote_zip = f"{remote_home}/naowalk_gestures/{package_name}.pkg"
sftp.put(local_zip, remote_zip)
os.remove(local_zip)
self.package_mgr.install(remote_zip)
finally:
sftp.close()
finally:
ssh.close()
behavior_name = f"{package_name}/{behavior_subdir}"
if not self.behavior_mgr.isBehaviorInstalled(behavior_name):
installed = self.behavior_mgr.getInstalledBehaviors()
candidates = [b for b in installed if b.endswith("/" + behavior_subdir)]
if not candidates:
raise RuntimeError(
f"'{behavior_name}' not found after install - installed behaviors: {installed}")
behavior_name = candidates[0]
self._installed[local_dir] = behavior_name
print(f"🕺 Running {behavior_name}...")
self.behavior_mgr.runBehavior(behavior_name) # blocks until it finishes
# ==========================================
# MAIN TELEOP
# ==========================================
class NaoTeleop:
HEAD_RATE = 1.4 # rad/sec - how fast a held IJKL key sweeps the head target angle
def __init__(self, session, nao_ip, audio_device_index=None, mic_gain=8.0, video_scale=2.0,
nao_ssh_user="nao", nao_ssh_password="nao", nao_ssh_port=22, relay_rate=16000,
media_relay=None, deploy_video_server=True, video_server_path=None,
gestures_dir=None):
self.session = session
self.nao_ip = nao_ip
self.nao_ssh_user = nao_ssh_user
self.nao_ssh_password = nao_ssh_password
self.nao_ssh_port = nao_ssh_port
self.motion = session.service("ALMotion")
self.posture = session.service("ALRobotPosture")
self.tts = session.service("ALTextToSpeech")
try:
self.battery = session.service("ALBattery")
except:
self.battery = None
self.music = NaoMusicPlayer(session, nao_ip, ssh_port=nao_ssh_port)
self.music_results = []
self.selected_result = 0
self.audio_device = None # only used for volume get/set now - video
# and both mic directions are relay-only
self.audio_device_index = audio_device_index
self.mic_gain = mic_gain
# One shared PyAudio context for the whole app - SoundReceiver (mic
# playback) and PushToTalk (push-to-talk capture) both use it,
# since two separate contexts crash PortAudio's PulseAudio backend.
self._pyaudio = pyaudio.PyAudio() if HAS_AUDIO else None
# The audio relay connection is bidirectional (see
# _audio_relay_loop) - this is the socket PushToTalk writes
# TYPE_TALK chunks into, guarded by a lock since the relay loop
# (on its own thread) can tear it down and reconnect at any time.
self._audio_sock = None
self._audio_sock_lock = threading.Lock()
self.mic_relay = PushToTalk(send_fn=self._send_talk_chunk, is_ready_fn=self._audio_relay_ready,
rate=relay_rate, pyaudio_instance=self._pyaudio)
self.ptt_active = False
self.running = True
self.head_yaw = 0.0
self.head_pitch = 0.0
self._last_head_update = time.time()
self.battery_level = 100
self.last_batt_check = 0
self.volume = 50
self.walk_speed = 1.0 # fixed, not adjustable at runtime - 0.60 was too slow
# for the walk engine to stay balanced and NAO kept
# tipping over
self.typing_mode = False
self.is_resting = False # tracks manual [DEL] crouch, toggled by idle_crouch()
self.chat_message = ""
self.music_search_mode = False
self.music_search_text = ""
self.music_search_pending = False
self.video_scale = video_scale
self.media_relay = media_relay
self.display_width = int(320 * video_scale)
self.display_height = int(240 * video_scale)
# Video and movement are decoupled onto their own threads (see
# _init_video / run) so a stalled network RPC to the robot
# never blocks keyboard input or the render loop.
self._latest_frame = None
self._latest_frame_ts = 0.0
self._video_stop_event = threading.Event()
self._desired_velocity = (0.0, 0.0, 0.0)
self._movement_stop_event = threading.Event()
self._movement_thread = None
# Deploy+launch nao_video_server.py over SSH now (before _init_video
# tries to connect to it) instead of requiring it to already be
# running manually. Everything - video, robot mic, push-to-talk -
# goes through this relay; there's no direct-qi fallback anymore.
self.video_server_launcher = None
if deploy_video_server:
self.video_server_launcher = VideoServerLauncher(
nao_ip, ssh_user=nao_ssh_user, ssh_password=nao_ssh_password,
ssh_port=nao_ssh_port, local_path=video_server_path)
self.video_server_launcher.deploy_and_start()
self._init_audio_device()
self._init_video()
if HAS_AUDIO:
self._init_mic_stream()
# Lets gesture scripts run full Choregraphe-authored behaviors
# (.xar, exported as a package folder) instead of hand-coding
# angleInterpolation calls - see ChoregrapheBehaviorRunner.
self.behavior_runner = ChoregrapheBehaviorRunner(
session, nao_ip, nao_ssh_user, nao_ssh_password, nao_ssh_port)
self.gestures_dir = gestures_dir or os.path.join(
os.path.dirname(os.path.abspath(__file__)), "gestures")
self._load_gestures()
pygame.init()
self.screen = pygame.display.set_mode((self.display_width, self.display_height))
pygame.display.set_caption("NAO Teleop + Music on Robot")
self.font = pygame.font.SysFont(None, 24)
self.clock = pygame.time.Clock()
self.update_title()
def _safe(self, label, fn):
try:
return fn()
except Exception as e:
print(f"{label} error: {e}")
return None
def _init_audio_device(self):
try:
self.audio_device = self.session.service("ALAudioDevice")
self.volume = self.audio_device.getOutputVolume()
print(f"🔊 NAO speaker volume: {self.volume}%")
except Exception as e:
print(f"⚠️ ALAudioDevice error: {e}")
self.audio_device = None
def _init_video(self):
print(f"✅ Camera via on-robot media relay: {self.media_relay}")
threading.Thread(target=self._video_relay_loop, daemon=True).start()
def _video_relay_loop(self):
"""Runs on its own thread for the app's lifetime, speaking the
lean length-framed TCP protocol from nao_video_server.py. The
render loop never touches the socket directly; it only ever
reads whatever self._latest_frame was last successfully set
to, so a stalled/reconnecting link can no longer freeze
keyboard input or movement. Video-only connection - a slow
video send could otherwise block audio behind it, so they're
split (see the audio port, media_relay port + 1, handled by
_audio_relay_loop)."""
host, _, port_str = self.media_relay.partition(":")
port = int(port_str) if port_str else 8000
while not self._video_stop_event.is_set():
sock = None
try:
sock = socket.create_connection((host, port), timeout=5)
sock.settimeout(5.0)
sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1)
print(f"✅ Video relay connected: {host}:{port}")
while not self._video_stop_event.is_set():
msg_type, length = struct.unpack(">BI", _recv_exact(sock, 5))
payload = _recv_exact(sock, length) if length else b""
if msg_type == TYPE_VIDEO:
arr = np.frombuffer(payload, dtype=np.uint8)
frame = cv2.imdecode(arr, cv2.IMREAD_COLOR)
if frame is not None:
self._latest_frame = cv2.resize(frame, (320, 240), interpolation=cv2.INTER_NEAREST)
self._latest_frame_ts = time.time()
# TYPE_HELLO: nothing to do with it yet, just skip.
except Exception as e:
print(f"video relay error: {e}")
finally:
if sock:
try: sock.close()
except: pass
if not self._video_stop_event.is_set():
time.sleep(0.5)
def _audio_relay_loop(self):
"""Mic audio's own connection - see _video_relay_loop for why
this is split off rather than sharing video's socket. Connects
to media_relay's host:port+1 by convention (matches AUDIO_PORT
= VIDEO_PORT + 1 in nao_video_server.py).
This connection is bidirectional: this loop reads TYPE_AUDIO
(robot's mic) off it, and _send_talk_chunk (called from the
PushToTalk thread) writes TYPE_TALK (push-to-talk audio) onto
the same socket - TCP is full-duplex, so neither direction
blocks the other. self._audio_sock is published here (and
cleared on disconnect) so _send_talk_chunk always has whatever
socket is currently live, or nothing to send to while
reconnecting."""
host, _, port_str = self.media_relay.partition(":")
video_port = int(port_str) if port_str else 8000
port = video_port + 1
while not self._video_stop_event.is_set():
sock = None
try:
sock = socket.create_connection((host, port), timeout=5)
sock.settimeout(5.0)
sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1)
with self._audio_sock_lock:
self._audio_sock = sock
print(f"✅ Audio relay connected: {host}:{port}")
while not self._video_stop_event.is_set():
msg_type, length = struct.unpack(">BI", _recv_exact(sock, 5))
payload = _recv_exact(sock, length) if length else b""
if msg_type == TYPE_AUDIO and getattr(self, "sound_receiver", None):
pcm = decode_ulaw(np.frombuffer(payload, dtype=np.uint8))
self.sound_receiver.feed(pcm.tobytes())
# TYPE_HELLO: nothing to do with it yet, just skip.
except Exception as e:
print(f"audio relay error: {e}")
finally:
with self._audio_sock_lock:
if self._audio_sock is sock:
self._audio_sock = None
if sock:
try: sock.close()
except: pass
if not self._video_stop_event.is_set():
time.sleep(0.5)
def _audio_relay_ready(self):
"""Whether the audio connection is currently up - PushToTalk
checks this before starting a press so it can fail fast with a
clear message instead of silently sending into the void."""
with self._audio_sock_lock:
return self._audio_sock is not None
def _send_talk_chunk(self, payload):
"""Called from the PushToTalk thread with one already mu-law-
encoded chunk. Grabs whatever socket _audio_relay_loop currently
has published and writes straight into it - if the link drops
mid-send this just raises/gets swallowed and the chunk is lost,
which is fine for live voice (the next chunk is a beat away)."""
with self._audio_sock_lock:
sock = self._audio_sock
if not sock:
return
try:
_send_framed(sock, TYPE_TALK, payload)
except Exception:
pass # _audio_relay_loop will notice the failure and reconnect
def get_image(self):
"""Non-blocking: returns whatever frame the background video
thread most recently captured (or None before the first one
arrives). Never makes a network call itself."""
return self._latest_frame
def _init_mic_stream(self):
try:
self.sound_receiver = SoundReceiver(self.audio_device_index, self.mic_gain,
pyaudio_instance=self._pyaudio)
print(f"✅ Mic stream ready (gain={self.mic_gain}x, via media relay - compressed, own connection)")
threading.Thread(target=self._audio_relay_loop, daemon=True).start()
except Exception as e:
print(f"❌ Mic init failed: {e}")
def _load_gestures(self):
"""Load number-key gestures from .py files in self.gestures_dir.
Each file is named '<digit>_<name>.py' (e.g. '4_gangnam_style.py')
and must define a module-level run(motion, tts) function - that's
the minimum contract, gets the already-connected ALMotion and
ALTextToSpeech proxies and does whatever it wants with them.
Gestures that need more - e.g. running a Choregraphe behavior via
self.behavior_runner, or resetting posture via
self.stand_for_walk() - can instead define run(motion, tts,
teleop) and get this NaoTeleop instance as a third argument; the
arity is detected here at load time so both signatures work. The
leading digit is which number key (0-9) triggers it; everything
after the underscore becomes its display name in the controls
hint, unless the file sets its own module-level NAME string.
"""
self.gestures = {} # pygame key constant -> (display_name, run_fn, n_params)
gestures_dir = self.gestures_dir
if not os.path.isdir(gestures_dir):
print(f"⚠️ No gestures folder at {gestures_dir}, skipping gesture load")
return
for fname in sorted(os.listdir(gestures_dir)):
if not fname.endswith(".py") or fname.startswith("_"):
continue
stem = fname[:-3]
digit_str, sep, rest = stem.partition("_")
if not (sep and len(digit_str) == 1 and digit_str.isdigit()):
print(f"⚠️ Skipping {fname}: name it '<digit>_name.py' (e.g. '4_gangnam_style.py')")
continue
path = os.path.join(gestures_dir, fname)
try:
spec = importlib.util.spec_from_file_location(f"nao_gesture_{stem}", path)
module = importlib.util.module_from_spec(spec)
spec.loader.exec_module(module)
run_fn = getattr(module, "run", None)
if not callable(run_fn):
print(f"⚠️ Skipping {fname}: no run(motion, tts) function defined")
continue
try:
n_params = len(inspect.signature(run_fn).parameters)
except (TypeError, ValueError):
n_params = 2
except Exception as e:
print(f"⚠️ Failed to load gesture {fname}: {e}")
continue
key = getattr(pygame, f"K_{digit_str}")
name = getattr(module, "NAME", rest.replace("_", " ").title() or stem)
if key in self.gestures:
print(f"⚠️ {fname} overrides number key {digit_str} (was {self.gestures[key][0]})")
self.gestures[key] = (name, run_fn, n_params)
print(f"✅ Loaded gesture: {digit_str}={name}")
def run_gesture(self, key):
if self.typing_mode or self.is_resting: return
entry = self.gestures.get(key)
if not entry: return
name, run_fn, n_params = entry
print(f"🕺 {name}...")
try:
if n_params >= 3:
run_fn(self.motion, self.tts, self)
else:
run_fn(self.motion, self.tts)
except Exception as e:
print(f"{name} gesture error: {e}")
def stand_for_walk(self):
"""Re-stiffen and drop straight into walk-ready posture -
mirrors the startup sequence in run() (stiffen -> walk arms ->
relax arm stiffness -> moveInit) rather than
goToPosture("Stand"), which goes through its own separate
stand-up animation instead of just being walk-ready. Shared by
the [DEL] wake-up toggle and by any gesture that needs to
reset his stance afterward (e.g. after a Choregraphe behavior
that moved the legs/hips)."""
self._safe("Restiffen", lambda: self.motion.setStiffnesses("Body", 1.0))
self._safe("Walk arms", lambda: self.motion.setMoveArmsEnabled(True, True))
self._safe("Relax arms", lambda: self.motion.setStiffnesses(["RArm", "LArm"], 0.3))
self._safe("moveInit", self.motion.moveInit)
self.is_resting = False
def idle_crouch(self):
"""The relaxed sit/crouch posture NAO settles into on shutdown -
toggled on demand with the Delete key. rest() cuts stiffness to
the whole body (that's the point of "rest"), so simply calling
it again does nothing on the second press - we have to
explicitly re-stiffen ourselves via stand_for_walk()."""
if not self.is_resting:
print("😴 Resting (idle crouch)...")
self._safe("stopMove", self.motion.stopMove)
self._safe("rest", self.motion.rest)
self.is_resting = True
else:
print("🧍 Waking up...")
self.stand_for_walk()
def handle_head_movement(self):
if self.typing_mode or self.music_search_mode: return
keys = pygame.key.get_pressed()
now = time.time()
dt = now - self._last_head_update
self._last_head_update = now
step = self.HEAD_RATE * dt
changed = False
if keys[pygame.K_i]: self.head_pitch = max(self.head_pitch - step, -0.5); changed = True
if keys[pygame.K_k]: self.head_pitch = min(self.head_pitch + step, 0.5); changed = True
if keys[pygame.K_j]: self.head_yaw = min(self.head_yaw + step, 2.0); changed = True
if keys[pygame.K_l]: self.head_yaw = max(self.head_yaw - step, -2.0); changed = True
if changed:
self._safe("Head move", lambda: self.motion.setAngles(
["HeadYaw", "HeadPitch"], [self.head_yaw, self.head_pitch], 0.18))
def reset_head(self):
self.head_yaw = 0.0
self.head_pitch = 0.0
self._safe("Head reset", lambda: self.motion.setAngles(["HeadYaw", "HeadPitch"], [0.0, 0.0], 0.4))
def update_title(self):
mode = "[TYPING] " if self.typing_mode else ""
pygame.display.set_caption(f"NAO Teleop | {mode}Battery: {self.battery_level}% | Vol: {self.volume}% | Speed: {self.walk_speed:.2f}")
def change_volume(self, delta):
self.volume = int(min(100, max(0, self.volume + delta)))
if self.audio_device:
self._safe("Volume set", lambda: self.audio_device.setOutputVolume(self.volume))
print(f"🔊 NAO volume: {self.volume}%")
self.update_title()
def check_battery(self):
if not self.battery: return
now = time.time()
if now - self.last_batt_check > 10:
try:
self.battery_level = self.battery.getBatteryCharge()
except: pass
self.update_title()
self.last_batt_check = now
def _movement_loop(self):
"""Runs on its own thread for the app's lifetime, sending whatever
velocity the render loop most recently recorded in
self._desired_velocity. moveToward/stopMove are blocking network
RPCs - if one hangs during a WiFi hiccup, only this thread stalls;
the render loop keeps reading keys and updating the target
velocity so the very next send (as soon as the network unblocks)
reflects your latest input, not a stale one from before the drop.
"""
stable_config = [
["StepHeight", 0.02],
["MaxStepFrequency", 0.5],
["TorsoWy", 0.08],
["MaxStepX", 0.03]
]
while not self._movement_stop_event.is_set():
x, y, theta = self._desired_velocity
try:
if abs(x) > 0.05 or abs(y) > 0.05 or abs(theta) > 0.05:
self.motion.moveToward(x, y, theta, stable_config)
else:
self.motion.stopMove()
except Exception as e:
print(f"Move error: {e}")
time.sleep(0.05) # ~20Hz send rate
def run(self):
print("🤖 Enabling stiffness (no wakeUp stand-up)...")
self.motion.setStiffnesses("Body", 1.0)
# Enable walk arms BEFORE forcing stiffness
print("💪 Enabling walk arms...")
self.motion.setMoveArmsEnabled(True, True)
# Force relaxed stiffness AFTER engine is active
self.motion.setStiffnesses(["RArm", "LArm"], 0.3)
print("🦶 Initializing walk engine...")
self._safe("moveInit", self.motion.moveInit)
print("🦶 Walk engine ready")
self._movement_thread = threading.Thread(target=self._movement_loop, daemon=True)
self._movement_thread.start()
gesture_hint = " | ".join(
f"{pygame.key.name(key)}={name}" for key, (name, *_rest) in sorted(self.gestures.items())
)
print(f"\n🎮 Controls: WASD/Arrows=Walk | IJKL=Head | SPACE=Reset head | {gesture_hint} | "
f"DEL=Toggle crouch/stand | M=Music search | X=Stop | -/==Volume | T=TTS | V=Push-to-talk (hold) | ESC=Quit")
try:
while self.running:
try:
self.check_battery()
for event in pygame.event.get():
if event.type == pygame.QUIT:
self.running = False
elif event.type == pygame.KEYDOWN:
if self.music_search_mode:
if event.key == pygame.K_RETURN:
if self.music_results:
entry = self.music_results[self.selected_result]
threading.Thread(
target=self.music.play,
args=(entry['url'], entry.get('title', 'Unknown')),
daemon=True
).start()
self.music_search_mode = False
elif self.music_search_text.strip() and not self.music_search_pending:
query = self.music_search_text
self.music_search_pending = True
def do_search(q=query):
results = self.music.search(q)
self.music_results = results
self.selected_result = 0
self.music_search_pending = False
if not results:
self.music_search_mode = False
threading.Thread(target=do_search, daemon=True).start()
else:
self.music_search_mode = False
elif event.key == pygame.K_ESCAPE:
self.music_search_mode = False
self.music_search_text = ""
self.music_results = []
elif event.key == pygame.K_UP:
if self.music_results:
self.selected_result = (self.selected_result - 1) % len(self.music_results[:5])
elif event.key == pygame.K_DOWN:
if self.music_results:
self.selected_result = (self.selected_result + 1) % len(self.music_results[:5])
elif event.key == pygame.K_BACKSPACE:
self.music_search_text = self.music_search_text[:-1]
self.music_results = []
elif event.key in (pygame.K_1, pygame.K_2, pygame.K_3, pygame.K_4, pygame.K_5):
idx = int(event.unicode) - 1
if idx < len(self.music_results):
entry = self.music_results[idx]
threading.Thread(
target=self.music.play,
args=(entry['url'], entry.get('title', 'Unknown')),
daemon=True
).start()
self.music_search_mode = False
elif event.unicode.isprintable():
self.music_search_text += event.unicode
self.music_results = []
elif self.typing_mode:
if event.key == pygame.K_RETURN:
if self.chat_message.strip():
threading.Thread(target=self.tts.say, args=(self.chat_message,), daemon=True).start()
self.chat_message = ""
self.typing_mode = False
elif event.key == pygame.K_ESCAPE:
self.chat_message = ""
self.typing_mode = False
elif event.key == pygame.K_BACKSPACE:
self.chat_message = self.chat_message[:-1]
else:
self.chat_message += event.unicode
else:
if event.key == pygame.K_ESCAPE:
self.running = False
elif event.key == pygame.K_DELETE:
self.idle_crouch()
elif event.key in self.gestures:
self.run_gesture(event.key)
elif event.key == pygame.K_SPACE:
self.reset_head()
elif event.key == pygame.K_m:
self.music_search_mode = True
self.music_search_text = ""
self.music_results = []
elif event.key == pygame.K_x:
self.music.stop()
elif event.key in (pygame.K_MINUS, pygame.K_KP_MINUS):
self.change_volume(-5)
elif event.key in (pygame.K_EQUALS, pygame.K_KP_PLUS):
self.change_volume(5)
elif event.key == pygame.K_t:
self.typing_mode = True
self._safe("stopMove", self.motion.stopMove)
if not self.typing_mode and not self.music_search_mode:
keys = pygame.key.get_pressed()
if keys[pygame.K_v] and not self.ptt_active:
self.ptt_active = True
self.mic_relay.start()
if getattr(self, "sound_receiver", None):
self.sound_receiver.muted = True
elif not keys[pygame.K_v] and self.ptt_active:
self.ptt_active = False
self.mic_relay.stop()
if getattr(self, "sound_receiver", None):
self.sound_receiver.muted = False
x = y = theta = 0.0
if keys[pygame.K_w] or keys[pygame.K_UP]: x = self.walk_speed
if keys[pygame.K_s] or keys[pygame.K_DOWN]: x = -self.walk_speed
if keys[pygame.K_a] or keys[pygame.K_LEFT]: theta = self.walk_speed * 0.8
if keys[pygame.K_d] or keys[pygame.K_RIGHT]: theta = -self.walk_speed * 0.8
# Just record intent here - the movement thread
# (see _movement_loop) is what actually calls
# moveToward/stopMove, so a stalled network RPC
# can't stop us reading the next key state.
self._desired_velocity = (x, y, theta)
self.handle_head_movement()
frame = self.get_image()
if frame is not None:
up = cv2.resize(frame, (self.display_width, self.display_height), interpolation=cv2.INTER_CUBIC)
rgb = cv2.cvtColor(up, cv2.COLOR_BGR2RGB)
surf = pygame.surfarray.make_surface(rgb.swapaxes(0,1))
self.screen.blit(surf, (0, 0))
if time.time() - self._latest_frame_ts > 1.0:
warn = self.font.render("⚠ VIDEO STALE - reconnecting...", True, (255, 60, 60))
self.screen.blit(warn, (10, self.display_height - 30))
else:
self.screen.fill((10, 10, 30))
# Music status
if self.music.is_loading:
text = self.font.render(f"{self.music.status}", True, (255, 200, 0))
self.screen.blit(text, (10, 10))
if self.music.progress > 0:
bar_w = 220
pct = min(self.music.progress, 100) / 100.0
pygame.draw.rect(self.screen, (70, 70, 70), (10, 32, bar_w, 8))
pygame.draw.rect(self.screen, (255, 200, 0), (10, 32, int(bar_w * pct), 8))
elif self.music.load_error:
text = self.font.render(f"♪ Error: {self.music.load_error}", True, (255, 80, 80))
self.screen.blit(text, (10, 10))
elif self.music.current_title:
text = self.font.render(f"{self.music.current_title[:45]}", True, (0, 255, 100))
self.screen.blit(text, (10, 10))
# Voice chat HUD (top-right) - only shown while V is held / connecting
relay_state = getattr(self.mic_relay, "state", "idle")
if relay_state == "connecting":
label, color = "VOICE: CONNECTING...", (255, 200, 0)
elif relay_state == "live":
label, color = "VOICE: LIVE - speak now", (0, 255, 100)
else:
label, color = None, None
if label:
text = self.font.render(label, True, color)
pad = 8
box_w, box_h = text.get_width() + pad * 2 + 18, text.get_height() + pad
box_x = self.display_width - box_w - 10
box_y = 10
s = pygame.Surface((box_w, box_h))
s.set_alpha(180)
s.fill((0, 0, 0))
self.screen.blit(s, (box_x, box_y))
# pulsing dot for "live" so it's easy to catch out of the corner of your eye
if relay_state == "live":
pulse = 0.5 + 0.5 * abs((pygame.time.get_ticks() % 1000) / 500.0 - 1.0)
dot_color = tuple(int(c * pulse) for c in color)
else:
dot_color = color
pygame.draw.circle(self.screen, dot_color, (box_x + pad + 5, box_y + box_h // 2), 5)
self.screen.blit(text, (box_x + pad + 18, box_y + pad // 2))
# Music search interface
if self.music_search_mode:
box_h = min(self.display_height - 80, 260)
s = pygame.Surface((self.display_width, box_h))
s.set_alpha(220)
s.fill((0, 0, 40))
self.screen.blit(s, (0, 80))
text = self.font.render(f"Search: {self.music_search_text}", True, (255, 255, 255))
self.screen.blit(text, (20, 100))
if self.music_results:
shown = self.music_results[:5]
for i, entry in enumerate(shown):
row_y = 140 + i * 30
if i == self.selected_result:
pygame.draw.rect(self.screen, (90, 90, 0), (16, row_y - 3, self.display_width - 32, 26))
color = (255, 255, 0) if i == self.selected_result else (200, 200, 200)
title = entry.get('title', 'No title')[:60]
line = self.font.render(f"{i+1}. {title}", True, color)
self.screen.blit(line, (20, row_y))
hint_y = 140 + len(shown) * 30 + 6
hint = self.font.render("\u2191\u2193 Select Enter Confirm 1-5 Quick pick Esc Cancel", True, (150, 150, 150))
self.screen.blit(hint, (20, hint_y))
elif self.music_search_pending:
text = self.font.render("Searching...", True, (200, 200, 200))
self.screen.blit(text, (20, 140))
else:
text = self.font.render("No results. Type to search again, Enter to search", True, (200, 200, 200))
self.screen.blit(text, (20, 140))
# TTS box
if self.typing_mode:
s = pygame.Surface((self.display_width, 40))
s.set_alpha(180)
s.fill((0,0,0))
self.screen.blit(s, (0, self.display_height-40))
text = self.font.render(f"Say: {self.chat_message}", True, (255,255,255))
self.screen.blit(text, (10, self.display_height-30))
pygame.display.flip()
self.clock.tick(0)
except Exception as e:
print(f"\u26a0\ufe0f Frame error (continuing): {e}")
finally:
print("🛑 Shutting down...")
self._movement_stop_event.set()
if self._movement_thread:
self._movement_thread.join(timeout=1.0)
self._video_stop_event.set()
self._safe("stopMove", self.motion.stopMove)
self._safe("music stop", self.music.stop)
if getattr(self, "mic_relay", None):
self._safe("mic relay stop", self.mic_relay.stop)
self._safe("mic relay close", self.mic_relay.close)
if getattr(self, "video_server_launcher", None):
self._safe("video server stop", self.video_server_launcher.close)
if getattr(self, "sound_receiver", None):
self._safe("sound receiver close", self.sound_receiver.close)
if getattr(self, "_pyaudio", None):
self._safe("pyaudio terminate", self._pyaudio.terminate)
self._safe("rest", self.motion.rest)
pygame.quit()
if __name__ == "__main__":
parser = argparse.ArgumentParser()
parser.add_argument("--ip", type=str, default="127.0.0.1")
parser.add_argument("--port", type=int, default=9559)
parser.add_argument("--audio-device-index", type=int, default=None)
parser.add_argument("--mic-gain", type=float, default=8.0)
parser.add_argument("--video-scale", type=float, default=2.0)
parser.add_argument("--media-relay", type=str, default=None,
help="host:port to reach the on-robot media relay (nao_video_server.py) "
"at, e.g. 127.0.0.1:8000 when tunneled through the phone. Setting "
"this also makes naowalk auto-scp nao_video_server.py to the "
"robot's home dir over SSH and launch it (see --no-deploy-video-"
"server to skip that and connect to one you started yourself). "
"BOTH video and mic audio come compressed over this one TCP "
"connection (JPEG + mu-law) instead of raw frames/PCM over the "
"qi session. To change the mic channel, set MIC_CHANNEL in "
"nao_video_server.py, since the mic subscription lives there now.")
parser.add_argument("--nao-ssh-user", type=str, default="nao",
help="SSH login for the live mic relay (push-to-talk)")
parser.add_argument("--nao-ssh-password", type=str, default="nao",
help="SSH password for the live mic relay (push-to-talk)")
parser.add_argument("--nao-ssh-port", type=int, default=22,
help="SSH port for the robot (music upload + push-to-talk). "
"Set this to your local forwarded port, e.g. 2222, when tunneling.")
parser.add_argument("--relay-rate", type=int, default=16000,
help="Sample rate for the live mic relay - try 48000 if audio sounds sped up/slow")
parser.add_argument("--no-deploy-video-server", action="store_true",
help="Don't auto-scp/launch nao_video_server.py on the robot when "
"--media-relay is set - use this if you're already running it "
"yourself and just want naowalk to connect to it.")
parser.add_argument("--video-server-path", type=str, default=None,
help="Local path to nao_video_server.py to deploy (default: the "
"copy sitting next to naowalk.py)")
parser.add_argument("--gestures-dir", type=str, default=None,
help="Folder of number-key gestures to load, each file named "
"'<digit>_name.py' with a run(motion, tts) function "
"(default: the 'gestures' folder next to naowalk.py)")
args = parser.parse_args()
session = qi.Session()
try:
session.connect(f"tcp://{args.ip}:{args.port}")
print(f"✅ Connected to {args.ip}")
except Exception as e:
print(f"❌ Connection failed: {e}")
sys.exit(1)
NaoTeleop(session, nao_ip=args.ip,
audio_device_index=args.audio_device_index,
mic_gain=args.mic_gain,
video_scale=args.video_scale,
media_relay=args.media_relay,
nao_ssh_user=args.nao_ssh_user,
nao_ssh_password=args.nao_ssh_password,
nao_ssh_port=args.nao_ssh_port,
relay_rate=args.relay_rate,
deploy_video_server=not args.no_deploy_video_server,
video_server_path=args.video_server_path,
gestures_dir=args.gestures_dir).run()