add tcp server
This commit is contained in:
parent
a1f5941c1a
commit
f86c55dd83
12 changed files with 308 additions and 99 deletions
7
Cargo.lock
generated
7
Cargo.lock
generated
|
@ -8,12 +8,6 @@ version = "0.1.4"
|
||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
checksum = "428f6ba17d0c927e57c15a86cf5d7d07a2f35b3fbf15b1eb36b7075459e150a3"
|
checksum = "428f6ba17d0c927e57c15a86cf5d7d07a2f35b3fbf15b1eb36b7075459e150a3"
|
||||||
|
|
||||||
[[package]]
|
|
||||||
name = "once_cell"
|
|
||||||
version = "1.17.1"
|
|
||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
|
||||||
checksum = "b7e5500299e16ebb147ae15a00a942af264cf3688f47923b8fc2cd5858f23ad3"
|
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "readformat"
|
name = "readformat"
|
||||||
version = "0.1.2"
|
version = "0.1.2"
|
||||||
|
@ -25,6 +19,5 @@ name = "spl"
|
||||||
version = "0.3.2"
|
version = "0.3.2"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"multicall",
|
"multicall",
|
||||||
"once_cell",
|
|
||||||
"readformat",
|
"readformat",
|
||||||
]
|
]
|
||||||
|
|
|
@ -9,5 +9,4 @@ authors = ["TudbuT"]
|
||||||
|
|
||||||
[dependencies]
|
[dependencies]
|
||||||
readformat = "0.1"
|
readformat = "0.1"
|
||||||
once_cell = "1.17"
|
|
||||||
multicall = "0.1"
|
multicall = "0.1"
|
||||||
|
|
|
@ -36,7 +36,10 @@ func dec {
|
||||||
pop
|
pop
|
||||||
}
|
}
|
||||||
func getms {
|
func getms {
|
||||||
0 time
|
0 time:unixms
|
||||||
|
}
|
||||||
|
func sleep {
|
||||||
|
time:sleep
|
||||||
}
|
}
|
||||||
|
|
||||||
func main { }
|
func main { }
|
||||||
|
|
63
server.spl
Normal file
63
server.spl
Normal file
|
@ -0,0 +1,63 @@
|
||||||
|
"#net.spl" import
|
||||||
|
"#stream.spl" import
|
||||||
|
|
||||||
|
"server" net:register
|
||||||
|
|
||||||
|
construct net:server namespace {
|
||||||
|
ServerStream
|
||||||
|
Type
|
||||||
|
_Types
|
||||||
|
types
|
||||||
|
;
|
||||||
|
Types { _Types | with this ;
|
||||||
|
this:_Types:new
|
||||||
|
}
|
||||||
|
register-type { | with id this ;
|
||||||
|
this:types id awrap aadd this:=types
|
||||||
|
id this:_Types dyn-def-field;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
0 anew net:server:=types
|
||||||
|
|
||||||
|
construct net:server:Type {
|
||||||
|
id
|
||||||
|
;
|
||||||
|
construct { this | with id this ;
|
||||||
|
id this:=id
|
||||||
|
this
|
||||||
|
}
|
||||||
|
create { ServerStream | with this ;
|
||||||
|
this:id net:server:ServerStream:new
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
construct net:server:_Types {
|
||||||
|
;
|
||||||
|
construct { this | with this ;
|
||||||
|
{ | with type ;
|
||||||
|
"type net:server:Type:new this:=<type>";
|
||||||
|
(type net:server:Type:new) (this ("=" type concat)) dyn-objcall
|
||||||
|
} net:server:types:foreach;
|
||||||
|
this
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
"tcp" net:server:register-type
|
||||||
|
|
||||||
|
construct net:server:ServerStream {
|
||||||
|
id
|
||||||
|
;
|
||||||
|
construct { this | with type this ;
|
||||||
|
type new-server-stream this:=id
|
||||||
|
this
|
||||||
|
}
|
||||||
|
accept { Stream | with this ;
|
||||||
|
def stream null Stream settype =stream
|
||||||
|
this:id accept-server-stream stream:=id
|
||||||
|
stream
|
||||||
|
}
|
||||||
|
close { | with this ;
|
||||||
|
this:id close-server-stream;
|
||||||
|
}
|
||||||
|
}
|
|
@ -72,6 +72,7 @@ pub struct Runtime {
|
||||||
types_by_id: HashMap<u32, AMType>,
|
types_by_id: HashMap<u32, AMType>,
|
||||||
next_stream_id: u128,
|
next_stream_id: u128,
|
||||||
streams: HashMap<u128, Arc<Mut<Stream>>>,
|
streams: HashMap<u128, Arc<Mut<Stream>>>,
|
||||||
|
server_streams: HashMap<u128, Arc<Mut<ServerStream>>>,
|
||||||
pub embedded_files: HashMap<&'static str, &'static str>,
|
pub embedded_files: HashMap<&'static str, &'static str>,
|
||||||
pub native_functions: HashMap<&'static str, (u32, FuncImpl)>,
|
pub native_functions: HashMap<&'static str, (u32, FuncImpl)>,
|
||||||
}
|
}
|
||||||
|
@ -101,6 +102,7 @@ impl Runtime {
|
||||||
types_by_id: HashMap::new(),
|
types_by_id: HashMap::new(),
|
||||||
next_stream_id: 0,
|
next_stream_id: 0,
|
||||||
streams: HashMap::new(),
|
streams: HashMap::new(),
|
||||||
|
server_streams: HashMap::new(),
|
||||||
embedded_files: HashMap::new(),
|
embedded_files: HashMap::new(),
|
||||||
native_functions: HashMap::new(),
|
native_functions: HashMap::new(),
|
||||||
};
|
};
|
||||||
|
@ -160,6 +162,23 @@ impl Runtime {
|
||||||
self.streams.remove(&id);
|
self.streams.remove(&id);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub fn register_server_stream(
|
||||||
|
&mut self,
|
||||||
|
stream: ServerStream,
|
||||||
|
) -> (u128, Arc<Mut<ServerStream>>) {
|
||||||
|
let id = (self.next_stream_id, self.next_stream_id += 1).0;
|
||||||
|
self.server_streams.insert(id, Arc::new(Mut::new(stream)));
|
||||||
|
(id, self.server_streams.get(&id).unwrap().clone())
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn get_server_stream(&self, id: u128) -> Option<Arc<Mut<ServerStream>>> {
|
||||||
|
self.server_streams.get(&id).cloned()
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn destroy_server_stream(&mut self, id: u128) {
|
||||||
|
self.server_streams.remove(&id);
|
||||||
|
}
|
||||||
|
|
||||||
pub fn load_native_function(&self, name: &str) -> &(u32, FuncImpl) {
|
pub fn load_native_function(&self, name: &str) -> &(u32, FuncImpl) {
|
||||||
self.native_functions.get(name).unwrap_or_else(|| {
|
self.native_functions.get(name).unwrap_or_else(|| {
|
||||||
panic!(
|
panic!(
|
||||||
|
|
|
@ -935,7 +935,7 @@ pub fn register(r: &mut Stack, o: Arc<Frame>) {
|
||||||
("write-sasm", write_sasm, 1),
|
("write-sasm", write_sasm, 1),
|
||||||
("write-file-sasm", write_file_sasm, 1),
|
("write-file-sasm", write_file_sasm, 1),
|
||||||
("fork", fork, 0),
|
("fork", fork, 0),
|
||||||
("time", time, 0),
|
("sleeptime", time, 0),
|
||||||
];
|
];
|
||||||
for f in fns {
|
for f in fns {
|
||||||
r.define_func(
|
r.define_func(
|
||||||
|
|
48
src/stream/mod.rs
Normal file
48
src/stream/mod.rs
Normal file
|
@ -0,0 +1,48 @@
|
||||||
|
mod rw;
|
||||||
|
mod server;
|
||||||
|
pub use rw::*;
|
||||||
|
pub use server::*;
|
||||||
|
|
||||||
|
use std::sync::{Arc, LazyLock};
|
||||||
|
|
||||||
|
use crate::{mutex::Mut, runtime::*};
|
||||||
|
|
||||||
|
static IS_INITIALIZED: LazyLock<Arc<Mut<bool>>> = LazyLock::new(|| Arc::new(Mut::new(false)));
|
||||||
|
|
||||||
|
pub fn register(r: &mut Stack, o: Arc<Frame>) {
|
||||||
|
if !*IS_INITIALIZED.lock_ro() {
|
||||||
|
register_stream_type("file", stream_file);
|
||||||
|
register_stream_type("tcp", stream_tcp);
|
||||||
|
register_stream_type("udp", stream_udp);
|
||||||
|
register_stream_type("cmd", stream_cmd);
|
||||||
|
register_server_stream_type("tcp", server_stream_tcp);
|
||||||
|
*IS_INITIALIZED.lock() = true;
|
||||||
|
}
|
||||||
|
|
||||||
|
type Fn = fn(&mut Stack) -> OError;
|
||||||
|
let fns: [(&str, Fn, u32); 10] = [
|
||||||
|
("new-stream", new_stream, 1),
|
||||||
|
("write-stream", write_stream, 1),
|
||||||
|
("write-all-stream", write_all_stream, 0),
|
||||||
|
("flush-stream", flush_stream, 0),
|
||||||
|
("read-stream", read_stream, 1),
|
||||||
|
("read-all-stream", read_all_stream, 0),
|
||||||
|
("close-stream", close_stream, 0),
|
||||||
|
("new-server-stream", new_server_stream, 1),
|
||||||
|
("accept-server-stream", accept_server_stream, 1),
|
||||||
|
("close-server-stream", close_server_stream, 0),
|
||||||
|
];
|
||||||
|
for f in fns {
|
||||||
|
r.define_func(
|
||||||
|
f.0.to_owned(),
|
||||||
|
AFunc::new(Func {
|
||||||
|
ret_count: f.2,
|
||||||
|
to_call: FuncImpl::Native(f.1),
|
||||||
|
run_as_base: false,
|
||||||
|
origin: o.clone(),
|
||||||
|
fname: None,
|
||||||
|
name: f.0.to_owned(),
|
||||||
|
}),
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
|
@ -1,20 +1,20 @@
|
||||||
use std::{
|
use std::{
|
||||||
|
borrow::BorrowMut,
|
||||||
collections::HashMap,
|
collections::HashMap,
|
||||||
fs::OpenOptions,
|
hint::black_box,
|
||||||
io::{Read, Write},
|
io::{Read, Write},
|
||||||
mem,
|
net::{TcpStream, UdpSocket},
|
||||||
net::{Shutdown, TcpStream, UdpSocket},
|
|
||||||
process::{self, Stdio},
|
process::{self, Stdio},
|
||||||
sync::Arc,
|
sync::{Arc, LazyLock},
|
||||||
};
|
};
|
||||||
|
|
||||||
use once_cell::sync::Lazy;
|
use crate::*;
|
||||||
|
|
||||||
use crate::{mutex::Mut, runtime::*, *};
|
use fs::OpenOptions;
|
||||||
|
use mutex::Mut;
|
||||||
|
|
||||||
static STREAM_TYPES: Lazy<Arc<Mut<HashMap<String, StreamType>>>> =
|
static STREAM_TYPES: LazyLock<Arc<Mut<HashMap<String, StreamType>>>> =
|
||||||
Lazy::new(|| Arc::new(Mut::new(HashMap::new())));
|
LazyLock::new(|| Arc::new(Mut::new(HashMap::new())));
|
||||||
static IS_INITIALIZED: Lazy<Arc<Mut<bool>>> = Lazy::new(|| Arc::new(Mut::new(false)));
|
|
||||||
|
|
||||||
/// Registers a custom stream type.
|
/// Registers a custom stream type.
|
||||||
pub fn register_stream_type(
|
pub fn register_stream_type(
|
||||||
|
@ -45,32 +45,37 @@ impl StreamType {
|
||||||
|
|
||||||
/// An SPL stream, holding a reader and a writer, and a function to close it.
|
/// An SPL stream, holding a reader and a writer, and a function to close it.
|
||||||
pub struct Stream {
|
pub struct Stream {
|
||||||
reader: Box<dyn Read + Send + Sync + 'static>,
|
pub(super) reader: Box<dyn Read + Send + Sync + 'static>,
|
||||||
writer: Box<dyn Write + Send + Sync + 'static>,
|
pub(super) _writer_storage: Option<Box<dyn Write + Send + Sync + 'static>>,
|
||||||
close: fn(&mut Self),
|
pub(super) writer: &'static mut (dyn Write + Send + Sync + 'static),
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Stream {
|
impl Stream {
|
||||||
pub fn new<T: Read + Write + Send + Sync + 'static>(main: T, close: fn(&mut Self)) -> Self {
|
pub fn new<T: Read + Write + Send + Sync + 'static>(main: T) -> Self {
|
||||||
let mut rw = Box::new(main);
|
let mut rw = Box::new(main);
|
||||||
Self {
|
Self {
|
||||||
// SAFETY: Because these are both in private fields on one object, they can not be
|
writer: unsafe {
|
||||||
// written to simultaneously or read from while writing due to the guards put in place
|
(rw.as_mut() as *mut (dyn Write + Send + Sync + 'static))
|
||||||
// by the borrow checker on the Stream.
|
.as_mut()
|
||||||
reader: Box::new(unsafe { mem::transmute::<&mut _, &mut T>(rw.as_mut()) }),
|
.unwrap()
|
||||||
writer: rw,
|
},
|
||||||
close,
|
_writer_storage: None,
|
||||||
|
reader: rw,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
pub fn new_split(
|
pub fn new_split(
|
||||||
reader: impl Read + Send + Sync + 'static,
|
reader: impl Read + Send + Sync + 'static,
|
||||||
writer: impl Write + Send + Sync + 'static,
|
writer: impl Write + Send + Sync + 'static,
|
||||||
close: fn(&mut Self),
|
|
||||||
) -> Self {
|
) -> Self {
|
||||||
|
let mut bx = Box::new(writer);
|
||||||
Self {
|
Self {
|
||||||
reader: Box::new(reader),
|
reader: Box::new(reader),
|
||||||
writer: Box::new(writer),
|
writer: unsafe {
|
||||||
close,
|
(bx.as_mut() as *mut (dyn Write + Send + Sync + 'static))
|
||||||
|
.as_mut()
|
||||||
|
.unwrap()
|
||||||
|
},
|
||||||
|
_writer_storage: Some(bx),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
@ -124,14 +129,14 @@ where
|
||||||
T: Fn(&mut Stack) -> Result<Stream, Error> + Sync + Send + 'static,
|
T: Fn(&mut Stack) -> Result<Stream, Error> + Sync + Send + 'static,
|
||||||
{
|
{
|
||||||
fn from(value: T) -> Self {
|
fn from(value: T) -> Self {
|
||||||
StreamType {
|
Self {
|
||||||
func: Arc::new(Box::new(value)),
|
func: Arc::new(Box::new(value)),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn new_stream(stack: &mut Stack) -> OError {
|
pub fn new_stream(stack: &mut Stack) -> OError {
|
||||||
require_on_stack!(s, Str, stack, "write-stream");
|
require_on_stack!(s, Str, stack, "new-stream");
|
||||||
let stream = get_stream_type(s.clone())
|
let stream = get_stream_type(s.clone())
|
||||||
.ok_or_else(|| stack.error(ErrorKind::VariableNotFound(format!("__stream-type-{s}"))))?
|
.ok_or_else(|| stack.error(ErrorKind::VariableNotFound(format!("__stream-type-{s}"))))?
|
||||||
.make_stream(stack)?;
|
.make_stream(stack)?;
|
||||||
|
@ -163,6 +168,7 @@ pub fn write_stream(stack: &mut Stack) -> OError {
|
||||||
)
|
)
|
||||||
.spl(),
|
.spl(),
|
||||||
);
|
);
|
||||||
|
black_box(&stream.lock_ro()._writer_storage);
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@ -184,6 +190,7 @@ pub fn write_all_stream(stack: &mut Stack) -> OError {
|
||||||
.lock()
|
.lock()
|
||||||
.write_all(&fixed[..])
|
.write_all(&fixed[..])
|
||||||
.map_err(|x| stack.error(ErrorKind::IO(format!("{x:?}"))))?;
|
.map_err(|x| stack.error(ErrorKind::IO(format!("{x:?}"))))?;
|
||||||
|
black_box(&stream.lock_ro()._writer_storage);
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@ -197,6 +204,7 @@ pub fn flush_stream(stack: &mut Stack) -> OError {
|
||||||
.lock()
|
.lock()
|
||||||
.flush()
|
.flush()
|
||||||
.map_err(|x| stack.error(ErrorKind::IO(format!("{x:?}"))))?;
|
.map_err(|x| stack.error(ErrorKind::IO(format!("{x:?}"))))?;
|
||||||
|
black_box(&stream.lock_ro()._writer_storage);
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@ -256,16 +264,11 @@ pub fn read_all_stream(stack: &mut Stack) -> OError {
|
||||||
|
|
||||||
pub fn close_stream(stack: &mut Stack) -> OError {
|
pub fn close_stream(stack: &mut Stack) -> OError {
|
||||||
require_on_stack!(id, Mega, stack, "close-stream");
|
require_on_stack!(id, Mega, stack, "close-stream");
|
||||||
if let Some(stream) = runtime(|rt| rt.get_stream(id as u128)) {
|
runtime_mut(|mut rt| rt.destroy_stream(id as u128));
|
||||||
let mut stream = stream.lock();
|
|
||||||
(stream.close)(&mut stream);
|
|
||||||
}
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
fn nop(_stream: &mut Stream) {}
|
pub(super) fn stream_file(stack: &mut Stack) -> Result<Stream, Error> {
|
||||||
|
|
||||||
fn stream_file(stack: &mut Stack) -> Result<Stream, Error> {
|
|
||||||
let truncate = stack.pop().lock_ro().is_truthy();
|
let truncate = stack.pop().lock_ro().is_truthy();
|
||||||
require_on_stack!(path, Str, stack, "FILE new-stream");
|
require_on_stack!(path, Str, stack, "FILE new-stream");
|
||||||
Ok(Stream::new(
|
Ok(Stream::new(
|
||||||
|
@ -276,34 +279,23 @@ fn stream_file(stack: &mut Stack) -> Result<Stream, Error> {
|
||||||
.truncate(truncate)
|
.truncate(truncate)
|
||||||
.open(path)
|
.open(path)
|
||||||
.map_err(|x| stack.error(ErrorKind::IO(x.to_string())))?,
|
.map_err(|x| stack.error(ErrorKind::IO(x.to_string())))?,
|
||||||
nop,
|
|
||||||
))
|
))
|
||||||
}
|
}
|
||||||
|
|
||||||
fn stream_tcp(stack: &mut Stack) -> Result<Stream, Error> {
|
pub(super) fn stream_tcp(stack: &mut Stack) -> Result<Stream, Error> {
|
||||||
require_int_on_stack!(port, stack, "TCP new-stream");
|
require_int_on_stack!(port, stack, "TCP new-stream");
|
||||||
require_on_stack!(ip, Str, stack, "TCP new-stream");
|
require_on_stack!(ip, Str, stack, "TCP new-stream");
|
||||||
fn close_tcp(stream: &mut Stream) {
|
|
||||||
unsafe {
|
|
||||||
let f = ((stream.reader.as_mut() as *mut dyn Read).cast() as *mut TcpStream)
|
|
||||||
.as_mut()
|
|
||||||
.unwrap();
|
|
||||||
let _ = f.shutdown(Shutdown::Both);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
Ok(Stream::new(
|
Ok(Stream::new(
|
||||||
TcpStream::connect((ip, port as u16))
|
TcpStream::connect((ip, port as u16))
|
||||||
.map_err(|x| stack.error(ErrorKind::IO(x.to_string())))?,
|
.map_err(|x| stack.error(ErrorKind::IO(x.to_string())))?,
|
||||||
close_tcp,
|
|
||||||
))
|
))
|
||||||
}
|
}
|
||||||
|
|
||||||
fn stream_udp(stack: &mut Stack) -> Result<Stream, Error> {
|
pub(super) fn stream_udp(stack: &mut Stack) -> Result<Stream, Error> {
|
||||||
require_int_on_stack!(port, stack, "UDP new-stream");
|
require_int_on_stack!(port, stack, "UDP new-stream");
|
||||||
require_on_stack!(ip, Str, stack, "UDP new-stream");
|
require_on_stack!(ip, Str, stack, "UDP new-stream");
|
||||||
require_int_on_stack!(self_port, stack, "UDP new-stream");
|
require_int_on_stack!(self_port, stack, "UDP new-stream");
|
||||||
require_on_stack!(self_ip, Str, stack, "UDP new-stream");
|
require_on_stack!(self_ip, Str, stack, "UDP new-stream");
|
||||||
fn close_udp(_stream: &mut Stream) {}
|
|
||||||
let sock = UdpSocket::bind((self_ip, self_port as u16))
|
let sock = UdpSocket::bind((self_ip, self_port as u16))
|
||||||
.map_err(|x| stack.error(ErrorKind::IO(x.to_string())))?;
|
.map_err(|x| stack.error(ErrorKind::IO(x.to_string())))?;
|
||||||
sock.connect((ip, port as u16))
|
sock.connect((ip, port as u16))
|
||||||
|
@ -323,12 +315,11 @@ fn stream_udp(stack: &mut Stack) -> Result<Stream, Error> {
|
||||||
self.0.recv(buf)
|
self.0.recv(buf)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
Ok(Stream::new(UdpRW(sock), close_udp))
|
Ok(Stream::new(UdpRW(sock)))
|
||||||
}
|
}
|
||||||
|
|
||||||
fn stream_cmd(stack: &mut Stack) -> Result<Stream, Error> {
|
pub(super) fn stream_cmd(stack: &mut Stack) -> Result<Stream, Error> {
|
||||||
require_on_stack!(a, Array, stack, "CMD new-stream");
|
require_on_stack!(a, Array, stack, "CMD new-stream");
|
||||||
fn close_cmd(_stream: &mut Stream) {}
|
|
||||||
let mut args = Vec::new();
|
let mut args = Vec::new();
|
||||||
for item in a.iter() {
|
for item in a.iter() {
|
||||||
if let Value::Str(ref s) = item.lock_ro().native {
|
if let Value::Str(ref s) = item.lock_ro().native {
|
||||||
|
@ -348,40 +339,5 @@ fn stream_cmd(stack: &mut Stack) -> Result<Stream, Error> {
|
||||||
Ok(Stream::new_split(
|
Ok(Stream::new_split(
|
||||||
command.stdout.take().unwrap(),
|
command.stdout.take().unwrap(),
|
||||||
command.stdin.take().unwrap(),
|
command.stdin.take().unwrap(),
|
||||||
close_cmd,
|
|
||||||
))
|
))
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn register(r: &mut Stack, o: Arc<Frame>) {
|
|
||||||
if !*IS_INITIALIZED.lock_ro() {
|
|
||||||
register_stream_type("file", stream_file);
|
|
||||||
register_stream_type("tcp", stream_tcp);
|
|
||||||
register_stream_type("udp", stream_udp);
|
|
||||||
register_stream_type("cmd", stream_cmd);
|
|
||||||
*IS_INITIALIZED.lock() = true;
|
|
||||||
}
|
|
||||||
|
|
||||||
type Fn = fn(&mut Stack) -> OError;
|
|
||||||
let fns: [(&str, Fn, u32); 7] = [
|
|
||||||
("new-stream", new_stream, 1),
|
|
||||||
("write-stream", write_stream, 1),
|
|
||||||
("write-all-stream", write_all_stream, 0),
|
|
||||||
("flush-stream", flush_stream, 0),
|
|
||||||
("read-stream", read_stream, 1),
|
|
||||||
("read-all-stream", read_all_stream, 0),
|
|
||||||
("close-stream", close_stream, 0),
|
|
||||||
];
|
|
||||||
for f in fns {
|
|
||||||
r.define_func(
|
|
||||||
f.0.to_owned(),
|
|
||||||
AFunc::new(Func {
|
|
||||||
ret_count: f.2,
|
|
||||||
to_call: FuncImpl::Native(f.1),
|
|
||||||
run_as_base: false,
|
|
||||||
origin: o.clone(),
|
|
||||||
fname: None,
|
|
||||||
name: f.0.to_owned(),
|
|
||||||
}),
|
|
||||||
);
|
|
||||||
}
|
|
||||||
}
|
|
105
src/stream/server.rs
Normal file
105
src/stream/server.rs
Normal file
|
@ -0,0 +1,105 @@
|
||||||
|
use std::{
|
||||||
|
collections::HashMap,
|
||||||
|
net::TcpListener,
|
||||||
|
sync::{Arc, LazyLock},
|
||||||
|
};
|
||||||
|
|
||||||
|
use crate::{mutex::Mut, *};
|
||||||
|
|
||||||
|
use super::Stream;
|
||||||
|
|
||||||
|
static SERVER_STREAM_TYPES: LazyLock<Arc<Mut<HashMap<String, ServerStreamType>>>> =
|
||||||
|
LazyLock::new(|| Arc::new(Mut::new(HashMap::new())));
|
||||||
|
|
||||||
|
/// Registers a custom stream type.
|
||||||
|
pub fn register_server_stream_type(
|
||||||
|
name: &str,
|
||||||
|
supplier: impl Fn(&mut Stack) -> Result<ServerStream, Error> + Sync + Send + 'static,
|
||||||
|
) {
|
||||||
|
SERVER_STREAM_TYPES
|
||||||
|
.lock()
|
||||||
|
.insert(name.to_owned(), ServerStreamType::from(supplier));
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Gets a stream type by name.
|
||||||
|
pub fn get_server_stream_type(name: String) -> Option<ServerStreamType> {
|
||||||
|
SERVER_STREAM_TYPES.lock_ro().get(&name).cloned()
|
||||||
|
}
|
||||||
|
|
||||||
|
/// An SPL stream type.
|
||||||
|
#[derive(Clone)]
|
||||||
|
pub struct ServerStreamType {
|
||||||
|
func: Arc<Box<dyn Fn(&mut Stack) -> Result<ServerStream, Error> + Sync + Send + 'static>>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl ServerStreamType {
|
||||||
|
pub fn make_stream(&self, stack: &mut Stack) -> Result<ServerStream, Error> {
|
||||||
|
(self.func)(stack)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl<T> From<T> for ServerStreamType
|
||||||
|
where
|
||||||
|
T: Fn(&mut Stack) -> Result<ServerStream, Error> + Sync + Send + 'static,
|
||||||
|
{
|
||||||
|
fn from(value: T) -> Self {
|
||||||
|
Self {
|
||||||
|
func: Arc::new(Box::new(value)),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// An SPL server stream, holding an acceptor and a function to close it.
|
||||||
|
pub struct ServerStream {
|
||||||
|
pub(super) acceptor: Box<dyn Fn(&mut Stack) -> Result<Stream, Error> + Sync + Send + 'static>,
|
||||||
|
}
|
||||||
|
impl ServerStream {
|
||||||
|
pub fn new(
|
||||||
|
acceptor: impl Fn(&mut Stack) -> Result<Stream, Error> + Sync + Send + 'static,
|
||||||
|
) -> Self {
|
||||||
|
Self {
|
||||||
|
acceptor: Box::new(acceptor),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn new_server_stream(stack: &mut Stack) -> OError {
|
||||||
|
require_on_stack!(s, Str, stack, "new-stream");
|
||||||
|
let stream = get_server_stream_type(s.clone())
|
||||||
|
.ok_or_else(|| stack.error(ErrorKind::VariableNotFound(format!("__stream-type-{s}"))))?
|
||||||
|
.make_stream(stack)?;
|
||||||
|
let stream = runtime_mut(move |mut rt| Ok(rt.register_server_stream(stream)))?;
|
||||||
|
stack.push(Value::Mega(stream.0 as i128).spl());
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn accept_server_stream(stack: &mut Stack) -> OError {
|
||||||
|
require_on_stack!(id, Mega, stack, "accept-server-stream");
|
||||||
|
let stream = runtime(|rt| {
|
||||||
|
rt.get_server_stream(id as u128)
|
||||||
|
.ok_or_else(|| stack.error(ErrorKind::VariableNotFound(format!("__stream-{id}"))))
|
||||||
|
})?;
|
||||||
|
let stream = (stream.lock_ro().acceptor)(stack)?;
|
||||||
|
stack.push((runtime_mut(move |mut rt| rt.register_stream(stream)).0 as i128).spl());
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn close_server_stream(stack: &mut Stack) -> OError {
|
||||||
|
require_on_stack!(id, Mega, stack, "close-server-stream");
|
||||||
|
runtime_mut(|mut rt| rt.destroy_server_stream(id as u128));
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) fn server_stream_tcp(stack: &mut Stack) -> Result<ServerStream, Error> {
|
||||||
|
require_int_on_stack!(port, stack, "TCP server-stream");
|
||||||
|
require_on_stack!(addr, Str, stack, "TCP server-stream");
|
||||||
|
let tcp = TcpListener::bind((addr, port as u16))
|
||||||
|
.map_err(|e| stack.error(ErrorKind::IO(format!("{e:?}"))))?;
|
||||||
|
let stream = ServerStream::new(move |stack| {
|
||||||
|
let socket = tcp
|
||||||
|
.accept()
|
||||||
|
.map_err(|e| stack.error(ErrorKind::IO(format!("{e:?}"))))?;
|
||||||
|
Ok(Stream::new(socket.0))
|
||||||
|
});
|
||||||
|
Ok(stream)
|
||||||
|
}
|
|
@ -66,7 +66,7 @@ construct StreamType {
|
||||||
|
|
||||||
def stream-types 0 anew =stream-types
|
def stream-types 0 anew =stream-types
|
||||||
|
|
||||||
construct _StreamType {
|
construct _StreamTypes {
|
||||||
;
|
;
|
||||||
construct { this | with this ;
|
construct { this | with this ;
|
||||||
{ | with type ;
|
{ | with type ;
|
||||||
|
@ -79,7 +79,7 @@ construct _StreamType {
|
||||||
|
|
||||||
func register-stream-type { | with id ;
|
func register-stream-type { | with id ;
|
||||||
stream-types [ id ] aadd =stream-types
|
stream-types [ id ] aadd =stream-types
|
||||||
id _StreamType dyn-def-field
|
id _StreamTypes dyn-def-field
|
||||||
}
|
}
|
||||||
|
|
||||||
"tcp" register-stream-type
|
"tcp" register-stream-type
|
||||||
|
@ -87,6 +87,6 @@ func register-stream-type { | with id ;
|
||||||
"file" register-stream-type
|
"file" register-stream-type
|
||||||
"cmd" register-stream-type
|
"cmd" register-stream-type
|
||||||
|
|
||||||
func StreamTypes { _StreamType |
|
func StreamTypes { _StreamTypes |
|
||||||
_StreamType:new
|
_StreamTypes:new
|
||||||
}
|
}
|
||||||
|
|
20
test.spl
20
test.spl
|
@ -2,6 +2,8 @@
|
||||||
"#stream.spl" import
|
"#stream.spl" import
|
||||||
"#http.spl" import
|
"#http.spl" import
|
||||||
"#messaging.spl" import
|
"#messaging.spl" import
|
||||||
|
"#server.spl" import
|
||||||
|
"#time.spl" import
|
||||||
|
|
||||||
"SPL tester" =program-name
|
"SPL tester" =program-name
|
||||||
|
|
||||||
|
@ -161,6 +163,24 @@ func main { int | with args ;
|
||||||
} fork
|
} fork
|
||||||
while { other-thread-done not } { "waiting for the other thread..." println }
|
while { other-thread-done not } { "waiting for the other thread..." println }
|
||||||
|
|
||||||
|
"" println
|
||||||
|
"testing tcp server" println
|
||||||
|
|
||||||
|
" starting server thread" println
|
||||||
|
{ |
|
||||||
|
def server "0.0.0.0" 4075 net:server:Types:tcp:create =server
|
||||||
|
while { 1 } {
|
||||||
|
def stream server:accept =stream
|
||||||
|
"Hello!" :to-bytes stream:write-exact;
|
||||||
|
stream:close;
|
||||||
|
}
|
||||||
|
} fork;
|
||||||
|
50 time:sleep;
|
||||||
|
" starting client" println;
|
||||||
|
def client "localhost" 4075 StreamTypes:tcp:create =client
|
||||||
|
1024 client:read-to-end:to-str println;
|
||||||
|
" ^ this should say 'Hello!'" println;
|
||||||
|
|
||||||
100
|
100
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
9
time.spl
9
time.spl
|
@ -1,6 +1,9 @@
|
||||||
|
|
||||||
func unixms { mega |
|
construct time namespace {
|
||||||
0 time
|
;
|
||||||
|
unixms { mega | with this ;
|
||||||
|
0 sleeptime
|
||||||
|
}
|
||||||
|
sleep { mega | with this ; _mega sleeptime }
|
||||||
}
|
}
|
||||||
|
|
||||||
func sleep { mega | time }
|
|
||||||
|
|
Loading…
Reference in a new issue