// Copyright (c) 2017-2024 Cloudflare, Inc. // Licensed under the Apache 2.0 license found in the LICENSE file or at: // https://opensource.org/licenses/Apache-2.0 #pragma once #include "basics.h" #include "events.h" #include "http.h" #include #include namespace workerd::api { using kj::uint; class Fetcher; class ReadableStream; class Response; // Implements the web standard EventSource API // https://developer.mozilla.org/en-US/docs/Web/API/EventSource class EventSource: public EventTarget { public: struct EventSourceInit { // We don't actually make use of the standard withCredentials option. If this is set to // any truthy value, we'll throw. jsg::Optional withCredentials; // This is a non-standard workers-specific extension that allows the EventSource to // use a custom Fetcher instance. jsg::Optional> fetcher; JSG_STRUCT(withCredentials, fetcher); }; enum class State { CONNECTING = 0, OPEN = 1, CLOSED = 2, }; EventSource(jsg::Lock& js, jsg::Url url, kj::Maybe init = kj::none); EventSource(jsg::Lock& js); static jsg::Ref constructor( jsg::Lock& js, kj::String url, jsg::Optional init); kj::ArrayPtr getUrl() const { KJ_IF_SOME(i, impl) { return i.url.getHref(); } return nullptr; } bool getWithCredentials() const { return false; } uint getReadyState() const { return static_cast(readyState); } void close(jsg::Lock& js); // A non-standard extension that creates an EventSource instance around a ReadableStream // instance. In this instance, automatic reconnection is disabled since there is no URL // or underlying fetch used. The ReadableStream instance must produce bytes. It will be // locked and disturbed, and will be read until it either ends or errors. Calling close() // will cause the stream to be canceled. static jsg::Ref from(jsg::Lock& js, jsg::Ref stream); kj::Maybe getOnOpen(jsg::Lock& js) { return onopenValue.map( [&](jsg::JsRef& ref) -> jsg::JsValue { return ref.getHandle(js); }); } void setOnOpen(jsg::Lock& js, jsg::JsValue value) { if (!value.isObject() && !value.isFunction()) { onopenValue = kj::none; } else { onopenValue = jsg::JsRef(js, value); } } kj::Maybe getOnMessage(jsg::Lock& js) { return onmessageValue.map( [&](jsg::JsRef& ref) -> jsg::JsValue { return ref.getHandle(js); }); } void setOnMessage(jsg::Lock& js, jsg::JsValue value) { if (!value.isObject() && !value.isFunction()) { onmessageValue = kj::none; } else { onmessageValue = jsg::JsRef(js, value); } } kj::Maybe getOnError(jsg::Lock& js) { return onerrorValue.map( [&](jsg::JsRef& ref) -> jsg::JsValue { return ref.getHandle(js); }); } void setOnError(jsg::Lock& js, jsg::JsValue value) { if (!value.isObject() && !value.isFunction()) { onerrorValue = kj::none; } else { onerrorValue = jsg::JsRef(js, value); } } JSG_RESOURCE_TYPE(EventSource) { JSG_INHERIT(EventTarget); JSG_METHOD(close); JSG_READONLY_PROTOTYPE_PROPERTY(url, getUrl); JSG_READONLY_PROTOTYPE_PROPERTY(withCredentials, getWithCredentials); JSG_READONLY_PROTOTYPE_PROPERTY(readyState, getReadyState); JSG_PROTOTYPE_PROPERTY(onopen, getOnOpen, setOnOpen); JSG_PROTOTYPE_PROPERTY(onmessage, getOnMessage, setOnMessage); JSG_PROTOTYPE_PROPERTY(onerror, getOnError, setOnError); JSG_STATIC_CONSTANT_NAMED(CONNECTING, static_cast(State::CONNECTING)); JSG_STATIC_CONSTANT_NAMED(OPEN, static_cast(State::OPEN)); JSG_STATIC_CONSTANT_NAMED(CLOSED, static_cast(State::CLOSED)); JSG_STATIC_METHOD(from); // EventSource is not defined by the spec as being disposable using ERM, but // it makes sense to do so. The dispose operation simply defers to close(). // This will enable `using eventsource = new EventSource(...)` JSG_DISPOSE(close); } struct PendingMessage { kj::Vector data; kj::Maybe event; kj::String id; }; // Called by the internal implementation to notify the EventSource about messages // received from the server. void enqueueMessages(kj::Array messages); // Called by the internal implementation to notify the EventSource that the server // has provided a new reconnection time. void setReconnectionTime(uint32_t time); // Called by the internal implementation to retrieve the last event id that was // specified by the server. kj::StringPtr getLastEventId(); // Called by the internal implementation to set the last event id that was specified // by the server. void setLastEventId(kj::String id); void visitForGc(jsg::GcVisitor& visitor); void visitForMemoryInfo(jsg::MemoryTracker& tracker) const; private: IoContext& context; struct FetchImpl { jsg::Url url; EventSourceInit options; // Indicates that the server previously responded with no content after a // successful connection. This is likely indicative of a bug on the server. // If this happens once, we'll try to reconnect. If it happens again, we'll // fail the connection. bool previousNoBody = false; }; // Used when the EventSource is created using the constructor. This // is the normal mode of operation, when the EventSource uses fetch // under the covers to connect, and reconnect, to the server. This // will be kj::none when the EventSource is created using the from() // method. kj::Maybe impl; jsg::Ref abortController; State readyState; kj::String lastEventId; // Indicates that the close method has been previously called. bool closeCalled = false; // The EventSource spec defines onopen, onmessage, and onerror as prototype // properties on the class. kj::Maybe> onopenValue; kj::Maybe> onmessageValue; kj::Maybe> onerrorValue; // The default reconnection wait time. This is fairly arbitrary and is left // entirely up to the implementation. The event stream can provide a new value. static constexpr auto DEFAULT_RECONNECTION_TIME = 2 * kj::SECONDS; static constexpr uint32_t MIN_RECONNECTION_TIME = 1000; static constexpr uint32_t MAX_RECONNECTION_TIME = 10 * 1000; kj::Duration reconnectionTime = DEFAULT_RECONNECTION_TIME; void notifyOpen(jsg::Lock& js); void notifyError(jsg::Lock& js, const jsg::JsValue& error, bool reconnecting = false); void notifyMessages(jsg::Lock& js, kj::Array messages); // The run() method handles the actual processing of the stream. void run(jsg::Lock& js, jsg::Ref stream, bool withReconnection = true, kj::Maybe> response = kj::none, kj::Maybe> fetcher = kj::none); // The start() method initializes the fetch and the processing of the // stream by calling run. void start(jsg::Lock& js); void reconnect(jsg::Lock& js); }; } // namespace workerd::api #define EW_EVENTSOURCE_ISOLATE_TYPES api::EventSource, api::EventSource::EventSourceInit