openai/openai-dotnet

Public

mirrored from https://github.com/openai/openai-dotnetAvailable

CodeCommitsIssuesPull requestsActionsInsightsSecurity
OpenAI_2.3.0

Branches

Tags

  • No tags available.
0Branches0Tags
Go to file
Add file
Code

Clone

HTTPS

Download ZIP

src/Utility/SseUpdateCollection.cs

203lines · modeblame

58f93c8dShivangiReja2 years ago1using System;
9f9f2936Jose Arriaga Maldonado2 years ago2using System.ClientModel;
3using System.ClientModel.Primitives;
4using System.Collections;
5using System.Collections.Generic;
674e0f77Stephen Toub2 years ago6using System.Net.ServerSentEvents;
9f9f2936Jose Arriaga Maldonado2 years ago7using System.Text.Json;
2ab1a942Jose Arriaga Maldonado1 years ago8using System.Threading;
9f9f2936Jose Arriaga Maldonado2 years ago9
10#nullable enable
11
0ca4c062Jose Arriaga Maldonado1 years ago12namespace OpenAI;
9f9f2936Jose Arriaga Maldonado2 years ago13
14/// <summary>
0ca4c062Jose Arriaga Maldonado1 years ago15/// Implementation of collection abstraction over streaming updates.
9f9f2936Jose Arriaga Maldonado2 years ago16/// </summary>
0ca4c062Jose Arriaga Maldonado1 years ago17internal class SseUpdateCollection<T> : CollectionResult<T>
9f9f2936Jose Arriaga Maldonado2 years ago18{
0ca4c062Jose Arriaga Maldonado1 years ago19private readonly Func<ClientResult> _sendRequestFunc;
20private readonly Func<SseItem<byte[]>, IEnumerable<T>> _eventDeserializerFunc;
2ab1a942Jose Arriaga Maldonado1 years ago21private readonly CancellationToken _cancellationToken;
9f9f2936Jose Arriaga Maldonado2 years ago22
5dce104aJose Arriaga Maldonado1 years ago23public List<Action> AdditionalDisposalActions { get; } = [];
24
0ca4c062Jose Arriaga Maldonado1 years ago25public SseUpdateCollection(
26Func<ClientResult> sendRequestFunc,
27Func<JsonElement, ModelReaderWriterOptions, IEnumerable<T>> jsonMultiDeserializerFunc,
28CancellationToken cancellationToken)
29: this(
30sendRequestFunc,
31AsyncSseUpdateCollection<T>.DeserializeSseToMultipleViaJson(jsonMultiDeserializerFunc),
32cancellationToken)
33
34{
35Argument.AssertNotNull(jsonMultiDeserializerFunc, nameof(jsonMultiDeserializerFunc));
36}
37
38public SseUpdateCollection(
39Func<ClientResult> sendRequestFunc,
40Func<JsonElement, ModelReaderWriterOptions, T> jsonSingleDeserializerFunc,
41CancellationToken cancellationToken)
42: this(
43sendRequestFunc,
44AsyncSseUpdateCollection<T>.DeserializeSseToSingleViaJson(jsonSingleDeserializerFunc),
45cancellationToken)
46{
47Argument.AssertNotNull(jsonSingleDeserializerFunc, nameof(jsonSingleDeserializerFunc));
48}
49
50public SseUpdateCollection(
51Func<ClientResult> sendRequestFunc,
52Func<SseItem<byte[]>, IEnumerable<T>> eventDeserializerFunc,
2ab1a942Jose Arriaga Maldonado1 years ago53CancellationToken cancellationToken)
9f9f2936Jose Arriaga Maldonado2 years ago54{
0ca4c062Jose Arriaga Maldonado1 years ago55Argument.AssertNotNull(sendRequestFunc, nameof(sendRequestFunc));
56Argument.AssertNotNull(eventDeserializerFunc, nameof(eventDeserializerFunc));
9f9f2936Jose Arriaga Maldonado2 years ago57
0ca4c062Jose Arriaga Maldonado1 years ago58_sendRequestFunc = sendRequestFunc;
59_eventDeserializerFunc = eventDeserializerFunc;
2ab1a942Jose Arriaga Maldonado1 years ago60_cancellationToken = cancellationToken;
9f9f2936Jose Arriaga Maldonado2 years ago61}
62
2ab1a942Jose Arriaga Maldonado1 years ago63public override ContinuationToken? GetContinuationToken(ClientResult page)
64// Continuation is not supported for SSE streams.
65=> null;
66
67public override IEnumerable<ClientResult> GetRawPages()
9f9f2936Jose Arriaga Maldonado2 years ago68{
2ab1a942Jose Arriaga Maldonado1 years ago69// We don't currently support resuming a dropped connection from the
70// last received event, so the response collection has a single element.
0ca4c062Jose Arriaga Maldonado1 years ago71yield return _sendRequestFunc();
2ab1a942Jose Arriaga Maldonado1 years ago72}
73
0ca4c062Jose Arriaga Maldonado1 years ago74protected override IEnumerable<T> GetValuesFromPage(ClientResult page)
2ab1a942Jose Arriaga Maldonado1 years ago75{
5dce104aJose Arriaga Maldonado1 years ago76using IEnumerator<T> enumerator = new SseUpdateEnumerator<T>(_eventDeserializerFunc, page, _cancellationToken, AdditionalDisposalActions);
2ab1a942Jose Arriaga Maldonado1 years ago77while (enumerator.MoveNext())
78{
79yield return enumerator.Current;
80}
9f9f2936Jose Arriaga Maldonado2 years ago81}
82
0ca4c062Jose Arriaga Maldonado1 years ago83private sealed class SseUpdateEnumerator<U> : IEnumerator<U>
9f9f2936Jose Arriaga Maldonado2 years ago84{
674e0f77Stephen Toub2 years ago85private static ReadOnlySpan<byte> TerminalData => "[DONE]"u8;
9f9f2936Jose Arriaga Maldonado2 years ago86
5dce104aJose Arriaga Maldonado1 years ago87private List<Action> _additionalDisposalActions;
88
2ab1a942Jose Arriaga Maldonado1 years ago89private readonly CancellationToken _cancellationToken;
90private readonly PipelineResponse _response;
9f9f2936Jose Arriaga Maldonado2 years ago91
92// These enumerators represent what is effectively a doubly-nested
93// loop over the outer event collection and the inner update collection,
94// i.e.:
95// foreach (var sse in _events) {
96// // get _updates from sse event
97// foreach (var update in _updates) { ... }
98// }
674e0f77Stephen Toub2 years ago99private IEnumerator<SseItem<byte[]>>? _events;
0ca4c062Jose Arriaga Maldonado1 years ago100private IEnumerator<U>? _updates;
101private readonly Func<SseItem<byte[]>, IEnumerable<U>> _eventDeserializerFunc;
9f9f2936Jose Arriaga Maldonado2 years ago102
0ca4c062Jose Arriaga Maldonado1 years ago103private U? _current;
9f9f2936Jose Arriaga Maldonado2 years ago104private bool _started;
105
0ca4c062Jose Arriaga Maldonado1 years ago106public SseUpdateEnumerator(
107Func<SseItem<byte[]>, IEnumerable<U>> eventDeserializerFunc,
108ClientResult page,
5dce104aJose Arriaga Maldonado1 years ago109CancellationToken cancellationToken,
110List<Action> additionalDisposalActions)
9f9f2936Jose Arriaga Maldonado2 years ago111{
0ca4c062Jose Arriaga Maldonado1 years ago112Argument.AssertNotNull(eventDeserializerFunc, nameof(eventDeserializerFunc));
2ab1a942Jose Arriaga Maldonado1 years ago113Argument.AssertNotNull(page, nameof(page));
9f9f2936Jose Arriaga Maldonado2 years ago114
0ca4c062Jose Arriaga Maldonado1 years ago115_eventDeserializerFunc = eventDeserializerFunc;
2ab1a942Jose Arriaga Maldonado1 years ago116_response = page.GetRawResponse();
117_cancellationToken = cancellationToken;
5dce104aJose Arriaga Maldonado1 years ago118_additionalDisposalActions = additionalDisposalActions;
9f9f2936Jose Arriaga Maldonado2 years ago119}
120
0ca4c062Jose Arriaga Maldonado1 years ago121U IEnumerator<U>.Current => _current!;
9f9f2936Jose Arriaga Maldonado2 years ago122
2ab1a942Jose Arriaga Maldonado1 years ago123object IEnumerator.Current => _current!;
9f9f2936Jose Arriaga Maldonado2 years ago124
125public bool MoveNext()
126{
127if (_events is null && _started)
128{
0ca4c062Jose Arriaga Maldonado1 years ago129throw new ObjectDisposedException(typeof(U).Name);
9f9f2936Jose Arriaga Maldonado2 years ago130}
131
2ab1a942Jose Arriaga Maldonado1 years ago132_cancellationToken.ThrowIfCancellationRequested();
9f9f2936Jose Arriaga Maldonado2 years ago133_events ??= CreateEventEnumerator();
134_started = true;
135
136if (_updates is not null && _updates.MoveNext())
137{
138_current = _updates.Current;
139return true;
140}
141
142if (_events.MoveNext())
143{
674e0f77Stephen Toub2 years ago144if (_events.Current.Data.AsSpan().SequenceEqual(TerminalData))
9f9f2936Jose Arriaga Maldonado2 years ago145{
146_current = default;
147return false;
148}
149
0ca4c062Jose Arriaga Maldonado1 years ago150_updates = _eventDeserializerFunc.Invoke(_events.Current).GetEnumerator();
9f9f2936Jose Arriaga Maldonado2 years ago151
152if (_updates.MoveNext())
153{
154_current = _updates.Current;
155return true;
156}
157}
158
159_current = default;
160return false;
161}
162
674e0f77Stephen Toub2 years ago163private IEnumerator<SseItem<byte[]>> CreateEventEnumerator()
9f9f2936Jose Arriaga Maldonado2 years ago164{
2ab1a942Jose Arriaga Maldonado1 years ago165if (_response.ContentStream is null)
9f9f2936Jose Arriaga Maldonado2 years ago166{
167throw new InvalidOperationException("Unable to create result from response with null ContentStream");
168}
169
2ab1a942Jose Arriaga Maldonado1 years ago170IEnumerable<SseItem<byte[]>> enumerable = SseParser.Create(_response.ContentStream, (_, bytes) => bytes.ToArray()).Enumerate();
9f9f2936Jose Arriaga Maldonado2 years ago171return enumerable.GetEnumerator();
172}
173
174public void Reset()
175{
176throw new NotSupportedException("Cannot seek back in an SSE stream.");
177}
178
179public void Dispose()
180{
181Dispose(true);
182GC.SuppressFinalize(this);
183}
184
185private void Dispose(bool disposing)
186{
187if (disposing && _events is not null)
188{
189_events.Dispose();
190_events = null;
191
2ab1a942Jose Arriaga Maldonado1 years ago192// Dispose the response so we don't leave the network connection open.
193_response?.Dispose();
9f9f2936Jose Arriaga Maldonado2 years ago194}
5dce104aJose Arriaga Maldonado1 years ago195
196foreach (Action additionalDisposalAction in _additionalDisposalActions ?? [])
197{
198additionalDisposalAction?.Invoke();
199}
200_additionalDisposalActions?.Clear();
9f9f2936Jose Arriaga Maldonado2 years ago201}
202}
203}