123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274 |
- // Copyright 2021 The Chromium Authors. All rights reserved.
- // Use of this source code is governed by a BSD-style license that can be
- // found in the LICENSE file.
- #include "components/reporting/client/report_queue_provider.h"
- #include <memory>
- #include <string>
- #include "base/bind.h"
- #include "base/callback.h"
- #include "base/feature_list.h"
- #include "base/memory/ptr_util.h"
- #include "base/memory/ref_counted.h"
- #include "base/memory/scoped_refptr.h"
- #include "base/sequence_checker.h"
- #include "base/task/bind_post_task.h"
- #include "base/task/sequenced_task_runner.h"
- #include "base/task/task_traits.h"
- #include "base/task/thread_pool.h"
- #include "components/reporting/client/report_queue.h"
- #include "components/reporting/client/report_queue_configuration.h"
- #include "components/reporting/client/report_queue_impl.h"
- #include "components/reporting/proto/synced/record_constants.pb.h"
- #include "components/reporting/storage/storage_module_interface.h"
- #include "components/reporting/util/status.h"
- #include "components/reporting/util/status_macros.h"
- #include "components/reporting/util/statusor.h"
- namespace reporting {
- using InitCompleteCallback = base::OnceCallback<void(Status)>;
- // Report queue creation request. Recorded in the `create_request_queue_` when
- // provider cannot create queues yet.
- class ReportQueueProvider::CreateReportQueueRequest {
- public:
- static void New(std::unique_ptr<ReportQueueConfiguration> config,
- CreateReportQueueCallback create_cb) {
- auto* const provider = GetInstance();
- DCHECK(provider)
- << "Provider must exist, otherwise it is an internal error";
- auto request = base::WrapUnique(
- new CreateReportQueueRequest(std::move(config), std::move(create_cb)));
- provider->sequenced_task_runner_->PostTask(
- FROM_HERE,
- base::BindOnce(
- [](base::WeakPtr<ReportQueueProvider> provider,
- std::unique_ptr<CreateReportQueueRequest> request) {
- if (!provider) {
- std::move(request->release_create_cb())
- .Run(Status(error::UNAVAILABLE,
- "Provider has been shut down"));
- return;
- }
- DCHECK_CALLED_ON_VALID_SEQUENCE(provider->sequence_checker_);
- provider->create_request_queue_.push(std::move(request));
- provider->CheckInitializationState();
- },
- provider->GetWeakPtr(), std::move(request)));
- }
- CreateReportQueueRequest(const CreateReportQueueRequest& other) = delete;
- CreateReportQueueRequest& operator=(const CreateReportQueueRequest& other) =
- delete;
- ~CreateReportQueueRequest() = default;
- std::unique_ptr<ReportQueueConfiguration> release_config() {
- DCHECK(config_) << "Can only be released once";
- return std::move(config_);
- }
- ReportQueueProvider::CreateReportQueueCallback release_create_cb() {
- DCHECK(create_cb_) << "Can only be released once";
- return std::move(create_cb_);
- }
- private:
- // Constructor is only called by `New` factory method above.
- CreateReportQueueRequest(std::unique_ptr<ReportQueueConfiguration> config,
- CreateReportQueueCallback create_cb)
- : config_(std::move(config)), create_cb_(std::move(create_cb)) {}
- std::unique_ptr<ReportQueueConfiguration> config_;
- CreateReportQueueCallback create_cb_;
- };
- // ReportQueueProvider core implementation.
- // static
- bool ReportQueueProvider::IsEncryptedReportingPipelineEnabled() {
- return base::FeatureList::IsEnabled(kEncryptedReportingPipeline);
- }
- // static
- const base::Feature ReportQueueProvider::kEncryptedReportingPipeline{
- "EncryptedReportingPipeline", base::FEATURE_ENABLED_BY_DEFAULT};
- ReportQueueProvider::ReportQueueProvider(
- StorageModuleCreateCallback storage_create_cb)
- : storage_create_cb_(storage_create_cb),
- sequenced_task_runner_(base::ThreadPool::CreateSequencedTaskRunner(
- {base::TaskPriority::BEST_EFFORT, base::MayBlock()})) {
- DETACH_FROM_SEQUENCE(sequence_checker_);
- }
- ReportQueueProvider::~ReportQueueProvider() = default;
- base::WeakPtr<ReportQueueProvider> ReportQueueProvider::GetWeakPtr() {
- return weak_ptr_factory_.GetWeakPtr();
- }
- scoped_refptr<StorageModuleInterface> ReportQueueProvider::storage() const {
- DCHECK_CALLED_ON_VALID_SEQUENCE(sequence_checker_);
- return storage_;
- }
- scoped_refptr<base::SequencedTaskRunner>
- ReportQueueProvider::sequenced_task_runner() const {
- return sequenced_task_runner_;
- }
- void ReportQueueProvider::CreateNewQueue(
- std::unique_ptr<ReportQueueConfiguration> config,
- CreateReportQueueCallback cb) {
- sequenced_task_runner_->PostTask(
- FROM_HERE,
- base::BindOnce(
- [](base::WeakPtr<ReportQueueProvider> provider,
- std::unique_ptr<ReportQueueConfiguration> config,
- CreateReportQueueCallback cb) {
- if (!provider) {
- std::move(cb).Run(
- Status(error::UNAVAILABLE, "Provider has been shut down"));
- return;
- }
- // Configure report queue config with an appropriate DM token and
- // proceed to create the queue if configuration was successful.
- DCHECK_CALLED_ON_VALID_SEQUENCE(provider->sequence_checker_);
- auto report_queue_configured_cb = base::BindOnce(
- [](scoped_refptr<StorageModuleInterface> storage,
- CreateReportQueueCallback cb,
- StatusOr<std::unique_ptr<ReportQueueConfiguration>>
- config_result) {
- // If configuration hit an error, we abort and
- // report this through the callback
- if (!config_result.ok()) {
- std::move(cb).Run(config_result.status());
- return;
- }
- // Proceed to create the queue on arbitrary thread.
- base::ThreadPool::PostTask(
- FROM_HERE,
- base::BindOnce(&ReportQueueImpl::Create,
- std::move(config_result.ValueOrDie()),
- storage, std::move(cb)));
- },
- provider->storage_, std::move(cb));
- provider->ConfigureReportQueue(
- std::move(config), std::move(report_queue_configured_cb));
- },
- GetWeakPtr(), std::move(config), std::move(cb)));
- }
- StatusOr<std::unique_ptr<ReportQueue, base::OnTaskRunnerDeleter>>
- ReportQueueProvider::CreateNewSpeculativeQueue() {
- return SpeculativeReportQueueImpl::Create();
- }
- void ReportQueueProvider::OnInitCompleted() {}
- // static
- void ReportQueueProvider::CreateQueue(
- std::unique_ptr<ReportQueueConfiguration> config,
- CreateReportQueueCallback create_cb) {
- if (!IsEncryptedReportingPipelineEnabled()) {
- Status not_enabled = Status(
- error::FAILED_PRECONDITION,
- "The Encrypted Reporting Pipeline is not enabled. Please enable it on "
- "the command line using --enable-features=EncryptedReportingPipeline");
- VLOG(1) << not_enabled;
- std::move(create_cb).Run(not_enabled);
- return;
- }
- CreateReportQueueRequest::New(std::move(config), std::move(create_cb));
- }
- // static
- StatusOr<std::unique_ptr<ReportQueue, base::OnTaskRunnerDeleter>>
- ReportQueueProvider::CreateSpeculativeQueue(
- std::unique_ptr<ReportQueueConfiguration> config) {
- if (!IsEncryptedReportingPipelineEnabled()) {
- Status not_enabled = Status(
- error::FAILED_PRECONDITION,
- "The Encrypted Reporting Pipeline is not enabled. Please enable it on "
- "the command line using --enable-features=EncryptedReportingPipeline");
- VLOG(1) << not_enabled;
- return not_enabled;
- }
- // Instantiate speculative queue, bail out in case of an error.
- ASSIGN_OR_RETURN(auto speculative_queue,
- GetInstance()->CreateNewSpeculativeQueue());
- // Initiate underlying queue creation.
- CreateReportQueueRequest::New(
- std::move(config), speculative_queue->PrepareToAttachActualQueue());
- return speculative_queue;
- }
- void ReportQueueProvider::CheckInitializationState() {
- DCHECK_CALLED_ON_VALID_SEQUENCE(sequence_checker_);
- if (!storage_) {
- // Provider not ready.
- DCHECK(!create_request_queue_.empty()) << "Request queue cannot be empty";
- if (create_request_queue_.size() > 1) {
- // More than one request in the queue - it means Storage creation has
- // already been started.
- return;
- }
- // Start Storage creation on an arbitrary thread. Upon completion resume on
- // sequenced task runner.
- base::ThreadPool::PostTask(
- FROM_HERE,
- base::BindOnce(
- [](StorageModuleCreateCallback storage_create_cb,
- OnStorageModuleCreatedCallback on_storage_created_cb) {
- storage_create_cb.Run(std::move(on_storage_created_cb));
- },
- storage_create_cb_,
- base::BindPostTask(
- sequenced_task_runner_,
- base::BindOnce(&ReportQueueProvider::OnStorageModuleConfigured,
- GetWeakPtr()))));
- return;
- }
- // Storage ready, create all report queues that were submitted.
- // Note that `CreateNewQueue` call offsets heavy work to arbitrary threads.
- while (!create_request_queue_.empty()) {
- auto& report_queue_request = create_request_queue_.front();
- CreateNewQueue(report_queue_request->release_config(),
- report_queue_request->release_create_cb());
- create_request_queue_.pop();
- }
- }
- void ReportQueueProvider::OnStorageModuleConfigured(
- StatusOr<scoped_refptr<StorageModuleInterface>> storage_result) {
- DCHECK_CALLED_ON_VALID_SEQUENCE(sequence_checker_);
- if (!storage_result.ok()) {
- // Storage creation failed, kill all requests.
- while (!create_request_queue_.empty()) {
- auto& report_queue_request = create_request_queue_.front();
- std::move(report_queue_request->release_create_cb())
- .Run(Status(error::UNAVAILABLE, "Unable to build a ReportQueue"));
- create_request_queue_.pop();
- }
- return;
- }
- // Storage ready, create all report queues that were submitted.
- // Note that `CreateNewQueue` call offsets heavy work to arbitrary threads.
- DCHECK(!storage_) << "Storage module already recorded";
- OnInitCompleted();
- storage_ = storage_result.ValueOrDie();
- while (!create_request_queue_.empty()) {
- auto& report_queue_request = create_request_queue_.front();
- CreateNewQueue(report_queue_request->release_config(),
- report_queue_request->release_create_cb());
- create_request_queue_.pop();
- }
- }
- } // namespace reporting
|