Skip to content
File

Blob: src/workerd/api/eventsource.h

cpp206 lines
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 
13namespace workerd::api {
14 
15using kj::uint;
16class Fetcher;
17class ReadableStream;
18class Response;
19 
20// Implements the web standard EventSource API
21// https://developer.mozilla.org/en-US/docs/Web/API/EventSource
22class 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