From dfbed0c850748345964227c3e48e7a1c46de3d9a Mon Sep 17 00:00:00 2001 From: Michael Davidsaver Date: Wed, 14 Jul 2021 11:05:25 -0700 Subject: [PATCH] server ExecOp timer --- src/evhelper.cpp | 78 ++++++++++++++++++++++++++++++++++++++++++++ src/evhelper.h | 22 +++++++++++++ src/instcounters.h | 1 + src/pvxs/srvcommon.h | 11 +++++++ src/pvxs/util.h | 25 ++++++++++++++ src/serverget.cpp | 9 +++++ test/testget.cpp | 40 ++++++++++++++++++++++- 7 files changed, 185 insertions(+), 1 deletion(-) diff --git a/src/evhelper.cpp b/src/evhelper.cpp index bd57eb2..ae32ec9 100644 --- a/src/evhelper.cpp +++ b/src/evhelper.cpp @@ -43,6 +43,7 @@ size_t min_slice_size = 1024u; namespace pvxs {namespace impl { DEFINE_LOGGER(logerr, "pvxs.loop"); +DEFINE_LOGGER(logtimer, "pvxs.timer"); namespace mdetail { VFunctor0::~VFunctor0() {} @@ -665,4 +666,81 @@ void to_evbuf(evbuffer *buf, const Header& H, bool be) } // namespace impl +Timer::~Timer() {} + +bool Timer::cancel() +{ + if(!pvt) + throw std::logic_error("NULL Timer"); + + auto P(std::move(pvt)); + + return P->cancel(); +} + +Timer::Pvt::~Pvt() { + log_debug_printf(logtimer, "Timer %p %s\n", this, __func__); + (void)cancel(); +} + +bool Timer::Pvt::cancel() +{ + bool ret = false; + decltype (cb) trash; + + log_debug_printf(logtimer, "Timer %p pcancel\n", this); + + base.call([this, &ret, &trash](){ + trash = std::move(cb); + if(auto T = std::move(timer)) { + log_debug_printf(logtimer, "Timer %p dispose %p\n", this, T.get()); + ret = event_pending(T.get(), EV_TIMEOUT, nullptr); + (void)event_del(T.get()); + } + }); + + return ret; +} + +static +void expire_cb(evutil_socket_t, short, void * raw) +{ + auto self(static_cast(raw)); + log_debug_printf(logtimer, "Timer %p expires\n", self); + assert(self->base.base); + + try { + self->cb(); + } catch(std::exception& e){ + log_exc_printf(logtimer, "Unhandled exception in Timer callback: %s\n", e.what()); + } +} + +Timer Timer::Pvt::buildOneShot(double delay, const evbase& base, std::function&& cb) +{ + if(!cb) + throw std::invalid_argument("NULL cb"); + + Timer ret; + ret.pvt = std::make_shared(base, std::move(cb)); + + base.call([&ret, &base, delay](){ + evevent timer(event_new(base.base, -1, EV_TIMEOUT, &expire_cb, ret.pvt.get())); + ret.pvt->timer = std::move(timer); + + auto timo(totv(delay)); + + if(event_add(ret.pvt->timer.get(), &timo)) + throw std::runtime_error("Unable to start oneshot timer"); + + log_debug_printf(logtimer, "Create timer %p as %p with delay %f and %s\n", + ret.pvt.get(), + ret.pvt->timer.get(), + delay, + ret.pvt->cb.target_type().name()); + }); + + return ret; +} + } // namespace pvxs diff --git a/src/evhelper.h b/src/evhelper.h index af32e8b..3ba095f 100644 --- a/src/evhelper.h +++ b/src/evhelper.h @@ -233,6 +233,28 @@ struct PVXS_API evsocket } // namespace impl +#ifdef PVXS_EXPERT_API_ENABLED + +struct Timer::Pvt { + const evbase base; + std::function cb; + evevent timer; + + Pvt(const evbase& base, std::function&& cb) + :base(base), cb(std::move(cb)) + {} + ~Pvt(); + + bool cancel(); + + static + Timer buildOneShot(double delay, const evbase &base, std::function&& cb); + + INST_COUNTER(Timer); +}; + +#endif // PVXS_EXPERT_API_ENABLED + } // namespace pvxs #endif /* EVHELPER_H */ diff --git a/src/instcounters.h b/src/instcounters.h index 291794e..80956de 100644 --- a/src/instcounters.h +++ b/src/instcounters.h @@ -11,6 +11,7 @@ CASE(StructTop); CASE(UDPListener); CASE(evbase); CASE(evbaseRunning); +CASE(Timer); CASE(GPROp); CASE(Connection); diff --git a/src/pvxs/srvcommon.h b/src/pvxs/srvcommon.h index 8aac58e..a6d4414 100644 --- a/src/pvxs/srvcommon.h +++ b/src/pvxs/srvcommon.h @@ -11,6 +11,7 @@ #include #include +#include #include namespace pvxs { @@ -98,6 +99,16 @@ public: const Value& pvRequest() const { return _pvRequest; } virtual ~ExecOp(); + +#ifdef PVXS_EXPERT_API_ENABLED + //! Create/start timer. cb runs on worker associated with Channel of this Operation. + //! @since UNRELEASED + Timer timerOneShot(double delay, std::function&& cb) { + return _timerOneShot(delay, std::move(cb)); + } +#endif // PVXS_EXPERT_API_ENABLED +private: + virtual Timer _timerOneShot(double delay, std::function&& cb) =0; }; }} // namespace pvxs::server diff --git a/src/pvxs/util.h b/src/pvxs/util.h index 26eb5e2..c322dc2 100644 --- a/src/pvxs/util.h +++ b/src/pvxs/util.h @@ -14,6 +14,7 @@ #include #include #include +#include #include #include @@ -295,6 +296,30 @@ public: } }; +struct Timer; + +#ifdef PVXS_EXPERT_API_ENABLED + +//! Timer associated with a client::Context or server::Server +//! @since UNRELEASED +struct PVXS_API Timer { + struct Pvt; + + //! dtor implicitly cancel()s + ~Timer(); + //! Explicit cancel. + //! @returns true if the timer was running, and now is not. + bool cancel(); + + explicit operator bool() const { return pvt.operator bool(); } + +private: + std::shared_ptr pvt; + friend struct Pvt; +}; + +#endif // PVXS_EXPERT_API_ENABLED + } // namespace pvxs #endif // PVXS_UTIL_H diff --git a/src/serverget.cpp b/src/serverget.cpp index 2497655..4f43d78 100644 --- a/src/serverget.cpp +++ b/src/serverget.cpp @@ -326,6 +326,15 @@ struct ServerGPRExec : public server::ExecOp }); } + virtual Timer _timerOneShot(double delay, std::function&& fn) override final + { + auto serv = server.lock(); + if(!serv) + throw std::logic_error("Can't start timer on deal server"); + + return Timer::Pvt::buildOneShot(delay, serv->acceptor_loop.internal(), std::move(fn)); + } + const std::weak_ptr server; const std::weak_ptr op; diff --git a/test/testget.cpp b/test/testget.cpp index 71487ea..2f5bafb 100644 --- a/test/testget.cpp +++ b/test/testget.cpp @@ -367,6 +367,43 @@ struct Tester { testShow()<wait(4.0); })<<" pvRequest selects no fields"; } + + void delayExec() + { + testShow()<<__func__; + + epicsEvent done; + Timer slowdown; + + mbox.onPut([this, &done, &slowdown](server::SharedPV& pv, std::unique_ptr&& rawop, Value&& rawval) { + // on server worker + std::shared_ptr op(std::move(rawop)); + auto val(std::move(rawval)); + testPass("In onPut"); + + slowdown = op->timerOneShot(0.01, [](){ + testFail("I should not run."); + }); + + testTrue(slowdown.cancel()); + + slowdown = op->timerOneShot(0.01, [this, &done, op, val](){ + testPass("I should run"); + done.signal(); + mbox.post(val); + op->reply(); + }); + + // op->reply() from timer + }); + + mbox.open(initial); + serv.start(); + + auto op = cli.put("mailbox") + .set("value", 42) + .exec()->wait(5.0); + } }; struct ErrorSource : public server::Source @@ -439,7 +476,7 @@ void testError(bool phase) MAIN(testget) { - testPlan(53); + testPlan(56); testSetup(); logger_config_env(); Tester().testConnector(); @@ -452,6 +489,7 @@ MAIN(testget) Tester().orphan(); Tester().manualExec(); Tester().badRequest(); + Tester().delayExec(); testError(false); testError(true); cleanup_for_valgrind();