using UnityEngine;
using BestHTTP;
using System.Text;
using System;
using System.Linq;
using System.Collections.Generic;

/// <summary>
/// 基于BestHTTP的流式请求处理组件
/// 功能:处理文本/音频流式数据,支持分片传输和超时检测
/// </summary>
public class BestHttpStream : MonoBehaviour
{
    //==================================================================
    // 配置区
    //==================================================================
    #region 配置参数

    private const float RequestTimeout = 30f;

    [Header("音频设置")]
    [Tooltip("音频包超时阈值(秒)")]
    public float audioTimeout = 0.5f;

    [Tooltip("最小有效音频包大小(字节)")]
    public int minAudioPacketSize = 128;

    [Tooltip("文本结束标记")]
    public string[] textEndMarkers = new[] { "。", "!", "?", "\n", "[END]" };

    #endregion

    //==================================================================
    // 运行时数据
    //==================================================================
    #region 运行时变量

    private ChunkedAudioPlayer _audioPlayer;
    private StringBuilder buffer = new StringBuilder();
    private Dictionary<string, StringBuilder> _textBuffers = new Dictionary<string, StringBuilder>();
    private Dictionary<string, List<byte[]>> _audioPackets = new Dictionary<string, List<byte[]>>();
    private Dictionary<string, float> _lastPacketTime = new Dictionary<string, float>();

    #endregion

    //==================================================================
    // 事件定义
    //==================================================================
    #region 事件系统

    public event Action<string, string> OnTextComplete;  // sessionId, content
    public event Action<string, byte[]> OnAudioComplete; // sessionId, pcmData
    public event Action<string, byte[]> OnAudioChunk;    // sessionId, chunkData

    #endregion

    //==================================================================
    // Unity生命周期方法
    //==================================================================
    #region Unity回调

    void Start()
    {
        InitializeAudioPlayer();
        // 示例请求(测试用)
        var requestData = CreateVoiceRequest("2", "介绍一下西安城墙");
        PostJsonStreamRequest(PathConfig.API_URL, requestData);
    }

    void Update()
    {
        CheckAudioTimeouts();
    }

    void OnDestroy()
    {
        // 清理事件订阅
        OnTextComplete = null;
        OnAudioComplete = null;
        OnAudioChunk = null;
    }

    #endregion

    //==================================================================
    // 核心请求方法
    //==================================================================
    #region 请求管理

    /// <summary>
    /// 创建语音请求数据结构
    /// </summary>
    private VoiceData CreateVoiceRequest(string type, string content)
    {
        return new VoiceData
        {
            question = content,
            userId = PathConfig.USER_ID,
            sessionId = "",
            type = type
        };
    }

    /// <summary>
    /// 发起JSON流式POST请求
    /// </summary>
    public void PostJsonStreamRequest(string url, object data)
    {
        try
        {
            if (string.IsNullOrEmpty(url) || !Uri.IsWellFormedUriString(url, UriKind.Absolute))
            {
                Debug.LogError($"无效的URL: {url}");
                return;
            }

            var request = new HTTPRequest(new Uri(url), HTTPMethods.Post, OnRequestFinished)
            {
                StreamChunksImmediately = true,
                RawData = Encoding.UTF8.GetBytes(JsonUtility.ToJson(data)),
                Timeout = TimeSpan.FromSeconds(RequestTimeout)
            };

            request.AddHeader("Content-Type", "application/json");
            request.AddHeader("Accept", "text/event-stream");

            request.OnStreamingData += OnStreamDataReceived;
            request.Send();
            Debug.Log($"流式请求已发送: {DateTime.Now}, URL: {url}");
        }
        catch (Exception ex)
        {
            Debug.LogError($"请求初始化异常: {ex.Message}");
        }
    }

    #endregion

    //==================================================================
    // 响应处理
    //==================================================================
    #region 响应处理

    /// <summary>
    /// 请求完成回调
    /// </summary>
    private void OnRequestFinished(HTTPRequest originalRequest, HTTPResponse response)
    {
        try
        {
            if (response == null)
            {
                Debug.LogError("请求失败:无响应");
                if (originalRequest.Exception != null)
                {
                    Debug.LogError($"底层异常: {originalRequest.Exception.Message}");
                }
                return;
            }

            Debug.Log($"[{DateTime.Now}] 请求完成,状态: {response.StatusCode}");

            if (!response.IsSuccess)
            {
                Debug.LogError($"请求失败: {response.Message} (状态码: {response.StatusCode})");
                if (!string.IsNullOrEmpty(response.DataAsText))
                {
                    Debug.LogError($"服务器错误信息: {response.DataAsText}");
                }
            }
        }
        catch (Exception ex)
        {
            Debug.LogError($"请求完成处理异常: {ex.Message}");
        }
    }

    /// <summary>
    /// 流式数据接收回调
    /// </summary>
    private bool OnStreamDataReceived(HTTPRequest originalRequest, HTTPResponse response, byte[] fragment, int dataFragmentLength)
    {
        try
        {
            if (fragment == null || dataFragmentLength <= 0)
            {
                Debug.LogWarning("接收到空数据片段");
                return true;
            }
            ProcessData(fragment, dataFragmentLength);
            return true;
        }
        catch (Exception ex)
        {
            Debug.LogError($"流数据处理异常: {ex.Message}\n{ex.StackTrace}");
            return false;
        }
    }

    #endregion

    //==================================================================
    // 数据处理流水线
    //==================================================================
    #region 数据处理

    /// <summary>
    /// 原始数据处理入口
    /// </summary>
    void ProcessData(byte[] fragment, int dataFragmentLength)
    {
        string fragmentText = Encoding.UTF8.GetString(fragment, 0, dataFragmentLength);
        buffer.Append(fragmentText);
        ProcessBuffer();
    }

    /// <summary>
    /// 缓冲区分行处理
    /// </summary>
    void ProcessBuffer()
    {
        while (true)
        {
            int newLineIndex = buffer.ToString().IndexOf('\n');
            if (newLineIndex == -1) break;

            string completeMessage = buffer.ToString(0, newLineIndex);
            buffer.Remove(0, newLineIndex + 1);

            ProcessStreamData(completeMessage);
        }
    }

    /// <summary>
    /// 流式数据内容处理
    /// </summary>
    private void ProcessStreamData(string rawData)
    {
        try
        {
            if (string.IsNullOrWhiteSpace(rawData))
                return;

            string[] lines = rawData.Split('\n');
            foreach (string line in lines.Where(l => !string.IsNullOrWhiteSpace(l)))
            {
                string trimmedLine = line.Trim();

                if (trimmedLine.StartsWith("data:"))
                {
                    string payload = trimmedLine.Substring(5).Trim();
                    HandleDataPayload(payload);
                }
            }
        }
        catch (Exception ex)
        {
            Debug.LogError($"数据内容处理异常: {ex.Message}");
        }
    }

    /// <summary>
    /// 有效负载分发处理
    /// </summary>
    private void HandleDataPayload(string payload)
    {
        try
        {
            //Debug.Log($"收到原始数据: {payload}");

            // 1. 优先尝试解析为 JSON
            if (payload.StartsWith("{") && payload.EndsWith("}"))
            {
                var response = JsonUtility.FromJson<BestResponseData>(payload);
                if (response != null && response.success)
                {
                    ProcessTypedData(response.data);
                    return;
                }
            }

            // 2. 如果不是 JSON,检查是否为裸 Base64 音频数据
            if (IsBase64String(payload))
            {
                HandleAudioPayload(payload);
                return;
            }

            // 3. 检查是否为十六进制二进制数据
            if (IsHexString(payload))
            {
                HandleBinaryPayload(payload);
                return;
            }

            // 4. 默认视为纯文本指令
            HandleTextCommand(payload);
        }
        catch (Exception ex)
        {
            Debug.LogError($"处理数据时发生异常: {ex.Message}\n{ex.StackTrace}");
        }
    }

    /// <summary>
    /// 类型化数据处理
    /// </summary>
    private void ProcessTypedData(ResponseData data)
    {
        switch (data.type)
        {
            case 1: // 文本消息
                HandleTextCommand(data.content);
                break;
            case 2: // 控制指令
                Debug.Log($"处理控制指令: {data.content}");
                break;
            case 3: // 音频(Base64)
                HandleAudioPayload(data.content);
                break;
            default:
                Debug.LogWarning($"未知的 type 类型: {data.type}");
                break;
        }
    }

    #endregion

    //==================================================================
    // 数据类型处理
    //==================================================================
    #region 数据类型处理

    /// <summary>
    /// 处理音频数据(Base64编码)
    /// </summary>
    private void HandleAudioPayload(string base64Payload)
    {
        try
        {
            byte[] audioData = Convert.FromBase64String(base64Payload);
            Debug.Log($"收到音频数据: {base64Payload}");

            if (audioData.Length % 2 == 0)
            {
                _audioPlayer.HandleAudioChunk(audioData);
            }
            else
            {
                Debug.LogWarning("非标准PCM音频数据,可能需要解码");
            }
        }
        catch (Exception ex)
        {
            Debug.LogError($"音频数据处理失败: {ex.Message}");
        }
    }

    /// <summary>
    /// 处理二进制数据(十六进制字符串)
    /// </summary>
    private void HandleBinaryPayload(string hexPayload)
    {
        try
        {
            byte[] binaryData = new byte[hexPayload.Length / 2];
            for (int i = 0; i < hexPayload.Length; i += 2)
            {
                binaryData[i / 2] = Convert.ToByte(hexPayload.Substring(i, 2), 16);
            }

            Debug.Log($"收到二进制数据: {binaryData.Length}字节");

            if (binaryData.Length > 4)
            {
                if (binaryData[0] == 0xFF && binaryData[1] == 0xD8)
                {
                    Debug.Log("检测到JPEG图像数据");
                }
                else if (Encoding.ASCII.GetString(binaryData, 0, 4) == "RIFF")
                {
                    Debug.Log("检测到WAV音频数据");
                    _audioPlayer.HandleAudioChunk(binaryData);
                }
            }
        }
        catch (Exception ex)
        {
            Debug.LogError($"二进制数据处理失败: {ex.Message}");
        }
    }

    /// <summary>
    /// 处理文本指令
    /// </summary>
    private void HandleTextCommand(string command)
    {
        command = command.Trim().ToUpper();
        Debug.Log($"收到文本指令: {command}");


    }

    #endregion

    //==================================================================
    // 音频管理
    //==================================================================
    #region 音频管理

    /// <summary>
    /// 初始化音频播放器
    /// </summary>
    private void InitializeAudioPlayer()
    {
        _audioPlayer = FindObjectOfType<ChunkedAudioPlayer>();
        if (_audioPlayer == null)
        {
            var audioPlayerObj = new GameObject("StreamingAudioPlayer");
            _audioPlayer = audioPlayerObj.AddComponent<ChunkedAudioPlayer>();
        }
    }

    /// <summary>
    /// 检查音频包超时
    /// </summary>
    private void CheckAudioTimeouts()
    {
        float now = Time.time;
        List<string> toRemove = new List<string>();

        foreach (var kv in _lastPacketTime)
        {
            if (now - kv.Value > audioTimeout)
            {
                toRemove.Add(kv.Key);
            }
        }

        foreach (var sessionId in toRemove)
        {
            FinalizeAudioSession(sessionId);
        }
    }

    /// <summary>
    /// 完成音频会话
    /// </summary>
    private void FinalizeAudioSession(string sessionId)
    {
        if (_audioPackets.TryGetValue(sessionId, out var packets))
        {
            byte[] fullAudio = CombinePackets(packets);
            OnAudioComplete?.Invoke(sessionId, fullAudio);

            _audioPackets.Remove(sessionId);
            _lastPacketTime.Remove(sessionId);
        }
    }

    /// <summary>
    /// 合并音频包
    /// </summary>
    private byte[] CombinePackets(List<byte[]> packets)
    {
        int totalSize = 0;
        foreach (var p in packets) totalSize += p.Length;

        byte[] result = new byte[totalSize];
        int offset = 0;

        foreach (var p in packets)
        {
            Buffer.BlockCopy(p, 0, result, offset, p.Length);
            offset += p.Length;
        }

        return result;
    }

    #endregion

    //==================================================================
    // 会话管理
    //==================================================================
    #region 会话管理

    /// <summary>
    /// 清理指定会话的数据
    /// </summary>
    public void ClearSession(string sessionId)
    {
        var textKeysToRemove = new List<string>();
        foreach (var kv in _textBuffers)
        {
            if (kv.Key.StartsWith(sessionId))
                textKeysToRemove.Add(kv.Key);
        }
        foreach (var key in textKeysToRemove)
            _textBuffers.Remove(key);

        _audioPackets.Remove(sessionId);
        _lastPacketTime.Remove(sessionId);
    }

    #endregion

    //==================================================================
    // 工具方法
    //==================================================================
    #region 工具方法

    /// <summary>
    /// 检查字符串是否为Base64编码
    /// </summary>
    private bool IsBase64String(string s)
    {
        if (string.IsNullOrEmpty(s) || s.Length % 4 != 0)
            return false;

        try
        {
            Convert.FromBase64String(s);
            return true;
        }
        catch
        {
            return false;
        }
    }

    /// <summary>
    /// 检查字符串是否为十六进制
    /// </summary>
    private bool IsHexString(string s)
    {
        if (string.IsNullOrEmpty(s))
            return false;

        foreach (char c in s)
        {
            if (!Uri.IsHexDigit(c))
                return false;
        }
        return true;
    }

    #endregion
}

Logo

魔乐社区(Modelers.cn) 是一个中立、公益的人工智能社区,提供人工智能工具、模型、数据的托管、展示与应用协同服务,为人工智能开发及爱好者搭建开放的学习交流平台。社区通过理事会方式运作,由全产业链共同建设、共同运营、共同享有,推动国产AI生态繁荣发展。

更多推荐