This commit is contained in:
Astro 2017-06-04 02:05:08 +02:00
commit 482bf77955
4 changed files with 36 additions and 42 deletions

View file

@ -6,4 +6,6 @@ authors = ["Astro <astro@spaceboyz.net>"]
[dependencies] [dependencies]
futures = "*" futures = "*"
tokio-core = "*" tokio-core = "*"
tokio-io = "*"
bytes = "*"
RustyXML = "*" RustyXML = "*"

View file

@ -1,9 +1,9 @@
#[macro_use] #[macro_use]
extern crate futures; extern crate futures;
extern crate tokio_core; extern crate tokio_core;
extern crate tokio_io;
extern crate bytes;
extern crate xml; extern crate xml;
extern crate rustls;
extern crate tokio_rustls;
mod xmpp_codec; mod xmpp_codec;

View file

@ -1,16 +1,12 @@
use std::fmt; use std::fmt;
use std::net::SocketAddr; use std::net::SocketAddr;
use std::net::ToSocketAddrs;
use std::sync::Arc;
use std::io::{Error, ErrorKind}; use std::io::{Error, ErrorKind};
use futures::{Future, BoxFuture, Sink, Poll, Async}; use futures::{Future, Sink, Poll, Async};
use futures::stream::{Stream, iter}; use futures::stream::Stream;
use futures::sink; use futures::sink;
use tokio_core::reactor::Handle; use tokio_core::reactor::Handle;
use tokio_core::io::Io; use tokio_io::AsyncRead;
use tokio_core::net::{TcpStream, TcpStreamNew}; use tokio_core::net::{TcpStream, TcpStreamNew};
use rustls::ClientConfig;
use tokio_rustls::ClientConfigExt;
use super::{XMPPStream, XMPPCodec, Packet}; use super::{XMPPStream, XMPPCodec, Packet};
@ -25,7 +21,6 @@ enum TcpClientState {
SendStart(sink::Send<XMPPStream<TcpStream>>), SendStart(sink::Send<XMPPStream<TcpStream>>),
RecvStart(Option<XMPPStream<TcpStream>>), RecvStart(Option<XMPPStream<TcpStream>>),
Established, Established,
Invalid,
} }
impl fmt::Debug for TcpClientState { impl fmt::Debug for TcpClientState {
@ -35,9 +30,9 @@ impl fmt::Debug for TcpClientState {
TcpClientState::SendStart(_) => "SendStart", TcpClientState::SendStart(_) => "SendStart",
TcpClientState::RecvStart(_) => "RecvStart", TcpClientState::RecvStart(_) => "RecvStart",
TcpClientState::Established => "Established", TcpClientState::Established => "Established",
TcpClientState::Invalid => "Invalid",
}; };
write!(fmt, "{}", s) try!(write!(fmt, "{}", s));
Ok(())
} }
} }
@ -58,7 +53,7 @@ impl Future for TcpClient {
let (new_state, result) = match self.state { let (new_state, result) = match self.state {
TcpClientState::Connecting(ref mut tcp_stream_new) => { TcpClientState::Connecting(ref mut tcp_stream_new) => {
let tcp_stream = try_ready!(tcp_stream_new.poll()); let tcp_stream = try_ready!(tcp_stream_new.poll());
let xmpp_stream = tcp_stream.framed(XMPPCodec::new()); let xmpp_stream = AsyncRead::framed(tcp_stream, XMPPCodec::new());
let send = xmpp_stream.send(Packet::StreamStart); let send = xmpp_stream.send(Packet::StreamStart);
let new_state = TcpClientState::SendStart(send); let new_state = TcpClientState::SendStart(send);
(new_state, Ok(Async::NotReady)) (new_state, Ok(Async::NotReady))
@ -82,7 +77,7 @@ impl Future for TcpClient {
let new_state = TcpClientState::Established; let new_state = TcpClientState::Established;
(new_state, Ok(Async::Ready(xmpp_stream))) (new_state, Ok(Async::Ready(xmpp_stream)))
}, },
TcpClientState::Established | TcpClientState::Invalid => TcpClientState::Established =>
unreachable!(), unreachable!(),
}; };

View file

@ -1,9 +1,11 @@
use std; use std;
use std::fmt::Write;
use std::str::from_utf8; use std::str::from_utf8;
use std::io::{Error, ErrorKind}; use std::io::{Error, ErrorKind};
use std::collections::HashMap; use std::collections::HashMap;
use tokio_core::io::{Codec, EasyBuf, Framed}; use tokio_io::codec::{Framed, Encoder, Decoder};
use xml; use xml;
use bytes::*;
const NS_XMLNS: &'static str = "http://www.w3.org/2000/xmlns/"; const NS_XMLNS: &'static str = "http://www.w3.org/2000/xmlns/";
const NS_STREAMS: &'static str = "http://etherx.jabber.org/streams"; const NS_STREAMS: &'static str = "http://etherx.jabber.org/streams";
@ -67,22 +69,20 @@ impl XMPPCodec {
} }
} }
impl Codec for XMPPCodec { impl Decoder for XMPPCodec {
type In = Packet; type Item = Packet;
type Out = Packet; type Error = Error;
fn decode(&mut self, buf: &mut EasyBuf) -> Result<Option<Self::In>, Error> { fn decode(&mut self, buf: &mut BytesMut) -> Result<Option<Self::Item>, Self::Error> {
println!("XMPPCodec.decode {:?}", buf.len()); println!("XMPPCodec.decode {:?}", buf.len());
let buf_len = buf.len(); match from_utf8(buf.take().as_ref()) {
let chunk = buf.drain_to(buf_len);
match from_utf8(chunk.as_slice()) {
Ok(s) => Ok(s) =>
self.parser.feed_str(s), self.parser.feed_str(s),
Err(e) => Err(e) =>
return Err(Error::new(ErrorKind::InvalidInput, e)), return Err(Error::new(ErrorKind::InvalidInput, e)),
} }
let mut new_root = None; let mut new_root: Option<XMPPRoot> = None;
let mut result = None; let mut result = None;
for event in &mut self.parser { for event in &mut self.parser {
match self.root { match self.root {
@ -128,29 +128,26 @@ impl Codec for XMPPCodec {
Ok(result) Ok(result)
} }
fn encode(&mut self, msg: Self::Out, buf: &mut Vec<u8>) -> Result<(), Error> { fn decode_eof(&mut self, buf: &mut BytesMut) -> Result<Option<Self::Item>, Error> {
match msg { self.decode(buf)
}
}
impl Encoder for XMPPCodec {
type Item = Packet;
type Error = Error;
fn encode(&mut self, item: Self::Item, dst: &mut BytesMut) -> Result<(), Self::Error> {
match item {
Packet::StreamStart => { Packet::StreamStart => {
let mut write = |s: &str| { write!(dst,
buf.extend_from_slice(s.as_bytes()); "<?xml version='1.0'?>\n
}; <stream:stream version='1.0' to='spaceboyz.net' xmlns='{}' xmlns:stream='{}'>\n",
NS_CLIENT, NS_STREAMS)
write("<?xml version='1.0'?>\n"); .map_err(|_| Error::from(ErrorKind::WriteZero))
write("<stream:stream");
write(" version='1.0'");
write(" to='spaceboyz.net'");
write(&format!(" xmlns='{}'", NS_CLIENT));
write(&format!(" xmlns:stream='{}'", NS_STREAMS));
write(">\n");
Ok(())
}, },
// TODO: Implement all // TODO: Implement all
_ => Ok(()) _ => Ok(())
} }
} }
fn decode_eof(&mut self, _buf: &mut EasyBuf) -> Result<Self::In, Error> {
Err(Error::from(ErrorKind::UnexpectedEof))
}
} }