| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
1 parent 07a6770 commit 45e28a8
9 files changed
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -181,7 +181,7 @@ void JSStream::DoAfterWrite(const FunctionCallbackInfo<Value>& args) { | |||
| 181 | 181 | ASSIGN_OR_RETURN_UNWRAP(&wrap, args.Holder()); | |
| 182 | 182 | ASSIGN_OR_RETURN_UNWRAP(&w, args[0].As<Object>()); | |
| 183 | 183 | ||
| 184 | - wrap->OnAfterWrite(w); | ||
| 184 | + w->Done(0); | ||
| 185 | 185 | } | |
| 186 | 186 | ||
| 187 | 187 | ||
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -975,9 +975,6 @@ inline void Http2Session::SetChunksSinceLastWrite(size_t n) { | |||
| 975 | 975 | ||
| 976 | 976 | WriteWrap* Http2Session::AllocateSend() { | |
| 977 | 977 | HandleScope scope(env()->isolate()); | |
| 978 | - auto AfterWrite = [](WriteWrap* req, int status) { | ||
| 979 | - req->Dispose(); | ||
| 980 | - }; | ||
| 981 | 978 | Local<Object> obj = | |
| 982 | 979 | env()->write_wrap_constructor_function() | |
| 983 | 980 | ->NewInstance(env()->context()).ToLocalChecked(); | |
@@ -987,7 +984,7 @@ WriteWrap* Http2Session::AllocateSend() { | |||
| 987 | 984 | session(), | |
| 988 | 985 | NGHTTP2_SETTINGS_MAX_FRAME_SIZE); | |
| 989 | 986 | // Max frame size + 9 bytes for the header | |
| 990 | - return WriteWrap::New(env(), obj, stream_, AfterWrite, size + 9); | ||
| 987 | + return WriteWrap::New(env(), obj, stream_, size + 9); | ||
| 991 | 988 | } | |
| 992 | 989 | ||
| 993 | 990 | void Http2Session::Send(WriteWrap* req, char* buf, size_t length) { | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -144,15 +144,19 @@ void StreamBase::JSMethod(const FunctionCallbackInfo<Value>& args) { | |||
| 144 | 144 | } | |
| 145 | 145 | ||
| 146 | 146 | ||
| 147 | + inline void ShutdownWrap::OnDone(int status) { | ||
| 148 | + stream()->AfterShutdown(this, status); | ||
| 149 | + } | ||
| 150 | + | ||
| 151 | + | ||
| 147 | 152 | WriteWrap* WriteWrap::New(Environment* env, | |
| 148 | 153 | Local<Object> obj, | |
| 149 | 154 | StreamBase* wrap, | |
| 150 | - DoneCb cb, | ||
| 151 | 155 | size_t extra) { | |
| 152 | 156 | size_t storage_size = ROUND_UP(sizeof(WriteWrap), kAlignSize) + extra; | |
| 153 | 157 | char* storage = new char[storage_size]; | |
| 154 | 158 | ||
| 155 | - return new(storage) WriteWrap(env, obj, wrap, cb, storage_size); | ||
| 159 | + return new(storage) WriteWrap(env, obj, wrap, storage_size); | ||
| 156 | 160 | } | |
| 157 | 161 | ||
| 158 | 162 | ||
@@ -172,6 +176,10 @@ size_t WriteWrap::ExtraSize() const { | |||
| 172 | 176 | return storage_size_ - ROUND_UP(sizeof(*this), kAlignSize); | |
| 173 | 177 | } | |
| 174 | 178 | ||
| 179 | + inline void WriteWrap::OnDone(int status) { | ||
| 180 | + stream()->AfterWrite(this, status); | ||
| 181 | + } | ||
| 182 | + | ||
| 175 | 183 | } // namespace node | |
| 176 | 184 | ||
| 177 | 185 | #endif // defined(NODE_WANT_INTERNALS) && NODE_WANT_INTERNALS | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -55,8 +55,7 @@ int StreamBase::Shutdown(const FunctionCallbackInfo<Value>& args) { | |||
| 55 | 55 | AsyncHooks::DefaultTriggerAsyncIdScope(env, wrap->get_async_id()); | |
| 56 | 56 | ShutdownWrap* req_wrap = new ShutdownWrap(env, | |
| 57 | 57 | req_wrap_obj, | |
| 58 | - this, | ||
| 59 | - AfterShutdown); | ||
| 58 | + this); | ||
| 60 | 59 | ||
| 61 | 60 | int err = DoShutdown(req_wrap); | |
| 62 | 61 | if (err) | |
@@ -66,7 +65,6 @@ int StreamBase::Shutdown(const FunctionCallbackInfo<Value>& args) { | |||
| 66 | 65 | ||
| 67 | 66 | ||
| 68 | 67 | void StreamBase::AfterShutdown(ShutdownWrap* req_wrap, int status) { | |
| 69 | - StreamBase* wrap = req_wrap->wrap(); | ||
| 70 | 68 | Environment* env = req_wrap->env(); | |
| 71 | 69 | ||
| 72 | 70 | // The wrap and request objects should still be there. | |
@@ -78,7 +76,7 @@ void StreamBase::AfterShutdown(ShutdownWrap* req_wrap, int status) { | |||
| 78 | 76 | Local<Object> req_wrap_obj = req_wrap->object(); | |
| 79 | 77 | Local<Value> argv[3] = { | |
| 80 | 78 | Integer::New(env->isolate(), status), | |
| 81 | - wrap->GetObject(), | ||
| 79 | + GetObject(), | ||
| 82 | 80 | req_wrap_obj | |
| 83 | 81 | }; | |
| 84 | 82 | ||
@@ -159,8 +157,7 @@ int StreamBase::Writev(const FunctionCallbackInfo<Value>& args) { | |||
| 159 | 157 | CHECK_NE(wrap, nullptr); | |
| 160 | 158 | AsyncHooks::DefaultTriggerAsyncIdScope trigger_scope(env, | |
| 161 | 159 | wrap->get_async_id()); | |
| 162 | - req_wrap = WriteWrap::New(env, req_wrap_obj, this, AfterWrite, | ||
| 163 | - storage_size); | ||
| 160 | + req_wrap = WriteWrap::New(env, req_wrap_obj, this, storage_size); | ||
| 164 | 161 | } | |
| 165 | 162 | ||
| 166 | 163 | offset = 0; | |
@@ -252,7 +249,7 @@ int StreamBase::WriteBuffer(const FunctionCallbackInfo<Value>& args) { | |||
| 252 | 249 | CHECK_NE(wrap, nullptr); | |
| 253 | 250 | AsyncHooks::DefaultTriggerAsyncIdScope trigger_scope(env, | |
| 254 | 251 | wrap->get_async_id()); | |
| 255 | - req_wrap = WriteWrap::New(env, req_wrap_obj, this, AfterWrite); | ||
| 252 | + req_wrap = WriteWrap::New(env, req_wrap_obj, this); | ||
| 256 | 253 | } | |
| 257 | 254 | ||
| 258 | 255 | err = DoWrite(req_wrap, bufs, count, nullptr); | |
@@ -338,8 +335,7 @@ int StreamBase::WriteString(const FunctionCallbackInfo<Value>& args) { | |||
| 338 | 335 | CHECK_NE(wrap, nullptr); | |
| 339 | 336 | AsyncHooks::DefaultTriggerAsyncIdScope trigger_scope(env, | |
| 340 | 337 | wrap->get_async_id()); | |
| 341 | - req_wrap = WriteWrap::New(env, req_wrap_obj, this, AfterWrite, | ||
| 342 | - storage_size); | ||
| 338 | + req_wrap = WriteWrap::New(env, req_wrap_obj, this, storage_size); | ||
| 343 | 339 | } | |
| 344 | 340 | ||
| 345 | 341 | data = req_wrap->Extra(); | |
@@ -401,7 +397,6 @@ int StreamBase::WriteString(const FunctionCallbackInfo<Value>& args) { | |||
| 401 | 397 | ||
| 402 | 398 | ||
| 403 | 399 | void StreamBase::AfterWrite(WriteWrap* req_wrap, int status) { | |
| 404 | - StreamBase* wrap = req_wrap->wrap(); | ||
| 405 | 400 | Environment* env = req_wrap->env(); | |
| 406 | 401 | ||
| 407 | 402 | HandleScope handle_scope(env->isolate()); | |
@@ -413,19 +408,19 @@ void StreamBase::AfterWrite(WriteWrap* req_wrap, int status) { | |||
| 413 | 408 | // Unref handle property | |
| 414 | 409 | Local<Object> req_wrap_obj = req_wrap->object(); | |
| 415 | 410 | req_wrap_obj->Delete(env->context(), env->handle_string()).FromJust(); | |
| 416 | - wrap->OnAfterWrite(req_wrap); | ||
| 411 | + OnAfterWrite(req_wrap, status); | ||
| 417 | 412 | ||
| 418 | 413 | Local<Value> argv[] = { | |
| 419 | 414 | Integer::New(env->isolate(), status), | |
| 420 | - wrap->GetObject(), | ||
| 415 | + GetObject(), | ||
| 421 | 416 | req_wrap_obj, | |
| 422 | 417 | Undefined(env->isolate()) | |
| 423 | 418 | }; | |
| 424 | 419 | ||
| 425 | - const char* msg = wrap->Error(); | ||
| 420 | + const char* msg = Error(); | ||
| 426 | 421 | if (msg != nullptr) { | |
| 427 | 422 | argv[3] = OneByteString(env->isolate(), msg); | |
| 428 | - wrap->ClearError(); | ||
| 423 | + ClearError(); | ||
| 429 | 424 | } | |
| 430 | 425 | ||
| 431 | 426 | if (req_wrap_obj->Has(env->context(), env->oncomplete_string()).FromJust()) | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -16,39 +16,37 @@ namespace node { | |||
| 16 | 16 | // Forward declarations | |
| 17 | 17 | class StreamBase; | |
| 18 | 18 | ||
| 19 | - template <class Req> | ||
| 19 | + template<typename Base> | ||
| 20 | 20 | class StreamReq { | |
| 21 | 21 | public: | |
| 22 | - typedef void (*DoneCb)(Req* req, int status); | ||
| 23 | - | ||
| 24 | - explicit StreamReq(DoneCb cb) : cb_(cb) { | ||
| 22 | + explicit StreamReq(StreamBase* stream) : stream_(stream) { | ||
| 25 | 23 | } | |
| 26 | 24 | ||
| 27 | 25 | inline void Done(int status, const char* error_str = nullptr) { | |
| 28 | - Req* req = static_cast<Req*>(this); | ||
| 26 | + Base* req = static_cast<Base*>(this); | ||
| 29 | 27 | Environment* env = req->env(); | |
| 30 | 28 | if (error_str != nullptr) { | |
| 31 | 29 | req->object()->Set(env->error_string(), | |
| 32 | 30 | OneByteString(env->isolate(), error_str)); | |
| 33 | 31 | } | |
| 34 | 32 | ||
| 35 | - cb_(req, status); | ||
| 33 | + req->OnDone(status); | ||
| 36 | 34 | } | |
| 37 | 35 | ||
| 36 | + inline StreamBase* stream() const { return stream_; } | ||
| 37 | + | ||
| 38 | 38 | private: | |
| 39 | - DoneCb cb_; | ||
| 39 | + StreamBase* const stream_; | ||
| 40 | 40 | }; | |
| 41 | 41 | ||
| 42 | 42 | class ShutdownWrap : public ReqWrap<uv_shutdown_t>, | |
| 43 | 43 | public StreamReq<ShutdownWrap> { | |
| 44 | 44 | public: | |
| 45 | 45 | ShutdownWrap(Environment* env, | |
| 46 | 46 | v8::Local<v8::Object> req_wrap_obj, | |
| 47 | - StreamBase* wrap, | ||
| 48 | - DoneCb cb) | ||
| 47 | + StreamBase* stream) | ||
| 49 | 48 | : ReqWrap(env, req_wrap_obj, AsyncWrap::PROVIDER_SHUTDOWNWRAP), | |
| 50 | - StreamReq<ShutdownWrap>(cb), | ||
| 51 | - wrap_(wrap) { | ||
| 49 | + StreamReq<ShutdownWrap>(stream) { | ||
| 52 | 50 | Wrap(req_wrap_obj, this); | |
| 53 | 51 | } | |
| 54 | 52 | ||
@@ -60,27 +58,22 @@ class ShutdownWrap : public ReqWrap<uv_shutdown_t>, | |||
| 60 | 58 | return ContainerOf(&ShutdownWrap::req_, req); | |
| 61 | 59 | } | |
| 62 | 60 | ||
| 63 | - inline StreamBase* wrap() const { return wrap_; } | ||
| 64 | 61 | size_t self_size() const override { return sizeof(*this); } | |
| 65 | 62 | ||
| 66 | - private: | ||
| 67 | - StreamBase* const wrap_; | ||
| 63 | + inline void OnDone(int status); // Just calls stream()->AfterShutdown() | ||
| 68 | 64 | }; | |
| 69 | 65 | ||
| 70 | - class WriteWrap: public ReqWrap<uv_write_t>, | ||
| 71 | - public StreamReq<WriteWrap> { | ||
| 66 | + class WriteWrap : public ReqWrap<uv_write_t>, | ||
| 67 | + public StreamReq<WriteWrap> { | ||
| 72 | 68 | public: | |
| 73 | 69 | static inline WriteWrap* New(Environment* env, | |
| 74 | 70 | v8::Local<v8::Object> obj, | |
| 75 | - StreamBase* wrap, | ||
| 76 | - DoneCb cb, | ||
| 71 | + StreamBase* stream, | ||
| 77 | 72 | size_t extra = 0); | |
| 78 | 73 | inline void Dispose(); | |
| 79 | 74 | inline char* Extra(size_t offset = 0); | |
| 80 | 75 | inline size_t ExtraSize() const; | |
| 81 | 76 | ||
| 82 | - inline StreamBase* wrap() const { return wrap_; } | ||
| 83 | - | ||
| 84 | 77 | size_t self_size() const override { return storage_size_; } | |
| 85 | 78 | ||
| 86 | 79 | static WriteWrap* from_req(uv_write_t* req) { | |
@@ -91,24 +84,22 @@ class WriteWrap: public ReqWrap<uv_write_t>, | |||
| 91 | 84 | ||
| 92 | 85 | WriteWrap(Environment* env, | |
| 93 | 86 | v8::Local<v8::Object> obj, | |
| 94 | - StreamBase* wrap, | ||
| 95 | - DoneCb cb) | ||
| 87 | + StreamBase* stream) | ||
| 96 | 88 | : ReqWrap(env, obj, AsyncWrap::PROVIDER_WRITEWRAP), | |
| 97 | - StreamReq<WriteWrap>(cb), | ||
| 98 | - wrap_(wrap), | ||
| 89 | + StreamReq<WriteWrap>(stream), | ||
| 99 | 90 | storage_size_(0) { | |
| 100 | 91 | Wrap(obj, this); | |
| 101 | 92 | } | |
| 102 | 93 | ||
| 94 | + inline void OnDone(int status); // Just calls stream()->AfterWrite() | ||
| 95 | + | ||
| 103 | 96 | protected: | |
| 104 | 97 | WriteWrap(Environment* env, | |
| 105 | 98 | v8::Local<v8::Object> obj, | |
| 106 | - StreamBase* wrap, | ||
| 107 | - DoneCb cb, | ||
| 99 | + StreamBase* stream, | ||
| 108 | 100 | size_t storage_size) | |
| 109 | 101 | : ReqWrap(env, obj, AsyncWrap::PROVIDER_WRITEWRAP), | |
| 110 | - StreamReq<WriteWrap>(cb), | ||
| 111 | - wrap_(wrap), | ||
| 102 | + StreamReq<WriteWrap>(stream), | ||
| 112 | 103 | storage_size_(storage_size) { | |
| 113 | 104 | Wrap(obj, this); | |
| 114 | 105 | } | |
@@ -129,7 +120,6 @@ class WriteWrap: public ReqWrap<uv_write_t>, | |||
| 129 | 120 | // WriteWrap. Ensure this never happens. | |
| 130 | 121 | void operator delete(void* ptr) { UNREACHABLE(); } | |
| 131 | 122 | ||
| 132 | - StreamBase* const wrap_; | ||
| 133 | 123 | const size_t storage_size_; | |
| 134 | 124 | }; | |
| 135 | 125 | ||
@@ -151,7 +141,7 @@ class StreamResource { | |||
| 151 | 141 | void* ctx; | |
| 152 | 142 | }; | |
| 153 | 143 | ||
| 154 | - typedef void (*AfterWriteCb)(WriteWrap* w, void* ctx); | ||
| 144 | + typedef void (*AfterWriteCb)(WriteWrap* w, int status, void* ctx); | ||
| 155 | 145 | typedef void (*AllocCb)(size_t size, uv_buf_t* buf, void* ctx); | |
| 156 | 146 | typedef void (*ReadCb)(ssize_t nread, | |
| 157 | 147 | const uv_buf_t* buf, | |
@@ -176,9 +166,9 @@ class StreamResource { | |||
| 176 | 166 | virtual void ClearError(); | |
| 177 | 167 | ||
| 178 | 168 | // Events | |
| 179 | - inline void OnAfterWrite(WriteWrap* w) { | ||
| 169 | + inline void OnAfterWrite(WriteWrap* w, int status) { | ||
| 180 | 170 | if (!after_write_cb_.is_empty()) | |
| 181 | - after_write_cb_.fn(w, after_write_cb_.ctx); | ||
| 171 | + after_write_cb_.fn(w, status, after_write_cb_.ctx); | ||
| 182 | 172 | } | |
| 183 | 173 | ||
| 184 | 174 | inline void OnAlloc(size_t size, uv_buf_t* buf) { | |
@@ -208,14 +198,12 @@ class StreamResource { | |||
| 208 | 198 | inline Callback<ReadCb> read_cb() { return read_cb_; } | |
| 209 | 199 | inline Callback<DestructCb> destruct_cb() { return destruct_cb_; } | |
| 210 | 200 | ||
| 211 | - private: | ||
| 201 | + protected: | ||
| 212 | 202 | Callback<AfterWriteCb> after_write_cb_; | |
| 213 | 203 | Callback<AllocCb> alloc_cb_; | |
| 214 | 204 | Callback<ReadCb> read_cb_; | |
| 215 | 205 | Callback<DestructCb> destruct_cb_; | |
| 216 | 206 | uint64_t bytes_read_; | |
| 217 | - | ||
| 218 | - friend class StreamBase; | ||
| 219 | 207 | }; | |
| 220 | 208 | ||
| 221 | 209 | class StreamBase : public StreamResource { | |
@@ -257,6 +245,10 @@ class StreamBase : public StreamResource { | |||
| 257 | 245 | v8::Local<v8::Object> buf, | |
| 258 | 246 | v8::Local<v8::Object> handle); | |
| 259 | 247 | ||
| 248 | + // These are called by the respective {Write,Shutdown}Wrap class. | ||
| 249 | + virtual void AfterShutdown(ShutdownWrap* req, int status); | ||
| 250 | + virtual void AfterWrite(WriteWrap* req, int status); | ||
| 251 | + | ||
| 260 | 252 | protected: | |
| 261 | 253 | explicit StreamBase(Environment* env) : env_(env), consumed_(false) { | |
| 262 | 254 | } | |
@@ -267,10 +259,6 @@ class StreamBase : public StreamResource { | |||
| 267 | 259 | virtual AsyncWrap* GetAsyncWrap() = 0; | |
| 268 | 260 | virtual v8::Local<v8::Object> GetObject(); | |
| 269 | 261 | ||
| 270 | - // Libuv callbacks | ||
| 271 | - static void AfterShutdown(ShutdownWrap* req, int status); | ||
| 272 | - static void AfterWrite(WriteWrap* req, int status); | ||
| 273 | - | ||
| 274 | 262 | // JS Methods | |
| 275 | 263 | int ReadStart(const v8::FunctionCallbackInfo<v8::Value>& args); | |
| 276 | 264 | int ReadStop(const v8::FunctionCallbackInfo<v8::Value>& args); | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -92,7 +92,6 @@ LibuvStreamWrap::LibuvStreamWrap(Environment* env, | |||
| 92 | 92 | provider), | |
| 93 | 93 | StreamBase(env), | |
| 94 | 94 | stream_(stream) { | |
| 95 | - set_after_write_cb({ OnAfterWriteImpl, this }); | ||
| 96 | 95 | set_alloc_cb({ OnAllocImpl, this }); | |
| 97 | 96 | set_read_cb({ OnReadImpl, this }); | |
| 98 | 97 | } | |
@@ -299,13 +298,13 @@ void LibuvStreamWrap::SetBlocking(const FunctionCallbackInfo<Value>& args) { | |||
| 299 | 298 | ||
| 300 | 299 | int LibuvStreamWrap::DoShutdown(ShutdownWrap* req_wrap) { | |
| 301 | 300 | int err; | |
| 302 | - err = uv_shutdown(req_wrap->req(), stream(), AfterShutdown); | ||
| 301 | + err = uv_shutdown(req_wrap->req(), stream(), AfterUvShutdown); | ||
| 303 | 302 | req_wrap->Dispatched(); | |
| 304 | 303 | return err; | |
| 305 | 304 | } | |
| 306 | 305 | ||
| 307 | 306 | ||
| 308 | - void LibuvStreamWrap::AfterShutdown(uv_shutdown_t* req, int status) { | ||
| 307 | + void LibuvStreamWrap::AfterUvShutdown(uv_shutdown_t* req, int status) { | ||
| 309 | 308 | ShutdownWrap* req_wrap = ShutdownWrap::from_req(req); | |
| 310 | 309 | CHECK_NE(req_wrap, nullptr); | |
| 311 | 310 | HandleScope scope(req_wrap->env()->isolate()); | |
@@ -360,9 +359,9 @@ int LibuvStreamWrap::DoWrite(WriteWrap* w, | |||
| 360 | 359 | uv_stream_t* send_handle) { | |
| 361 | 360 | int r; | |
| 362 | 361 | if (send_handle == nullptr) { | |
| 363 | - r = uv_write(w->req(), stream(), bufs, count, AfterWrite); | ||
| 362 | + r = uv_write(w->req(), stream(), bufs, count, AfterUvWrite); | ||
| 364 | 363 | } else { | |
| 365 | - r = uv_write2(w->req(), stream(), bufs, count, send_handle, AfterWrite); | ||
| 364 | + r = uv_write2(w->req(), stream(), bufs, count, send_handle, AfterUvWrite); | ||
| 366 | 365 | } | |
| 367 | 366 | ||
| 368 | 367 | if (!r) { | |
@@ -383,7 +382,7 @@ int LibuvStreamWrap::DoWrite(WriteWrap* w, | |||
| 383 | 382 | } | |
| 384 | 383 | ||
| 385 | 384 | ||
| 386 | - void LibuvStreamWrap::AfterWrite(uv_write_t* req, int status) { | ||
| 385 | + void LibuvStreamWrap::AfterUvWrite(uv_write_t* req, int status) { | ||
| 387 | 386 | WriteWrap* req_wrap = WriteWrap::from_req(req); | |
| 388 | 387 | CHECK_NE(req_wrap, nullptr); | |
| 389 | 388 | HandleScope scope(req_wrap->env()->isolate()); | |
@@ -392,9 +391,9 @@ void LibuvStreamWrap::AfterWrite(uv_write_t* req, int status) { | |||
| 392 | 391 | } | |
| 393 | 392 | ||
| 394 | 393 | ||
| 395 | - void LibuvStreamWrap::OnAfterWriteImpl(WriteWrap* w, void* ctx) { | ||
| 396 | - LibuvStreamWrap* wrap = static_cast<LibuvStreamWrap*>(ctx); | ||
| 397 | - wrap->UpdateWriteQueueSize(); | ||
| 394 | + void LibuvStreamWrap::AfterWrite(WriteWrap* w, int status) { | ||
| 395 | + StreamBase::AfterWrite(w, status); | ||
| 396 | + UpdateWriteQueueSize(); | ||
| 398 | 397 | } | |
| 399 | 398 | ||
| 400 | 399 | } // namespace node | |
| Back | FazBrowse Home | New Git URL |
0 commit comments