EventSource.cs 32 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955
  1. #if !BESTHTTP_DISABLE_SERVERSENT_EVENTS
  2. using BestHTTP.Core;
  3. using BestHTTP.Extensions;
  4. using BestHTTP.Logger;
  5. using BestHTTP.PlatformSupport.Memory;
  6. using System;
  7. using System.Collections.Concurrent;
  8. using System.Collections.Generic;
  9. using System.Linq;
  10. using System.Text;
  11. #if UNITY_WEBGL && !UNITY_EDITOR
  12. using System.Runtime.InteropServices;
  13. #endif
  14. namespace BestHTTP.ServerSentEvents
  15. {
  16. /// <summary>
  17. /// Possible states of an EventSource object.
  18. /// </summary>
  19. public enum States
  20. {
  21. Initial,
  22. Connecting,
  23. Open,
  24. Retrying,
  25. Closing,
  26. Closed
  27. }
  28. public delegate void OnGeneralEventDelegate(EventSource eventSource);
  29. public delegate void OnMessageDelegate(EventSource eventSource, BestHTTP.ServerSentEvents.Message message);
  30. public delegate void OnErrorDelegate(EventSource eventSource, string error);
  31. public delegate bool OnRetryDelegate(EventSource eventSource);
  32. public delegate void OnEventDelegate(EventSource eventSource, BestHTTP.ServerSentEvents.Message message);
  33. public delegate void OnStateChangedDelegate(EventSource eventSource, States oldState, States newState);
  34. #if !UNITY_WEBGL || UNITY_EDITOR
  35. public delegate void OnCommentDelegate(EventSource eventSource, string comment);
  36. #endif
  37. #if UNITY_WEBGL && !UNITY_EDITOR
  38. delegate void OnWebGLEventSourceOpenDelegate(uint id);
  39. delegate void OnWebGLEventSourceMessageDelegate(uint id, string eventStr, string data, string eventId, int retry);
  40. delegate void OnWebGLEventSourceErrorDelegate(uint id, string reason);
  41. #endif
  42. /// <summary>
  43. /// http://www.w3.org/TR/eventsource/
  44. /// </summary>
  45. public class EventSource : IProtocol
  46. #if !UNITY_WEBGL || UNITY_EDITOR
  47. , IHeartbeat
  48. #endif
  49. {
  50. #region Public Properties
  51. /// <summary>
  52. /// Uri of the remote endpoint.
  53. /// </summary>
  54. public Uri Uri { get; private set; }
  55. /// <summary>
  56. /// Current state of the EventSource object.
  57. /// </summary>
  58. public States State
  59. {
  60. get
  61. {
  62. return _state;
  63. }
  64. private set
  65. {
  66. States oldState = _state;
  67. _state = value;
  68. if (OnStateChanged != null)
  69. {
  70. try
  71. {
  72. OnStateChanged(this, oldState, _state);
  73. }
  74. catch(Exception ex)
  75. {
  76. HTTPManager.Logger.Exception("EventSource", "OnStateChanged", ex);
  77. }
  78. }
  79. }
  80. }
  81. private States _state;
  82. /// <summary>
  83. /// Time to wait to do a reconnect attempt. Default to 2 sec. The server can overwrite this setting.
  84. /// </summary>
  85. public TimeSpan ReconnectionTime { get; set; }
  86. /// <summary>
  87. /// The last successfully received event's id.
  88. /// </summary>
  89. public string LastEventId { get; private set; }
  90. public HostConnectionKey ConnectionKey { get; private set; }
  91. public bool IsClosed { get { return this.State == States.Closed; } }
  92. public LoggingContext LoggingContext { get; private set; }
  93. #if !UNITY_WEBGL || UNITY_EDITOR
  94. /// <summary>
  95. /// The internal request object of the EventSource.
  96. /// </summary>
  97. public HTTPRequest InternalRequest { get; private set; }
  98. #else
  99. public bool WithCredentials { get; set; }
  100. #endif
  101. #endregion
  102. #region Public Events
  103. /// <summary>
  104. /// Called when successfully connected to the server.
  105. /// </summary>
  106. public event OnGeneralEventDelegate OnOpen;
  107. /// <summary>
  108. /// Called on every message received from the server.
  109. /// </summary>
  110. public event OnMessageDelegate OnMessage;
  111. /// <summary>
  112. /// Called when an error occurs.
  113. /// </summary>
  114. public event OnErrorDelegate OnError;
  115. #if !UNITY_WEBGL || UNITY_EDITOR
  116. /// <summary>
  117. /// Called when the EventSource will try to do a retry attempt. If this function returns with false, it will cancel the attempt.
  118. /// </summary>
  119. public event OnRetryDelegate OnRetry;
  120. /// <summary>
  121. /// This event is called for comments received from the server.
  122. /// </summary>
  123. public event OnCommentDelegate OnComment;
  124. #endif
  125. /// <summary>
  126. /// Called when the EventSource object closed.
  127. /// </summary>
  128. public event OnGeneralEventDelegate OnClosed;
  129. /// <summary>
  130. /// Called every time when the State property changed.
  131. /// </summary>
  132. public event OnStateChangedDelegate OnStateChanged;
  133. #endregion
  134. #region Privates
  135. /// <summary>
  136. /// A dictionary to store eventName => delegate mapping.
  137. /// </summary>
  138. private Dictionary<string, OnEventDelegate> EventTable;
  139. #if !UNITY_WEBGL || UNITY_EDITOR
  140. /// <summary>
  141. /// Number of retry attempts made.
  142. /// </summary>
  143. private byte RetryCount;
  144. /// <summary>
  145. /// When we called the Retry function. We will delay the Open call from here.
  146. /// </summary>
  147. private DateTime RetryCalled;
  148. /// <summary>
  149. /// Buffer for the read data.
  150. /// </summary>
  151. private byte[] LineBuffer;
  152. /// <summary>
  153. /// Buffer position.
  154. /// </summary>
  155. private int LineBufferPos = 0;
  156. /// <summary>
  157. /// The currently receiving and parsing message
  158. /// </summary>
  159. private BestHTTP.ServerSentEvents.Message CurrentMessage;
  160. /// <summary>
  161. /// Completed messages that waiting to be dispatched
  162. /// </summary>
  163. //private List<BestHTTP.ServerSentEvents.Message> CompletedMessages = new List<BestHTTP.ServerSentEvents.Message>();
  164. private ConcurrentQueue<BestHTTP.ServerSentEvents.Message> CompletedMessages = new ConcurrentQueue<Message>();
  165. #else
  166. private static Dictionary<uint, EventSource> EventSources = new Dictionary<uint, EventSource>();
  167. private uint Id;
  168. #endif
  169. #endregion
  170. public EventSource(Uri uri, int readBufferSizeOverride = 0)
  171. {
  172. this.Uri = uri;
  173. this.LoggingContext = new LoggingContext(this);
  174. this.ReconnectionTime = TimeSpan.FromMilliseconds(2000);
  175. this.ConnectionKey = new HostConnectionKey(this.Uri.Host, HostDefinition.GetKeyFor(this.Uri
  176. #if !BESTHTTP_DISABLE_PROXY && (!UNITY_WEBGL || UNITY_EDITOR)
  177. , HTTPManager.Proxy
  178. #endif
  179. ));
  180. #if !UNITY_WEBGL || UNITY_EDITOR
  181. this.InternalRequest = new HTTPRequest(Uri, HTTPMethods.Get, true, true, OnRequestFinished);
  182. // Set headers
  183. this.InternalRequest.SetHeader("Accept", "text/event-stream");
  184. this.InternalRequest.SetHeader("Cache-Control", "no-cache");
  185. this.InternalRequest.SetHeader("Accept-Encoding", "identity");
  186. this.InternalRequest.StreamChunksImmediately = true;
  187. this.InternalRequest.ReadBufferSizeOverride = readBufferSizeOverride;
  188. this.InternalRequest.OnStreamingData = OnData;
  189. // Disable internal retry
  190. this.InternalRequest.MaxRetries = 0;
  191. this.InternalRequest.Context.Add("EventSource", this.LoggingContext);
  192. #else
  193. if (!ES_IsSupported())
  194. throw new NotSupportedException("This browser isn't support the EventSource protocol!");
  195. this.Id = ES_Create(this.Uri.ToString(), WithCredentials, OnOpenCallback, OnMessageCallback, OnErrorCallback);
  196. EventSources.Add(this.Id, this);
  197. #endif
  198. }
  199. #region Public Functions
  200. /// <summary>
  201. /// Start to connect to the remote server.
  202. /// </summary>
  203. public void Open()
  204. {
  205. if (this.State != States.Initial &&
  206. this.State != States.Retrying &&
  207. this.State != States.Closed)
  208. return;
  209. this.State = States.Connecting;
  210. #if !UNITY_WEBGL || UNITY_EDITOR
  211. if (!string.IsNullOrEmpty(this.LastEventId))
  212. this.InternalRequest.SetHeader("Last-Event-ID", this.LastEventId);
  213. this.InternalRequest.Send();
  214. #endif
  215. }
  216. /// <summary>
  217. /// Start to close the connection.
  218. /// </summary>
  219. public void Close()
  220. {
  221. if (this.State == States.Closing ||
  222. this.State == States.Closed)
  223. return;
  224. this.State = States.Closing;
  225. #if !UNITY_WEBGL || UNITY_EDITOR
  226. if (this.InternalRequest != null)
  227. this.CancellationRequested();
  228. else
  229. this.State = States.Closed;
  230. #else
  231. ES_Close(this.Id);
  232. SetClosed("Close");
  233. EventSources.Remove(this.Id);
  234. ES_Release(this.Id);
  235. #endif
  236. }
  237. /// <summary>
  238. /// With this function an event handler can be subscribed for an event name.
  239. /// </summary>
  240. public void On(string eventName, OnEventDelegate action)
  241. {
  242. if (EventTable == null)
  243. EventTable = new Dictionary<string, OnEventDelegate>();
  244. EventTable[eventName] = action;
  245. #if UNITY_WEBGL && !UNITY_EDITOR
  246. ES_AddEventHandler(this.Id, eventName);
  247. #endif
  248. }
  249. /// <summary>
  250. /// With this function the event handler can be removed for the given event name.
  251. /// </summary>
  252. /// <param name="eventName"></param>
  253. public void Off(string eventName)
  254. {
  255. if (eventName == null || EventTable == null)
  256. return;
  257. EventTable.Remove(eventName);
  258. }
  259. #endregion
  260. #region Private Helper Functions
  261. private void CallOnError(string error, string msg)
  262. {
  263. HTTPManager.Logger.Verbose(nameof(EventSource), $"CallOnError({this.State}, {error}, {msg})", this.LoggingContext);
  264. if (OnError != null)
  265. {
  266. try
  267. {
  268. OnError(this, error);
  269. }
  270. catch (Exception ex)
  271. {
  272. HTTPManager.Logger.Exception("EventSource", msg + " - OnError", ex, this.LoggingContext);
  273. }
  274. }
  275. }
  276. #if !UNITY_WEBGL || UNITY_EDITOR
  277. private bool CallOnRetry()
  278. {
  279. if (OnRetry != null)
  280. {
  281. try
  282. {
  283. return OnRetry(this);
  284. }
  285. catch(Exception ex)
  286. {
  287. HTTPManager.Logger.Exception("EventSource", "CallOnRetry", ex, this.LoggingContext);
  288. }
  289. }
  290. return true;
  291. }
  292. #endif
  293. private void SetClosed(string msg)
  294. {
  295. HTTPManager.Logger.Information("EventSource", $"SetClosed({this.State}, {msg})", this.LoggingContext);
  296. this.State = States.Closed;
  297. if (OnClosed != null)
  298. {
  299. try
  300. {
  301. OnClosed(this);
  302. }
  303. catch (Exception ex)
  304. {
  305. HTTPManager.Logger.Exception("EventSource", msg + " - OnClosed", ex, this.LoggingContext);
  306. }
  307. }
  308. }
  309. #if !UNITY_WEBGL || UNITY_EDITOR
  310. private void Retry()
  311. {
  312. HTTPManager.Logger.Information("EventSource", $"Retry({this.State})", this.LoggingContext);
  313. if (RetryCount > 0 ||
  314. !CallOnRetry())
  315. {
  316. SetClosed("Retry");
  317. return;
  318. }
  319. RetryCount++;
  320. RetryCalled = DateTime.UtcNow;
  321. HTTPManager.Heartbeats.Subscribe(this);
  322. this.State = States.Retrying;
  323. }
  324. #endif
  325. #endregion
  326. #region HTTP Request Implementation
  327. #if !UNITY_WEBGL || UNITY_EDITOR
  328. private void OnRequestFinished(HTTPRequest req, HTTPResponse resp)
  329. {
  330. HTTPManager.Logger.Information("EventSource", string.Format("OnRequestFinished - State: {0}, StatusCode: {1}", this.State, resp != null ? resp.StatusCode : 0), req.Context);
  331. if (this.State == States.Closed)
  332. return;
  333. if (this.State == States.Closing || req.IsCancellationRequested)
  334. {
  335. SetClosed("OnRequestFinished");
  336. return;
  337. }
  338. string reason = string.Empty;
  339. // In some cases retry is prohibited
  340. bool canRetry = true;
  341. switch (req.State)
  342. {
  343. // The request finished without any problem.
  344. case HTTPRequestStates.Finished:
  345. // HTTP 200 OK responses that have a Content-Type specifying an unsupported type, or that have no Content-Type at all, must cause the user agent to fail the connection.
  346. if (resp.StatusCode == 200)
  347. {
  348. // https://github.com/Benedicht/BestHTTP-Issues/issues/168
  349. // https://html.spec.whatwg.org/multipage/iana.html#text/event-stream
  350. // Content-Type's value must be "text/event-stream", but it might have an optional charset parameter
  351. var contentTypeValue = resp.GetFirstHeaderValue("content-type");
  352. if (contentTypeValue == null)
  353. {
  354. canRetry = false;
  355. reason = "No Content-Type header found!";
  356. break;
  357. }
  358. HeaderParser contentType = new HeaderParser(contentTypeValue);
  359. bool eventStream = contentType?.Values?.FirstOrDefault()?.Key == "text/event-stream";
  360. if (!eventStream)
  361. {
  362. reason = $"No Content-Type header with value 'text/event-stream' present. Got '{contentTypeValue}'";
  363. canRetry = false;
  364. break;
  365. }
  366. canRetry = false;
  367. }
  368. // HTTP 500 Internal Server Error, 502 Bad Gateway, 503 Service Unavailable, and 504 Gateway Timeout responses, and any network error that prevents the connection
  369. // from being established in the first place (e.g. DNS errors), must cause the user agent to asynchronously reestablish the connection.
  370. // Any other HTTP response code not listed here must cause the user agent to fail the connection.
  371. if (canRetry &&
  372. resp.StatusCode != 500 &&
  373. resp.StatusCode != 502 &&
  374. resp.StatusCode != 503 &&
  375. resp.StatusCode != 504)
  376. {
  377. canRetry = false;
  378. reason = string.Format("Request Finished Successfully, but the server sent an error. Status Code: {0}-{1} Message: {2}",
  379. resp.StatusCode,
  380. resp.Message,
  381. resp.DataAsText);
  382. }
  383. break;
  384. // The request finished with an unexpected error. The request's Exception property may contain more info about the error.
  385. case HTTPRequestStates.Error:
  386. reason = "Request Finished with Error! " + (req.Exception != null ? (req.Exception.Message + "\n" + req.Exception.StackTrace) : "No Exception");
  387. break;
  388. // The request aborted, initiated by the user.
  389. case HTTPRequestStates.Aborted:
  390. // If the state is Closing, then it's a normal behaviour, and we close the EventSource
  391. reason = "OnRequestFinished - Aborted without request. EventSource's State: " + this.State;
  392. break;
  393. // Connecting to the server is timed out.
  394. case HTTPRequestStates.ConnectionTimedOut:
  395. reason = "Connection Timed Out!";
  396. break;
  397. // The request didn't finished in the given time.
  398. case HTTPRequestStates.TimedOut:
  399. reason = "Processing the request Timed Out!";
  400. break;
  401. }
  402. // If we are not closing the EventSource, then we will try to reconnect.
  403. if (this.State < States.Closing)
  404. {
  405. if (!string.IsNullOrEmpty(reason))
  406. CallOnError(reason, "OnRequestFinished");
  407. if (canRetry)
  408. Retry();
  409. else
  410. SetClosed("OnRequestFinished");
  411. }
  412. else
  413. SetClosed("OnRequestFinished");
  414. }
  415. private bool OnData(HTTPRequest request, HTTPResponse response, byte[] dataFragment, int dataFragmentLength)
  416. {
  417. if (HTTPManager.Logger.Level == Loglevels.All)
  418. HTTPManager.Logger.Information("EventSource", $"OnData({dataFragment.AsBuffer(dataFragmentLength)})", this.LoggingContext);
  419. if (this.State == States.Connecting)
  420. {
  421. string contentType = response.GetFirstHeaderValue("content-type");
  422. bool IsUpgraded = response.StatusCode == 200 &&
  423. !string.IsNullOrEmpty(contentType) &&
  424. contentType.ToLower().StartsWith("text/event-stream");
  425. if (IsUpgraded)
  426. {
  427. ProtocolEventHelper.AddProtocol(this);
  428. if (this.OnOpen != null)
  429. {
  430. try
  431. {
  432. this.OnOpen(this);
  433. }
  434. catch (Exception ex)
  435. {
  436. HTTPManager.Logger.Exception("EventSource", "OnOpen", ex, request.Context);
  437. }
  438. }
  439. this.RetryCount = 0;
  440. this.State = States.Open;
  441. }
  442. else
  443. {
  444. this.State = States.Closing;
  445. request.Abort();
  446. }
  447. }
  448. if (this.State == States.Closing)
  449. return true;
  450. if (FeedData(dataFragment, dataFragmentLength))
  451. ProtocolEventHelper.EnqueueProtocolEvent(new ProtocolEventInfo(this));
  452. return true;
  453. }
  454. #region Data Parsing
  455. public bool FeedData(byte[] buffer, int count)
  456. {
  457. if (HTTPManager.Logger.Level == Loglevels.All)
  458. HTTPManager.Logger.Information("EventSource", $"FeedData({buffer.AsBuffer(count)})", this.LoggingContext);
  459. if (count == -1)
  460. count = buffer.Length;
  461. if (count == 0)
  462. return false;
  463. if (LineBuffer == null)
  464. LineBuffer = BufferPool.Get(1024, true);
  465. int newlineIdx;
  466. int pos = 0;
  467. bool hasMessageToSend = false;
  468. do
  469. {
  470. newlineIdx = -1;
  471. int skipCount = 1; // to skip CR and/or LF
  472. for (int i = pos; i < count && newlineIdx == -1; ++i)
  473. {
  474. // Lines must be separated by either a U+000D CARRIAGE RETURN U+000A LINE FEED (CRLF) character pair, a single U+000A LINE FEED (LF) character, or a single U+000D CARRIAGE RETURN (CR) character.
  475. if (buffer[i] == HTTPResponse.CR)
  476. {
  477. if (i + 1 < count && buffer[i + 1] == HTTPResponse.LF)
  478. skipCount = 2;
  479. newlineIdx = i;
  480. }
  481. else if (buffer[i] == HTTPResponse.LF)
  482. newlineIdx = i;
  483. }
  484. int copyIndex = newlineIdx == -1 ? count : newlineIdx;
  485. if (LineBuffer.Length < LineBufferPos + (copyIndex - pos))
  486. {
  487. int newSize = LineBufferPos + (copyIndex - pos);
  488. BufferPool.Resize(ref LineBuffer, newSize, true, false);
  489. }
  490. Array.Copy(buffer, pos, LineBuffer, LineBufferPos, copyIndex - pos);
  491. LineBufferPos += copyIndex - pos;
  492. if (newlineIdx == -1)
  493. return hasMessageToSend;
  494. hasMessageToSend |= ParseLine(LineBuffer, LineBufferPos);
  495. LineBufferPos = 0;
  496. //pos += newlineIdx + skipCount;
  497. pos = newlineIdx + skipCount;
  498. } while (newlineIdx != -1 && pos < count);
  499. return hasMessageToSend;
  500. }
  501. bool ParseLine(byte[] buffer, int count)
  502. {
  503. // If the line is empty (a blank line) => Dispatch the event
  504. if (count == 0)
  505. {
  506. if (CurrentMessage != null)
  507. {
  508. if (HTTPManager.Logger.Level == Loglevels.All)
  509. HTTPManager.Logger.Information("EventSource", $"ParseLine - event dispatch ({this.State}, {CurrentMessage})", this.LoggingContext);
  510. CompletedMessages.Enqueue(CurrentMessage);
  511. CurrentMessage = null;
  512. return true;
  513. }
  514. return false;
  515. }
  516. // If the line starts with a U+003A COLON character (:) => Ignore the line.
  517. if (buffer[0] == 0x3A)
  518. {
  519. var msg = new Message() { IsComment = true, Data = Encoding.UTF8.GetString(buffer, 1, count - 1) };
  520. if (HTTPManager.Logger.Level == Loglevels.All)
  521. HTTPManager.Logger.Information("EventSource", $"ParseLine - comment (':') found ({this.State}, {msg})", this.LoggingContext);
  522. this.CompletedMessages.Enqueue(msg);
  523. return true;
  524. }
  525. //If the line contains a U+003A COLON character (:)
  526. int colonIdx = -1;
  527. for (int i = 0; i < count && colonIdx == -1; ++i)
  528. if (buffer[i] == 0x3A)
  529. colonIdx = i;
  530. string field;
  531. string value;
  532. if (colonIdx != -1)
  533. {
  534. // Collect the characters on the line before the first U+003A COLON character (:), and let field be that string.
  535. field = Encoding.UTF8.GetString(buffer, 0, colonIdx);
  536. //Collect the characters on the line after the first U+003A COLON character (:), and let value be that string. If value starts with a U+0020 SPACE character, remove it from value.
  537. if (colonIdx + 1 < count && buffer[colonIdx + 1] == 0x20)
  538. colonIdx++;
  539. colonIdx++;
  540. // discarded because it is not followed by a blank line
  541. if (colonIdx >= count)
  542. return false;
  543. value = Encoding.UTF8.GetString(buffer, colonIdx, count - colonIdx);
  544. }
  545. else
  546. {
  547. // Otherwise, the string is not empty but does not contain a U+003A COLON character (:) =>
  548. // Process the field using the whole line as the field name, and the empty string as the field value.
  549. field = Encoding.UTF8.GetString(buffer, 0, count);
  550. value = string.Empty;
  551. }
  552. if (CurrentMessage == null)
  553. CurrentMessage = new BestHTTP.ServerSentEvents.Message();
  554. switch (field)
  555. {
  556. // If the field name is "id" => Set the last event ID buffer to the field value.
  557. case "id":
  558. CurrentMessage.Id = value;
  559. break;
  560. // If the field name is "event" => Set the event type buffer to field value.
  561. case "event":
  562. CurrentMessage.Event = value;
  563. break;
  564. // If the field name is "data" => Append the field value to the data buffer, then append a single U+000A LINE FEED (LF) character to the data buffer.
  565. case "data":
  566. // Append a new line if we already have some data. This way we can skip step 3.) in the EventSource's OnMessageReceived.
  567. // We do only null check, because empty string can be valid payload
  568. if (CurrentMessage.Data != null)
  569. CurrentMessage.Data += Environment.NewLine;
  570. CurrentMessage.Data += value;
  571. break;
  572. // If the field name is "retry" => If the field value consists of only ASCII digits, then interpret the field value as an integer in base ten,
  573. // and set the event stream's reconnection time to that integer. Otherwise, ignore the field.
  574. case "retry":
  575. int result;
  576. if (int.TryParse(value, out result))
  577. CurrentMessage.Retry = TimeSpan.FromMilliseconds(result);
  578. break;
  579. // Otherwise: The field is ignored.
  580. default:
  581. break;
  582. }
  583. return false;
  584. }
  585. #endregion
  586. #endif
  587. #endregion
  588. #region EventStreamResponse Event Handlers
  589. private void OnMessageReceived(BestHTTP.ServerSentEvents.Message message)
  590. {
  591. if (HTTPManager.Logger.Level == Loglevels.All)
  592. HTTPManager.Logger.Information("EventSource", $"OnMessageReceived({this.State}, {message})", this.LoggingContext);
  593. if (this.State >= States.Closing)
  594. return;
  595. // 1.) Set the last event ID string of the event source to value of the last event ID buffer.
  596. // The buffer does not get reset, so the last event ID string of the event source remains set to this value until the next time it is set by the server.
  597. // We check here only for null, because it can be a non-null but empty string.
  598. if (message.Id != null)
  599. this.LastEventId = message.Id;
  600. if (message.Retry.TotalMilliseconds > 0)
  601. this.ReconnectionTime = message.Retry;
  602. // 2.) If the data buffer is an empty string, set the data buffer and the event type buffer to the empty string and abort these steps.
  603. if (string.IsNullOrEmpty(message.Data))
  604. return;
  605. // 3.) If the data buffer's last character is a U+000A LINE FEED (LF) character, then remove the last character from the data buffer.
  606. // This step can be ignored. We constructed the string to be able to skip this step.
  607. if (OnMessage != null && !message.IsComment)
  608. {
  609. try
  610. {
  611. OnMessage(this, message);
  612. }
  613. catch (Exception ex)
  614. {
  615. HTTPManager.Logger.Exception("EventSource", "OnMessageReceived - OnMessage", ex, this.LoggingContext);
  616. }
  617. }
  618. #if !UNITY_WEBGL || UNITY_EDITOR
  619. else if (message.IsComment && this.OnComment != null)
  620. {
  621. try
  622. {
  623. this.OnComment(this, message.Data);
  624. }
  625. catch (Exception ex)
  626. {
  627. HTTPManager.Logger.Exception("EventSource", "OnMessageReceived - OnComment", ex, this.LoggingContext);
  628. }
  629. }
  630. #endif
  631. if (EventTable != null && !string.IsNullOrEmpty(message.Event))
  632. {
  633. OnEventDelegate action;
  634. if (EventTable.TryGetValue(message.Event, out action))
  635. {
  636. if (action != null)
  637. {
  638. try
  639. {
  640. action(this, message);
  641. }
  642. catch(Exception ex)
  643. {
  644. HTTPManager.Logger.Exception("EventSource", "OnMessageReceived - action", ex, this.LoggingContext);
  645. }
  646. }
  647. }
  648. }
  649. }
  650. public void HandleEvents()
  651. {
  652. #if !UNITY_WEBGL || UNITY_EDITOR
  653. if (this.State == States.Open)
  654. {
  655. BestHTTP.ServerSentEvents.Message message;
  656. while (this.CompletedMessages.TryDequeue(out message))
  657. OnMessageReceived(message);
  658. }
  659. #endif
  660. }
  661. public void CancellationRequested()
  662. {
  663. #if !UNITY_WEBGL || UNITY_EDITOR
  664. if (this.InternalRequest != null)
  665. this.InternalRequest.Abort();
  666. #else
  667. Close();
  668. #endif
  669. }
  670. public void Dispose()
  671. {
  672. }
  673. #endregion
  674. #region IHeartbeat Implementation
  675. #if !UNITY_WEBGL || UNITY_EDITOR
  676. void IHeartbeat.OnHeartbeatUpdate(TimeSpan dif)
  677. {
  678. if (this.State != States.Retrying)
  679. {
  680. HTTPManager.Heartbeats.Unsubscribe(this);
  681. return;
  682. }
  683. if (DateTime.UtcNow - RetryCalled >= ReconnectionTime)
  684. {
  685. Open();
  686. if (this.State != States.Connecting)
  687. SetClosed("OnHeartbeatUpdate");
  688. HTTPManager.Heartbeats.Unsubscribe(this);
  689. }
  690. }
  691. #endif
  692. #endregion
  693. #region WebGL Static Callbacks
  694. #if UNITY_WEBGL && !UNITY_EDITOR
  695. [AOT.MonoPInvokeCallback(typeof(OnWebGLEventSourceOpenDelegate))]
  696. static void OnOpenCallback(uint id)
  697. {
  698. EventSource es;
  699. if (EventSources.TryGetValue(id, out es))
  700. {
  701. if (es.OnOpen != null)
  702. {
  703. try
  704. {
  705. es.OnOpen(es);
  706. }
  707. catch(Exception ex)
  708. {
  709. HTTPManager.Logger.Exception("EventSource", "OnOpen", ex, es.LoggingContext);
  710. }
  711. }
  712. es.State = States.Open;
  713. }
  714. else
  715. HTTPManager.Logger.Warning("EventSource", "OnOpenCallback - No EventSource found for id: " + id.ToString());
  716. }
  717. [AOT.MonoPInvokeCallback(typeof(OnWebGLEventSourceMessageDelegate))]
  718. static void OnMessageCallback(uint id, string eventStr, string data, string eventId, int retry)
  719. {
  720. EventSource es;
  721. if (EventSources.TryGetValue(id, out es))
  722. {
  723. var msg = new BestHTTP.ServerSentEvents.Message();
  724. msg.Id = eventId;
  725. msg.Data = data;
  726. msg.Event = eventStr;
  727. msg.Retry = TimeSpan.FromSeconds(retry);
  728. es.OnMessageReceived(msg);
  729. }
  730. }
  731. [AOT.MonoPInvokeCallback(typeof(OnWebGLEventSourceErrorDelegate))]
  732. static void OnErrorCallback(uint id, string reason)
  733. {
  734. EventSource es;
  735. if (EventSources.TryGetValue(id, out es))
  736. {
  737. es.CallOnError(reason, "OnErrorCallback");
  738. es.SetClosed("OnError");
  739. EventSources.Remove(id);
  740. }
  741. try
  742. {
  743. ES_Release(id);
  744. }
  745. catch (Exception ex)
  746. {
  747. HTTPManager.Logger.Exception("EventSource", "ES_Release", ex);
  748. }
  749. }
  750. #endif
  751. #endregion
  752. #region WebGL Interface
  753. #if UNITY_WEBGL && !UNITY_EDITOR
  754. [DllImport("__Internal")]
  755. static extern bool ES_IsSupported();
  756. [DllImport("__Internal")]
  757. static extern uint ES_Create(string url, bool withCred, OnWebGLEventSourceOpenDelegate onOpen, OnWebGLEventSourceMessageDelegate onMessage, OnWebGLEventSourceErrorDelegate onError);
  758. [DllImport("__Internal")]
  759. static extern void ES_AddEventHandler(uint id, string eventName);
  760. [DllImport("__Internal")]
  761. static extern void ES_Close(uint id);
  762. [DllImport("__Internal")]
  763. static extern void ES_Release(uint id);
  764. #endif
  765. #endregion
  766. }
  767. }
  768. #endif