// Copyright 2022-2023 Google LLC // // Licensed under the Apache License, Version 2.0 (the "License"); // you may not use this file except in compliance with the License. // You may obtain a copy of the License at // // https://www.apache.org/licenses/LICENSE-2.0 // // Unless required by applicable law or agreed to in writing, software // distributed under the License is distributed on an "AS IS" BASIS, // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. // See the License for the specific language governing permissions and // limitations under the License. #ifndef THIRD_PARTY_NEARBY_SHARING_INCOMING_FRAMES_READER_H_ #define THIRD_PARTY_NEARBY_SHARING_INCOMING_FRAMES_READER_H_ #include #include #include #include #include #include #include #include "absl/base/thread_annotations.h" #include "absl/synchronization/mutex.h" #include "absl/time/time.h" #include "internal/platform/task_runner.h" #include "sharing/nearby_connection.h" #include "sharing/proto/wire_format.pb.h" #include "sharing/thread_timer.h" namespace nearby { namespace sharing { // Helper class to read incoming frames from Nearby devices. class IncomingFramesReader : public std::enable_shared_from_this { public: IncomingFramesReader(TaskRunner& service_thread, NearbyConnection* connection); virtual ~IncomingFramesReader(); IncomingFramesReader(const IncomingFramesReader&) = delete; IncomingFramesReader& operator=(IncomingFramesReader&) = delete; // Reads an incoming frame from connection. `callback` is called // with the frame read from connection or nullopt if connection socket is // closed or timeout has occurred. If timeout has occurred, the `is_timeout` // parameter will be true. Set `timeout` to absl::ZeroDuration() to disable // timeout. // // Note: Callers are expected wait for `callback` to be run before scheduling // subsequent calls to ReadFrame(..). virtual void ReadFrame( std::function< void(bool is_timeout, std::optional)> callback, absl::Duration timeout) ABSL_LOCKS_EXCLUDED(mutex_); // Reads a frame of type `frame_type` from `connection`. `callback` is called // with the frame read from connection or nullopt if connection socket is // closed or `timeout` units of time have passed. If timeout has occurred, // the `is_timeout` parameter will be true. Set `timeout` to // absl::ZeroDuration() to disable timeout. // // Note: Callers are expected wait for |callback| to be run before scheduling // subsequent calls to ReadFrame(..). virtual void ReadFrame( nearby::sharing::service::proto::V1Frame::FrameType frame_type, std::function< void(bool is_timeout, std::optional)> callback, absl::Duration timeout) ABSL_LOCKS_EXCLUDED(mutex_); std::weak_ptr GetWeakPtr() { return this->weak_from_this(); } private: struct ReadFrameInfo { std::optional frame_type = std::nullopt; std::function)> callback = nullptr; absl::Duration timeout = absl::ZeroDuration(); }; void ProcessReadRequest( std::optional frame_type, std::function< void(bool is_timeout, std::optional)> callback, absl::Duration timeout) ABSL_LOCKS_EXCLUDED(mutex_); void CloseAllPendingReads(bool is_timeout) ABSL_LOCKS_EXCLUDED(mutex_); void ReadNextFrame() ABSL_LOCKS_EXCLUDED(mutex_); void OnDataReadFromConnection(const std::vector& bytes) ABSL_LOCKS_EXCLUDED(mutex_); void OnTimeout(); void Done(std::unique_ptr frame) ABSL_LOCKS_EXCLUDED(mutex_); std::unique_ptr PopCachedFrame( std::optional frame_type) ABSL_EXCLUSIVE_LOCKS_REQUIRED(mutex_); TaskRunner& service_thread_; NearbyConnection* const connection_; absl::Mutex mutex_; std::queue read_frame_info_queue_ ABSL_GUARDED_BY(mutex_); // Caches frames read from NearbyConnection which are not used immediately. std::list> cached_frames_ ABSL_GUARDED_BY(mutex_); std::unique_ptr timeout_timer_ ABSL_GUARDED_BY(mutex_); }; } // namespace sharing } // namespace nearby #endif // THIRD_PARTY_NEARBY_SHARING_INCOMING_FRAMES_READER_H_