File
Blob: src/workerd/api/eventsource.h
| 1 | // Copyright (c) 2017-2024 Cloudflare, Inc. |
| 2 | // Licensed under the Apache 2.0 license found in the LICENSE file or at: |
| 3 | // https://opensource.org/licenses/Apache-2.0 |
| 4 | |
| 5 | #pragma once |
| 6 | #include "basics.h" |
| 7 | #include "events.h" |
| 8 | #include "http.h" |
| 9 | |
| 10 | #include <workerd/jsg/jsg.h> |
| 11 | #include <workerd/jsg/url.h> |
| 12 | |
| 13 | namespace workerd::api { |
| 14 | |
| 15 | using kj::uint; |
| 16 | class Fetcher; |
| 17 | class ReadableStream; |
| 18 | class Response; |
| 19 | |
| 20 | // Implements the web standard EventSource API |
| 21 | // https://developer.mozilla.org/en-US/docs/Web/API/EventSource |
| 22 | class EventSource: public EventTarget { |
| 23 | public: |
| 24 | struct EventSourceInit { |
| 25 | // We don't actually make use of the standard withCredentials option. If this is set to |
| 26 | // any truthy value, we'll throw. |
| 27 | jsg::Optional<bool> withCredentials; |
| 28 | |
| 29 | // This is a non-standard workers-specific extension that allows the EventSource to |
| 30 | // use a custom Fetcher instance. |
| 31 | jsg::Optional<jsg::Ref<Fetcher>> fetcher; |
| 32 | JSG_STRUCT(withCredentials, fetcher); |
| 33 | }; |
| 34 | |
| 35 | enum class State { |
| 36 | CONNECTING = 0, |
| 37 | OPEN = 1, |
| 38 | CLOSED = 2, |
| 39 | }; |
| 40 | |
| 41 | EventSource(jsg::Lock& js, jsg::Url url, kj::Maybe<EventSourceInit> init = kj::none); |
| 42 | |
| 43 | EventSource(jsg::Lock& js); |
| 44 | |
| 45 | static jsg::Ref<EventSource> constructor( |
| 46 | jsg::Lock& js, kj::String url, jsg::Optional<EventSourceInit> init); |
| 47 | |
| 48 | kj::ArrayPtr<const char> getUrl() const { |
| 49 | KJ_IF_SOME(i, impl) { |
| 50 | return i.url.getHref(); |
| 51 | } |
| 52 | return nullptr; |
| 53 | } |
| 54 | bool getWithCredentials() const { |
| 55 | return false; |
| 56 | } |
| 57 | uint getReadyState() const { |
| 58 | return static_cast<uint>(readyState); |
| 59 | } |
| 60 | |
| 61 | void close(jsg::Lock& js); |
| 62 | |
| 63 | // A non-standard extension that creates an EventSource instance around a ReadableStream |
| 64 | // instance. In this instance, automatic reconnection is disabled since there is no URL |
| 65 | // or underlying fetch used. The ReadableStream instance must produce bytes. It will be |
| 66 | // locked and disturbed, and will be read until it either ends or errors. Calling close() |
| 67 | // will cause the stream to be canceled. |
| 68 | static jsg::Ref<EventSource> from(jsg::Lock& js, jsg::Ref<ReadableStream> stream); |
| 69 | |
| 70 | kj::Maybe<jsg::JsValue> getOnOpen(jsg::Lock& js) { |
| 71 | return onopenValue.map( |
| 72 | [&](jsg::JsRef<jsg::JsValue>& ref) -> jsg::JsValue { return ref.getHandle(js); }); |
| 73 | } |
| 74 | void setOnOpen(jsg::Lock& js, jsg::JsValue value) { |
| 75 | if (!value.isObject() && !value.isFunction()) { |
| 76 | onopenValue = kj::none; |
| 77 | } else { |
| 78 | onopenValue = jsg::JsRef<jsg::JsValue>(js, value); |
| 79 | } |
| 80 | } |
| 81 | kj::Maybe<jsg::JsValue> getOnMessage(jsg::Lock& js) { |
| 82 | return onmessageValue.map( |
| 83 | [&](jsg::JsRef<jsg::JsValue>& ref) -> jsg::JsValue { return ref.getHandle(js); }); |
| 84 | } |
| 85 | void setOnMessage(jsg::Lock& js, jsg::JsValue value) { |
| 86 | if (!value.isObject() && !value.isFunction()) { |
| 87 | onmessageValue = kj::none; |
| 88 | } else { |
| 89 | onmessageValue = jsg::JsRef<jsg::JsValue>(js, value); |
| 90 | } |
| 91 | } |
| 92 | kj::Maybe<jsg::JsValue> getOnError(jsg::Lock& js) { |
| 93 | return onerrorValue.map( |
| 94 | [&](jsg::JsRef<jsg::JsValue>& ref) -> jsg::JsValue { return ref.getHandle(js); }); |
| 95 | } |
| 96 | void setOnError(jsg::Lock& js, jsg::JsValue value) { |
| 97 | if (!value.isObject() && !value.isFunction()) { |
| 98 | onerrorValue = kj::none; |
| 99 | } else { |
| 100 | onerrorValue = jsg::JsRef<jsg::JsValue>(js, value); |
| 101 | } |
| 102 | } |
| 103 | |
| 104 | JSG_RESOURCE_TYPE(EventSource) { |
| 105 | JSG_INHERIT(EventTarget); |
| 106 | JSG_METHOD(close); |
| 107 | JSG_READONLY_PROTOTYPE_PROPERTY(url, getUrl); |
| 108 | JSG_READONLY_PROTOTYPE_PROPERTY(withCredentials, getWithCredentials); |
| 109 | JSG_READONLY_PROTOTYPE_PROPERTY(readyState, getReadyState); |
| 110 | JSG_PROTOTYPE_PROPERTY(onopen, getOnOpen, setOnOpen); |
| 111 | JSG_PROTOTYPE_PROPERTY(onmessage, getOnMessage, setOnMessage); |
| 112 | JSG_PROTOTYPE_PROPERTY(onerror, getOnError, setOnError); |
| 113 | JSG_STATIC_CONSTANT_NAMED(CONNECTING, static_cast<uint>(State::CONNECTING)); |
| 114 | JSG_STATIC_CONSTANT_NAMED(OPEN, static_cast<uint>(State::OPEN)); |
| 115 | JSG_STATIC_CONSTANT_NAMED(CLOSED, static_cast<uint>(State::CLOSED)); |
| 116 | JSG_STATIC_METHOD(from); |
| 117 | |
| 118 | // EventSource is not defined by the spec as being disposable using ERM, but |
| 119 | // it makes sense to do so. The dispose operation simply defers to close(). |
| 120 | // This will enable `using eventsource = new EventSource(...)` |
| 121 | JSG_DISPOSE(close); |
| 122 | } |
| 123 | |
| 124 | struct PendingMessage { |
| 125 | kj::Vector<kj::String> data; |
| 126 | kj::Maybe<kj::String> event; |
| 127 | kj::String id; |
| 128 | }; |
| 129 | |
| 130 | // Called by the internal implementation to notify the EventSource about messages |
| 131 | // received from the server. |
| 132 | void enqueueMessages(kj::Array<PendingMessage> messages); |
| 133 | |
| 134 | // Called by the internal implementation to notify the EventSource that the server |
| 135 | // has provided a new reconnection time. |
| 136 | void setReconnectionTime(uint32_t time); |
| 137 | |
| 138 | // Called by the internal implementation to retrieve the last event id that was |
| 139 | // specified by the server. |
| 140 | kj::StringPtr getLastEventId(); |
| 141 | |
| 142 | // Called by the internal implementation to set the last event id that was specified |
| 143 | // by the server. |
| 144 | void setLastEventId(kj::String id); |
| 145 | |
| 146 | void visitForGc(jsg::GcVisitor& visitor); |
| 147 | void visitForMemoryInfo(jsg::MemoryTracker& tracker) const; |
| 148 | |
| 149 | private: |
| 150 | IoContext& context; |
| 151 | struct FetchImpl { |
| 152 | jsg::Url url; |
| 153 | EventSourceInit options; |
| 154 | // Indicates that the server previously responded with no content after a |
| 155 | // successful connection. This is likely indicative of a bug on the server. |
| 156 | // If this happens once, we'll try to reconnect. If it happens again, we'll |
| 157 | // fail the connection. |
| 158 | bool previousNoBody = false; |
| 159 | }; |
| 160 | // Used when the EventSource is created using the constructor. This |
| 161 | // is the normal mode of operation, when the EventSource uses fetch |
| 162 | // under the covers to connect, and reconnect, to the server. This |
| 163 | // will be kj::none when the EventSource is created using the from() |
| 164 | // method. |
| 165 | kj::Maybe<FetchImpl> impl; |
| 166 | jsg::Ref<AbortController> abortController; |
| 167 | State readyState; |
| 168 | kj::String lastEventId; |
| 169 | |
| 170 | // Indicates that the close method has been previously called. |
| 171 | bool closeCalled = false; |
| 172 | |
| 173 | // The EventSource spec defines onopen, onmessage, and onerror as prototype |
| 174 | // properties on the class. |
| 175 | kj::Maybe<jsg::JsRef<jsg::JsValue>> onopenValue; |
| 176 | kj::Maybe<jsg::JsRef<jsg::JsValue>> onmessageValue; |
| 177 | kj::Maybe<jsg::JsRef<jsg::JsValue>> onerrorValue; |
| 178 | |
| 179 | // The default reconnection wait time. This is fairly arbitrary and is left |
| 180 | // entirely up to the implementation. The event stream can provide a new value. |
| 181 | static constexpr auto DEFAULT_RECONNECTION_TIME = 2 * kj::SECONDS; |
| 182 | static constexpr uint32_t MIN_RECONNECTION_TIME = 1000; |
| 183 | static constexpr uint32_t MAX_RECONNECTION_TIME = 10 * 1000; |
| 184 | |
| 185 | kj::Duration reconnectionTime = DEFAULT_RECONNECTION_TIME; |
| 186 | |
| 187 | void notifyOpen(jsg::Lock& js); |
| 188 | void notifyError(jsg::Lock& js, const jsg::JsValue& error, bool reconnecting = false); |
| 189 | void notifyMessages(jsg::Lock& js, kj::Array<PendingMessage> messages); |
| 190 | |
| 191 | // The run() method handles the actual processing of the stream. |
| 192 | void run(jsg::Lock& js, |
| 193 | jsg::Ref<ReadableStream> stream, |
| 194 | bool withReconnection = true, |
| 195 | kj::Maybe<jsg::Ref<Response>> response = kj::none, |
| 196 | kj::Maybe<jsg::Ref<Fetcher>> fetcher = kj::none); |
| 197 | // The start() method initializes the fetch and the processing of the |
| 198 | // stream by calling run. |
| 199 | void start(jsg::Lock& js); |
| 200 | void reconnect(jsg::Lock& js); |
| 201 | }; |
| 202 | |
| 203 | } // namespace workerd::api |
| 204 | |
| 205 | #define EW_EVENTSOURCE_ISOLATE_TYPES api::EventSource, api::EventSource::EventSourceInit |