|
@@ -18,32 +18,19 @@ NodeStreamLoader::NodeStreamLoader(network::ResourceResponseHead head,
|
|
|
network::mojom::URLLoaderClientPtr client,
|
|
|
v8::Isolate* isolate,
|
|
|
v8::Local<v8::Object> emitter)
|
|
|
- : binding_(this),
|
|
|
+ : binding_(this, std::move(loader)),
|
|
|
client_(std::move(client)),
|
|
|
isolate_(isolate),
|
|
|
emitter_(isolate, emitter),
|
|
|
weak_factory_(this) {
|
|
|
- auto weak = weak_factory_.GetWeakPtr();
|
|
|
- binding_.Bind(std::move(loader));
|
|
|
binding_.set_connection_error_handler(
|
|
|
- base::BindOnce(&NodeStreamLoader::OnConnectionError, weak));
|
|
|
-
|
|
|
- mojo::ScopedDataPipeConsumerHandle consumer;
|
|
|
- MojoResult rv = mojo::CreateDataPipe(nullptr, &producer_, &consumer);
|
|
|
- if (rv != MOJO_RESULT_OK) {
|
|
|
- OnError(nullptr);
|
|
|
- return;
|
|
|
- }
|
|
|
-
|
|
|
- client_->OnReceiveResponse(head);
|
|
|
- client_->OnStartLoadingResponseBody(std::move(consumer));
|
|
|
+ base::BindOnce(&NodeStreamLoader::NotifyComplete,
|
|
|
+ weak_factory_.GetWeakPtr(), net::ERR_FAILED));
|
|
|
|
|
|
- On("end", base::BindRepeating(&NodeStreamLoader::OnEnd, weak));
|
|
|
- On("error", base::BindRepeating(&NodeStreamLoader::OnError, weak));
|
|
|
- // Since every node::MakeCallback call has a micro scope itself, we have to
|
|
|
- // subscribe |data| at last otherwise |end|'s listener won't be called when
|
|
|
- // it is emitted in the same tick.
|
|
|
- On("data", base::BindRepeating(&NodeStreamLoader::OnData, weak));
|
|
|
+ // PostTask since it might destruct.
|
|
|
+ base::SequencedTaskRunnerHandle::Get()->PostTask(
|
|
|
+ FROM_HERE, base::BindOnce(&NodeStreamLoader::Start,
|
|
|
+ weak_factory_.GetWeakPtr(), std::move(head)));
|
|
|
}
|
|
|
|
|
|
NodeStreamLoader::~NodeStreamLoader() {
|
|
@@ -58,62 +45,98 @@ NodeStreamLoader::~NodeStreamLoader() {
|
|
|
node::MakeCallback(isolate_, emitter_.Get(isolate_), "removeListener",
|
|
|
node::arraysize(args), args, {0, 0});
|
|
|
}
|
|
|
+
|
|
|
+ // Release references.
|
|
|
+ emitter_.Reset();
|
|
|
+ buffer_.Reset();
|
|
|
}
|
|
|
|
|
|
-void NodeStreamLoader::On(const char* event, EventCallback callback) {
|
|
|
- v8::Locker locker(isolate_);
|
|
|
- v8::Isolate::Scope isolate_scope(isolate_);
|
|
|
- v8::HandleScope handle_scope(isolate_);
|
|
|
+void NodeStreamLoader::Start(network::ResourceResponseHead head) {
|
|
|
+ mojo::ScopedDataPipeProducerHandle producer;
|
|
|
+ mojo::ScopedDataPipeConsumerHandle consumer;
|
|
|
+ MojoResult rv = mojo::CreateDataPipe(nullptr, &producer, &consumer);
|
|
|
+ if (rv != MOJO_RESULT_OK) {
|
|
|
+ NotifyComplete(net::ERR_INSUFFICIENT_RESOURCES);
|
|
|
+ return;
|
|
|
+ }
|
|
|
|
|
|
- // emitter.on(event, callback)
|
|
|
- v8::Local<v8::Value> args[] = {
|
|
|
- mate::StringToV8(isolate_, event),
|
|
|
- mate::CallbackToV8(isolate_, std::move(callback)),
|
|
|
- };
|
|
|
- node::MakeCallback(isolate_, emitter_.Get(isolate_), "on",
|
|
|
- node::arraysize(args), args, {0, 0});
|
|
|
+ producer_ =
|
|
|
+ std::make_unique<mojo::StringDataPipeProducer>(std::move(producer));
|
|
|
|
|
|
- handlers_[event].Reset(isolate_, args[1]);
|
|
|
+ client_->OnReceiveResponse(head);
|
|
|
+ client_->OnStartLoadingResponseBody(std::move(consumer));
|
|
|
+
|
|
|
+ auto weak = weak_factory_.GetWeakPtr();
|
|
|
+ On("end",
|
|
|
+ base::BindRepeating(&NodeStreamLoader::NotifyComplete, weak, net::OK));
|
|
|
+ On("error", base::BindRepeating(&NodeStreamLoader::NotifyComplete, weak,
|
|
|
+ net::ERR_FAILED));
|
|
|
+ On("readable", base::BindRepeating(&NodeStreamLoader::ReadMore, weak));
|
|
|
}
|
|
|
|
|
|
-void NodeStreamLoader::OnData(mate::Arguments* args) {
|
|
|
- v8::Local<v8::Value> buffer;
|
|
|
- args->GetNext(&buffer);
|
|
|
- if (!node::Buffer::HasInstance(buffer)) {
|
|
|
- args->ThrowError("data must be Buffer");
|
|
|
+void NodeStreamLoader::NotifyComplete(int result) {
|
|
|
+ // Wait until write finishes or fails.
|
|
|
+ if (is_writing_) {
|
|
|
+ ended_ = true;
|
|
|
+ result_ = result;
|
|
|
return;
|
|
|
}
|
|
|
|
|
|
- size_t ssize = node::Buffer::Length(buffer);
|
|
|
- uint32_t size = base::saturated_cast<uint32_t>(ssize);
|
|
|
- MojoResult result = producer_->WriteData(node::Buffer::Data(buffer), &size,
|
|
|
- MOJO_WRITE_DATA_FLAG_NONE);
|
|
|
- if (result != MOJO_RESULT_OK || size < ssize) {
|
|
|
- OnError(nullptr);
|
|
|
- return;
|
|
|
- }
|
|
|
+ client_->OnComplete(network::URLLoaderCompletionStatus(result));
|
|
|
+ delete this;
|
|
|
}
|
|
|
|
|
|
-void NodeStreamLoader::OnEnd(mate::Arguments* args) {
|
|
|
- client_->OnComplete(network::URLLoaderCompletionStatus(net::OK));
|
|
|
- client_.reset();
|
|
|
- MaybeDeleteSelf();
|
|
|
-}
|
|
|
+void NodeStreamLoader::ReadMore() {
|
|
|
+ // buffer = emitter.read()
|
|
|
+ v8::MaybeLocal<v8::Value> ret = node::MakeCallback(
|
|
|
+ isolate_, emitter_.Get(isolate_), "read", 0, nullptr, {0, 0});
|
|
|
|
|
|
-void NodeStreamLoader::OnError(mate::Arguments* args) {
|
|
|
- client_->OnComplete(network::URLLoaderCompletionStatus(net::ERR_FAILED));
|
|
|
- client_.reset();
|
|
|
- MaybeDeleteSelf();
|
|
|
+ // If there is no buffer read, wait until |readable| is emitted again.
|
|
|
+ v8::Local<v8::Value> buffer;
|
|
|
+ if (!ret.ToLocal(&buffer) || !node::Buffer::HasInstance(buffer))
|
|
|
+ return;
|
|
|
+
|
|
|
+ // Hold the buffer until the write is done.
|
|
|
+ buffer_.Reset(isolate_, buffer);
|
|
|
+
|
|
|
+ // Write buffer to mojo pipe asyncronously.
|
|
|
+ is_writing_ = true;
|
|
|
+ producer_->Write(
|
|
|
+ base::StringPiece(node::Buffer::Data(buffer),
|
|
|
+ node::Buffer::Length(buffer)),
|
|
|
+ mojo::StringDataPipeProducer::AsyncWritingMode::
|
|
|
+ STRING_STAYS_VALID_UNTIL_COMPLETION,
|
|
|
+ base::BindOnce(&NodeStreamLoader::DidWrite, weak_factory_.GetWeakPtr()));
|
|
|
}
|
|
|
|
|
|
-void NodeStreamLoader::OnConnectionError() {
|
|
|
- binding_.Close();
|
|
|
- MaybeDeleteSelf();
|
|
|
+void NodeStreamLoader::DidWrite(MojoResult result) {
|
|
|
+ is_writing_ = false;
|
|
|
+ // We were told to end streaming.
|
|
|
+ if (ended_) {
|
|
|
+ NotifyComplete(result_);
|
|
|
+ return;
|
|
|
+ }
|
|
|
+
|
|
|
+ if (result == MOJO_RESULT_OK)
|
|
|
+ ReadMore();
|
|
|
+ else
|
|
|
+ NotifyComplete(net::ERR_FAILED);
|
|
|
}
|
|
|
|
|
|
-void NodeStreamLoader::MaybeDeleteSelf() {
|
|
|
- if (!binding_.is_bound() && !client_.is_bound())
|
|
|
- delete this;
|
|
|
+void NodeStreamLoader::On(const char* event, EventCallback callback) {
|
|
|
+ v8::Locker locker(isolate_);
|
|
|
+ v8::Isolate::Scope isolate_scope(isolate_);
|
|
|
+ v8::HandleScope handle_scope(isolate_);
|
|
|
+
|
|
|
+ // emitter.on(event, callback)
|
|
|
+ v8::Local<v8::Value> args[] = {
|
|
|
+ mate::StringToV8(isolate_, event),
|
|
|
+ mate::CallbackToV8(isolate_, std::move(callback)),
|
|
|
+ };
|
|
|
+ handlers_[event].Reset(isolate_, args[1]);
|
|
|
+ node::MakeCallback(isolate_, emitter_.Get(isolate_), "on",
|
|
|
+ node::arraysize(args), args, {0, 0});
|
|
|
+ // No more code bellow, as this class may destruct when subscribing.
|
|
|
}
|
|
|
|
|
|
} // namespace atom
|