This commit is contained in:
Brian Smith 2023-09-27 19:24:08 -05:00 committed by roytam1
commit 3979e4847c
2 changed files with 182 additions and 4 deletions

View file

@ -8,6 +8,9 @@
#include "nsITransport.h"
#include "nsIStreamTransportService.h"
#include "nsProxyRelease.h"
#include "WorkerPrivate.h"
#include "WorkerRunnable.h"
#include "Workers.h"
#include "mozilla/dom/DOMError.h"
@ -19,7 +22,67 @@ static NS_DEFINE_CID(kStreamTransportServiceCID,
namespace mozilla {
namespace dom {
NS_IMPL_ISUPPORTS(FetchStream, nsIInputStreamCallback)
using namespace workers;
namespace {
class FetchStreamWorkerHolder final : public WorkerHolder
{
public:
explicit FetchStreamWorkerHolder(FetchStream* aStream)
: WorkerHolder()
, mStream(aStream)
, mWasNotified(false)
{}
bool Notify(Status aStatus) override
{
if (!mWasNotified) {
mWasNotified = true;
mStream->Close();
}
return true;
}
WorkerPrivate* GetWorkerPrivate() const
{
return mWorkerPrivate;
}
private:
RefPtr<FetchStream> mStream;
bool mWasNotified;
};
class FetchStreamWorkerHolderShutdown final : public WorkerControlRunnable
{
public:
FetchStreamWorkerHolderShutdown(WorkerPrivate* aWorkerPrivate,
UniquePtr<WorkerHolder>&& aHolder,
nsCOMPtr<nsIGlobalObject>&& aGlobal)
: WorkerControlRunnable(aWorkerPrivate)
, mHolder(Move(aHolder))
, mGlobal(Move(aGlobal))
{}
bool
WorkerRun(JSContext* aCx, WorkerPrivate* aWorkerPrivate) override
{
mHolder = nullptr;
mGlobal = nullptr;
return true;
}
private:
UniquePtr<WorkerHolder> mHolder;
nsCOMPtr<nsIGlobalObject> mGlobal;
};
} // anonymous
NS_IMPL_ISUPPORTS(FetchStream, nsIInputStreamCallback, nsIObserver,
nsISupportsWeakReference)
/* static */ JSObject*
FetchStream::Create(JSContext* aCx, nsIGlobalObject* aGlobal,
@ -30,6 +93,35 @@ FetchStream::Create(JSContext* aCx, nsIGlobalObject* aGlobal,
RefPtr<FetchStream> stream = new FetchStream(aGlobal, aInputStream);
if (NS_IsMainThread()) {
nsCOMPtr<nsIObserverService> os = mozilla::services::GetObserverService();
if (NS_WARN_IF(!os)) {
aRv.Throw(NS_ERROR_FAILURE);
return nullptr;
}
aRv = os->AddObserver(stream, DOM_WINDOW_DESTROYED_TOPIC, true);
if (NS_WARN_IF(aRv.Failed())) {
return nullptr;
}
} else {
WorkerPrivate* workerPrivate = GetWorkerPrivateFromContext(aCx);
MOZ_ASSERT(workerPrivate);
UniquePtr<FetchStreamWorkerHolder> holder(
new FetchStreamWorkerHolder(stream));
if (NS_WARN_IF(!holder->HoldWorker(workerPrivate, Closing))) {
aRv.Throw(NS_ERROR_DOM_INVALID_STATE_ERR);
return nullptr;
}
// Note, this will create a ref-cycle between the holder and the stream.
// The cycle is broken when the stream is closed or the worker begins
// shutting down.
stream->mWorkerHolder = Move(holder);
}
if (!JS::HasReadableStreamCallbacks(aCx)) {
JS::SetReadableStreamCallbacks(aCx,
&FetchStream::RequestDataCallback,
@ -228,11 +320,16 @@ FetchStream::FinalizeCallback(void* aUnderlyingSource, uint8_t aFlags)
MOZ_DIAGNOSTIC_ASSERT(aUnderlyingSource);
MOZ_DIAGNOSTIC_ASSERT(aFlags == FETCH_STREAM_FLAG);
// This can be called in any thread.
RefPtr<FetchStream> stream =
dont_AddRef(static_cast<FetchStream*>(aUnderlyingSource));
stream->mState = eClosed;
stream->mReadableStream = nullptr;
if (stream->mState == eClosed) {
return;
}
stream->CloseAndReleaseObjects();
}
FetchStream::FetchStream(nsIGlobalObject* aGlobal,
@ -248,7 +345,6 @@ FetchStream::FetchStream(nsIGlobalObject* aGlobal,
FetchStream::~FetchStream()
{
NS_ProxyRelease(mOwningEventTarget, mGlobal.forget());
}
void
@ -332,5 +428,70 @@ FetchStream::OnInputStreamReady(nsIAsyncInputStream* aStream)
return NS_OK;
}
void
FetchStream::Close()
{
if (mState == eClosed) {
return;
}
AutoJSAPI jsapi;
if (NS_WARN_IF(!jsapi.Init(mGlobal))) {
return;
}
JSContext* cx = jsapi.cx();
JS::Rooted<JSObject*> stream(cx, mReadableStream);
JS::ReadableStreamClose(cx, stream);
CloseAndReleaseObjects();
}
void
FetchStream::CloseAndReleaseObjects()
{
MOZ_DIAGNOSTIC_ASSERT(mState != eClosed);
mState = eClosed;
if (mWorkerHolder) {
RefPtr<FetchStreamWorkerHolderShutdown> r =
new FetchStreamWorkerHolderShutdown(
static_cast<FetchStreamWorkerHolder*>(mWorkerHolder.get())->GetWorkerPrivate(),
Move(mWorkerHolder), Move(mGlobal));
r->Dispatch();
} else {
RefPtr<FetchStream> self = this;
RefPtr<Runnable> r = NS_NewRunnableFunction(
[self] () {
nsCOMPtr<nsIObserverService> os = mozilla::services::GetObserverService();
if (os) {
os->RemoveObserver(self, DOM_WINDOW_DESTROYED_TOPIC);
}
self->mGlobal = nullptr;
});
NS_DispatchToMainThread(r);
}
}
// nsIObserver
// -----------
NS_IMETHODIMP
FetchStream::Observe(nsISupports* aSubject, const char* aTopic,
const char16_t* aData)
{
AssertIsOnMainThread();
MOZ_ASSERT(strcmp(aTopic, DOM_WINDOW_DESTROYED_TOPIC) == 0);
nsCOMPtr<nsPIDOMWindowInner> window = do_QueryInterface(mGlobal);
if (SameCOMIdentity(aSubject, window)) {
Close();
}
return NS_OK;
}
} // dom namespace
} // mozilla namespace