-
Notifications
You must be signed in to change notification settings - Fork 3
Expand file tree
/
Copy pathISStreamer.cs
More file actions
321 lines (280 loc) · 12.7 KB
/
Copy pathISStreamer.cs
File metadata and controls
321 lines (280 loc) · 12.7 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
/*Copyright 2019 Tektronix Inc.
*
*Licensed under the Apache License, Version 2.0 (the "License");
*you may not use this file except in compliance with the License.
*You may obtain a copy of the License at
*
* https://www.apache.org/licenses/LICENSE-2.0
*
*Unless required by applicable law or agreed to in writing, software
*distributed under the License is distributed on an "AS IS" BASIS,
*WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
*See the License for the specific language governing permissions and
*limitations under the License.
*/
using System;
using System.Collections.Generic;
using System.Linq;
using System.Text;
using System.Threading.Tasks;
using System.Net.Http;
using System.Collections;
namespace InitialState.Streaming
{
public struct StreamResponse
{
private bool success;
public System.Net.HttpStatusCode StatusCode;
public int RateLimit;
public int RateLimitRemaining;
public DateTime RateLimitReset;
public bool Success { get => success; internal set => success = value; }
}
/// <summary>
/// An object that serves to stream event data to an Initial State Event Data Stream.
/// </summary>
public class ISStreamer : IDisposable
{
/// <summary>
/// An indicator of the status of a <see cref="CreateBucket(string, string, string)"/> call.
/// </summary>
public enum CreateBucketStatus
{
Success,
AlreadyExists,
Error
}
// Private Fields
//---------------
private HttpClient _httpClient = null;
private StringBuilder _jsonStrBuilder = null;
private readonly DateTime _epochDateTime = new DateTime(1970, 1, 1, 0, 0, 0, 0);
// Properties
//-----------
/// <summary>
/// The base address for the Initial State API endpoint.
/// </summary>
public Uri ApiBaseAddress { get; private set; }
/// <summary>
/// The access key that will be used to access the Initial State API.
/// </summary>
public string AccessKey { get; private set; } = null;
/// <summary>
/// The bucket key that identifies the Initial State event data bucket to which data will be streamed.
/// </summary>
public string BucketKey { get; private set; } = null;
/// <summary>
/// Holds a collection of <see cref="ISEventData"/> that will be streamed to Initial State when <see cref="Stream"/>() is called.
/// </summary>
public ISEventDataCollection EventData { get; set; }
// Constructors
//-------------
/// <summary>
/// Initializes a new instance of the <see cref="ISStreamer"/> class.
/// </summary>
/// <remarks>Uses the default API base address https://groker.init.st/api/</remarks>
public ISStreamer() : this("https://groker.init.st/api/") { }
/// <summary>
/// Initializes a new instance of the <see cref="ISStreamer"/> class with a specified API base address.
/// </summary>
/// <param name="apiBaseAddress">The URL where InitalState API calls will be made.</param>
public ISStreamer(string apiBaseAddress)
{
this.ApiBaseAddress = new Uri(apiBaseAddress);
this.EventData = new ISEventDataCollection();
this._jsonStrBuilder = new StringBuilder(8192);
}
// Public Methods
//---------------
/// <summary>
/// Connects the <see cref="ISStreamer"/> to an Event Data Bucket for streaming.
/// </summary>
/// <param name="accessKey">The access key that will be used to access the Initial State API.</param>
/// <param name="bucketKey">The bucket key that identifies the Initial State Event Data bucket to which event data will be streamed.</param>
public void ConnectBucket(string accessKey, string bucketKey)
{
//return CreateBucket(accessKey, bucketKey, null);
this.AccessKey = accessKey;
this.BucketKey = bucketKey;
closeHttpClient();
this._httpClient = new HttpClient { BaseAddress = this.ApiBaseAddress };
this._httpClient.DefaultRequestHeaders.Clear();
this._httpClient.DefaultRequestHeaders.TryAddWithoutValidation("Accept-Version", "~0");
this._httpClient.DefaultRequestHeaders.TryAddWithoutValidation("X-Is-AccessKey", this.AccessKey);
this._httpClient.DefaultRequestHeaders.TryAddWithoutValidation("X-Is-BucketKey", this.BucketKey);
}
/// <summary>
/// Creates and connects the <see cref="ISStreamer"/> to a new Event Data Bucket for streaming.
/// </summary>
/// <param name="accessKey">The access key that will be used to access the Initial State API.</param>
/// <param name="bucketKey">The bucket key that identifies the Initial State event data bucket to which data will be streamed.</param>
/// <param name="bucketName">A name for the newly created event data bucket.</param>
/// <returns>A status indicating the result of the <see cref="CreateBucket(string, string, string)"/> call.</returns>
/// <remarks>If the bucket already exists, then the method will return <see cref="CreateBucketStatus.AlreadyExists"/>.</remarks>
public CreateBucketStatus CreateBucket(string accessKey, string bucketKey, string bucketName)
{
closeHttpClient();
this._httpClient = new HttpClient { BaseAddress = this.ApiBaseAddress };
this._httpClient.DefaultRequestHeaders.Clear();
this._httpClient.DefaultRequestHeaders.TryAddWithoutValidation("Accept-Version", "~0");
this._httpClient.DefaultRequestHeaders.TryAddWithoutValidation("X-IS-AccessKey", accessKey);
string json;
if (bucketName != null)
json = $"{{ \"bucketKey\": \"{bucketKey}\", \"bucketName\": \"{bucketName}\"}}";
else
json = $"{{ \"bucketKey\": \"{bucketKey}\"}}";
// Connect/Create Bucket
try
{
using (var content = new StringContent(json,
System.Text.Encoding.Default,
"application/json"))
{
using (var responseTask = _httpClient.PostAsync("buckets", content))
{
responseTask.Wait(5000);
// Keys accepted so store them
if (responseTask.Result.StatusCode == System.Net.HttpStatusCode.Created ||
responseTask.Result.StatusCode == System.Net.HttpStatusCode.NoContent)
{
this.AccessKey = accessKey;
this.BucketKey = bucketKey;
}
if (responseTask.Result.StatusCode == System.Net.HttpStatusCode.Created) // Then bucket was created
{
return CreateBucketStatus.Success;
}
else if (responseTask.Result.StatusCode == System.Net.HttpStatusCode.NoContent) // Then bucket already exists
{
return CreateBucketStatus.AlreadyExists;
}
else
{
return CreateBucketStatus.Error;
}
}
}
}
catch
{
closeHttpClient();
throw;
}
}
/// <summary>
/// Streams the events contained in <see cref="EventData"/> to Initial State. This call blocks and will not return until streaming of the data is complete.
/// </summary>
/// /// <returns>A <see cref="StreamResponse?"/> containing the response for the stream event.</returns>
/// <remarks>Upon succesful streaming the the event data, <see cref="EventData"/> will be cleared.</remarks>
public StreamResponse? Stream()
{
return Stream(-1);
}
/// <summary>
/// Streams the events contained in <see cref="EventData"/> to Initial State. This call blocks and will not return until streaming of the data is complete or timeout occurs.
/// </summary>
/// <returns>A <see cref="StreamResponse?"/> containing the response for the stream event.</returns>
/// <remarks>Upon succesful streaming the the event data, <see cref="EventData"/> will be cleared.</remarks>
public StreamResponse? Stream(int milliSecondsTimeout)
{
var streamTask = StreamAsync();
streamTask.Wait(milliSecondsTimeout);
if (streamTask.Exception != null)
{
throw streamTask.Exception;
}
return streamTask.Result;
}
/// <summary>
/// Streams the event data contained in <see cref="EventData"/> to Initial State.
/// </summary>s
/// <returns>A <see cref="StreamResponse?"/> containing the response for the stream event.</returns>
/// <remarks>Upon succesful streaming the the event data, <see cref="EventData"/> will be cleared.</remarks>
public async Task<StreamResponse?> StreamAsync()
{
if (_httpClient == null)
{
throw new InvalidOperationException("Cannot post data. Stream is not connected to a bucket.");
}
if (this.EventData.Count == 0)
{
// Nothing to send so just return
return null;
}
StreamResponse sr = new StreamResponse();
_jsonStrBuilder.Clear();
_jsonStrBuilder.Append("[\r\n");
foreach (ISEventData entry in this.EventData)
{
_jsonStrBuilder.Append(entry.ToJsonString() + ",\r\n");
}
_jsonStrBuilder.Remove(_jsonStrBuilder.Length - 3, 3); // Remove the last ",\r\n"
_jsonStrBuilder.Append("\r\n]");
using (var response = await sendToApi(_jsonStrBuilder.ToString()))
{
sr.Success = response.StatusCode == System.Net.HttpStatusCode.NoContent;
sr.StatusCode = response.StatusCode;
int.TryParse(response.Headers.GetValues("X-RateLimit-Limit").First(), out sr.RateLimit);
int.TryParse(response.Headers.GetValues("X-RateLimit-Remaining").First(), out sr.RateLimitRemaining);
int.TryParse(response.Headers.GetValues("X-RateLimit-Reset").First(), out int epochTimestamp);
sr.RateLimitReset = _epochDateTime.AddSeconds(epochTimestamp).ToLocalTime();
// Data is sent so clear out the buffer
if (sr.Success)
{
this.EventData.Clear();
}
}
return sr;
}
/// <summary>
/// Closes the <see cref="ISStreamer"/>'s connection to the event data bucket.
/// </summary>
public void Close()
{
closeHttpClient();
AccessKey = null;
BucketKey = null;
}
// Private Methods
//----------------
private async Task<HttpResponseMessage> sendToApi(string jsonData)
{
using (var content = new StringContent(jsonData, Encoding.UTF8, "application/json"))
{
var response = await _httpClient.PostAsync("events", content);
return response;
}
}
private void closeHttpClient()
{
if (this._httpClient != null)
{
this._httpClient.Dispose();
this._httpClient = null;
}
}
#region IDisposable Support
private bool disposedValue = false;
protected virtual void Dispose(bool disposing)
{
if (!disposedValue)
{
if (disposing)
{
closeHttpClient();
if (EventData != null)
{
EventData.Clear();
EventData = null;
}
}
disposedValue = true;
}
}
public void Dispose()
{
Dispose(true);
}
#endregion
}
}