using Dpz.Core.Service.Network.Models;

namespace Dpz.Core.Service.Network;

/// <summary>
/// 流式响应解析器
/// 负责将SSE格式的流式数据解析为完整的ChatCompletionResponse
/// </summary>
[Obsolete("use IOpenAiService")]
internal class StreamResponseParser(ILogger<ChatCompletionService> logger)
{
    private readonly JsonSerializerOptions _jsonOptions = new()
    {
        PropertyNameCaseInsensitive = true,
    };

    /// <summary>
    /// 解析SSE格式的流式响应行
    /// </summary>
    /// <param name="lines">响应行集合</param>
    /// <returns>解析后的ChatCompletionResponse,如果解析失败返回null</returns>
    public ChatCompletionResponse? Parse(string[] lines)
    {
        var accumulator = new StreamAccumulator();

        foreach (var line in lines)
        {
            if (!TryParseLine(line, accumulator)) { }
        }

        return accumulator.BuildFinalResponse();
    }

    /// <summary>
    /// 解析SSE格式的流式响应行,并通过回调推送增量内容
    /// </summary>
    /// <param name="lines">响应行集合</param>
    /// <param name="onDeltaContent">当收到新增内容时的回调,参数为新增的内容片段</param>
    /// <returns>解析后的ChatCompletionResponse,如果解析失败返回null</returns>
    public ChatCompletionResponse? ParseWithCallback(string[] lines, Action<string> onDeltaContent)
    {
        var accumulator = new StreamAccumulator();

        foreach (var line in lines)
        {
            if (!TryParseLineWithCallback(line, accumulator, onDeltaContent)) { }
        }

        return accumulator.BuildFinalResponse();
    }

    /// <summary>
    /// 解析SSE格式的流式响应行,并通过异步回调推送增量内容
    /// </summary>
    public async Task<ChatCompletionResponse?> ParseWithCallbackAsync(
        string[] lines,
        Func<string, Task> onDeltaContent
    )
    {
        var accumulator = new StreamAccumulator();

        foreach (var line in lines)
        {
            if (!await TryParseLineWithCallbackAsync(line, accumulator, onDeltaContent)) { }
        }

        return accumulator.BuildFinalResponse();
    }

    /// <summary>
    /// 尝试解析单行数据
    /// </summary>
    private bool TryParseLine(string line, StreamAccumulator accumulator)
    {
        if (string.IsNullOrWhiteSpace(line) || !line.StartsWith("data:"))
        {
            return false;
        }

        var jsonStr = line["data:".Length..].Trim();

        // SSE规范中的结束标记
        if (jsonStr == "[DONE]")
        {
            return false;
        }

        if (string.IsNullOrEmpty(jsonStr))
        {
            return false;
        }

        return TryDeserializeChunk(jsonStr, accumulator);
    }

    /// <summary>
    /// 尝试解析单行数据并推送增量内容
    /// </summary>
    private bool TryParseLineWithCallback(
        string line,
        StreamAccumulator accumulator,
        Action<string> onDeltaContent
    )
    {
        if (string.IsNullOrWhiteSpace(line) || !line.StartsWith("data:"))
        {
            return false;
        }

        var jsonStr = line["data:".Length..].Trim();

        // SSE规范中的结束标记
        if (jsonStr == "[DONE]")
        {
            return false;
        }

        if (string.IsNullOrEmpty(jsonStr))
        {
            return false;
        }

        return TryDeserializeChunk(jsonStr, accumulator, onDeltaContent);
    }

    /// <summary>
    /// 尝试解析单行数据并推送增量内容(异步版本)
    /// </summary>
    private async Task<bool> TryParseLineWithCallbackAsync(
        string line,
        StreamAccumulator accumulator,
        Func<string, Task> onDeltaContent
    )
    {
        if (string.IsNullOrWhiteSpace(line) || !line.StartsWith("data:"))
        {
            return false;
        }

        var jsonStr = line["data:".Length..].Trim();

        // SSE规范中的结束标记
        if (jsonStr == "[DONE]")
        {
            return false;
        }

        if (string.IsNullOrEmpty(jsonStr))
        {
            return false;
        }

        return await TryDeserializeChunkAsync(jsonStr, accumulator, onDeltaContent);
    }

    /// <summary>
    /// 尝试反序列化流式块
    /// </summary>
    private bool TryDeserializeChunk(
        string jsonStr,
        StreamAccumulator accumulator,
        Action<string>? onDeltaContent = null
    )
    {
        try
        {
            var chunk = JsonSerializer.Deserialize<StreamCompletionChunk>(jsonStr, _jsonOptions);
            if (chunk == null)
            {
                return false;
            }

            // 初始化响应(使用第一个块的元数据)
            if (!accumulator.IsInitialized)
            {
                accumulator.Initialize(chunk);
            }

            // 累积内容和使用情况,并提取增量内容
            var deltaContent = accumulator.AccumulateChunk(chunk);

            // 如果有新增内容且提供了回调,则调用回调
            if (!string.IsNullOrEmpty(deltaContent) && onDeltaContent != null)
            {
                onDeltaContent(deltaContent);
            }

            return true;
        }
        catch (JsonException ex)
        {
            var truncatedJson = jsonStr[..Math.Min(100, jsonStr.Length)];
            logger.LogWarning(ex, "Failed to deserialize stream chunk: {JsonStr}", truncatedJson);
            return false;
        }
    }

    /// <summary>
    /// 尝试反序列化流式块(异步版本)
    /// </summary>
    private async Task<bool> TryDeserializeChunkAsync(
        string jsonStr,
        StreamAccumulator accumulator,
        Func<string, Task> onDeltaContent
    )
    {
        try
        {
            var chunk = JsonSerializer.Deserialize<StreamCompletionChunk>(jsonStr, _jsonOptions);
            if (chunk == null)
            {
                return false;
            }

            // 初始化响应(使用第一个块的元数据)
            if (!accumulator.IsInitialized)
            {
                accumulator.Initialize(chunk);
            }

            // 累积内容和使用情况,并提取增量内容
            var deltaContent = accumulator.AccumulateChunk(chunk);

            // 如果有新增内容,则异步调用回调
            if (!string.IsNullOrEmpty(deltaContent))
            {
                await onDeltaContent(deltaContent);
            }

            return true;
        }
        catch (JsonException ex)
        {
            var truncatedJson = jsonStr[..Math.Min(100, jsonStr.Length)];
            logger.LogWarning(ex, "Failed to deserialize stream chunk: {JsonStr}", truncatedJson);
            return false;
        }
    }
}

/// <summary>
/// 流式响应累积器
/// 负责累积多个流式块的数据,最终组合为完整响应
/// </summary>
[Obsolete("use IOpenAiService")]
internal class StreamAccumulator
{
    private ChatCompletionResponse? _response;
    private readonly StringBuilder _contentBuilder = new();
    private readonly StringBuilder _reasoningBuilder = new();
    private bool _initialized;

    public bool IsInitialized => _initialized;

    /// <summary>
    /// 使用第一个块的元数据初始化响应
    /// </summary>
    public void Initialize(StreamCompletionChunk chunk)
    {
        _response = new ChatCompletionResponse
        {
            Id = chunk.Id,
            Object = chunk.Object,
            Created = chunk.Created,
            Model = chunk.Model,
            Provider = chunk.Provider,
            Choices = BuildInitialChoices(chunk),
            Usage = chunk.Usage,
        };
        _initialized = true;
    }

    /// <summary>
    /// 从流式块构建初始的选择项
    /// </summary>
    private static List<CompletionChoice> BuildInitialChoices(StreamCompletionChunk chunk)
    {
        return chunk
            .Choices.Select(c => new CompletionChoice
            {
                Index = c.Index,
                Message = new CompletionMessage
                {
                    Role = c.Delta.Role ?? "assistant",
                    Content = string.Empty,
                },
                LogProbs = c.LogProbs,
                FinishReason = c.FinishReason ?? string.Empty,
            })
            .ToList();
    }

    /// <summary>
    /// 累积单个流式块的数据,返回本次新增的内容
    /// </summary>
    /// <returns>本次累积产生的新增内容(delta)</returns>
    public string AccumulateChunk(StreamCompletionChunk chunk)
    {
        if (_response == null || _response.Choices.Count == 0)
        {
            return string.Empty;
        }

        if (chunk.Choices.Count == 0)
        {
            return string.Empty;
        }

        var choice = chunk.Choices[0];
        var deltaContent = AccumulateContent(choice);
        AccumulateFinishReason(choice);
        AccumulateUsage(chunk);

        return deltaContent;
    }

    /// <summary>
    /// 累积内容和推理数据,返回本次新增内容
    /// </summary>
    private string AccumulateContent(StreamCompletionChoice choice)
    {
        var delta = new StringBuilder();

        if (choice.Delta.Content != null)
        {
            _contentBuilder.Append(choice.Delta.Content);
            delta.Append(choice.Delta.Content);
        }

        if (choice.Delta.Reasoning != null)
        {
            _reasoningBuilder.Append(choice.Delta.Reasoning);
        }

        return delta.ToString();
    }

    /// <summary>
    /// 更新完成原因
    /// </summary>
    private void AccumulateFinishReason(StreamCompletionChoice choice)
    {
        if (!string.IsNullOrEmpty(choice.FinishReason) && _response != null)
        {
            _response.Choices[choice.Index].FinishReason = choice.FinishReason;
        }
    }

    /// <summary>
    /// 更新使用统计信息
    /// </summary>
    private void AccumulateUsage(StreamCompletionChunk chunk)
    {
        if (chunk.Usage != null && _response != null)
        {
            _response.Usage = chunk.Usage;
        }
    }

    /// <summary>
    /// 构建最终响应
    /// </summary>
    public ChatCompletionResponse? BuildFinalResponse()
    {
        if (_response?.Choices.Count > 0)
        {
            _response.Choices[0].Message.Content = _contentBuilder.ToString();
        }

        return _response;
    }
}
⚠⚠    以下内容为AI分析的结果,请根据实际情况进行判断。

这是一个用于解析流式响应的代码文件,主要包含两个核心类:StreamResponseParserStreamAccumulator。让我详细解释一下:

主要功能

这个代码是用来处理 SSE(Server-Sent Events)格式的流式 AI 对话响应,比如 OpenAI 的 ChatGPT 流式响应。它能够将分块传输的数据重新组装成完整的对话响应。

StreamResponseParser 类

这是主要的解析器类,负责:

核心方法:

  1. Parse(string[] lines) - 基础解析方法,将SSE格式的行数据解析为完整响应
  2. ParseWithCallback(string[] lines, Action<string> onDeltaContent) - 解析时提供同步回调,实时获取新增内容
  3. ParseWithCallbackAsync(...) - 异步版本的回调解析

解析流程:

// 1. 检查行格式是否以 "data:" 开头
if (!line.StartsWith("data:")) return false;

// 2. 提取JSON数据
var jsonStr = line["data:".Length..].Trim();

// 3. 检查结束标记
if (jsonStr == "[DONE]") return false;

// 4. 反序列化为流式块对象
var chunk = JsonSerializer.Deserialize<StreamCompletionChunk>(jsonStr);

StreamAccumulator 类

这是数据累积器,负责:

主要功能:

  1. 初始化响应结构 - 使用第一个数据块的元数据创建基础响应对象
  2. 累积内容 - 逐步收集每个块的内容片段
  3. 构建最终响应 - 将所有累积的数据组装成完整的 ChatCompletionResponse

工作原理:

// 初始化(第一次)
public void Initialize(StreamCompletionChunk chunk)
{
    _response = new ChatCompletionResponse
    {
        Id = chunk.Id,
        Model = chunk.Model,
        // ... 其他元数据
    };
}

// 累积每个块的数据
public string AccumulateChunk(StreamCompletionChunk chunk)
{
    // 累积文本内容
    if (choice.Delta.Content != null)
    {
        _contentBuilder.Append(choice.Delta.Content);
        delta.Append(choice.Delta.Content);
    }
    
    // 返回本次新增的内容
    return delta.ToString();
}

使用场景

这个解析器主要用于处理类似以下格式的流式数据:

data: {"id":"chatcmpl-123","choices":[{"delta":{"content":"Hello"}}]}
data: {"id":"chatcmpl-123","choices":[{"delta":{"content":" world"}}]}
data: {"id":"chatcmpl-123","choices":[{"delta":{"content":"!"}}]}
data: [DONE]

最终会组装成包含完整内容 "Hello world!" 的 ChatCompletionResponse 对象。

注意事项

  • 代码标记为 [Obsolete],表明这是一个过时的实现,建议使用 IOpenAiService
  • 包含错误处理机制,当JSON反序列化失败时会记录警告日志
  • 支持异步处理,可以在解析过程中实时推送增量内容给调用者

这种设计模式在处理大型语言模型的流式输出时非常有用,能够让用户看到实时的响应生成过程,而不是等待完整响应完成后才显示。

评论加载中...