Files
naowalk/nao_video_server.py

434 lines
16 KiB
Python

# -*- coding: utf-8 -*-
"""
nao_video_server.py - RUNS ON THE ROBOT, under its own Python 2.7.
Opens a qi session to NAOqi over localhost (zero network hop - it's
the same box) and serves TWO persistent, length-framed TCP streams -
one for compressed video (JPEG), one for compressed audio (mu-law) -
each on its own port/connection. They used to share one connection,
but a single big video frame mid-send could block audio behind it for
its whole transfer time, which is a bad tradeoff for something as
latency-sensitive as audio. Splitting them means a slow video frame
can never delay audio.
The audio connection is bidirectional: this side keeps writing the
robot's own mic (TYPE_AUDIO) to the client, while a reader thread on
the same socket listens for TYPE_TALK chunks pushed the other way
(push-to-talk audio from the client's mic) and plays them on the
robot's speaker via TalkPlayer. Video stays one-way, server->client.
Why not HTTP: the phone tunnel only forwards specific TCP ports, so
this still rides over TCP - but a 5-byte binary header has none of
HTTP's chunked/multipart-boundary parsing overhead, and it's trivial
to keep alive and reconnect on a flaky link.
Why mic capture lives here now: it used to be shipped raw (16-bit PCM,
uncompressed) straight through NAOqi's own pub/sub over the qi RPC
session the client holds - i.e. across the tunnel, uncompressed, the
whole time. Subscribing to ALAudioDevice locally (same zero-hop trick
as the camera) means only mu-law-compressed bytes ever cross the
tunnel, half the size of raw PCM.
Wire format, identical on both connections (all integers big-endian,
via `struct`):
1 byte type 0x00 HELLO 0x01 VIDEO 0x02 AUDIO 0x03 TALK
4 bytes length
N bytes payload
HELLO payload: JSON, informational only
VIDEO payload: JPEG bytes (video connection, server->client only)
AUDIO payload: mu-law 8-bit samples, mono (audio connection, server->client - robot's mic)
TALK payload: mu-law 8-bit samples, mono (audio connection, client->server - push-to-talk)
Under load (slow/lossy tunnel), video quality and frame rate step
down automatically (see QUALITY_TIERS) and recover once sends are
fast again - and since encoding is throttled to match, this also cuts
CPU use on a congested link instead of just wasting it on frames that
would've been discarded anyway.
Needs the SDK's own env vars, since the bindings aren't on a login
shell's default path:
export PYTHONPATH=/opt/aldebaran/lib/python2.7/site-packages
export LD_LIBRARY_PATH=/opt/aldebaran/lib
python2 nao_video_server.py
Stdlib only (BaseHTTPServer/Flask are gone entirely now) plus qi,
cv2, and numpy, which are already confirmed present on the robot.
"""
import collections
import json
import socket
import struct
import subprocess
import threading
import time
import cv2
import numpy as np
import qi
# ---- tunables ----
NAOQI_PORT = 9561 # local qi session port (confirmed working)
VIDEO_PORT = 8000 # forward both of these through the phone tunnel
AUDIO_PORT = 8001 # kept on its own connection so a slow video
# frame can never block/delay audio delivery
CAMERA_RESOLUTION = 1 # 0=kQQVGA(160x120) 1=kQVGA(320x240) 2=kVGA(640x480)
CAMERA_COLORSPACE = 11 # kRGBColorSpace
CAMERA_FPS = 20
AUDIO_RATE = 16000
AUDIO_CHANNELS = 1
MIC_CHANNEL = 1 # 1=left 2=right 3=front 4=rear (ALAudioDevice channel index)
AUDIO_QUEUE_MAX = 60 # chunks buffered before we start dropping the oldest
# (jpeg_quality, min_seconds_between_frames) - best link first, worst last.
# The server steps down a tier when a video send is slow, and back up
# after a run of fast ones.
QUALITY_TIERS = [
(70, 0.05),
(55, 0.07),
(40, 0.10),
(25, 0.15),
]
SLOW_SEND_THRESHOLD = 0.20 # seconds - a send this slow means congestion
FAST_STREAK_TO_RECOVER = 30 # consecutive quick sends before stepping back up
TYPE_HELLO = 0x00
TYPE_VIDEO = 0x01
TYPE_AUDIO = 0x02
TYPE_TALK = 0x03
def encode_ulaw(pcm16):
"""int16 numpy array -> uint8 numpy array, mu-law companded.
Must match decode_ulaw() in naowalk.py exactly, or the far end
just hears noise - keep the two in sync if you touch this."""
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 decode_ulaw(q):
"""uint8 mu-law samples -> int16 PCM. Inverse of encode_ulaw()
above, used to play back incoming push-to-talk (TYPE_TALK) chunks.
Must match naowalk.py's encode_ulaw() exactly, or push-to-talk
just produces noise on the robot's speaker."""
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)
class FrameGrabber(object):
"""Background thread: keeps the single most recent JPEG-encoded
frame available, at whatever quality tier the link currently
supports. A stalled capture never blocks a client's read."""
def __init__(self, session):
self.video = session.service("ALVideoDevice")
self.client = self.video.subscribeCamera(
"MediaServer", 0, CAMERA_RESOLUTION, CAMERA_COLORSPACE, CAMERA_FPS)
self._lock = threading.Lock()
self._jpeg = None
self._tier = 0
self._fast_streak = 0
self._stop = threading.Event()
self._thread = threading.Thread(target=self._loop)
self._thread.daemon = True
self._thread.start()
def _loop(self):
last_encode = 0.0
while not self._stop.is_set():
try:
image = self.video.getImageRemote(self.client)
if not image or len(image) < 7:
continue
# getImageRemote() blocks until the next camera frame is
# ready (~1/CAMERA_FPS), so this loop is already paced -
# but under a congested tier we send far less often than
# that, so skip the (expensive) convert+encode work for
# frames we're just going to overwrite unsent anyway.
quality, min_interval = self.quality_and_interval()
now = time.time()
if now - last_encode < min_interval:
continue
w, h = image[0], image[1]
arr = np.frombuffer(bytearray(image[6]), dtype=np.uint8).reshape((h, w, 3))
bgr = cv2.cvtColor(arr, cv2.COLOR_RGB2BGR)
ok, buf = cv2.imencode(".jpg", bgr, [cv2.IMWRITE_JPEG_QUALITY, quality])
if ok:
with self._lock:
self._jpeg = buf.tobytes()
last_encode = now
except Exception as e:
print("frame grab error: %s" % e)
time.sleep(0.2)
def quality_and_interval(self):
with self._lock:
return QUALITY_TIERS[self._tier]
def report_send_time(self, elapsed):
"""Called by the writer thread after each video send so we can
adapt to how the link is actually behaving right now."""
with self._lock:
if elapsed > SLOW_SEND_THRESHOLD:
self._fast_streak = 0
if self._tier < len(QUALITY_TIERS) - 1:
self._tier += 1
print("link looks congested - dropping to quality tier %d" % self._tier)
else:
self._fast_streak += 1
if self._fast_streak >= FAST_STREAK_TO_RECOVER and self._tier > 0:
self._tier -= 1
self._fast_streak = 0
print("link recovered - raising to quality tier %d" % self._tier)
def latest(self):
with self._lock:
return self._jpeg
def close(self):
self._stop.set()
try:
self.video.unsubscribe(self.client)
except Exception:
pass
class MicGrabber(object):
"""Subscribes to the robot's own mic locally (zero-hop, same trick
as FrameGrabber) and mu-law-encodes each chunk as it arrives. This
replaces the old design where raw PCM was pushed to the client
directly through NAOqi's pub/sub over the (tunneled) qi session."""
def __init__(self, session, channel=MIC_CHANNEL):
self.audio = session.service("ALAudioDevice")
self._svc_name = "MicRelay"
self._lock = threading.Lock()
self._chunks = collections.deque(maxlen=AUDIO_QUEUE_MAX)
session.registerService(self._svc_name, self)
time.sleep(0.5) # let NAOqi's service directory propagate the
# registration before ALAudioDevice looks it up
# by name, or subscribe() fails to find it
self.audio.setClientPreferences(self._svc_name, AUDIO_RATE, channel, 0)
self.audio.subscribe(self._svc_name)
def processRemote(self, nbOfChannels, nbOfSamplesByChannel, timeStamp, buffer):
try:
pcm = np.frombuffer(bytes(buffer), dtype=np.int16)
chunk = encode_ulaw(pcm).tobytes()
with self._lock:
self._chunks.append(chunk)
except Exception as e:
print("mic encode error: %s" % e)
def pop_all(self):
"""Drain everything queued right now, oldest first."""
with self._lock:
out = list(self._chunks)
self._chunks.clear()
return out
def close(self):
try:
self.audio.unsubscribe(self._svc_name)
except Exception:
pass
class TalkPlayer(object):
"""Decodes incoming TYPE_TALK chunks (push-to-talk audio pushed
from the client) and plays them on the robot's speaker.
Uses a persistent `aplay` subprocess fed raw PCM over its stdin -
the same playback mechanism the old SSH `nc -l | aplay` relay
used, just sourced from decoded network chunks instead of a raw
pipe. The process is kept open across calls (silence between
presses is fine); it's only recreated if it exited or hasn't
been started yet."""
def __init__(self, rate=AUDIO_RATE):
self.rate = rate
self._proc = None
self._lock = threading.Lock()
def _ensure_proc(self):
if self._proc is None or self._proc.poll() is not None:
self._proc = subprocess.Popen(
["aplay", "-q", "-t", "raw", "-f", "S16_LE",
"-c", str(AUDIO_CHANNELS), "-r", str(self.rate)],
stdin=subprocess.PIPE)
def play(self, ulaw_payload):
pcm = decode_ulaw(np.frombuffer(ulaw_payload, dtype=np.uint8))
with self._lock:
try:
self._ensure_proc()
self._proc.stdin.write(pcm.tobytes())
self._proc.stdin.flush()
except Exception as e:
print("talk playback error: %s" % e)
self._proc = None
def close(self):
with self._lock:
if self._proc is not None:
try:
self._proc.stdin.close()
self._proc.terminate()
except Exception:
pass
self._proc = None
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).
Mirror of naowalk.py's _recv_exact()."""
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(conn, msg_type, payload):
conn.sendall(struct.pack(">BI", msg_type, len(payload)))
conn.sendall(payload)
def serve_video(grabber, host, port):
"""Video gets its own connection/thread - a slow sendall() here
(a big JPEG frame over a bad link) must never be able to delay
audio, which used to share this same socket and paid for it in
latency."""
sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
sock.bind((host, port))
sock.listen(1)
print("video server listening on :%d" % port)
while True:
conn, addr = sock.accept()
conn.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1)
conn.settimeout(5.0)
print("video client connected: %s" % (addr,))
last_jpeg = None
last_send = 0.0
try:
_send(conn, TYPE_HELLO, json.dumps({"stream": "video"}).encode("utf-8"))
while True:
jpeg = grabber.latest()
_, min_interval = grabber.quality_and_interval()
now = time.time()
if jpeg is not None and jpeg is not last_jpeg and now - last_send >= min_interval:
t0 = time.time()
_send(conn, TYPE_VIDEO, jpeg)
grabber.report_send_time(time.time() - t0)
last_jpeg = jpeg
last_send = now
else:
time.sleep(0.01)
except Exception as e:
print("video client disconnected (%s)" % e)
finally:
try:
conn.close()
except Exception:
pass
def _talk_reader(conn, player):
"""Runs alongside the audio writer loop on the same connection -
the audio socket is bidirectional now, so this thread's the read
half: it pulls TYPE_TALK frames pushed from the client (push-to-
talk) and hands each decoded chunk to TalkPlayer. conn has a
socket timeout set by the caller, so a quiet stretch between
button presses just means occasional socket.timeout - not a
disconnect - and the loop keeps going."""
try:
while True:
try:
header = _recv_exact(conn, 5)
except socket.timeout:
continue
msg_type, length = struct.unpack(">BI", header)
payload = _recv_exact(conn, length)
if msg_type == TYPE_TALK:
player.play(payload)
except Exception as e:
print("talk reader stopped (%s)" % e)
def serve_audio(mic, player, host, port):
"""Audio's own connection/thread. Checked eagerly (short sleep,
not paced like video) since gaps are far more noticeable than a
video frame being a beat stale."""
sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
sock.bind((host, port))
sock.listen(1)
print("audio server listening on :%d" % port)
while True:
conn, addr = sock.accept()
conn.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1)
conn.settimeout(5.0)
print("audio client connected: %s" % (addr,))
reader_thread = threading.Thread(target=_talk_reader, args=(conn, player))
reader_thread.daemon = True
reader_thread.start()
try:
_send(conn, TYPE_HELLO,
json.dumps({"rate": AUDIO_RATE, "channels": AUDIO_CHANNELS}).encode("utf-8"))
while True:
chunks = mic.pop_all() if mic is not None else []
if chunks:
for chunk in chunks:
_send(conn, TYPE_AUDIO, chunk)
else:
time.sleep(0.005)
except Exception as e:
print("audio client disconnected (%s)" % e)
finally:
try:
conn.close()
except Exception:
pass
def main():
session = qi.Session()
session.connect("tcp://127.0.0.1:%d" % NAOQI_PORT)
print("connected to local NAOqi on port %d" % NAOQI_PORT)
grabber = FrameGrabber(session)
mic = None
try:
mic = MicGrabber(session)
print("mic capture ready (mu-law, %dHz)" % AUDIO_RATE)
except Exception as e:
print("mic capture unavailable (%s) - video only" % e)
video_thread = threading.Thread(target=serve_video, args=(grabber, "0.0.0.0", VIDEO_PORT))
video_thread.daemon = True
video_thread.start()
player = TalkPlayer()
try:
serve_audio(mic, player, "0.0.0.0", AUDIO_PORT)
except KeyboardInterrupt:
pass
finally:
grabber.close()
if mic:
mic.close()
player.close()
if __name__ == "__main__":
main()