From f51af0f3751151725035fdb7c3c49bf718591115 Mon Sep 17 00:00:00 2001 From: Jude Kwashie Date: Wed, 9 Sep 2026 11:50:48 +0000 Subject: [PATCH 1/3] test(storage): add concurrent putData e2e for Windows crash Sequential putData tests never overlap progress callbacks, so they miss the native access violation in TaskStateListener::OnProgress. --- .../integration_test/reference_e2e.dart | 47 +++++++++++++++++++ 1 file changed, 47 insertions(+) diff --git a/packages/firebase_storage/firebase_storage/example/integration_test/reference_e2e.dart b/packages/firebase_storage/firebase_storage/example/integration_test/reference_e2e.dart index 5b382fb03767..42524fdc4020 100644 --- a/packages/firebase_storage/firebase_storage/example/integration_test/reference_e2e.dart +++ b/packages/firebase_storage/firebase_storage/example/integration_test/reference_e2e.dart @@ -369,6 +369,53 @@ void setupReferenceTests() { expect(complete.metadata?.contentType, 'image/jpeg'); }, ); + + // Overlapping putData calls can native-crash Windows (0xC0000005 in + // TaskStateListener::OnProgress). Sequential putData tests do not hit + // that path. See https://github.com/firebase/flutterfire/issues/18664. + test( + 'uploads many small files concurrently and reads them back', + () async { + final bytes = Uint8List.fromList( + utf8.encode( + jsonEncode({ + 'id': List.generate(300, (i) => i), + 'label': List.generate(300, (i) => 'Synthetic customer $i'), + }), + ), + ); + const rounds = 3; + const uploadsPerRound = 13; + final prefix = + 'flutter-tests/concurrent-put-data/${DateTime.now().microsecondsSinceEpoch}'; + + for (var round = 0; round < rounds; round++) { + final refs = List.generate( + uploadsPerRound, + (i) => storage.ref('$prefix/round-$round-$i.json'), + ); + + final snapshots = await Future.wait( + refs.map( + (ref) => ref.putData( + bytes, + SettableMetadata(contentType: 'application/json'), + ), + ), + ); + + expect(snapshots, hasLength(uploadsPerRound)); + for (final snapshot in snapshots) { + expect(snapshot.state, TaskState.success); + } + + for (final ref in refs) { + expect(await ref.getData(), bytes); + } + } + }, + timeout: const Timeout(Duration(minutes: 2)), + ); }, ); From adc3af736c30fc321b4bab14d94414274ff6635f Mon Sep 17 00:00:00 2001 From: Jude Kwashie Date: Wed, 9 Sep 2026 14:37:04 +0000 Subject: [PATCH 2/3] fix(storage,windows): dispatch task events on the platform thread Concurrent putData crashed Windows because progress callbacks called EventSink::Success from an SDK worker thread and the stream handler never owned the sink. --- .../integration_test/reference_e2e.dart | 7 +- .../windows/firebase_storage_plugin.cpp | 356 ++++++++++++------ 2 files changed, 255 insertions(+), 108 deletions(-) diff --git a/packages/firebase_storage/firebase_storage/example/integration_test/reference_e2e.dart b/packages/firebase_storage/firebase_storage/example/integration_test/reference_e2e.dart index 42524fdc4020..6e8bf83e6c41 100644 --- a/packages/firebase_storage/firebase_storage/example/integration_test/reference_e2e.dart +++ b/packages/firebase_storage/firebase_storage/example/integration_test/reference_e2e.dart @@ -370,9 +370,10 @@ void setupReferenceTests() { }, ); - // Overlapping putData calls can native-crash Windows (0xC0000005 in - // TaskStateListener::OnProgress). Sequential putData tests do not hit - // that path. See https://github.com/firebase/flutterfire/issues/18664. + // Overlapping putData calls previously native-crashed Windows by + // sending task events off the platform thread. Sequential putData + // tests do not hit that path. + // See https://github.com/firebase/flutterfire/issues/18664. test( 'uploads many small files concurrently and reads them back', () async { diff --git a/packages/firebase_storage/firebase_storage/windows/firebase_storage_plugin.cpp b/packages/firebase_storage/firebase_storage/windows/firebase_storage_plugin.cpp index 5fcac57d17fe..a88dd1d940f2 100644 --- a/packages/firebase_storage/firebase_storage/windows/firebase_storage_plugin.cpp +++ b/packages/firebase_storage/firebase_storage/windows/firebase_storage_plugin.cpp @@ -27,14 +27,15 @@ #include // #include +#include #include #include #include +#include +#include #include #include #include -// #include -#include #include using ::firebase::App; using ::firebase::Future; @@ -53,9 +54,138 @@ static std::string kLibraryName = "flutter-fire-gcs"; static std::string kStorageMethodChannelName = "plugins.flutter.io/firebase_storage"; static std::string kStorageTaskEventName = "taskEvent"; + +namespace { + +constexpr wchar_t kTaskRunnerWindowClassName[] = + L"FirebaseStorageWindowsTaskRunnerWindow"; +constexpr UINT kTaskRunnerWindowMessage = WM_APP + 0x4674; + +class PlatformThreadDispatcher { + public: + static PlatformThreadDispatcher& GetInstance() { + static PlatformThreadDispatcher instance; + return instance; + } + + void Initialize() { + std::lock_guard lock(mutex_); + if (window_ != nullptr) { + return; + } + + platform_thread_id_ = GetCurrentThreadId(); + + WNDCLASSW window_class = {}; + window_class.lpfnWndProc = PlatformThreadDispatcher::WindowProc; + window_class.hInstance = GetModuleHandle(nullptr); + window_class.lpszClassName = kTaskRunnerWindowClassName; + + RegisterClassW(&window_class); + window_ = + CreateWindowExW(0, kTaskRunnerWindowClassName, L"", 0, 0, 0, 0, 0, + HWND_MESSAGE, nullptr, window_class.hInstance, this); + } + + void Post(std::function task) { + if (window_ == nullptr) { + task(); + return; + } + if (GetCurrentThreadId() == platform_thread_id_) { + task(); + return; + } + + { + std::lock_guard lock(mutex_); + tasks_.push(std::move(task)); + } + PostMessageW(window_, kTaskRunnerWindowMessage, 0, 0); + } + + private: + static LRESULT CALLBACK WindowProc(HWND window, UINT message, WPARAM wparam, + LPARAM lparam) { + if (message == WM_NCCREATE) { + auto create_struct = reinterpret_cast(lparam); + SetWindowLongPtr( + window, GWLP_USERDATA, + reinterpret_cast(create_struct->lpCreateParams)); + return TRUE; + } + + auto dispatcher = reinterpret_cast( + GetWindowLongPtr(window, GWLP_USERDATA)); + if (dispatcher != nullptr && message == kTaskRunnerWindowMessage) { + dispatcher->ProcessTasks(); + return 0; + } + + return DefWindowProc(window, message, wparam, lparam); + } + + void ProcessTasks() { + std::queue> tasks; + { + std::lock_guard lock(mutex_); + tasks.swap(tasks_); + } + + while (!tasks.empty()) { + tasks.front()(); + tasks.pop(); + } + } + + PlatformThreadDispatcher() = default; + + HWND window_ = nullptr; + DWORD platform_thread_id_ = 0; + std::mutex mutex_; + std::queue> tasks_; +}; + +struct EventSinkState { + std::mutex mutex; + std::unique_ptr> events; + bool active = true; +}; + +void SendSuccessOnPlatformThread(std::shared_ptr state, + flutter::EncodableValue value) { + if (!state) { + return; + } + + PlatformThreadDispatcher::GetInstance().Post( + [state, value = std::move(value)]() mutable { + std::lock_guard lock(state->mutex); + if (state->active && state->events) { + state->events->Success(value); + } + }); +} + +void DeactivateSinkOnPlatformThread(std::shared_ptr state) { + if (!state) { + return; + } + + PlatformThreadDispatcher::GetInstance().Post([state]() { + std::lock_guard lock(state->mutex); + state->active = false; + state->events.reset(); + }); +} + +} // namespace + // static void FirebaseStoragePlugin::RegisterWithRegistrar( flutter::PluginRegistrarWindows* registrar) { + PlatformThreadDispatcher::GetInstance().Initialize(); + auto plugin = std::make_unique(); messenger_ = registrar->messenger(); FirebaseStorageHostApi::SetUp(registrar->messenger(), plugin.get()); @@ -523,10 +653,10 @@ std::string kErrorName = "error"; class TaskStateListener : public Listener { public: - TaskStateListener(flutter::EventSink* events) { - events_ = events; - } - virtual void OnProgress(firebase::storage::Controller* controller) { + explicit TaskStateListener(std::shared_ptr events_state) + : events_state_(std::move(events_state)) {} + + void OnProgress(firebase::storage::Controller* controller) override { flutter::EncodableMap event = flutter::EncodableMap(); event[kTaskStateName] = static_cast(InternalStorageTaskState::kRunning); @@ -537,10 +667,10 @@ class TaskStateListener : public Listener { snapshot[kTaskSnapshotBytesTransferred] = controller->bytes_transferred(); event[kTaskSnapshotName] = snapshot; - events_->Success(event); + SendSuccessOnPlatformThread(events_state_, flutter::EncodableValue(event)); } - virtual void OnPaused(firebase::storage::Controller* controller) { + void OnPaused(firebase::storage::Controller* controller) override { flutter::EncodableMap event = flutter::EncodableMap(); event[kTaskStateName] = static_cast(InternalStorageTaskState::kPaused); event[kTaskAppName] = controller->GetReference().storage()->app()->name(); @@ -550,12 +680,38 @@ class TaskStateListener : public Listener { snapshot[kTaskSnapshotBytesTransferred] = controller->bytes_transferred(); event[kTaskSnapshotName] = snapshot; - events_->Success(event); + SendSuccessOnPlatformThread(events_state_, flutter::EncodableValue(event)); } - flutter::EventSink* events_; + private: + std::shared_ptr events_state_; }; +void SendMetadataTaskResult(std::shared_ptr events_state, + const std::string& app_name, + const Future& data_result) { + if (data_result.error() == firebase::storage::kErrorNone) { + flutter::EncodableMap event = flutter::EncodableMap(); + event[kTaskStateName] = + static_cast(InternalStorageTaskState::kSuccess); + event[kTaskAppName] = app_name; + flutter::EncodableMap snapshot = flutter::EncodableMap(); + snapshot[kTaskSnapshotPath] = data_result.result()->path(); + snapshot[kTaskSnapshotTotalBytes] = data_result.result()->size_bytes(); + snapshot[kTaskSnapshotBytesTransferred] = + data_result.result()->size_bytes(); + snapshot[kMetadataName] = ConvertMedadataToPigeon(data_result.result()); + event[kTaskSnapshotName] = snapshot; + SendSuccessOnPlatformThread(std::move(events_state), + flutter::EncodableValue(event)); + } else { + SendSuccessOnPlatformThread( + std::move(events_state), + flutter::EncodableValue(FirebaseStoragePlugin::ErrorStreamEvent( + data_result, app_name))); + } +} + class PutDataStreamHandler : public flutter::StreamHandler { public: @@ -576,52 +732,41 @@ class PutDataStreamHandler const flutter::EncodableValue* arguments, std::unique_ptr>&& events) override { - events_ = std::move(events); - - TaskStateListener* putStringListener = new TaskStateListener(events_.get()); - StorageReference reference = storage_->GetReference(reference_path_); + events_state_ = std::make_shared(); + events_state_->events = std::move(events); + listener_ = std::make_shared(events_state_); + reference_ = std::make_shared( + storage_->GetReference(reference_path_)); Metadata* storage_metadata = FirebaseStoragePlugin::CreateStorageMetadataFromPigeon(&meta_data_); Future future_result; if (storage_metadata) { future_result = - reference.PutBytes(data_.data(), data_.size(), *storage_metadata, - putStringListener, controller_); + reference_->PutBytes(data_.data(), data_.size(), *storage_metadata, + listener_.get(), controller_); } else { - future_result = reference.PutBytes(data_.data(), data_.size(), - putStringListener, controller_); + future_result = reference_->PutBytes(data_.data(), data_.size(), + listener_.get(), controller_); } ::Sleep(1); // timing for c++ sdk grabbing a mutex - future_result.OnCompletion([this](const Future& data_result) { - if (data_result.error() == firebase::storage::kErrorNone) { - flutter::EncodableMap event = flutter::EncodableMap(); - event[kTaskStateName] = - static_cast(InternalStorageTaskState::kSuccess); - event[kTaskAppName] = std::string(storage_->app()->name()); - flutter::EncodableMap snapshot = flutter::EncodableMap(); - snapshot[kTaskSnapshotPath] = data_result.result()->path(); - snapshot[kTaskSnapshotTotalBytes] = data_result.result()->size_bytes(); - snapshot[kTaskSnapshotBytesTransferred] = - data_result.result()->size_bytes(); - snapshot[kMetadataName] = ConvertMedadataToPigeon(data_result.result()); - event[kTaskSnapshotName] = snapshot; - - events_->Success(event); - } else { - flutter::EncodableMap map = FirebaseStoragePlugin::ErrorStreamEvent( - data_result, storage_->app()->name()); - - events_->Success(map); - } - }); + future_result.OnCompletion( + [events_state = events_state_, + app_name = std::string(storage_->app()->name()), + reference = reference_, + listener = listener_](const Future& data_result) { + (void)reference; + (void)listener; + SendMetadataTaskResult(events_state, app_name, data_result); + }); return nullptr; } std::unique_ptr> OnCancelInternal(const flutter::EncodableValue* arguments) override { + DeactivateSinkOnPlatformThread(events_state_); return nullptr; } @@ -631,8 +776,9 @@ class PutDataStreamHandler std::vector data_; InternalSettableMetadata meta_data_; Controller* controller_; - std::unique_ptr>&& events_ = - nullptr; + std::shared_ptr events_state_; + std::shared_ptr reference_; + std::shared_ptr listener_; }; class PutFileStreamHandler @@ -654,10 +800,11 @@ class PutFileStreamHandler const flutter::EncodableValue* arguments, std::unique_ptr>&& events) override { - events_ = std::move(events); - - TaskStateListener* putFileListener = new TaskStateListener(events_.get()); - StorageReference reference = storage_->GetReference(reference_path_); + events_state_ = std::make_shared(); + events_state_->events = std::move(events); + listener_ = std::make_shared(events_state_); + reference_ = std::make_shared( + storage_->GetReference(reference_path_)); Metadata* storage_metadata = FirebaseStoragePlugin::CreateStorageMetadataFromPigeon( @@ -665,41 +812,29 @@ class PutFileStreamHandler Future future_result; if (storage_metadata) { - future_result = reference.PutFile(file_path_.c_str(), *storage_metadata, - putFileListener, controller_); + future_result = reference_->PutFile(file_path_.c_str(), *storage_metadata, + listener_.get(), controller_); } else { - future_result = - reference.PutFile(file_path_.c_str(), putFileListener, controller_); + future_result = reference_->PutFile(file_path_.c_str(), listener_.get(), + controller_); } ::Sleep(1); // timing for c++ sdk grabbing a mutex - future_result.OnCompletion([this](const Future& data_result) { - if (data_result.error() == firebase::storage::kErrorNone) { - flutter::EncodableMap event = flutter::EncodableMap(); - event[kTaskStateName] = - static_cast(InternalStorageTaskState::kSuccess); - event[kTaskAppName] = std::string(storage_->app()->name()); - flutter::EncodableMap snapshot = flutter::EncodableMap(); - snapshot[kTaskSnapshotPath] = data_result.result()->path(); - snapshot[kTaskSnapshotTotalBytes] = data_result.result()->size_bytes(); - snapshot[kTaskSnapshotBytesTransferred] = - data_result.result()->size_bytes(); - snapshot[kMetadataName] = ConvertMedadataToPigeon(data_result.result()); - event[kTaskSnapshotName] = snapshot; - - events_->Success(event); - } else { - flutter::EncodableMap map = FirebaseStoragePlugin::ErrorStreamEvent( - data_result, storage_->app()->name()); - - events_->Success(map); - } - }); + future_result.OnCompletion( + [events_state = events_state_, + app_name = std::string(storage_->app()->name()), + reference = reference_, + listener = listener_](const Future& data_result) { + (void)reference; + (void)listener; + SendMetadataTaskResult(events_state, app_name, data_result); + }); return nullptr; } std::unique_ptr> OnCancelInternal(const flutter::EncodableValue* arguments) override { + DeactivateSinkOnPlatformThread(events_state_); return nullptr; } @@ -708,9 +843,10 @@ class PutFileStreamHandler std::string reference_path_; std::string file_path_; Controller* controller_; - std::unique_ptr>&& events_ = - nullptr; std::unique_ptr meta_data_; + std::shared_ptr events_state_; + std::shared_ptr reference_; + std::shared_ptr listener_; }; class GetFileStreamHandler @@ -729,45 +865,54 @@ class GetFileStreamHandler const flutter::EncodableValue* arguments, std::unique_ptr>&& events) override { - events_ = std::move(events); + events_state_ = std::make_shared(); + events_state_->events = std::move(events); std::unique_lock lock(mtx_); - TaskStateListener* getFileListener = new TaskStateListener(events_.get()); - StorageReference reference = storage_->GetReference(reference_path_); - Future future_result = - reference.GetFile(file_path_.c_str(), getFileListener, controller_); + listener_ = std::make_shared(events_state_); + reference_ = std::make_shared( + storage_->GetReference(reference_path_)); + Future future_result = reference_->GetFile( + file_path_.c_str(), listener_.get(), controller_); ::Sleep(1); // timing for c++ sdk grabbing a mutex - future_result.OnCompletion([this](const Future& data_result) { - if (data_result.error() == firebase::storage::kErrorNone) { - flutter::EncodableMap event = flutter::EncodableMap(); - event[kTaskStateName] = - static_cast(InternalStorageTaskState::kSuccess); - event[kTaskAppName] = std::string(storage_->app()->name()); - flutter::EncodableMap snapshot = flutter::EncodableMap(); - size_t data_size = *data_result.result(); - snapshot[kTaskSnapshotTotalBytes] = - flutter::EncodableValue(static_cast(data_size)); - snapshot[kTaskSnapshotBytesTransferred] = - flutter::EncodableValue(static_cast(data_size)); - snapshot[kTaskSnapshotPath] = EncodableValue(reference_path_); - event[kTaskSnapshotName] = snapshot; - - events_->Success(event); - } else { - flutter::EncodableMap map = FirebaseStoragePlugin::ErrorStreamEvent( - data_result, storage_->app()->name()); - - events_->Success(map); - } - }); + future_result.OnCompletion( + [events_state = events_state_, + app_name = std::string(storage_->app()->name()), + path = reference_path_, reference = reference_, + listener = listener_](const Future& data_result) { + (void)reference; + (void)listener; + if (data_result.error() == firebase::storage::kErrorNone) { + flutter::EncodableMap event = flutter::EncodableMap(); + event[kTaskStateName] = + static_cast(InternalStorageTaskState::kSuccess); + event[kTaskAppName] = app_name; + flutter::EncodableMap snapshot = flutter::EncodableMap(); + size_t data_size = *data_result.result(); + snapshot[kTaskSnapshotTotalBytes] = + flutter::EncodableValue(static_cast(data_size)); + snapshot[kTaskSnapshotBytesTransferred] = + flutter::EncodableValue(static_cast(data_size)); + snapshot[kTaskSnapshotPath] = EncodableValue(path); + event[kTaskSnapshotName] = snapshot; + + SendSuccessOnPlatformThread(events_state, + flutter::EncodableValue(event)); + } else { + SendSuccessOnPlatformThread( + events_state, + flutter::EncodableValue(FirebaseStoragePlugin::ErrorStreamEvent( + data_result, app_name))); + } + }); return nullptr; } std::unique_ptr> OnCancelInternal(const flutter::EncodableValue* arguments) override { std::unique_lock lock(mtx_); - + DeactivateSinkOnPlatformThread(events_state_); return nullptr; } @@ -777,8 +922,9 @@ class GetFileStreamHandler std::string file_path_; Controller* controller_; std::mutex mtx_; - std::unique_ptr>&& events_ = - nullptr; + std::shared_ptr events_state_; + std::shared_ptr reference_; + std::shared_ptr listener_; }; void FirebaseStoragePlugin::ReferencePutData( From 8feedf919a175775393296d15ca287ecf0f9ba81 Mon Sep 17 00:00:00 2001 From: Jude Kwashie Date: Wed, 9 Sep 2026 14:50:14 +0000 Subject: [PATCH 3/3] style(storage,windows): apply clang-format to task event dispatcher --- .../windows/firebase_storage_plugin.cpp | 46 +++++++++---------- 1 file changed, 22 insertions(+), 24 deletions(-) diff --git a/packages/firebase_storage/firebase_storage/windows/firebase_storage_plugin.cpp b/packages/firebase_storage/firebase_storage/windows/firebase_storage_plugin.cpp index a88dd1d940f2..9704203b23b2 100644 --- a/packages/firebase_storage/firebase_storage/windows/firebase_storage_plugin.cpp +++ b/packages/firebase_storage/firebase_storage/windows/firebase_storage_plugin.cpp @@ -707,8 +707,8 @@ void SendMetadataTaskResult(std::shared_ptr events_state, } else { SendSuccessOnPlatformThread( std::move(events_state), - flutter::EncodableValue(FirebaseStoragePlugin::ErrorStreamEvent( - data_result, app_name))); + flutter::EncodableValue( + FirebaseStoragePlugin::ErrorStreamEvent(data_result, app_name))); } } @@ -752,15 +752,14 @@ class PutDataStreamHandler ::Sleep(1); // timing for c++ sdk grabbing a mutex - future_result.OnCompletion( - [events_state = events_state_, - app_name = std::string(storage_->app()->name()), - reference = reference_, - listener = listener_](const Future& data_result) { - (void)reference; - (void)listener; - SendMetadataTaskResult(events_state, app_name, data_result); - }); + future_result.OnCompletion([events_state = events_state_, + app_name = std::string(storage_->app()->name()), + reference = reference_, listener = listener_]( + const Future& data_result) { + (void)reference; + (void)listener; + SendMetadataTaskResult(events_state, app_name, data_result); + }); return nullptr; } @@ -815,20 +814,19 @@ class PutFileStreamHandler future_result = reference_->PutFile(file_path_.c_str(), *storage_metadata, listener_.get(), controller_); } else { - future_result = reference_->PutFile(file_path_.c_str(), listener_.get(), - controller_); + future_result = + reference_->PutFile(file_path_.c_str(), listener_.get(), controller_); } ::Sleep(1); // timing for c++ sdk grabbing a mutex - future_result.OnCompletion( - [events_state = events_state_, - app_name = std::string(storage_->app()->name()), - reference = reference_, - listener = listener_](const Future& data_result) { - (void)reference; - (void)listener; - SendMetadataTaskResult(events_state, app_name, data_result); - }); + future_result.OnCompletion([events_state = events_state_, + app_name = std::string(storage_->app()->name()), + reference = reference_, listener = listener_]( + const Future& data_result) { + (void)reference; + (void)listener; + SendMetadataTaskResult(events_state, app_name, data_result); + }); return nullptr; } @@ -872,8 +870,8 @@ class GetFileStreamHandler listener_ = std::make_shared(events_state_); reference_ = std::make_shared( storage_->GetReference(reference_path_)); - Future future_result = reference_->GetFile( - file_path_.c_str(), listener_.get(), controller_); + Future future_result = + reference_->GetFile(file_path_.c_str(), listener_.get(), controller_); ::Sleep(1); // timing for c++ sdk grabbing a mutex future_result.OnCompletion(