|
11 | 11 | // See the License for the specific language governing permissions and |
12 | 12 | // limitations under the License. |
13 | 13 |
|
| 14 | +using System.Runtime.CompilerServices; |
| 15 | + |
14 | 16 | namespace Synapse.Api.Client.Services; |
15 | 17 |
|
16 | 18 | /// <summary> |
@@ -105,60 +107,89 @@ public virtual async Task<IAsyncEnumerable<TResource>> ListAsync(IEnumerable<Lab |
105 | 107 | } |
106 | 108 |
|
107 | 109 | /// <inheritdoc/> |
108 | | - public virtual async Task<IAsyncEnumerable<IResourceWatchEvent<TResource>>> WatchAsync(string? @namespace = null, IEnumerable<LabelSelector>? labelSelectors = null, CancellationToken cancellationToken = default) |
| 110 | + public virtual async IAsyncEnumerable<IResourceWatchEvent<TResource>> WatchAsync(string? @namespace = null, IEnumerable<LabelSelector>? labelSelectors = null, [EnumeratorCancellation] CancellationToken cancellationToken = default) |
109 | 111 | { |
110 | 112 | var resource = new TResource(); |
111 | | - var uri = string.IsNullOrWhiteSpace(@namespace) ? $"/api/{resource.Definition.Version}/{resource.Definition.Plural}/watch" : $"/api/{resource.Definition.Version}/{resource.Definition.Plural}/{@namespace}/watch"; |
| 113 | + var uri = string.IsNullOrWhiteSpace(@namespace) ? $"/api/{resource.Definition.Version}/{resource.Definition.Plural}/watch/sse" : $"/api/{resource.Definition.Version}/{resource.Definition.Plural}/{@namespace}/watch"; |
112 | 114 | var queryStringArguments = new Dictionary<string, string>(); |
113 | 115 | if (labelSelectors?.Any() == true) queryStringArguments.Add("labelSelector", labelSelectors.Select(s => s.ToString()).Join(',')); |
114 | 116 | if (queryStringArguments.Count != 0) uri += $"?{queryStringArguments.Select(kvp => $"{kvp.Key}={kvp.Value}").Join('&')}"; |
115 | 117 | using var request = await this.ProcessRequestAsync(new HttpRequestMessage(HttpMethod.Get, uri), cancellationToken).ConfigureAwait(false); |
116 | | - request.EnableWebAssemblyStreamingResponse(); |
117 | 118 | var response = await this.HttpClient.SendAsync(request, HttpCompletionOption.ResponseHeadersRead, cancellationToken); |
118 | 119 | var responseStream = await response.Content.ReadAsStreamAsync(cancellationToken).ConfigureAwait(false); |
119 | | - return this.JsonSerializer.DeserializeAsyncEnumerable<ResourceWatchEvent<TResource>>(responseStream, cancellationToken)!; |
| 120 | + using var streamReader = new StreamReader(await response.Content.ReadAsStreamAsync(cancellationToken).ConfigureAwait(false)); |
| 121 | + while (!streamReader.EndOfStream) |
| 122 | + { |
| 123 | + var sseMessage = await streamReader.ReadLineAsync(cancellationToken).ConfigureAwait(false); |
| 124 | + if (string.IsNullOrWhiteSpace(sseMessage)) continue; |
| 125 | + var json = sseMessage["data: ".Length..].Trim(); |
| 126 | + var e = JsonSerializer.Deserialize<ResourceWatchEvent<TResource>>(json)!; |
| 127 | + yield return e; |
| 128 | + } |
120 | 129 | } |
121 | 130 |
|
122 | 131 | /// <inheritdoc/> |
123 | | - public virtual async Task<IAsyncEnumerable<IResourceWatchEvent<TResource>>> WatchAsync(IEnumerable<LabelSelector>? labelSelectors = null, CancellationToken cancellationToken = default) |
| 132 | + public virtual async IAsyncEnumerable<IResourceWatchEvent<TResource>> WatchAsync(IEnumerable<LabelSelector>? labelSelectors = null, [EnumeratorCancellation] CancellationToken cancellationToken = default) |
124 | 133 | { |
125 | 134 | var resource = new TResource(); |
126 | | - var uri = $"/api/{resource.Definition.Version}/{resource.Definition.Plural}/watch"; |
| 135 | + var uri = $"/api/{resource.Definition.Version}/{resource.Definition.Plural}/watch/sse"; |
127 | 136 | var queryStringArguments = new Dictionary<string, string>(); |
128 | 137 | if (labelSelectors?.Any() == true) queryStringArguments.Add("labelSelector", labelSelectors.Select(s => s.ToString()).Join(',')); |
129 | 138 | if (queryStringArguments.Count != 0) uri += $"?{queryStringArguments.Select(kvp => $"{kvp.Key}={kvp.Value}").Join('&')}"; |
130 | 139 | using var request = await this.ProcessRequestAsync(new HttpRequestMessage(HttpMethod.Get, uri), cancellationToken).ConfigureAwait(false); |
131 | 140 | request.EnableWebAssemblyStreamingResponse(); |
132 | 141 | var response = await this.HttpClient.SendAsync(request, HttpCompletionOption.ResponseHeadersRead, cancellationToken).ConfigureAwait(false); |
133 | 142 | var responseStream = await response.Content.ReadAsStreamAsync(cancellationToken).ConfigureAwait(false); |
134 | | - return this.JsonSerializer.DeserializeAsyncEnumerable<ResourceWatchEvent<TResource>>(responseStream, cancellationToken)!; |
| 143 | + using var streamReader = new StreamReader(await response.Content.ReadAsStreamAsync(cancellationToken).ConfigureAwait(false)); |
| 144 | + while (!streamReader.EndOfStream) |
| 145 | + { |
| 146 | + var sseMessage = await streamReader.ReadLineAsync(cancellationToken).ConfigureAwait(false); |
| 147 | + if (string.IsNullOrWhiteSpace(sseMessage)) continue; |
| 148 | + var json = sseMessage["data: ".Length..].Trim(); |
| 149 | + var e = JsonSerializer.Deserialize<ResourceWatchEvent<TResource>>(json)!; |
| 150 | + yield return e; |
| 151 | + } |
135 | 152 | } |
136 | 153 |
|
137 | 154 | /// <inheritdoc/> |
138 | | - public virtual async Task<IAsyncEnumerable<IResourceWatchEvent<TResource>>> MonitorAsync(string name, string @namespace, CancellationToken cancellationToken = default) |
| 155 | + public virtual async IAsyncEnumerable<IResourceWatchEvent<TResource>> MonitorAsync(string name, string @namespace, [EnumeratorCancellation]CancellationToken cancellationToken = default) |
139 | 156 | { |
140 | 157 | ArgumentException.ThrowIfNullOrWhiteSpace(name); |
141 | 158 | ArgumentException.ThrowIfNullOrWhiteSpace(@namespace); |
142 | 159 | var resource = new TResource(); |
143 | | - var uri = $"/api/{resource.Definition.Version}/{resource.Definition.Plural}/{@namespace}/{name}/monitor"; |
| 160 | + var uri = $"/api/{resource.Definition.Version}/{resource.Definition.Plural}/{@namespace}/{name}/monitor/sse"; |
144 | 161 | using var request = await this.ProcessRequestAsync(new HttpRequestMessage(HttpMethod.Get, uri), cancellationToken).ConfigureAwait(false); |
145 | | - request.EnableWebAssemblyStreamingResponse(); |
146 | 162 | var response = await this.HttpClient.SendAsync(request, HttpCompletionOption.ResponseHeadersRead, cancellationToken).ConfigureAwait(false); |
147 | 163 | var responseStream = await response.Content.ReadAsStreamAsync(cancellationToken).ConfigureAwait(false); |
148 | | - return this.JsonSerializer.DeserializeAsyncEnumerable<ResourceWatchEvent<TResource>>(responseStream, cancellationToken)!; |
| 164 | + using var streamReader = new StreamReader(await response.Content.ReadAsStreamAsync(cancellationToken).ConfigureAwait(false)); |
| 165 | + while (!streamReader.EndOfStream) |
| 166 | + { |
| 167 | + var sseMessage = await streamReader.ReadLineAsync(cancellationToken).ConfigureAwait(false); |
| 168 | + if (string.IsNullOrWhiteSpace(sseMessage)) continue; |
| 169 | + var json = sseMessage["data: ".Length..].Trim(); |
| 170 | + var e = JsonSerializer.Deserialize<ResourceWatchEvent<TResource>>(json)!; |
| 171 | + yield return e; |
| 172 | + } |
149 | 173 | } |
150 | 174 |
|
151 | 175 | /// <inheritdoc/> |
152 | | - public virtual async Task<IAsyncEnumerable<IResourceWatchEvent<TResource>>> MonitorAsync(string name, CancellationToken cancellationToken = default) |
| 176 | + public virtual async IAsyncEnumerable<IResourceWatchEvent<TResource>> MonitorAsync(string name, [EnumeratorCancellation]CancellationToken cancellationToken = default) |
153 | 177 | { |
154 | 178 | ArgumentException.ThrowIfNullOrWhiteSpace(name); |
155 | 179 | var resource = new TResource(); |
156 | | - var uri = $"/api/{resource.Definition.Version}/{resource.Definition.Plural}/{name}/monitor"; |
| 180 | + var uri = $"/api/{resource.Definition.Version}/{resource.Definition.Plural}/{name}/monitor/sse"; |
157 | 181 | using var request = await this.ProcessRequestAsync(new HttpRequestMessage(HttpMethod.Get, uri), cancellationToken).ConfigureAwait(false); |
158 | | - request.EnableWebAssemblyStreamingResponse(); |
159 | 182 | var response = await this.HttpClient.SendAsync(request, HttpCompletionOption.ResponseHeadersRead, cancellationToken).ConfigureAwait(false); |
160 | 183 | var responseStream = await response.Content.ReadAsStreamAsync(cancellationToken).ConfigureAwait(false); |
161 | | - return this.JsonSerializer.DeserializeAsyncEnumerable<ResourceWatchEvent<TResource>>(responseStream, cancellationToken)!; |
| 184 | + using var streamReader = new StreamReader(await response.Content.ReadAsStreamAsync(cancellationToken).ConfigureAwait(false)); |
| 185 | + while (!streamReader.EndOfStream) |
| 186 | + { |
| 187 | + var sseMessage = await streamReader.ReadLineAsync(cancellationToken).ConfigureAwait(false); |
| 188 | + if (string.IsNullOrWhiteSpace(sseMessage)) continue; |
| 189 | + var json = sseMessage["data: ".Length..].Trim(); |
| 190 | + var e = JsonSerializer.Deserialize<ResourceWatchEvent<TResource>>(json)!; |
| 191 | + yield return e; |
| 192 | + } |
162 | 193 | } |
163 | 194 |
|
164 | 195 | /// <inheritdoc/> |
|
0 commit comments