diff --git a/Cargo.lock b/Cargo.lock index 6d51290..ba049c8 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -733,9 +733,9 @@ checksum = "0717cef1bc8b636c6e1c1bbdefc09e6322da8a9321966e8928ef80d20f7f770f" [[package]] name = "llvm-sys" -version = "201.0.1" +version = "221.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9bb947e8b79254ca10d496d0798a9ba1287dcf68e50a92b016fec1cc45bef447" +checksum = "2e52e36cd9eef0d5c4ba35c751c252f207642293009bb774185f84f678c7683c" dependencies = [ "anyhow", "cc", diff --git a/Cargo.toml b/Cargo.toml index de3484b..d8705f6 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -28,7 +28,7 @@ io-uring = "0.7.10" enum_dispatch = "0.3.13" pest = "2.8.1" pest_derive = "2.8.1" -llvm-sys = "201.0.1" +llvm-sys = "221.1.0" docopt = "1.1.1" signal-hook = "0.3.18" diff --git a/Containerfile b/Containerfile index 66c2afe..79ca86c 100644 --- a/Containerfile +++ b/Containerfile @@ -1,4 +1,4 @@ -FROM registry.fedoraproject.org/fedora:43 AS builder +FROM registry.fedoraproject.org/fedora:44 AS builder ARG RUST_VERSION=stable diff --git a/src/main.rs b/src/main.rs index f177741..9b1f241 100644 --- a/src/main.rs +++ b/src/main.rs @@ -85,15 +85,15 @@ fn run_script(script_path: String) -> Vec<(i32, u64)> { unreachable!() }; - let workers: u32 = + let workers: usize = args.get("workers").cloned().unwrap().parse().unwrap(); let duration: u64 = args.get("duration").cloned().unwrap().parse().unwrap(); (0..workers) - .filter_map(|_| { - let worker = new_script_worker(node.clone()); + .filter_map(|i| { + let worker = new_script_worker(node.clone(), i); match fork() { Ok(Fork::Parent(child)) => { @@ -290,7 +290,7 @@ mod tests { let ast: Vec = parse_instructions(input).unwrap(); assert_eq!(ast.len(), 1); - new_script_worker(ast[0].clone()).run_payload().unwrap(); + new_script_worker(ast[0].clone(), 0).run_payload().unwrap(); } #[test] @@ -305,7 +305,7 @@ mod tests { assert_eq!(nodes.len(), 1); let prepared_nodes = apply_rules(nodes); - new_script_worker(prepared_nodes[0].clone()) + new_script_worker(prepared_nodes[0].clone(), 0) .run_payload() .unwrap(); } @@ -322,7 +322,7 @@ mod tests { assert_eq!(nodes.len(), 1); let prepared_nodes = apply_rules(nodes); - new_script_worker(prepared_nodes[0].clone()) + new_script_worker(prepared_nodes[0].clone(), 0) .run_payload() .unwrap(); } @@ -339,7 +339,7 @@ mod tests { assert_eq!(nodes.len(), 1); let prepared_nodes = apply_rules(nodes); - new_script_worker(prepared_nodes[0].clone()) + new_script_worker(prepared_nodes[0].clone(), 0) .run_payload() .unwrap(); } @@ -367,8 +367,8 @@ mod tests { let _ = apply(vec![&ast[0]]); // run workers - new_script_worker(ast[1].clone()).run_payload().unwrap(); - new_script_worker(ast[2].clone()).run_payload().unwrap(); + new_script_worker(ast[1].clone(), 0).run_payload().unwrap(); + new_script_worker(ast[2].clone(), 0).run_payload().unwrap(); } #[test] @@ -412,4 +412,40 @@ mod tests { assert_eq!(args.get("workers").cloned().unwrap(), "2".to_string()); assert_eq!(args.get("duration").cloned().unwrap(), "10".to_string()); } + + #[test] + fn test_listen_with_zipf() { + let input = r#" + main (workers = 2, duration = 10) { + listen(8081, zipf(10, 1.4)); + } + "#; + + let nodes: Vec = parse_instructions(input).unwrap(); + assert_eq!(nodes.len(), 1); + + let prepared_nodes = apply_rules(nodes); + + new_script_worker(prepared_nodes[0].clone(), 0) + .run_payload() + .unwrap(); + } + + #[test] + fn test_sleep() { + let input = r#" + main (workers = 2, duration = 10) { + sleep(0.01); + } + "#; + + let nodes: Vec = parse_instructions(input).unwrap(); + assert_eq!(nodes.len(), 1); + + let prepared_nodes = apply_rules(nodes); + + new_script_worker(prepared_nodes[0].clone(), 0) + .run_payload() + .unwrap(); + } } diff --git a/src/script/ast.rs b/src/script/ast.rs index 92e7fe3..85d8956 100644 --- a/src/script/ast.rs +++ b/src/script/ast.rs @@ -1,12 +1,19 @@ use std::collections::HashMap; +#[derive(Debug, Clone, PartialEq)] +pub enum ConstType { + Text(String), + Int(u64), + Double(f64), +} + #[derive(Debug, Clone, PartialEq)] pub enum Arg { /// Null constant Null, /// Simple constant - Const { text: String }, + Const { value: ConstType }, /// Variable available at runtime Var { name: String }, @@ -28,6 +35,12 @@ pub enum Instruction { /// Send a message to a server at specified address Ping { server: Arg }, + + /// Listen on a specified number of endpoints from the lower boundary + Listen { lower: Arg, n: Arg }, + + /// Sleep for specified amount of time + Sleep { interval: Arg }, } #[derive(Debug, Clone, PartialEq)] diff --git a/src/script/grammar.peg b/src/script/grammar.peg index 51e02fa..04b2442 100644 --- a/src/script/grammar.peg +++ b/src/script/grammar.peg @@ -4,10 +4,11 @@ COMMENT = _{"//" ~ (!NEWLINE ~ ANY)*} ident_char = {ASCII_ALPHA | "_" | "$"} ident = @{ident_char ~ (ASCII_DIGIT | ident_char)*} -constant = { - "\"" ~ value ~ "\"" - | ASCII_DIGIT+ -} +constant = { text | double | int } + +text = { "\"" ~ value ~ "\"" } +int = { ASCII_DIGIT+ } +double = { ASCII_DIGIT+ ~ "." ~ ASCII_DIGIT+ } randomPath = { "random_path" } randomString = { "random_string" } @@ -15,6 +16,7 @@ randomString = { "random_string" } dynamicName = { randomPath | randomString + | zipf } dynamic = {dynamicName ~ args} @@ -38,6 +40,8 @@ port = { "port" } open = { "open" } ping = { "ping" } debug = { "debug" } +listen = { "listen" } +sleep = { "sleep" } funcName = { task @@ -46,6 +50,8 @@ funcName = { | open | ping | debug + | listen + | sleep } exp = { "exp" } diff --git a/src/script/parser.rs b/src/script/parser.rs index 6fe5873..3f2ebae 100644 --- a/src/script/parser.rs +++ b/src/script/parser.rs @@ -2,7 +2,9 @@ use log::trace; use pest::{self, Parser, error::Error}; use std::collections::HashMap; -use crate::script::ast::{Arg, Dist, Instruction, MachineInstruction, Node}; +use crate::script::ast::{ + Arg, ConstType, Dist, Instruction, MachineInstruction, Node, +}; #[derive(Debug)] pub enum ParseError { @@ -136,6 +138,52 @@ fn build_ast_from_function( } } +fn pair_to_arg(pair: pest::iterators::Pair) -> Arg { + match pair.as_rule() { + Rule::constant => { + let value = first_nested_pair(pair); + match value.as_rule() { + Rule::text => Arg::Const { + value: ConstType::Text(pair_to_string(first_nested_pair( + value, + ))), + }, + Rule::int => Arg::Const { + value: ConstType::Int(pair_to_int(value)), + }, + Rule::double => Arg::Const { + value: ConstType::Double(pair_to_double(value)), + }, + unknown => { + panic!("Unknown constant type {unknown:?}") + } + } + } + Rule::ident => Arg::Var { + name: pair_to_string(pair), + }, + Rule::dynamic => { + let mut inner = pair.into_inner(); + let name = inner.next().expect("No argument name"); + let args_pair = inner.next().expect("No argument value"); + + let args: Vec = args_pair + .into_inner() + .map(|arg| { + let value = first_nested_pair(arg); + pair_to_arg(value) + }) + .collect(); + + Arg::Dynamic { + name: pair_to_string(name), + args, + } + } + unknown => panic!("Unknown arg type {unknown:?}"), + } +} + fn build_ast_from_instr( pairs: pest::iterators::Pairs, ) -> Vec { @@ -149,40 +197,7 @@ fn build_ast_from_instr( let args: Vec = args_pair .into_inner() - .map(|arg| { - let a = first_nested_pair(arg); - match a.as_rule() { - Rule::constant => Arg::Const { - text: pair_to_string(first_nested_pair(a)), - }, - Rule::ident => Arg::Var { - name: pair_to_string(a), - }, - Rule::dynamic => { - let mut inner = a.into_inner(); - let name = inner.next().expect("No argument name"); - let args_pair = - inner.next().expect("No argument value"); - - let args: Vec = args_pair - .into_inner() - .map(|arg| { - let a = - first_nested_pair(first_nested_pair(arg)); - Arg::Const { - text: pair_to_string(a), - } - }) - .collect(); - - Arg::Dynamic { - name: pair_to_string(name), - args, - } - } - unknown => panic!("Unknown arg type {unknown:?}"), - } - }) + .map(|arg| pair_to_arg(first_nested_pair(arg))) .collect(); match first_nested_pair(name).as_rule() { @@ -211,6 +226,17 @@ fn build_ast_from_instr( server: args[0].clone(), }); } + Rule::listen => { + instr.push(Instruction::Listen { + lower: args[0].clone(), + n: args[1].clone(), + }); + } + Rule::sleep => { + instr.push(Instruction::Sleep { + interval: args[0].clone(), + }); + } unknown => panic!("Unknown instruction type {unknown:?}"), } } @@ -253,9 +279,9 @@ fn build_ast_from_dist(pair: pest::iterators::Pair) -> Dist { fn string_from_pair(pair: pest::iterators::Pair) -> String { assert!(matches!(pair.as_rule(), Rule::constant | Rule::ident)); - // Extract "value" (Constants) or "name" (Identifier) + // Extract "value" (text Constants) or "name" (Identifier) // and convert it to String - pair_to_string(first_nested_pair(pair)) + pair_to_string(first_nested_pair(first_nested_pair(pair))) } fn string_from_argument( @@ -276,10 +302,23 @@ fn pair_to_string(pair: pest::iterators::Pair) -> String { pair.as_span().as_str().to_string() } +fn pair_to_int(pair: pest::iterators::Pair) -> u64 { + pair.as_span().as_str().to_string().parse().unwrap() +} + +fn pair_to_double(pair: pest::iterators::Pair) -> f64 { + pair.as_span().as_str().to_string().parse().unwrap() +} + fn first_nested_pair( pair: pest::iterators::Pair, ) -> pest::iterators::Pair { - pair.into_inner().next().expect("Cannot get first pair") + let mut inner = pair.clone().into_inner(); + if inner.is_empty() { + pair + } else { + inner.next().expect("Cannot get first pair") + } } #[cfg(test)] @@ -310,7 +349,7 @@ mod tests { instructions[0], Instruction::Open { path: Arg::Const { - text: "/tmp/test".to_string() + value: ConstType::Text("/tmp/test".to_string()) } } ); @@ -346,7 +385,7 @@ mod tests { path: Arg::Dynamic { name: "random_path".to_string(), args: vec![Arg::Const { - text: "/tmp".to_string() + value: ConstType::Text("/tmp".to_string()) }], } } @@ -383,7 +422,7 @@ mod tests { instructions[0], Instruction::Debug { text: Arg::Const { - text: "run task stub".to_string(), + value: ConstType::Text("run task stub".to_string()), } } ); @@ -431,7 +470,7 @@ mod tests { instructions[0], Instruction::Debug { text: Arg::Const { - text: "ping server".to_string(), + value: ConstType::Text("ping server".to_string()), } } ); @@ -440,7 +479,7 @@ mod tests { instructions[1], Instruction::Ping { server: Arg::Const { - text: "127.0.0.1:8080".to_string(), + value: ConstType::Text("127.0.0.1:8080".to_string()), }, } ); diff --git a/src/worker/mod.rs b/src/worker/mod.rs index 0cc74d3..e5bf229 100644 --- a/src/worker/mod.rs +++ b/src/worker/mod.rs @@ -74,6 +74,6 @@ pub fn new_worker( } } -pub fn new_script_worker(node: Node) -> Box { - Box::new(ScriptWorker::new(node)) +pub fn new_script_worker(node: Node, worker: usize) -> Box { + Box::new(ScriptWorker::new(node, worker)) } diff --git a/src/worker/script.rs b/src/worker/script.rs index 9af1d35..80e588d 100644 --- a/src/worker/script.rs +++ b/src/worker/script.rs @@ -7,9 +7,13 @@ use std::{ fs::OpenOptions, io::Write, io::prelude::*, - net::{Shutdown, TcpStream}, + net::{Shutdown, TcpListener, TcpStream}, + os::fd::{AsRawFd, RawFd}, process::Command, - sync::Arc, + sync::{ + Arc, + atomic::{AtomicUsize, Ordering}, + }, thread, time, }; @@ -17,7 +21,7 @@ use std::sync::LazyLock; use log::{Level, debug, log_enabled, trace}; use rand::{Rng, distributions::Alphanumeric, thread_rng}; -use rand_distr::Exp; +use rand_distr::{Exp, Zipf}; use llvm::core::*; use llvm::execution_engine::*; @@ -29,12 +33,13 @@ use std::mem; use crate::{Worker, WorkerError}; -use crate::script::ast::{Arg, Dist, Instruction, Node}; +use crate::script::ast::{Arg, ConstType, Dist, Instruction, Node}; #[derive(Debug, Clone)] enum RuntimeType { Int, Pointer, + Double, } #[derive(Debug, Clone)] @@ -74,7 +79,6 @@ pub unsafe extern "C" fn debug(text: *const i8) -> u64 { /// terminated C-string. #[unsafe(no_mangle)] pub unsafe extern "C" fn open_file(path: *const i8) -> u64 { - //let path = unsafe { CString::from_raw(path as *mut i8) }; let path = unsafe { CStr::from_ptr(path) }; debug!("Open path {:?}", path); let mut file = OpenOptions::new() @@ -138,8 +142,60 @@ pub unsafe extern "C" fn task(name: *const i8, args: *const i8) -> u64 { .unwrap() } +/// Listen on a specified number of ports starting from the lower boundary. +/// Open connections will live until the end of the block of work and will be +/// shutdown in cleanup instruction. Note that there will be no concurrency +/// issues, since threads for listening are registered in sychronous manner. +/// +/// TODO: Decrement max_ports on batch exit. +/// +/// # Safety +#[unsafe(no_mangle)] +pub unsafe extern "C" fn listen_on_ports(lower: u64, n: u64) -> u64 { + debug!("Listen {lower} {n}"); + let max_ports = MAX_PORTS.fetch_add(n as usize, Ordering::Relaxed); + + let start_port = lower + max_ports as u64; + for port in start_port..start_port + n { + let addr = format!("0.0.0.0:{port}"); + let listener = TcpListener::bind(&addr) + .expect("Couldn't listen on the specified address"); + let fd = listener.as_raw_fd(); + + trace!("Listen {addr}, fd {fd}"); + SOCKETS.with(|socks| socks.borrow_mut().push(fd)); + + thread::spawn(move || for _stream in listener.incoming() {}); + } + + 0 +} + thread_local! { static POINTERS: RefCell> = const { RefCell::new(vec![]) }; + static SOCKETS: RefCell> = const { RefCell::new(vec![]) }; +} + +static MAX_PORTS: AtomicUsize = AtomicUsize::new(0); + +/// Return a random integer from zipf distribution with specified +/// size and exponent. +/// +/// # Safety +#[unsafe(no_mangle)] +pub unsafe extern "C" fn zipf(size: u64, exp: f64) -> u64 { + debug!("zipf {size} {exp}"); + thread_rng().sample(Zipf::new(size, exp).unwrap()) as u64 +} + +/// Sleeps for specified amount of time. +/// +/// # Safety +#[unsafe(no_mangle)] +pub unsafe extern "C" fn sleep(interval: f64) -> u64 { + debug!("Sleep {interval}"); + thread::sleep(time::Duration::from_secs_f64(interval)); + 0 } /// Return a randomly generated string. @@ -184,15 +240,29 @@ pub unsafe extern "C" fn random_path(base: *const i8) -> *const i8 { /// # Safety #[unsafe(no_mangle)] pub unsafe extern "C" fn cleanup(_: *const i8) -> u64 { + debug!("Cleanup"); POINTERS.with(|ps| { let mut vec = ps.borrow_mut(); for p in vec.as_slice() { + trace!("Cleanup {:?}", p); let _ = unsafe { CString::from_raw(*p) }; } vec.clear(); }); + SOCKETS.with(|socks| { + let mut vec = socks.borrow_mut(); + for fd in vec.as_slice() { + trace!("Shutdown {fd}"); + unsafe { + libc::shutdown(*fd, libc::SHUT_RD); + } + } + + vec.clear(); + }); + 0 } @@ -245,6 +315,24 @@ pub static RUNTIME: LazyLock> = return_type: RuntimeType::Int, }, ), + ( + "listen".to_string(), + RuntimeFunc { + func: listen_on_ports as *const () as usize, + param_count: 2, + param_types: &[RuntimeType::Int, RuntimeType::Int], + return_type: RuntimeType::Int, + }, + ), + ( + "sleep".to_string(), + RuntimeFunc { + func: sleep as *const () as usize, + param_count: 1, + param_types: &[RuntimeType::Double], + return_type: RuntimeType::Int, + }, + ), // dynamic values ( "random_path".to_string(), @@ -264,6 +352,15 @@ pub static RUNTIME: LazyLock> = return_type: RuntimeType::Pointer, }, ), + ( + "zipf".to_string(), + RuntimeFunc { + func: zipf as *const () as usize, + param_count: 2, + param_types: &[RuntimeType::Int, RuntimeType::Double], + return_type: RuntimeType::Int, + }, + ), // utils ( "cleanup".to_string(), @@ -308,19 +405,32 @@ impl ScriptWorker { let iptr = LLVMIntPtrTypeInContext(ctx.context, td); LLVMConstNull(iptr) }, - Arg::Const { text } => unsafe { - // The name of all constants created this way will be "const", - // which is ugly, but not a problem as LLVM modifies this to - // make sure uniqueness, i.e. they will be: - // - // @const, @const.1, @const.2, ... - // - // in the jited code. - LLVMBuildGlobalString( - ctx.builder, - format!("{text}\0").as_ptr() as *const _, - c"const".as_ptr() as *const _, - ) + Arg::Const { value } => unsafe { + match value { + ConstType::Text(text) => { + // The name of all constants created this way will be + // "const", which is ugly, but + // not a problem as LLVM modifies this to + // make sure uniqueness, i.e. they will be: + // + // @const, @const.1, @const.2, ... + // + // in the jited code. + LLVMBuildGlobalString( + ctx.builder, + format!("{text}\0").as_ptr() as *const _, + c"const".as_ptr() as *const _, + ) + } + ConstType::Int(value) => { + let i64t = LLVMInt64TypeInContext(ctx.context); + LLVMConstInt(i64t, value, 0) + } + ConstType::Double(value) => { + let double = LLVMDoubleTypeInContext(ctx.context); + LLVMConstReal(double, value) + } + } }, Arg::Var { name } => { *ctx.module_state.get(&name).expect("No variable") @@ -359,18 +469,23 @@ impl ScriptWorker { *func, args_ptr.as_mut_ptr(), args.len().try_into().unwrap(), - c"{name}".as_ptr() as *const _, + c"const".as_ptr() as *const _, ) } } } } - pub fn new(node: Node) -> Self { + pub fn new(node: Node, worker: usize) -> Self { let mut module_runtime: HashMap = HashMap::new(); let mut module_state: HashMap = HashMap::new(); + // FIXME: It's an ugly temporary hack to split port space with a + // hardcoded constant. Do this better via deriving in applu_rules how + // listen arguments look like and how the ranges should be. + MAX_PORTS.fetch_add(worker * 1000, Ordering::Relaxed); + unsafe { // Set up a context, module and builder in that context. let context = LLVMContextCreate(); @@ -414,6 +529,7 @@ impl ScriptWorker { // get a type for main function let i64t = LLVMInt64TypeInContext(context); let boolt = LLVMInt1TypeInContext(context); + let double = LLVMDoubleTypeInContext(context); let iptr = LLVMIntPtrTypeInContext(context, td); // Insert runtime functions into the module @@ -428,6 +544,7 @@ impl ScriptWorker { .map(|t| match t { RuntimeType::Pointer => iptr, RuntimeType::Int => i64t, + RuntimeType::Double => double, }) .collect::>(); @@ -435,6 +552,7 @@ impl ScriptWorker { match f.return_type { RuntimeType::Int => i64t, RuntimeType::Pointer => iptr, + RuntimeType::Double => double, }, function_args.as_mut_ptr(), f.param_count, @@ -533,6 +651,16 @@ impl ScriptWorker { Self::jit_instruction(c"debug", vec![text], &ctx); "debug" } + + Instruction::Listen { lower, n } => { + Self::jit_instruction(c"listen", vec![lower, n], &ctx); + "listen" + } + + Instruction::Sleep { interval } => { + Self::jit_instruction(c"sleep", vec![interval], &ctx); + "sleep" + } }; // Populate the global mapping with observed runtime functions diff --git a/workloads/example.ber b/workloads/example.ber index 7b5a690..af89835 100644 --- a/workloads/example.ber +++ b/workloads/example.ber @@ -1,6 +1,5 @@ machine { server(8080); - profile("bpf"); } // Named work block @@ -11,8 +10,11 @@ main (workers = 2, duration = 10) { // open(path) -- open file by path, create if needed and write something to it debug("run task stub"); task(stub, random_string()); + debug(random_string()); debug("open file /tmp/test"); open("/tmp/test"); + debug("listen on random ports starting from 8081"); + listen(8081, zipf(10, 1.4)); } : exp { // If no distribution provided, do the unit only once. rate = 10.0; diff --git a/workloads/example.short.ber b/workloads/example.short.ber index b4fcafa..0b58070 100644 --- a/workloads/example.short.ber +++ b/workloads/example.short.ber @@ -10,4 +10,6 @@ main (workers = 1) { open("/tmp/test"); debug("ping server"); ping("127.0.0.1:8080"); + debug("listen on random ports starting from 8081"); + listen(8081, zipf(200, 1.4)); }