Skip to content

Commit

Permalink
add eventloop implementations for call_* and call_when_*_completed fu…
Browse files Browse the repository at this point in the history
…nctions
  • Loading branch information
philoinovsky committed Jun 26, 2021
1 parent 68fa9dc commit e6b292f
Show file tree
Hide file tree
Showing 6 changed files with 197 additions and 1 deletion.
1 change: 1 addition & 0 deletions build/Jamfile
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,7 @@ lib boost_python
import.cpp
exec.cpp
object/function_doc_signature.cpp
eventloop.cpp
: # requirements
<link>static:<define>BOOST_PYTHON_STATIC_LIB
<define>BOOST_PYTHON_SOURCE
Expand Down
1 change: 1 addition & 0 deletions include/boost/python.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@
# include <boost/python/docstring_options.hpp>
# include <boost/python/enum.hpp>
# include <boost/python/errors.hpp>
# include <boost/python/eventloop.hpp>
# include <boost/python/exception_translator.hpp>
# include <boost/python/exec.hpp>
# include <boost/python/extract.hpp>
Expand Down
111 changes: 111 additions & 0 deletions include/boost/python/eventloop.hpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,111 @@
// Copyright Pan Yue 2021.
// Distributed under the Boost Software License, Version 1.0. (See
// accompanying file LICENSE_1_0.txt or copy at
// http://www.boost.org/LICENSE_1_0.txt)

// TODO:
// 1. posix::stream_descriptor need windows version
// 2. call_* need return async.Handle
# ifndef EVENT_LOOP_PY2021_H_
# define EVENT_LOOP_PY2021_H_

#include <unordered_map>
#include <boost/asio.hpp>
#include <boost/python.hpp>

namespace a = boost::asio;
namespace c = std::chrono;
namespace py = boost::python;

namespace boost { namespace python { namespace eventloop {

class EventLoop
{
private:
int64_t _timer_id = 0;
a::io_context::strand _strand;
std::unordered_map<int, std::unique_ptr<a::steady_timer>> _id_to_timer_map;
// read: key = fd * 2 + 0, write: key = fd * 2 + 1
std::unordered_map<int, std::unique_ptr<a::posix::stream_descriptor>> _descriptor_map;
std::chrono::steady_clock::time_point _created_time;

void _add_reader_or_writer(int fd, py::object f, int key);
void _remove_reader_or_writer(int key);

public:
EventLoop(a::io_context& ctx):
_strand{ctx}, _created_time{std::chrono::steady_clock::now()}
{
}

// TODO: An instance of asyncio.Handle is returned, which can be used later to cancel the callback.
inline void call_soon(py::object f)
{
_strand.post([f, loop=this] {
f(boost::ref(*loop));
});
return;
}

// TODO: implement this
inline void call_soon_thread_safe(py::object f) {};

// Schedule callback to be called after the given delay number of seconds
// TODO: An instance of asyncio.Handle is returned, which can be used later to cancel the callback.
void call_later(double delay, py::object f);

void call_at(double when, py::object f);

inline double time()
{
return static_cast<std::chrono::duration<double>>(std::chrono::steady_clock::now() - _created_time).count();
}

// week 2 ......start......

inline void add_reader(int fd, py::object f)
{
_add_reader_or_writer(fd, f, fd * 2);
}

inline void remove_reader(int fd)
{
_remove_reader_or_writer(fd * 2);
}

inline void add_writer(int fd, py::object f)
{
_add_reader_or_writer(fd, f, fd * 2 + 1);
}

inline void remove_writer(int fd)
{
_remove_reader_or_writer(fd * 2 + 1);
}


void sock_recv(py::object sock, int bytes);

void sock_recv_into(py::object sock, py::object buffer);

void sock_sendall(py::object sock, py::object data);

void sock_connect(py::object sock, py::object address);

void sock_accept(py::object sock);

void sock_sendfile(py::object sock, py::object file, int offset = 0, int count = 0, bool fallback = true);

// week 2 ......end......

void run()
{
_strand.context().run();
}
};


}}}


# endif
81 changes: 81 additions & 0 deletions src/eventloop.cpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,81 @@
// Copyright Pan Yue 2021.
// Distributed under the Boost Software License, Version 1.0. (See
// accompanying file LICENSE_1_0.txt or copy at
// http://www.boost.org/LICENSE_1_0.txt)

// TODO:
// 1. posix::stream_descriptor need windows version
// 2. call_* need return async.Handle

#include <boost/asio.hpp>
#include <boost/bind.hpp>
#include <boost/python.hpp>

namespace a = boost::asio;
namespace c = std::chrono;
namespace py = boost::python;

namespace boost { namespace python { namespace eventloop {

void EventLoop::_add_reader_or_writer(int fd, py::object f, int key)
{
// add descriptor
if (_descriptor_map.find(key) == _descriptor_map.end())
{
_descriptor_map.emplace(key,
std::move(std::make_unique<a::posix::stream_descriptor>(_strand.context(), fd))
);
}

_descriptor_map.find(key)->second->async_wait(a::posix::descriptor::wait_type::wait_read,
a::bind_executor(_strand, [key, f, loop=this] (const boost::system::error_code& ec)
{
// move descriptor
auto iter = loop->_descriptor_map.find(key);
if (iter != loop->_descriptor_map.end())
{
iter->second->release();
loop->_descriptor_map.erase(iter);
}
loop->call_soon(f);
}));
return;
}

void EventLoop::_remove_reader_or_writer(int key)
{
auto iter = _descriptor_map.find(key);
if (iter != _descriptor_map.end())
{
iter->second->release();
_descriptor_map.erase(iter);
}
}

void EventLoop::call_later(double delay, py::object f)
{
// add timer
_id_to_timer_map.emplace(_timer_id,
std::move(std::make_unique<a::steady_timer>(_strand.context(),
std::chrono::steady_clock::now() + std::chrono::nanoseconds(int64_t(delay * 1e9))))
);

_id_to_timer_map.find(_timer_id)->second->async_wait(
// remove timer
a::bind_executor(_strand, [id=_timer_id, f, loop=this] (const boost::system::error_code& ec)
{
loop->_id_to_timer_map.erase(id);
loop->call_soon(f);
}));
_timer_id++;
}

void EventLoop::call_at(double when, py::object f)
{
double diff = when - time();
if (diff > 0)
return call_later(diff, f);
return call_soon(f);
}

}}}
3 changes: 2 additions & 1 deletion src/fabscript
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,8 @@ bpl = library('boost_python' + root.py_suffix,
'wrapper.cpp',
'import.cpp',
'exec.cpp',
'object/function_doc_signature.cpp'],
'object/function_doc_signature.cpp',
'eventloop.cpp'],
dependencies=root.config,
features=features + define('BOOST_PYTHON_SOURCE'))

Expand Down
1 change: 1 addition & 0 deletions test/Jamfile
Original file line number Diff line number Diff line change
Expand Up @@ -84,6 +84,7 @@ bpl-test crossmod_exception
: crossmod_exception.py crossmod_exception_a.cpp crossmod_exception_b.cpp
]

[ bpl-test eventloop ]
[ bpl-test injected ]
[ bpl-test properties ]
[ bpl-test return_arg ]
Expand Down

0 comments on commit e6b292f

Please sign in to comment.