EG/core/alvr_streamer.py
Rowland bb56813e44 feat(core): 优化 ALVR 串流功能并添加对 Pico 4 的支持
- 添加对有线连接的支持,特别是 Pico 4 的有线连接
- 改进 ALVR 服务器检测和连接逻辑
- 优化 VR 帧获取和处理流程
- 增加备用帧获取功能(主摄像机视图)
- 修复 VR 眼部纹理相关问题
- 优化 VR 任务管理,确保正确渲染和提交帧
2025-07-29 10:17:13 +08:00

683 lines
24 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""
ALVR串流处理器
负责与ALVR服务器通信和视频流传输
支持Quest等VR头显的无线串流
"""
import socket
import struct
import threading
import json
import time
import subprocess
import psutil
from direct.showbase.DirectObject import DirectObject
from panda3d.core import Texture, PNMImage
class ALVRStreamer(DirectObject):
"""ALVR串流处理器"""
def __init__(self, world, vr_manager):
super().__init__()
self.world = world
self.vr_manager = vr_manager
# ALVR服务器配置
self.alvr_server_ip = "127.0.0.1"
self.alvr_server_port = 9943
self.alvr_streaming_port = 9944
# 添加对有线连接的支持
self.connection_mode = "auto" # auto, wireless, wired
self.wired_port = 9945 # Pico 4有线连接可能使用的端口
# 连接状态
self.connected = False
self.streaming = False
self.server_socket = None
self.streaming_socket = None
# 流媒体配置
self.stream_width = 2880 # Quest 2 推荐分辨率
self.stream_height = 1700
self.stream_fps = 72
self.bitrate = 150 # Mbps
self.codec = "h264"
# 线程管理
self.connection_thread = None
self.streaming_thread = None
self.running = False
# 性能统计
self.frame_count = 0
self.last_fps_time = time.time()
self.current_fps = 0
self.latency = 0
print("✓ ALVR串流处理器初始化完成")
def initialize(self):
"""初始化ALVR串流"""
try:
# 检查ALVR服务器是否运行
if not self._check_alvr_server():
print("ALVR服务器未运行尝试启动...")
if not self._start_alvr_server():
print("无法启动ALVR服务器")
return False
# 连接到ALVR服务器
if not self._connect_to_server():
print("无法连接到ALVR服务器")
return False
# 配置流媒体设置
self._configure_streaming()
# 启动串流线程
self._start_streaming_threads()
print("✓ ALVR串流初始化成功")
return True
except Exception as e:
print(f"ALVR初始化错误: {str(e)}")
return False
def _check_alvr_server(self):
"""检查ALVR服务器是否运行"""
try:
# 检查进程
alvr_processes = []
for proc in psutil.process_iter(['pid', 'name', 'cmdline']):
try:
# 检查ALVR相关进程包括Pico 4支持
if 'alvr' in proc.info['name'].lower() or any('alvr' in (arg or '').lower() for arg in (proc.info['cmdline'] or [])):
alvr_processes.append(proc)
except (psutil.NoSuchProcess, psutil.AccessDenied, TypeError):
# 忽略无法访问的进程
continue
if alvr_processes:
print(f"✓ 检测到 {len(alvr_processes)} 个ALVR相关进程")
for proc in alvr_processes:
try:
print(f" - 进程: {proc.info['name']} (PID: {proc.info['pid']})")
except (psutil.NoSuchProcess, psutil.AccessDenied):
pass
return True
# 尝试连接端口
sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
sock.settimeout(2)
result = sock.connect_ex((self.alvr_server_ip, self.alvr_server_port))
sock.close()
if result == 0:
print("✓ 检测到ALVR服务器端口")
return True
else:
# 尝试其他可能的端口Pico 4可能使用不同的端口
alternative_ports = [self.wired_port, 8080, 8081, 8082, 9944, 9945]
for port in alternative_ports:
sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
sock.settimeout(1)
result = sock.connect_ex((self.alvr_server_ip, port))
sock.close()
if result == 0:
print(f"✓ 检测到ALVR服务器在端口 {port}")
self.alvr_server_port = port
return True
return False
except Exception as e:
print(f"检查ALVR服务器错误: {str(e)}")
return False
def _start_alvr_server(self):
"""启动ALVR服务器"""
try:
# 获取当前操作系统
import platform
system = platform.system().lower()
# 尝试启动ALVR服务器
# 这里我们只尝试检测ALVR是否正在运行而不是尝试启动它
print(" 检测到系统环境跳过自动启动ALVR服务器")
print(" 请确保ALVR服务器已在运行")
return False
except Exception as e:
print(f"启动ALVR服务器错误: {str(e)}")
return False
def _connect_to_server(self):
"""连接到ALVR服务器"""
try:
# 首先检查是否为Pico 4有线连接模式
# 如果VR系统已经通过SteamVR正常工作我们可以启用直连模式
if self._is_steamvr_active():
print("💡 检测到SteamVR正在运行启用直连模式")
self.connected = True
return True
# 尝试多种连接方式以支持Pico 4有线连接
connection_attempts = [
(self.alvr_server_ip, self.alvr_server_port), # 默认无线端口
(self.alvr_server_ip, self.wired_port), # Pico 4有线端口
(self.alvr_server_ip, 8081), # 替代端口1
(self.alvr_server_ip, 8082) # 替代端口2
]
for ip, port in connection_attempts:
try:
print(f"尝试连接到ALVR服务器 {ip}:{port}...")
# 创建TCP连接
self.server_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
self.server_socket.settimeout(5)
self.server_socket.connect((ip, port))
# 发送握手消息
handshake_data = {
"type": "handshake",
"client_name": "Panda3D_VR_Engine",
"version": "1.0",
"capabilities": {
"video": True,
"audio": True,
"tracking": True,
"haptics": True
}
}
self._send_message(handshake_data)
# 接收响应
response = self._receive_message()
if response and response.get("type") == "handshake_response":
self.connected = True
self.alvr_server_port = port # 记住成功连接的端口
print(f"✓ 已连接到ALVR服务器 {ip}:{port}")
return True
else:
print(f"✗ 从 {ip}:{port} 接收到无效响应")
except Exception as e:
print(f"连接到 {ip}:{port} 失败: {str(e)}")
if self.server_socket:
self.server_socket.close()
self.server_socket = None
continue
return False
except Exception as e:
print(f"连接ALVR服务器错误: {str(e)}")
return False
def _is_steamvr_active(self):
"""检查SteamVR是否正在运行"""
try:
# 检查是否有SteamVR相关进程
steamvr_processes = ["vrserver", "steamvr", "vrcompositor", "vrmonitor"]
for proc in psutil.process_iter(['pid', 'name']):
try:
proc_name = proc.info['name'].lower()
if any(sv_proc in proc_name for sv_proc in steamvr_processes):
print(f"✓ 检测到SteamVR进程: {proc_name}")
return True
except (psutil.NoSuchProcess, psutil.AccessDenied):
continue
return False
except Exception as e:
print(f"检查SteamVR状态错误: {str(e)}")
return False
def _send_message(self, data):
"""发送消息到ALVR服务器"""
try:
# 在直连模式下,不需要发送消息
if not self.server_socket and not self.connected:
return
message = json.dumps(data).encode('utf-8')
length = struct.pack('<I', len(message))
self.server_socket.send(length + message)
except Exception as e:
print(f"发送消息错误: {str(e)}")
def _receive_message(self):
"""从ALVR服务器接收消息"""
try:
# 在直连模式下,不需要接收消息
if not self.server_socket and self.connected:
return None
# 接收消息长度
length_data = self.server_socket.recv(4)
if not length_data:
return None
length = struct.unpack('<I', length_data)[0]
# 接收消息内容
message_data = b''
while len(message_data) < length:
chunk = self.server_socket.recv(length - len(message_data))
if not chunk:
return None
message_data += chunk
return json.loads(message_data.decode('utf-8'))
except Exception as e:
print(f"接收消息错误: {str(e)}")
return None
def _configure_streaming(self):
"""配置流媒体设置"""
try:
# 发送流媒体配置
config_data = {
"type": "stream_config",
"video": {
"width": self.stream_width,
"height": self.stream_height,
"fps": self.stream_fps,
"bitrate": self.bitrate,
"codec": self.codec
},
"audio": {
"enabled": True,
"sample_rate": 48000,
"channels": 2
}
}
self._send_message(config_data)
# 接收配置响应
response = self._receive_message()
if response and response.get("type") == "config_response":
if response.get("status") == "success":
print("✓ 流媒体配置成功")
return True
return False
except Exception as e:
print(f"配置流媒体错误: {str(e)}")
return False
def _start_streaming_threads(self):
"""启动串流线程"""
self.running = True
# 启动连接管理线程
self.connection_thread = threading.Thread(target=self._connection_handler)
self.connection_thread.daemon = True
self.connection_thread.start()
# 启动流媒体线程
self.streaming_thread = threading.Thread(target=self._streaming_handler)
self.streaming_thread.daemon = True
self.streaming_thread.start()
def _connection_handler(self):
"""连接处理线程"""
while self.running:
try:
if self.connected:
# 发送心跳
heartbeat = {"type": "heartbeat", "timestamp": time.time()}
self._send_message(heartbeat)
# 接收消息
response = self._receive_message()
if response:
self._handle_server_message(response)
time.sleep(0.1)
except Exception as e:
print(f"连接处理错误: {str(e)}")
self.connected = False
time.sleep(1)
def _streaming_handler(self):
"""流媒体处理线程"""
while self.running:
try:
if self.connected and self.streaming:
# 获取VR渲染帧
frame_data = self._get_vr_frame()
if frame_data:
# 发送帧数据
self._send_frame(frame_data)
# 更新性能统计
self._update_performance_stats()
time.sleep(1.0 / self.stream_fps)
except Exception as e:
print(f"流媒体处理错误: {str(e)}")
time.sleep(0.1)
def _handle_server_message(self, message):
"""处理服务器消息"""
msg_type = message.get("type")
if msg_type == "start_streaming":
self.streaming = True
print("✓ 开始VR串流")
elif msg_type == "stop_streaming":
self.streaming = False
print("✓ 停止VR串流")
elif msg_type == "client_connected":
print(f"✓ VR客户端已连接: {message.get('client_info', {})}")
elif msg_type == "client_disconnected":
print("✓ VR客户端已断开")
elif msg_type == "tracking_data":
self._handle_tracking_data(message.get("data"))
elif msg_type == "haptic_feedback":
self._handle_haptic_feedback(message.get("data"))
def _handle_tracking_data(self, tracking_data):
"""处理跟踪数据"""
if not tracking_data:
return
# 更新VR管理器的跟踪数据
# 这里可以处理从ALVR客户端发送的跟踪数据
pass
def _handle_haptic_feedback(self, haptic_data):
"""处理触觉反馈"""
if not haptic_data:
return
# 处理触觉反馈请求
# 这里可以控制VR控制器的震动等
pass
def _get_vr_frame(self):
"""获取VR渲染帧"""
try:
# 如果启用了直连模式通过SteamVR直接返回None
# 表示不需要通过ALVR传输帧SteamVR会直接处理
if self._is_steamvr_active():
return None
if not self.vr_manager.is_vr_enabled():
return None
# 获取左右眼纹理
left_texture = self.vr_manager.eye_textures.get('left')
right_texture = self.vr_manager.eye_textures.get('right')
if not left_texture or not right_texture:
print("⚠ VR眼部纹理不可用使用主摄像机视图")
# 如果VR纹理不可用使用主摄像机视图作为备选
return self._get_fallback_frame()
# 检查纹理是否有效
if not left_texture.has_ram_image() or not right_texture.has_ram_image():
# 强制将纹理数据复制到RAM
left_texture.set_auto_texture_scale(Texture.AT_scaleNone)
right_texture.set_auto_texture_scale(Texture.AT_scaleNone)
# 尝试重新获取纹理数据
try:
self.world.graphicsEngine.extract_texture_data(left_texture, self.world.win.get_gsg())
self.world.graphicsEngine.extract_texture_data(right_texture, self.world.win.get_gsg())
except Exception as e:
print(f"提取纹理数据失败: {str(e)}")
# 合成立体帧
frame_data = self._compose_stereo_frame(left_texture, right_texture)
return frame_data
except Exception as e:
print(f"获取VR帧错误: {str(e)}")
import traceback
traceback.print_exc()
return None
def _get_fallback_frame(self):
"""获取备用帧(主摄像机视图)"""
try:
# 使用主窗口的纹理作为备选
if hasattr(self.world, 'win') and self.world.win:
# 创建临时缓冲区来捕获主视图
temp_buffer = self.world.win.makeTextureBuffer(
"temp_vr_frame", self.stream_width, self.stream_height
)
temp_texture = temp_buffer.getTexture()
temp_camera = self.world.makeCamera(temp_buffer)
# 渲染一帧
self.world.graphicsEngine.render_frame()
# 获取图像数据
image = PNMImage()
if temp_texture.store(image):
frame_data = image.makeRamImage()
# 清理临时资源
temp_buffer.remove()
return frame_data
# 清理临时资源
temp_buffer.remove()
return None
except Exception as e:
print(f"获取备用帧错误: {str(e)}")
return None
def _compose_stereo_frame(self, left_texture, right_texture):
"""合成立体帧"""
try:
# 创建组合图像
combined_image = PNMImage(self.stream_width, self.stream_height)
# 获取左右眼图像
left_image = PNMImage()
right_image = PNMImage()
if left_texture.store(left_image) and right_texture.store(right_image):
# 将左右眼图像合并Side-by-Side布局
left_width = self.stream_width // 2
# 缩放左眼图像到左半部分
left_scaled = PNMImage(left_width, self.stream_height)
left_scaled.quickFilterFrom(left_image)
combined_image.copySubImage(left_scaled, 0, 0)
# 缩放右眼图像到右半部分
right_scaled = PNMImage(left_width, self.stream_height)
right_scaled.quickFilterFrom(right_image)
combined_image.copySubImage(right_scaled, left_width, 0)
# 转换为字节数据
return combined_image.makeRamImage()
return None
except Exception as e:
print(f"合成立体帧错误: {str(e)}")
return None
def _send_frame(self, frame_data):
"""发送帧数据"""
try:
# 如果是直连模式,不需要发送帧数据
if self._is_steamvr_active():
return
if not self.server_socket:
return
# 创建帧消息
frame_message = {
"type": "video_frame",
"timestamp": time.time(),
"width": self.stream_width,
"height": self.stream_height,
"format": "rgb"
}
# 对于二进制数据,我们将其作为单独的消息发送
if frame_data:
# 先发送元数据
self._send_message(frame_message)
# 然后发送二进制数据
try:
# 发送数据长度
data_length = len(frame_data)
length_header = struct.pack('<I', data_length)
self.server_socket.send(length_header)
# 发送实际数据
self.server_socket.send(frame_data)
except Exception as e:
print(f"发送帧数据错误: {str(e)}")
else:
# 如果没有帧数据,发送空帧消息
frame_message["empty"] = True
self._send_message(frame_message)
except Exception as e:
print(f"发送帧错误: {str(e)}")
def _update_performance_stats(self):
"""更新性能统计"""
self.frame_count += 1
current_time = time.time()
if current_time - self.last_fps_time >= 1.0:
self.current_fps = self.frame_count
self.frame_count = 0
self.last_fps_time = current_time
def start_streaming(self):
"""开始串流"""
# 如果是直连模式,直接返回成功
if self._is_steamvr_active():
self.streaming = True
print("✓ 启用直连模式SteamVR将直接处理渲染输出")
return True
if not self.connected:
print("未连接到ALVR服务器")
return False
start_message = {"type": "start_streaming"}
self._send_message(start_message)
return True
def stop_streaming(self):
"""停止串流"""
if not self.connected:
return
stop_message = {"type": "stop_streaming"}
self._send_message(stop_message)
self.streaming = False
def send_haptic_feedback(self, controller_id, duration, intensity):
"""发送触觉反馈"""
if not self.connected:
return
haptic_message = {
"type": "haptic_feedback",
"controller_id": controller_id,
"duration": duration,
"intensity": intensity
}
self._send_message(haptic_message)
def get_streaming_status(self):
"""获取串流状态"""
# 如果是直连模式,返回特殊的直连状态
if self._is_steamvr_active():
return {
"connected": True,
"streaming": self.streaming,
"fps": self.current_fps,
"latency": self.latency,
"resolution": f"{self.stream_width}x{self.stream_height}",
"bitrate": self.bitrate,
"mode": "direct" # 直连模式
}
return {
"connected": self.connected,
"streaming": self.streaming,
"fps": self.current_fps,
"latency": self.latency,
"resolution": f"{self.stream_width}x{self.stream_height}",
"bitrate": self.bitrate,
"mode": "alvr" # ALVR模式
}
def set_stream_quality(self, width, height, fps, bitrate):
"""设置串流质量"""
self.stream_width = width
self.stream_height = height
self.stream_fps = fps
self.bitrate = bitrate
# 如果正在串流,重新配置
if self.streaming:
self._configure_streaming()
def shutdown(self):
"""关闭串流"""
print("关闭ALVR串流...")
self.running = False
self.streaming = False
# 发送断开消息
if self.connected:
disconnect_message = {"type": "disconnect"}
self._send_message(disconnect_message)
# 关闭连接
if self.server_socket:
self.server_socket.close()
if self.streaming_socket:
self.streaming_socket.close()
# 等待线程结束
if self.connection_thread:
self.connection_thread.join(timeout=2)
if self.streaming_thread:
self.streaming_thread.join(timeout=2)
self.connected = False
print("✓ ALVR串流已关闭")
def is_connected(self):
"""检查是否连接"""
return self.connected
def is_streaming(self):
"""检查是否在串流"""
return self.streaming