Unverified Commit 58f1d373 authored by Joseph Noir's avatar Joseph Noir Committed by GitHub

Merge pull request #904

Add scheduled_send for sending with absolute timeout
parents db299ded 1b579742
......@@ -64,11 +64,11 @@ public:
group(intrusive_ptr<abstract_group> gptr);
inline explicit operator bool() const noexcept {
explicit operator bool() const noexcept {
return static_cast<bool>(ptr_);
}
inline bool operator!() const noexcept {
bool operator!() const noexcept {
return !ptr_;
}
......@@ -76,7 +76,7 @@ public:
intptr_t compare(const group& other) const noexcept;
inline intptr_t compare(const invalid_group_t&) const noexcept {
intptr_t compare(const invalid_group_t&) const noexcept {
return ptr_ ? 1 : 0;
}
......@@ -98,12 +98,16 @@ public:
friend error inspect(deserializer&, group&);
inline abstract_group* get() const noexcept {
abstract_group* get() const noexcept {
return ptr_.get();
}
/// @cond PRIVATE
actor_system& system() const {
return ptr_->system();
}
template <class... Ts>
void eq_impl(message_id mid, strong_actor_ptr sender,
execution_unit* ctx, Ts&&... xs) const {
......@@ -113,27 +117,25 @@ public:
make_message(std::forward<Ts>(xs)...), ctx);
}
inline bool subscribe(strong_actor_ptr who) const {
bool subscribe(strong_actor_ptr who) const {
if (!ptr_)
return false;
return ptr_->subscribe(std::move(who));
}
inline void unsubscribe(const actor_control_block* who) const {
void unsubscribe(const actor_control_block* who) const {
if (ptr_)
ptr_->unsubscribe(who);
}
/// CAF's messaging primitives assume a non-null guarantee. A group
/// object indirects pointer-like access to a group to prevent UB.
inline const group* operator->() const noexcept {
const group* operator->() const noexcept {
return this;
}
/// @endcond
private:
inline abstract_group* release() noexcept {
abstract_group* release() noexcept {
return ptr_.release();
}
......@@ -150,7 +152,7 @@ std::string to_string(const group& x);
namespace std {
template <>
struct hash<caf::group> {
inline size_t operator()(const caf::group& x) const {
size_t operator()(const caf::group& x) const {
// groups are singleton objects, the address is thus the best possible hash
return !x ? 0 : reinterpret_cast<size_t>(x.get());
}
......
This diff is collapsed.
......@@ -19,17 +19,19 @@
#pragma once
#include "caf/actor.hpp"
#include "caf/message.hpp"
#include "caf/actor_cast.hpp"
#include "caf/actor_addr.hpp"
#include "caf/message_id.hpp"
#include "caf/typed_actor.hpp"
#include "caf/actor_cast.hpp"
#include "caf/check_typed_input.hpp"
#include "caf/is_message_sink.hpp"
#include "caf/local_actor.hpp"
#include "caf/mailbox_element.hpp"
#include "caf/message.hpp"
#include "caf/message_id.hpp"
#include "caf/message_priority.hpp"
#include "caf/no_stages.hpp"
#include "caf/response_type.hpp"
#include "caf/system_messages.hpp"
#include "caf/is_message_sink.hpp"
#include "caf/message_priority.hpp"
#include "caf/check_typed_input.hpp"
#include "caf/typed_actor.hpp"
namespace caf {
......@@ -69,6 +71,7 @@ void send_as(const Source& src, const Dest& dest, Ts&&... xs) {
template <message_priority P = message_priority::normal, class Source,
class Dest, class... Ts>
void unsafe_send_as(Source* src, const Dest& dest, Ts&&... xs) {
static_assert(sizeof...(Ts) > 0, "no message to send");
if (dest)
actor_cast<abstract_actor*>(dest)->eq_impl(make_message_id(P),
src->ctrl(), src->context(),
......@@ -79,6 +82,7 @@ template <class... Ts>
void unsafe_response(local_actor* self, strong_actor_ptr src,
std::vector<strong_actor_ptr> stages, message_id mid,
Ts&&... xs) {
static_assert(sizeof...(Ts) > 0, "no message to send");
strong_actor_ptr next;
if (stages.empty()) {
next = src;
......@@ -108,6 +112,38 @@ void anon_send(const Dest& dest, Ts&&... xs) {
std::forward<Ts>(xs)...);
}
template <message_priority P = message_priority::normal, class Dest = actor,
class Rep = int, class Period = std::ratio<1>, class... Ts>
detail::enable_if_t<!std::is_same<Dest, group>::value>
delayed_anon_send(const Dest& dest, std::chrono::duration<Rep, Period> rtime,
Ts&&... xs) {
static_assert(sizeof...(Ts) > 0, "no message to send");
using token = detail::type_list<typename detail::implicit_conversions<
typename std::decay<Ts>::type>::type...>;
static_assert(response_type_unbox<signatures_of_t<Dest>, token>::valid,
"receiver does not accept given message");
if (dest) {
auto& clock = dest->home_system().clock();
auto timeout = clock.now() + rtime;
clock.schedule_message(timeout, actor_cast<strong_actor_ptr>(dest),
make_mailbox_element(nullptr, make_message_id(P),
no_stages,
std::forward<Ts>(xs)...));
}
}
template <class Rep = int, class Period = std::ratio<1>, class... Ts>
void delayed_anon_send(const group& dest,
std::chrono::duration<Rep, Period> rtime, Ts&&... xs) {
static_assert(sizeof...(Ts) > 0, "no message to send");
if (dest) {
auto& clock = dest->system().clock();
auto timeout = clock.now() + rtime;
clock.schedule_message(timeout, dest,
make_message(std::forward<Ts>(xs)...));
}
}
/// Anonymously sends `dest` an exit message.
template <class Dest>
void anon_send_exit(const Dest& dest, exit_reason reason) {
......
......@@ -48,7 +48,11 @@ public:
return &view_;
}
inline explicit operator bool() const {
const typed_actor_view<Sigs...>* operator->() const {
return &view_;
}
explicit operator bool() const {
return static_cast<bool>(view_.internal_ptr());
}
......
......@@ -77,6 +77,10 @@ public:
return self_->system();
}
actor_system& home_system() const {
return self_->home_system();
}
void quit(exit_reason reason = exit_reason::normal) {
self_->quit(reason);
}
......
......@@ -212,7 +212,7 @@ class delayed_testee : public composable_behavior<delayed_testee_actor> {
public:
result<void> operator()(int x) override {
CAF_CHECK_EQUAL(x, 42);
self->delayed_anon_send(self, std::chrono::milliseconds(10), true);
delayed_anon_send(self, std::chrono::milliseconds(10), true);
return unit;
}
......
......@@ -5,7 +5,7 @@
* | |___ / ___ \| _| Framework *
* \____/_/ \_|_| *
* *
* Copyright 2011-2018 Dominik Charousset *
* Copyright 2011-2019 Dominik Charousset *
* *
* Distributed under the terms and conditions of the BSD 3-Clause License or *
* (at your option) under the terms and conditions of the Boost Software *
......@@ -16,16 +16,14 @@
* http://www.boost.org/LICENSE_1_0.txt. *
******************************************************************************/
#define CAF_SUITE delayed_send
#define CAF_SUITE sender
#include <chrono>
#include "caf/actor_system.hpp"
#include "caf/behavior.hpp"
#include "caf/event_based_actor.hpp"
#include "caf/mixin/sender.hpp"
#include "caf/test/dsl.hpp"
#include <chrono>
using namespace caf;
using std::chrono::seconds;
......@@ -34,32 +32,51 @@ namespace {
behavior testee_impl(event_based_actor* self) {
self->set_default_handler(drop);
return {
[] {
return {[] {
// nop
}
};
}};
}
} // namespace <anonymous>
struct fixture : test_coordinator_fixture<> {
group grp;
actor testee;
fixture() {
grp = sys.groups().anonymous();
testee = sys.spawn_in_group(grp, testee_impl);
}
~fixture() {
anon_send_exit(testee, exit_reason::user_shutdown);
}
};
} // namespace
CAF_TEST_FIXTURE_SCOPE(request_timeout_tests, test_coordinator_fixture<>)
CAF_TEST_FIXTURE_SCOPE(sender_tests, fixture)
CAF_TEST(delayed actor message) {
auto testee = sys.spawn(testee_impl);
self->delayed_send(testee, seconds(1), "hello world");
sched.trigger_timeout();
expect((std::string), from(self).to(testee).with("hello world"));
}
CAF_TEST(delayed group message) {
auto grp = sys.groups().anonymous();
auto testee = sys.spawn_in_group(grp, testee_impl);
self->delayed_send(grp, seconds(1), "hello world");
sched.trigger_timeout();
expect((std::string), from(self).to(testee).with("hello world"));
// The group keeps a reference, so we need to shutdown 'manually'.
anon_send_exit(testee, exit_reason::user_shutdown);
}
CAF_TEST(scheduled actor message) {
self->scheduled_send(testee, self->clock().now() + seconds(1), "hello world");
sched.trigger_timeout();
expect((std::string), from(self).to(testee).with("hello world"));
}
CAF_TEST(scheduled group message) {
self->scheduled_send(grp, self->clock().now() + seconds(1), "hello world");
sched.trigger_timeout();
expect((std::string), from(self).to(testee).with("hello world"));
}
CAF_TEST_FIXTURE_SCOPE_END()
......@@ -57,7 +57,7 @@ timer::behavior_type timer_impl(timer::stateful_pointer<timer_state> self) {
timer::behavior_type timer_impl2(timer::pointer self) {
auto had_reset = std::make_shared<bool>(false);
self->delayed_anon_send(self, ms(100), reset_atom::value);
delayed_anon_send(self, ms(100), reset_atom::value);
return {
[=](reset_atom) {
CAF_MESSAGE("timer reset");
......
Markdown is supported
0%
or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment