From 961ca1d7da818be534a728791ee0afd5142f607e Mon Sep 17 00:00:00 2001 From: Oleg Utkin Date: Wed, 10 Nov 2021 09:46:49 +0000 Subject: [PATCH 01/17] add pop_many func for executor --- src/txapi.rs | 23 ++++++++++++++++++++++- 1 file changed, 22 insertions(+), 1 deletion(-) diff --git a/src/txapi.rs b/src/txapi.rs index 0785235..977498c 100644 --- a/src/txapi.rs +++ b/src/txapi.rs @@ -7,7 +7,7 @@ use std::os::unix::io::{AsRawFd, RawFd}; use std::io; use mlua::Lua; -type Task = Box Result<(), ChannelError> + Send>; +pub type Task = Box Result<(), ChannelError> + Send>; type TaskSender = async_channel::Sender; type TaskReceiver = async_channel::Receiver; @@ -105,6 +105,27 @@ impl Executor { } } + pub fn pop_many(&self, max_tasks: usize, coio_timeout: f64) -> Result, ChannelError> { + if self.task_rx.is_empty() { + let _ = self.eventfd.coio_read(coio_timeout); + } + + let mut tasks = Vec::with_capacity(max_tasks); + for _ in 0..max_tasks { + match self.task_rx.try_recv() { + Ok(func) => tasks.push(func), + Err(TryRecvError::Empty) => break, + Err(TryRecvError::Closed) => return Err(ChannelError::RXChannelClosed), + }; + + if self.task_rx.len() <= 1 { + break; + } + } + + Ok(tasks) + } + pub fn try_clone(&self) -> io::Result { Ok(Self { task_rx: self.task_rx.clone(), From aa83810482281a3a1aad4165db8463df56142aa6 Mon Sep 17 00:00:00 2001 From: Oleg Utkin Date: Fri, 12 Nov 2021 13:54:52 +0000 Subject: [PATCH 02/17] wip --- Cargo.toml | 1 + examples/bench/init.lua | 6 +- examples/bench/src/lib.rs | 14 ++++- src/config.rs | 6 +- src/fiber_pool.rs | 128 ++++++++++++++++++++++++++++++++++++++ src/lib.rs | 34 ++-------- 6 files changed, 155 insertions(+), 34 deletions(-) create mode 100644 src/fiber_pool.rs diff --git a/Cargo.toml b/Cargo.toml index 6eafbe1..12a4926 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -13,6 +13,7 @@ thiserror = "1.0" libc = "0.2" crossbeam-utils = "0.8" serde = "1.0" +crossbeam-channel = "0.5" [dependencies.mlua] version = "0.6" diff --git a/examples/bench/init.lua b/examples/bench/init.lua index 18c42b0..046c4ed 100644 --- a/examples/bench/init.lua +++ b/examples/bench/init.lua @@ -1,5 +1,9 @@ local txapi = require('bench') -txapi.start({}) -- will use default +txapi.start({ + fibers = 4, + max_batch = 8, + runtime = { type = "cur_thread" }, +}) -- will use default -- txapi.start({buffer = 128 }) -- txapi.start({buffer = 128, runtime = { type = "cur_thread" }}) -- txapi.start({buffer = 128, runtime = { type = "multi_thread" }}) -- default diff --git a/examples/bench/src/lib.rs b/examples/bench/src/lib.rs index 134e68f..cd3a286 100644 --- a/examples/bench/src/lib.rs +++ b/examples/bench/src/lib.rs @@ -5,7 +5,18 @@ use tokio::time::Instant; use xtm_rust::{run_module, Dispatcher, ModuleConfig}; async fn module_main(dispatcher: Dispatcher) { - let iterations = 10_000_000; + tokio::spawn({ + let dispatcher = dispatcher.try_clone().unwrap(); + async move { + let mut interval = tokio::time::interval(std::time::Duration::from_secs(1)); + loop { + println!("task_queue: {:>3}", dispatcher.len()); + interval.tick().await; + } + } + }); + + let iterations = 1_000_000; let worker_n = 6; let iterations_per_worker = iterations / worker_n; @@ -57,6 +68,7 @@ fn bench(lua: &Lua) -> LuaResult { "start", lua.create_function_mut(|lua, (config,): (LuaValue,)| { let config: ModuleConfig = lua.from_value(config)?; + println!("{:?}", config); run_module(module_main, config, lua).map_err(LuaError::external) })?, diff --git a/src/config.rs b/src/config.rs index ca6f775..9def20a 100644 --- a/src/config.rs +++ b/src/config.rs @@ -1,12 +1,13 @@ use serde::{Deserialize, Serialize}; use tokio::runtime; -#[derive(Debug, Serialize, Deserialize)] +#[derive(Debug, Serialize, Deserialize, Clone)] #[serde(default)] pub struct ModuleConfig { pub buffer: usize, pub fibers: usize, pub max_recv_retries: usize, + pub max_batch: usize, pub coio_timeout: f64, pub runtime: RuntimeConfig, } @@ -17,13 +18,14 @@ impl Default for ModuleConfig { buffer: 128, fibers: 16, max_recv_retries: 100, + max_batch: 16, coio_timeout: 1.0, runtime: RuntimeConfig::default(), } } } -#[derive(Debug, Serialize, Deserialize)] +#[derive(Debug, Serialize, Deserialize, Clone)] #[serde(tag = "type")] pub enum RuntimeConfig { #[serde(rename(deserialize = "cur_thread"))] diff --git a/src/fiber_pool.rs b/src/fiber_pool.rs new file mode 100644 index 0000000..12322fb --- /dev/null +++ b/src/fiber_pool.rs @@ -0,0 +1,128 @@ +use std::{collections::LinkedList, rc::Rc}; + +use crossbeam_channel::{unbounded, TryRecvError}; +use mlua::Lua; +use tarantool::fiber; + +use crate::{ChannelError, Executor, ModuleConfig, Task}; + +struct SchedulerArgs<'a> { + lua: &'a Lua, + executor: Executor, + config: ModuleConfig, +} +fn scheduler_f(args: Box) -> i32 { + let SchedulerArgs { + lua, + executor, + config, + } = *args; + let ModuleConfig { + max_batch, + coio_timeout, + fibers, + .. + } = config; + + let cond = Rc::new(fiber::Cond::new()); + let (tx, rx) = unbounded::(); + + let mut workers = LinkedList::new(); + for _ in 0..fibers { + let mut worker = fiber::Fiber::new("worker", &mut worker_f); + worker.start(WorkerArgs { + cond: cond.clone(), + lua, + rx: rx.clone(), + }); + workers.push_back(worker); + } + + loop { + let tasks = match executor.pop_many(max_batch, coio_timeout) { + Ok(tasks) => tasks, + Err(ChannelError::TXChannelClosed) => break, + Err(ChannelError::RXChannelClosed) => return 0, + Err(_err) => return -1, + }; + + for task in tasks { + tx.send(task).unwrap(); + cond.signal(); + } + } + + for worker in workers { + worker.join(); + } + + 0 +} + +struct WorkerArgs<'a> { + cond: Rc, + lua: &'a Lua, + rx: crossbeam_channel::Receiver, +} +fn worker_f(args: Box) -> i32 { + let WorkerArgs { cond, lua, rx } = *args; + + let thread_func = lua + .create_function(move |lua, _: ()| { + loop { + match rx.try_recv() { + Ok(task) => task(lua), + Err(TryRecvError::Disconnected) => return Ok(()), + Err(TryRecvError::Empty) => { + cond.wait(); + // let ok = cond.wait_timeout(std::time::Duration::from_secs(1)); + // if !ok { + // kill fiber + // return Ok(()); + // } + Ok(()) + } + } + .unwrap(); + } + }) + .unwrap(); + let thread = lua.create_thread(thread_func.clone()).unwrap(); + let _: () = thread.resume(()).unwrap(); + + 0 +} + +pub(crate) struct FiberPool<'a> { + lua: &'a Lua, + executor: Executor, + config: ModuleConfig, + + scheduler: fiber::Fiber<'a, SchedulerArgs<'a>>, +} + +impl<'a> FiberPool<'a> { + pub fn new(lua: &'a Lua, executor: Executor, config: ModuleConfig) -> Self { + let mut scheduler = fiber::Fiber::new("scheduler", &mut scheduler_f); + scheduler.set_joinable(true); + Self { + lua, + executor, + config, + scheduler, + } + } + + pub fn run(&mut self) { + self.scheduler.start(SchedulerArgs { + lua: self.lua, + executor: self.executor.try_clone().unwrap(), + config: self.config.clone(), + }); + } + + // join will exit when all Dispatchers die + pub fn join(&self) { + self.scheduler.join(); + } +} diff --git a/src/lib.rs b/src/lib.rs index caf50df..a0d2123 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -4,12 +4,11 @@ use std::future::Future; use crossbeam_utils::thread; use mlua::Lua; use tokio::runtime; - pub use txapi::*; -use tarantool::fiber::Fiber; mod eventfd; mod txapi; +mod fiber_pool; mod config; pub use config::*; @@ -27,31 +26,8 @@ where { let (dispatcher, executor) = channel(config.buffer)?; - let executor_loop = &mut |args: Box<(&Lua, Executor)>| { - let (lua, executor) = *args; - - let thread_func = lua.create_function(move |lua, _: ()| { - Ok(loop { - match executor.exec(lua, config.max_recv_retries, config.coio_timeout) { - Ok(_) => continue, - Err(ChannelError::TXChannelClosed) => continue, - Err(ChannelError::RXChannelClosed) => break 0, - Err(_err) => break -1, - } - }) - }).unwrap(); - let thread = lua.create_thread(thread_func).unwrap(); - thread.resume(()).unwrap() - }; - - // UNSAFE: fibers must die inside the current function - let mut fibers = Vec::with_capacity(config.fibers); - for _ in 0..config.fibers { - let mut fiber = Fiber::new("xtm", executor_loop); - fiber.set_joinable(true); - fiber.start((lua, executor.try_clone()?)); - fibers.push(fiber); - } + let mut fiber_pool = fiber_pool::FiberPool::new(lua, executor, config.clone()); + fiber_pool.run(); let result = thread::scope(|scope| { let module_thread = scope @@ -67,9 +43,7 @@ where }) .unwrap(); - for fiber in &fibers { - fiber.join(); - } + fiber_pool.join(); module_thread.join().unwrap().unwrap() }) .unwrap(); From 06fce204efaf756e6155efb4f6f2fadc318d479c Mon Sep 17 00:00:00 2001 From: Oleg Utkin Date: Fri, 12 Nov 2021 15:30:49 +0000 Subject: [PATCH 03/17] remove retries --- src/config.rs | 2 -- src/txapi.rs | 9 +-------- 2 files changed, 1 insertion(+), 10 deletions(-) diff --git a/src/config.rs b/src/config.rs index 9def20a..3a8f884 100644 --- a/src/config.rs +++ b/src/config.rs @@ -6,7 +6,6 @@ use tokio::runtime; pub struct ModuleConfig { pub buffer: usize, pub fibers: usize, - pub max_recv_retries: usize, pub max_batch: usize, pub coio_timeout: f64, pub runtime: RuntimeConfig, @@ -17,7 +16,6 @@ impl Default for ModuleConfig { Self { buffer: 128, fibers: 16, - max_recv_retries: 100, max_batch: 16, coio_timeout: 1.0, runtime: RuntimeConfig::default(), diff --git a/src/txapi.rs b/src/txapi.rs index 977498c..a22a9f8 100644 --- a/src/txapi.rs +++ b/src/txapi.rs @@ -86,7 +86,7 @@ impl Executor { Self { task_rx, eventfd } } - pub fn exec(&self, lua: &Lua, max_recv_retries: usize, coio_timeout: f64) -> Result<(), ChannelError> { + pub fn exec(&self, lua: &Lua, coio_timeout: f64) -> Result<(), ChannelError> { loop { match self.task_rx.try_recv() { Ok(func) => return func(lua), @@ -94,13 +94,6 @@ impl Executor { Err(TryRecvError::Closed) => return Err(ChannelError::RXChannelClosed), }; - for _ in 0..max_recv_retries { - match self.task_rx.try_recv() { - Ok(func) => return func(lua), - Err(TryRecvError::Empty) => tarantool::fiber::sleep(0.), - Err(TryRecvError::Closed) => return Err(ChannelError::RXChannelClosed), - }; - } let _ = self.eventfd.coio_read(coio_timeout); } } From 32ac69bfd5432a8862234681fa0881c35c984342 Mon Sep 17 00:00:00 2001 From: Oleg Utkin Date: Sun, 14 Nov 2021 14:04:13 +0000 Subject: [PATCH 04/17] handle executor clone error --- src/fiber_pool.rs | 5 +++-- src/lib.rs | 2 +- 2 files changed, 4 insertions(+), 3 deletions(-) diff --git a/src/fiber_pool.rs b/src/fiber_pool.rs index 12322fb..d05c1c3 100644 --- a/src/fiber_pool.rs +++ b/src/fiber_pool.rs @@ -113,12 +113,13 @@ impl<'a> FiberPool<'a> { } } - pub fn run(&mut self) { + pub fn run(&mut self) -> std::io::Result<()> { self.scheduler.start(SchedulerArgs { lua: self.lua, - executor: self.executor.try_clone().unwrap(), + executor: self.executor.try_clone()?, config: self.config.clone(), }); + Ok(()) } // join will exit when all Dispatchers die diff --git a/src/lib.rs b/src/lib.rs index a0d2123..4def8f6 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -27,7 +27,7 @@ where let (dispatcher, executor) = channel(config.buffer)?; let mut fiber_pool = fiber_pool::FiberPool::new(lua, executor, config.clone()); - fiber_pool.run(); + fiber_pool.run()?; let result = thread::scope(|scope| { let module_thread = scope From f3b5c074ad89b92ed142fac213711484f863f8f5 Mon Sep 17 00:00:00 2001 From: Oleg Utkin Date: Sun, 14 Nov 2021 15:31:35 +0000 Subject: [PATCH 05/17] add fiber_standby_timeout parameter --- src/config.rs | 4 +++- src/fiber_pool.rs | 29 +++++++++++++++++++---------- 2 files changed, 22 insertions(+), 11 deletions(-) diff --git a/src/config.rs b/src/config.rs index 3a8f884..9de646a 100644 --- a/src/config.rs +++ b/src/config.rs @@ -8,6 +8,7 @@ pub struct ModuleConfig { pub fibers: usize, pub max_batch: usize, pub coio_timeout: f64, + pub fiber_standby_timeout: f64, pub runtime: RuntimeConfig, } @@ -17,7 +18,8 @@ impl Default for ModuleConfig { buffer: 128, fibers: 16, max_batch: 16, - coio_timeout: 1.0, + coio_timeout: 0.1, + fiber_standby_timeout: 1.0, runtime: RuntimeConfig::default(), } } diff --git a/src/fiber_pool.rs b/src/fiber_pool.rs index d05c1c3..a71154f 100644 --- a/src/fiber_pool.rs +++ b/src/fiber_pool.rs @@ -15,14 +15,15 @@ fn scheduler_f(args: Box) -> i32 { let SchedulerArgs { lua, executor, - config, - } = *args; - let ModuleConfig { + config: + ModuleConfig { max_batch, coio_timeout, fibers, + fiber_standby_timeout, .. - } = config; + }, + } = *args; let cond = Rc::new(fiber::Cond::new()); let (tx, rx) = unbounded::(); @@ -34,6 +35,7 @@ fn scheduler_f(args: Box) -> i32 { cond: cond.clone(), lua, rx: rx.clone(), + fiber_standby_timeout, }); workers.push_back(worker); } @@ -63,9 +65,15 @@ struct WorkerArgs<'a> { cond: Rc, lua: &'a Lua, rx: crossbeam_channel::Receiver, + fiber_standby_timeout: f64, } fn worker_f(args: Box) -> i32 { - let WorkerArgs { cond, lua, rx } = *args; + let WorkerArgs { + cond, + lua, + rx, + fiber_standby_timeout, + } = *args; let thread_func = lua .create_function(move |lua, _: ()| { @@ -74,11 +82,12 @@ fn worker_f(args: Box) -> i32 { Ok(task) => task(lua), Err(TryRecvError::Disconnected) => return Ok(()), Err(TryRecvError::Empty) => { - cond.wait(); - // let ok = cond.wait_timeout(std::time::Duration::from_secs(1)); - // if !ok { - // kill fiber - // return Ok(()); + let signaled = cond.wait_timeout(std::time::Duration::from_secs_f64( + fiber_standby_timeout, + )); + // if !signaled { + // // kill fiber + // return Ok(()); // } Ok(()) } From f153c6dde62e92551a6b538f5b16553b44761445 Mon Sep 17 00:00:00 2001 From: Oleg Utkin Date: Sun, 14 Nov 2021 16:16:00 +0000 Subject: [PATCH 06/17] remove redundant clone --- src/fiber_pool.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/fiber_pool.rs b/src/fiber_pool.rs index a71154f..44a42d6 100644 --- a/src/fiber_pool.rs +++ b/src/fiber_pool.rs @@ -95,8 +95,8 @@ fn worker_f(args: Box) -> i32 { .unwrap(); } }) - .unwrap(); - let thread = lua.create_thread(thread_func.clone()).unwrap(); + .unwrap(); // TODO: fix TXChannelClosed (oneshot receiver was dropped) + let thread = lua.create_thread(thread_func).unwrap(); let _: () = thread.resume(()).unwrap(); 0 From befd794d1ed83078f1d05c1a9cbfe458a0fe70c4 Mon Sep 17 00:00:00 2001 From: Oleg Utkin Date: Sun, 14 Nov 2021 19:13:43 +0000 Subject: [PATCH 07/17] fix scheduler dead lock --- src/fiber_pool.rs | 29 ++++++++++++++++++----------- 1 file changed, 18 insertions(+), 11 deletions(-) diff --git a/src/fiber_pool.rs b/src/fiber_pool.rs index 44a42d6..d41d61c 100644 --- a/src/fiber_pool.rs +++ b/src/fiber_pool.rs @@ -17,11 +17,11 @@ fn scheduler_f(args: Box) -> i32 { executor, config: ModuleConfig { - max_batch, - coio_timeout, - fibers, + max_batch, + coio_timeout, + fibers, fiber_standby_timeout, - .. + .. }, } = *args; @@ -31,6 +31,7 @@ fn scheduler_f(args: Box) -> i32 { let mut workers = LinkedList::new(); for _ in 0..fibers { let mut worker = fiber::Fiber::new("worker", &mut worker_f); + worker.set_joinable(true); worker.start(WorkerArgs { cond: cond.clone(), lua, @@ -40,25 +41,31 @@ fn scheduler_f(args: Box) -> i32 { workers.push_back(worker); } - loop { + let result = loop { let tasks = match executor.pop_many(max_batch, coio_timeout) { Ok(tasks) => tasks, - Err(ChannelError::TXChannelClosed) => break, - Err(ChannelError::RXChannelClosed) => return 0, - Err(_err) => return -1, + Err(ChannelError::RXChannelClosed) => break Ok(()), + Err(err) => break Err(err), }; for task in tasks { - tx.send(task).unwrap(); + tx.send(task).unwrap(); // TODO: add error handling cond.signal(); } - } + }; + + // gracefully kill fibers + drop(tx); + cond.broadcast(); for worker in workers { worker.join(); } - 0 + match result { + Ok(_) => 0, + Err(_) => -1, + } } struct WorkerArgs<'a> { From 0326cdb5b749e91b99424e5c113993c7315caec0 Mon Sep 17 00:00:00 2001 From: Oleg Utkin Date: Mon, 15 Nov 2021 11:00:12 +0000 Subject: [PATCH 08/17] change bench config --- examples/bench/init.lua | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/examples/bench/init.lua b/examples/bench/init.lua index 046c4ed..d965c6c 100644 --- a/examples/bench/init.lua +++ b/examples/bench/init.lua @@ -1,7 +1,7 @@ local txapi = require('bench') txapi.start({ - fibers = 4, - max_batch = 8, + fibers = 16, + max_batch = 16, runtime = { type = "cur_thread" }, }) -- will use default -- txapi.start({buffer = 128 }) From f34943bfb3d4a94eb42a443928f7f32899c07215 Mon Sep 17 00:00:00 2001 From: Oleg Utkin Date: Fri, 26 Nov 2021 15:26:50 +0000 Subject: [PATCH 09/17] fix TXChannelClosed (oneshot receiver was dropped) --- src/fiber_pool.rs | 28 +++++++++++++++------------- 1 file changed, 15 insertions(+), 13 deletions(-) diff --git a/src/fiber_pool.rs b/src/fiber_pool.rs index d41d61c..30387f4 100644 --- a/src/fiber_pool.rs +++ b/src/fiber_pool.rs @@ -1,4 +1,4 @@ -use std::{collections::LinkedList, rc::Rc}; +use std::{collections::LinkedList, rc::Rc, time::Duration}; use crossbeam_channel::{unbounded, TryRecvError}; use mlua::Lua; @@ -81,32 +81,34 @@ fn worker_f(args: Box) -> i32 { rx, fiber_standby_timeout, } = *args; + let fiber_standby_timeout = Duration::from_secs_f64(fiber_standby_timeout); let thread_func = lua .create_function(move |lua, _: ()| { loop { match rx.try_recv() { - Ok(task) => task(lua), - Err(TryRecvError::Disconnected) => return Ok(()), + Ok(task) => match task(lua) { + Ok(()) => (), + Err(ChannelError::TXChannelClosed) => continue, + Err(err) => break Err(mlua::Error::external(err)), + }, + Err(TryRecvError::Disconnected) => break Ok(()), Err(TryRecvError::Empty) => { - let signaled = cond.wait_timeout(std::time::Duration::from_secs_f64( - fiber_standby_timeout, - )); + let signaled = cond.wait_timeout(fiber_standby_timeout); // if !signaled { // // kill fiber - // return Ok(()); + // break Ok(()); // } - Ok(()) } } - .unwrap(); } }) - .unwrap(); // TODO: fix TXChannelClosed (oneshot receiver was dropped) + .unwrap(); let thread = lua.create_thread(thread_func).unwrap(); - let _: () = thread.resume(()).unwrap(); - - 0 + match thread.resume(()) { + Ok(()) => 0, + Err(_) => -1, + } } pub(crate) struct FiberPool<'a> { From f3d71936fd8028264fffc0afd10c504d0ef52fb9 Mon Sep 17 00:00:00 2001 From: Oleg Utkin Date: Fri, 26 Nov 2021 15:27:04 +0000 Subject: [PATCH 10/17] change bench parameters --- examples/bench/src/lib.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/examples/bench/src/lib.rs b/examples/bench/src/lib.rs index cd3a286..d5ac83b 100644 --- a/examples/bench/src/lib.rs +++ b/examples/bench/src/lib.rs @@ -16,9 +16,9 @@ async fn module_main(dispatcher: Dispatcher) { } }); - let iterations = 1_000_000; + let iterations = 4_000_000; - let worker_n = 6; + let worker_n = 16; let iterations_per_worker = iterations / worker_n; let mut workers = Vec::new(); From 2e3418f1e197b96ab149715f8f0c374b607ee078 Mon Sep 17 00:00:00 2001 From: Alexey Andreev Date: Tue, 26 Oct 2021 13:45:04 +0300 Subject: [PATCH 11/17] added notifier parameter --- src/lib.rs | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/src/lib.rs b/src/lib.rs index 4def8f6..60d6589 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -13,13 +13,17 @@ mod fiber_pool; mod config; pub use config::*; +use tokio::sync::Notify; +use std::sync::Arc; + pub fn run_module( module_main: Func, config: ModuleConfig, lua: &Lua, + notifier : Arc, ) -> io::Result where - Func: FnOnce(Dispatcher) -> Fut, + Func: FnOnce(Dispatcher, Arc) -> Fut, Func: Send, Fut: Future, Fut::Output: Send, @@ -39,7 +43,7 @@ where .enable_time() .build()?; - Ok(rt.block_on(module_main(dispatcher))) + Ok(rt.block_on(module_main(dispatcher, notifier))) }) .unwrap(); From 500b568f8f3f598244b8fad8493287abb169ab7b Mon Sep 17 00:00:00 2001 From: Oleg Utkin Date: Fri, 22 Oct 2021 13:36:24 +0000 Subject: [PATCH 12/17] wip --- src/txapi.rs | 19 ++++++++++++++----- 1 file changed, 14 insertions(+), 5 deletions(-) diff --git a/src/txapi.rs b/src/txapi.rs index a22a9f8..1d35de6 100644 --- a/src/txapi.rs +++ b/src/txapi.rs @@ -6,6 +6,7 @@ use crate::eventfd; use std::os::unix::io::{AsRawFd, RawFd}; use std::io; use mlua::Lua; +use std::sync::{Arc, atomic}; pub type Task = Box Result<(), ChannelError> + Send>; type TaskSender = async_channel::Sender; @@ -26,17 +27,19 @@ pub enum ChannelError { pub struct Dispatcher { pub(crate) task_tx: TaskSender, pub(crate) eventfd: eventfd::EventFd, + pub(crate) waiters: Arc, } impl Dispatcher { - pub fn new(task_tx: TaskSender, eventfd: eventfd::EventFd) -> Self { - Self { task_tx, eventfd } + pub fn new(task_tx: TaskSender, eventfd: eventfd::EventFd, waiters: Arc) -> Self { + Self { task_tx, eventfd, waiters } } pub fn try_clone(&self) -> std::io::Result { Ok(Self { task_tx: self.task_tx.clone(), eventfd: self.eventfd.try_clone()?, + waiters: self.waiters.clone(), }) } @@ -79,11 +82,12 @@ impl Dispatcher { pub struct Executor { task_rx: TaskReceiver, eventfd: eventfd::EventFd, + waiters: Arc, } impl Executor { - pub fn new(task_rx: TaskReceiver, eventfd: eventfd::EventFd) -> Self { - Self { task_rx, eventfd } + pub fn new(task_rx: TaskReceiver, eventfd: eventfd::EventFd, waiters: Arc) -> Self { + Self { task_rx, eventfd, waiters } } pub fn exec(&self, lua: &Lua, coio_timeout: f64) -> Result<(), ChannelError> { @@ -123,6 +127,7 @@ impl Executor { Ok(Self { task_rx: self.task_rx.clone(), eventfd: self.eventfd.try_clone()?, + waiters: self.waiters.clone(), }) } @@ -140,6 +145,10 @@ impl AsRawFd for Executor { pub fn channel(buffer: usize) -> io::Result<(Dispatcher, Executor)> { let (task_tx, task_rx) = async_channel::bounded(buffer); let efd = eventfd::EventFd::new(0, false)?; + let waiters = Arc::new(atomic::AtomicUsize::new(0)); - Ok((Dispatcher::new(task_tx, efd.try_clone()?), Executor::new(task_rx, efd))) + Ok(( + Dispatcher::new(task_tx, efd.try_clone()?, waiters.clone()), + Executor::new(task_rx, efd, waiters), + )) } From 2c6ea0b7ae2ee97c0eef87a5fdc63de568276e74 Mon Sep 17 00:00:00 2001 From: Oleg Utkin Date: Fri, 22 Oct 2021 14:10:46 +0000 Subject: [PATCH 13/17] fix waiters counter decrement --- src/txapi.rs | 1 + 1 file changed, 1 insertion(+) diff --git a/src/txapi.rs b/src/txapi.rs index 1d35de6..344bdce 100644 --- a/src/txapi.rs +++ b/src/txapi.rs @@ -99,6 +99,7 @@ impl Executor { }; let _ = self.eventfd.coio_read(coio_timeout); + self.waiters.fetch_sub(1, atomic::Ordering::Relaxed); } } From 4dca3ec7989a1d905b717299b30be1017bf93a8f Mon Sep 17 00:00:00 2001 From: Oleg Utkin Date: Fri, 22 Oct 2021 14:19:56 +0000 Subject: [PATCH 14/17] add waiters method for dispatcher and executor --- src/txapi.rs | 8 ++++++++ 1 file changed, 8 insertions(+) diff --git a/src/txapi.rs b/src/txapi.rs index 344bdce..31d81e1 100644 --- a/src/txapi.rs +++ b/src/txapi.rs @@ -76,6 +76,10 @@ impl Dispatcher { pub fn len(&self) -> usize { self.task_tx.len() } + + pub fn waiters(&self) -> usize { + self.waiters.load(atomic::Ordering::Relaxed) + } } @@ -135,6 +139,10 @@ impl Executor { pub fn len(&self) -> usize { self.task_rx.len() } + + pub fn waiters(&self) -> usize { + self.waiters.load(atomic::Ordering::Relaxed) + } } impl AsRawFd for Executor { From 72650ddee12d28bd6343096ce938e362c63ede2c Mon Sep 17 00:00:00 2001 From: Alexey Andreev Date: Fri, 17 Jun 2022 13:20:55 +0300 Subject: [PATCH 15/17] Revert "Merge remote-tracking branch 'origin/notification_optimization' into vendored/stoppable_module" This reverts commit 4e404dea6479c086719811b08cfa24cbe1ed763f, reversing changes made to fea6c0614f649da121bc0bd511103fcd73e9d4fb. --- src/txapi.rs | 28 +++++----------------------- 1 file changed, 5 insertions(+), 23 deletions(-) diff --git a/src/txapi.rs b/src/txapi.rs index 31d81e1..7a3d822 100644 --- a/src/txapi.rs +++ b/src/txapi.rs @@ -6,7 +6,6 @@ use crate::eventfd; use std::os::unix::io::{AsRawFd, RawFd}; use std::io; use mlua::Lua; -use std::sync::{Arc, atomic}; pub type Task = Box Result<(), ChannelError> + Send>; type TaskSender = async_channel::Sender; @@ -27,19 +26,17 @@ pub enum ChannelError { pub struct Dispatcher { pub(crate) task_tx: TaskSender, pub(crate) eventfd: eventfd::EventFd, - pub(crate) waiters: Arc, } impl Dispatcher { - pub fn new(task_tx: TaskSender, eventfd: eventfd::EventFd, waiters: Arc) -> Self { - Self { task_tx, eventfd, waiters } + pub fn new(task_tx: TaskSender, eventfd: eventfd::EventFd) -> Self { + Self { task_tx, eventfd } } pub fn try_clone(&self) -> std::io::Result { Ok(Self { task_tx: self.task_tx.clone(), eventfd: self.eventfd.try_clone()?, - waiters: self.waiters.clone(), }) } @@ -86,12 +83,11 @@ impl Dispatcher { pub struct Executor { task_rx: TaskReceiver, eventfd: eventfd::EventFd, - waiters: Arc, } impl Executor { - pub fn new(task_rx: TaskReceiver, eventfd: eventfd::EventFd, waiters: Arc) -> Self { - Self { task_rx, eventfd, waiters } + pub fn new(task_rx: TaskReceiver, eventfd: eventfd::EventFd) -> Self { + Self { task_rx, eventfd } } pub fn exec(&self, lua: &Lua, coio_timeout: f64) -> Result<(), ChannelError> { @@ -103,7 +99,6 @@ impl Executor { }; let _ = self.eventfd.coio_read(coio_timeout); - self.waiters.fetch_sub(1, atomic::Ordering::Relaxed); } } @@ -132,17 +127,8 @@ impl Executor { Ok(Self { task_rx: self.task_rx.clone(), eventfd: self.eventfd.try_clone()?, - waiters: self.waiters.clone(), }) } - - pub fn len(&self) -> usize { - self.task_rx.len() - } - - pub fn waiters(&self) -> usize { - self.waiters.load(atomic::Ordering::Relaxed) - } } impl AsRawFd for Executor { @@ -154,10 +140,6 @@ impl AsRawFd for Executor { pub fn channel(buffer: usize) -> io::Result<(Dispatcher, Executor)> { let (task_tx, task_rx) = async_channel::bounded(buffer); let efd = eventfd::EventFd::new(0, false)?; - let waiters = Arc::new(atomic::AtomicUsize::new(0)); - Ok(( - Dispatcher::new(task_tx, efd.try_clone()?, waiters.clone()), - Executor::new(task_rx, efd, waiters), - )) + Ok((Dispatcher::new(task_tx, efd.try_clone()?), Executor::new(task_rx, efd))) } From 007a04429dbb9a1f7c4e8aef50d803ed807e222b Mon Sep 17 00:00:00 2001 From: Alexey Andreev Date: Fri, 17 Jun 2022 13:43:03 +0300 Subject: [PATCH 16/17] Revert "add waiters method for dispatcher and executor" This reverts commit 3da4b373f44187ba43e4e8b0f66933a3a1ba137a. --- src/txapi.rs | 4 ---- 1 file changed, 4 deletions(-) diff --git a/src/txapi.rs b/src/txapi.rs index 7a3d822..dcd5361 100644 --- a/src/txapi.rs +++ b/src/txapi.rs @@ -73,10 +73,6 @@ impl Dispatcher { pub fn len(&self) -> usize { self.task_tx.len() } - - pub fn waiters(&self) -> usize { - self.waiters.load(atomic::Ordering::Relaxed) - } } From f06b8a391c7f46aab1b57bcfd9765cfc0a420446 Mon Sep 17 00:00:00 2001 From: Alexey Andreev Date: Tue, 5 Jul 2022 13:31:48 +0300 Subject: [PATCH 17/17] tracing --- Cargo.toml | 6 ++++ src/fiber_pool.rs | 16 +++++++--- src/txapi.rs | 81 +++++++++++++++++++++++++++++++++-------------- 3 files changed, 75 insertions(+), 28 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index 12a4926..1aa7a47 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -14,6 +14,12 @@ libc = "0.2" crossbeam-utils = "0.8" serde = "1.0" crossbeam-channel = "0.5" +opentelemetry = { version = "0.17", features = ["rt-tokio"] } +opentelemetry-zipkin = { version = "0.15", features = ["reqwest-client"], default-features = false } +opentelemetry-jaeger = { version = "0.16.0", features = ["rt-tokio", "reqwest_collector_client"] } +tracing-opentelemetry = "0.17.3" +tracing = "0.1.35" +tracing-subscriber = "0.3.11" [dependencies.mlua] version = "0.6" diff --git a/src/fiber_pool.rs b/src/fiber_pool.rs index 30387f4..694d7d8 100644 --- a/src/fiber_pool.rs +++ b/src/fiber_pool.rs @@ -3,8 +3,10 @@ use std::{collections::LinkedList, rc::Rc, time::Duration}; use crossbeam_channel::{unbounded, TryRecvError}; use mlua::Lua; use tarantool::fiber; +use tracing_opentelemetry::OpenTelemetrySpanExt; +// use tracing; -use crate::{ChannelError, Executor, ModuleConfig, Task}; +use crate::{ChannelError, Executor, ModuleConfig, Task, InstrumentedTask}; struct SchedulerArgs<'a> { lua: &'a Lua, @@ -26,7 +28,7 @@ fn scheduler_f(args: Box) -> i32 { } = *args; let cond = Rc::new(fiber::Cond::new()); - let (tx, rx) = unbounded::(); + let (tx, rx) = unbounded::(); let mut workers = LinkedList::new(); for _ in 0..fibers { @@ -71,7 +73,7 @@ fn scheduler_f(args: Box) -> i32 { struct WorkerArgs<'a> { cond: Rc, lua: &'a Lua, - rx: crossbeam_channel::Receiver, + rx: crossbeam_channel::Receiver, fiber_standby_timeout: f64, } fn worker_f(args: Box) -> i32 { @@ -87,7 +89,13 @@ fn worker_f(args: Box) -> i32 { .create_function(move |lua, _: ()| { loop { match rx.try_recv() { - Ok(task) => match task(lua) { + Ok((func, span_ctx)) => match { + let span = tracing::span!(tracing::Level::TRACE, "fiber pool: exec"); + span.set_parent(span_ctx); + let _ = span.enter(); + + func(lua, span.context()) + } { Ok(()) => (), Err(ChannelError::TXChannelClosed) => continue, Err(err) => break Err(mlua::Error::external(err)), diff --git a/src/txapi.rs b/src/txapi.rs index dcd5361..083515b 100644 --- a/src/txapi.rs +++ b/src/txapi.rs @@ -1,15 +1,19 @@ -use tokio::sync::oneshot; +use crate::eventfd; use async_channel; use async_channel::TryRecvError; -use thiserror::Error; -use crate::eventfd; -use std::os::unix::io::{AsRawFd, RawFd}; -use std::io; use mlua::Lua; +use tracing_opentelemetry::OpenTelemetrySpanExt; +use std::io; +use std::os::unix::io::{AsRawFd, RawFd}; +use thiserror::Error; +use tokio::sync::oneshot; +use tracing::Instrument; +use opentelemetry::Context; -pub type Task = Box Result<(), ChannelError> + Send>; -type TaskSender = async_channel::Sender; -type TaskReceiver = async_channel::Receiver; +pub type Task = Box Result<(), ChannelError> + Send>; +pub type InstrumentedTask = (Task, Context); +type TaskSender = async_channel::Sender; +type TaskReceiver = async_channel::Receiver; #[derive(Error, Debug)] pub enum ChannelError { @@ -40,24 +44,33 @@ impl Dispatcher { }) } + #[tracing::instrument(level = "trace", skip_all)] pub async fn call(&self, func: Func) -> Result - where - Ret: Send + 'static, - Func: FnOnce(&Lua) -> Ret, - Func: Send + 'static, + where + Ret: Send + 'static, + Func: FnOnce(&Lua, Context) -> Ret, + Func: Send + 'static, { + tracing::event!(tracing::Level::TRACE, "bass drop begin"); + let (result_tx, result_rx) = oneshot::channel(); - let handler_func: Task = Box::new(move |lua| { + let result_rx = result_rx + .instrument(tracing::span!(tracing::Level::TRACE, "result_rx")); + let result_rx_span_ctx = result_rx.span().context(); + + let handler_func: Task = Box::new(move |lua, exec_ctx| { if result_tx.is_closed() { - return Err(ChannelError::TXChannelClosed) + return Err(ChannelError::TXChannelClosed); }; - let result = func(lua); - result_tx.send(result).or(Err(ChannelError::TXChannelClosed)) + let result = func(lua, exec_ctx); + result_tx + .send(result) + .or(Err(ChannelError::TXChannelClosed)) }); - + let task_tx_len = self.task_tx.len(); - if let Err(_channel_closed) = self.task_tx.send(handler_func).await { + if let Err(_channel_closed) = self.task_tx.send((handler_func, result_rx_span_ctx)).await { return Err(ChannelError::TXChannelClosed); } @@ -67,7 +80,11 @@ impl Dispatcher { } } - result_rx.await.or(Err(ChannelError::RXChannelClosed)) + tracing::event!(tracing::Level::TRACE, "bass drop end"); + + result_rx + .await + .or(Err(ChannelError::RXChannelClosed)) } pub fn len(&self) -> usize { @@ -75,7 +92,6 @@ impl Dispatcher { } } - pub struct Executor { task_rx: TaskReceiver, eventfd: eventfd::EventFd, @@ -86,10 +102,23 @@ impl Executor { Self { task_rx, eventfd } } + // #[tracing::instrument(level = "trace", skip_all)] pub fn exec(&self, lua: &Lua, coio_timeout: f64) -> Result<(), ChannelError> { loop { - match self.task_rx.try_recv() { - Ok(func) => return func(lua), + // tracing::event!(tracing::Level::TRACE, "exec: iteration"); + + // let _ = tracing::trace_span!("executing task").enter(); + match self.task_rx.try_recv() + { + Ok((func, span_ctx)) => { + // tracing::event!(tracing::Level::TRACE, "task: start"); + println!("{:?}", span_ctx); + + let res = func(lua, span_ctx); + + // tracing::event!(tracing::Level::TRACE, "task: finish"); + return res; + } Err(TryRecvError::Empty) => (), Err(TryRecvError::Closed) => return Err(ChannelError::RXChannelClosed), }; @@ -98,7 +127,8 @@ impl Executor { } } - pub fn pop_many(&self, max_tasks: usize, coio_timeout: f64) -> Result, ChannelError> { + // #[tracing::instrument(level = "trace", skip_all)] + pub fn pop_many(&self, max_tasks: usize, coio_timeout: f64) -> Result, ChannelError> { if self.task_rx.is_empty() { let _ = self.eventfd.coio_read(coio_timeout); } @@ -137,5 +167,8 @@ pub fn channel(buffer: usize) -> io::Result<(Dispatcher, Executor)> { let (task_tx, task_rx) = async_channel::bounded(buffer); let efd = eventfd::EventFd::new(0, false)?; - Ok((Dispatcher::new(task_tx, efd.try_clone()?), Executor::new(task_rx, efd))) + Ok(( + Dispatcher::new(task_tx, efd.try_clone()?), + Executor::new(task_rx, efd), + )) }