Add unit tests for BufferedIOStream

This commit is contained in:
Simon Fels 2016-12-21 16:27:40 +01:00
commit dbbb8ab795
9 changed files with 177 additions and 24 deletions

View file

@ -19,8 +19,6 @@
#include <stdlib.h>
#include <stdio.h>
#include "ErrorLog.h"
class IOStream {
public:
IOStream(size_t bufSize) {
@ -40,19 +38,15 @@ public:
unsigned char *alloc(size_t len) {
if (m_buf && len > m_free) {
if (flush() < 0) {
ERR("Failed to flush in alloc\n");
if (flush() < 0)
return NULL; // we failed to flush so something is wrong
}
}
if (!m_buf || len > m_bufsize) {
int allocLen = m_bufsize < len ? len : m_bufsize;
m_buf = (unsigned char *)allocBuffer(allocLen);
if (!m_buf) {
ERR("Alloc (%u bytes) failed\n", allocLen);
if (!m_buf)
return NULL;
}
m_bufsize = m_free = allocLen;
}

View file

@ -1,19 +1,16 @@
/*
* Copyright (C) 2016 Simon Fels <morphis@gravedo.de>
*
* This program is free software: you can redistribute it and/or modify it
* under the terms of the GNU General Public License version 3, as published
* by the Free Software Foundation.
*
* This program is distributed in the hope that it will be useful, but
* WITHOUT ANY WARRANTY; without even the implied warranties of
* MERCHANTABILITY, SATISFACTORY QUALITY, or FITNESS FOR A PARTICULAR
* PURPOSE. See the GNU General Public License for more details.
*
* You should have received a copy of the GNU General Public License along
* with this program. If not, see <http://www.gnu.org/licenses/>.
*
*/
// Copyright (C) 2016 The Android Open Source Project
//
// 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
//
// http://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 ANBOX_GRAPHICS_BUFFER_QUEUE_H_
#define ANBOX_GRAPHICS_BUFFER_QUEUE_H_
@ -32,6 +29,10 @@ class BufferQueue {
public:
BufferQueue(size_t capacity);
bool can_push_locked() const { return !closed_ && (count_ < capacity_); }
bool can_pop_locked() const { return count_ > 0U; }
bool is_closed_locked() const { return closed_; }
int try_push_locked(Buffer &&buffer);
int push_locked(Buffer &&buffer, std::unique_lock<std::mutex> &lock);
int try_pop_locked(Buffer *buffer);

View file

@ -105,6 +105,11 @@ void BufferedIOStream::post_data(Buffer &&data) {
in_queue_.push_locked(std::move(data), l);
}
bool BufferedIOStream::needs_data() {
std::unique_lock<std::mutex> l(lock_);
return !in_queue_.can_pop_locked();
}
void BufferedIOStream::thread_main() {
while (true) {
std::unique_lock<std::mutex> l(out_lock_);

View file

@ -44,6 +44,8 @@ class BufferedIOStream : public IOStream {
void forceStop() override;
void post_data(Buffer &&data);
bool needs_data();
private:
void thread_main();

View file

@ -22,6 +22,8 @@
#include "OpenGLESDispatch/EGLDispatch.h"
#include "ErrorLog.h"
#include <stdio.h>
namespace {

View file

@ -27,6 +27,8 @@
#include "anbox/logger.h"
#include "ErrorLog.h"
#include <stdio.h>
#include <glm/glm.hpp>

View file

@ -1,5 +1,7 @@
include_directories(
${Boost_INCLUDE_DIRS}
${CMAKE_SOURCE_DIR}
${CMAKE_SOURCE_DIR}/external/android-emugl/host/include
${CMAKE_SOURCE_DIR}/src
)

View file

@ -1 +1,2 @@
ANBOX_ADD_TEST(buffer_queue_tests buffer_queue_tests.cpp)
ANBOX_ADD_TEST(buffered_io_stream_tests buffered_io_stream_tests.cpp)

View file

@ -0,0 +1,144 @@
/*
* Copyright (C) 2016 Simon Fels <morphis@gravedo.de>
*
* This program is free software: you can redistribute it and/or modify it
* under the terms of the GNU General Public License version 3, as published
* by the Free Software Foundation.
*
* This program is distributed in the hope that it will be useful, but
* WITHOUT ANY WARRANTY; without even the implied warranties of
* MERCHANTABILITY, SATISFACTORY QUALITY, or FITNESS FOR A PARTICULAR
* PURPOSE. See the GNU General Public License for more details.
*
* You should have received a copy of the GNU General Public License along
* with this program. If not, see <http://www.gnu.org/licenses/>.
*
*/
#include "anbox/graphics/buffered_io_stream.h"
#include <chrono>
#include <gtest/gtest.h>
#include <gmock/gmock.h>
using namespace ::testing;
namespace {
class MockSocketMessenger :
public anbox::network::SocketMessenger {
public:
// anbox::network::SocketMessenger
MOCK_CONST_METHOD0(creds, anbox::network::Credentials());
MOCK_CONST_METHOD0(local_port, unsigned short());
MOCK_METHOD0(set_no_delay, void());
MOCK_METHOD0(close, void());
// anbox::network::MessageSender
MOCK_METHOD2(send, void(char const*, size_t));
MOCK_METHOD2(send_raw, ssize_t(char const*, size_t));
// anbox::network::MessageReceiver
MOCK_METHOD2(async_receive_msg, void(AnboxReadHandler const&, boost::asio::mutable_buffers_1 const&));
MOCK_METHOD1(receive_msg, boost::system::error_code(boost::asio::mutable_buffers_1 const&));
MOCK_METHOD0(available_bytes, size_t());
};
}
namespace anbox {
namespace graphics {
TEST(BufferedIOStream, CommitBufferWritesOutToMessenger) {
auto messenger = std::make_shared<MockSocketMessenger>();
BufferedIOStream stream(messenger);
const auto buffer_size{1000};
// We will write out the data we get in two junks of half the size
// the original buffer has.
EXPECT_CALL(*messenger, send_raw(_, buffer_size))
.Times(1)
.WillOnce(Return(buffer_size/2));
EXPECT_CALL(*messenger, send_raw(_, buffer_size/2))
.Times(1)
.WillOnce(Return(buffer_size/2));
char *ptr = static_cast<char*>(stream.allocBuffer(buffer_size));
ASSERT_NE(ptr, nullptr);
ASSERT_EQ(stream.commitBuffer(buffer_size), buffer_size);
// The BufferedIOStream class works internally with a thread to
// write out the actual data to the messenger. As it blocks in
// its d'tor for the writer thread to quit we can safely expect
// that the messenger will see all data it has to.
}
TEST(BufferedIOStream, WriterContinuesWhenSocketIsBusy) {
auto messenger = std::make_shared<MockSocketMessenger>();
BufferedIOStream stream(messenger);
const auto buffer_size{1000};
const auto first_chunk_size{100};
// The writer will check the error code of the send function
// and will retry writing the next chunk when it doesn't get
// EAGAIN anymore from the sender.
EXPECT_CALL(*messenger, send_raw(_, buffer_size))
.Times(1)
.WillOnce(Return(first_chunk_size));
EXPECT_CALL(*messenger, send_raw(_, buffer_size - first_chunk_size))
.Times(2)
.WillOnce(DoAll(Invoke([](char const*, size_t) { errno = EAGAIN; }), Return(-EAGAIN)))
.WillOnce(Return(buffer_size - first_chunk_size));
char *ptr = static_cast<char*>(stream.allocBuffer(buffer_size));
ASSERT_NE(ptr, nullptr);
ASSERT_EQ(stream.commitBuffer(buffer_size), buffer_size);
}
TEST(BufferedIOStream, ReadWhenEnoughDataAvailable) {
auto messenger = std::make_shared<MockSocketMessenger>();
BufferedIOStream stream(messenger);
Buffer buffer;
buffer.push_back(0x12);
buffer.push_back(0x34);
stream.post_data(std::move(buffer));
std::uint8_t read_data[1] = {0x0};
size_t size = 1;
EXPECT_NE(nullptr, stream.read(read_data, &size));
EXPECT_EQ(1, size);
EXPECT_EQ(0x12, read_data[0]);
EXPECT_NE(nullptr, stream.read(read_data, &size));
EXPECT_EQ(1, size);
EXPECT_EQ(0x34, read_data[0]);
}
TEST(BufferedIOStream, ReadWithNoDataAvailable) {
auto messenger = std::make_shared<MockSocketMessenger>();
BufferedIOStream stream(messenger);
bool stopped = false;
std::thread producer([&](){
while (!stopped) {
if (stream.needs_data()) {
Buffer buffer;
buffer.push_back(0x12);
buffer.push_back(0x34);
stream.post_data(std::move(buffer));
}
std::this_thread::sleep_for(std::chrono::milliseconds{10});
}
});
size_t size{10};
std::uint8_t read_data[size] = {0x0};
EXPECT_NE(nullptr, stream.read(read_data, &size));
EXPECT_EQ(2, size);
EXPECT_EQ(0x12, read_data[0]);
EXPECT_EQ(0x34, read_data[1]);
stopped = true;
producer.join();
}
} // namespace graphics
} // namespace anbox