minidom, tokio-xmpp: switch xml parsing to rxml
This commit is contained in:
parent
3901068717
commit
d4a5a8247b
33 changed files with 465 additions and 811 deletions
|
|
@ -203,30 +203,28 @@ impl Stream for Client {
|
|||
self.poll_next(cx)
|
||||
}
|
||||
ClientState::Disconnected => Poll::Ready(None),
|
||||
ClientState::Connecting(mut connect) => {
|
||||
match Pin::new(&mut connect).poll(cx) {
|
||||
Poll::Ready(Ok(Ok(stream))) => {
|
||||
let bound_jid = stream.jid.clone();
|
||||
self.state = ClientState::Connected(stream);
|
||||
Poll::Ready(Some(Event::Online {
|
||||
bound_jid,
|
||||
resumed: false,
|
||||
}))
|
||||
}
|
||||
Poll::Ready(Ok(Err(e))) => {
|
||||
self.state = ClientState::Disconnected;
|
||||
return Poll::Ready(Some(Event::Disconnected(e.into())));
|
||||
}
|
||||
Poll::Ready(Err(e)) => {
|
||||
self.state = ClientState::Disconnected;
|
||||
panic!("connect task: {}", e);
|
||||
}
|
||||
Poll::Pending => {
|
||||
self.state = ClientState::Connecting(connect);
|
||||
Poll::Pending
|
||||
}
|
||||
ClientState::Connecting(mut connect) => match Pin::new(&mut connect).poll(cx) {
|
||||
Poll::Ready(Ok(Ok(stream))) => {
|
||||
let bound_jid = stream.jid.clone();
|
||||
self.state = ClientState::Connected(stream);
|
||||
Poll::Ready(Some(Event::Online {
|
||||
bound_jid,
|
||||
resumed: false,
|
||||
}))
|
||||
}
|
||||
}
|
||||
Poll::Ready(Ok(Err(e))) => {
|
||||
self.state = ClientState::Disconnected;
|
||||
return Poll::Ready(Some(Event::Disconnected(e.into())));
|
||||
}
|
||||
Poll::Ready(Err(e)) => {
|
||||
self.state = ClientState::Disconnected;
|
||||
panic!("connect task: {}", e);
|
||||
}
|
||||
Poll::Pending => {
|
||||
self.state = ClientState::Connecting(connect);
|
||||
Poll::Pending
|
||||
}
|
||||
},
|
||||
ClientState::Connected(mut stream) => {
|
||||
// Poll sink
|
||||
match Pin::new(&mut stream).poll_ready(cx) {
|
||||
|
|
|
|||
|
|
@ -5,7 +5,6 @@ use std::borrow::Cow;
|
|||
use std::error::Error as StdError;
|
||||
use std::fmt;
|
||||
use std::io::Error as IoError;
|
||||
use std::str::Utf8Error;
|
||||
#[cfg(feature = "tls-rust")]
|
||||
use tokio_rustls::rustls::client::InvalidDnsNameError;
|
||||
#[cfg(feature = "tls-rust")]
|
||||
|
|
@ -106,44 +105,6 @@ impl From<InvalidDnsNameError> for Error {
|
|||
}
|
||||
}
|
||||
|
||||
/// Causes for stream parsing errors
|
||||
#[derive(Debug)]
|
||||
pub enum ParserError {
|
||||
/// Encoding error
|
||||
Utf8(Utf8Error),
|
||||
/// XML parse error
|
||||
Parse(ParseError),
|
||||
/// Illegal `</>`
|
||||
ShortTag,
|
||||
/// Required by `impl Decoder`
|
||||
Io(IoError),
|
||||
}
|
||||
|
||||
impl fmt::Display for ParserError {
|
||||
fn fmt(&self, fmt: &mut fmt::Formatter) -> fmt::Result {
|
||||
match self {
|
||||
ParserError::Utf8(e) => write!(fmt, "UTF-8 error: {}", e),
|
||||
ParserError::Parse(e) => write!(fmt, "parse error: {}", e),
|
||||
ParserError::ShortTag => write!(fmt, "short tag"),
|
||||
ParserError::Io(e) => write!(fmt, "IO error: {}", e),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl StdError for ParserError {}
|
||||
|
||||
impl From<IoError> for ParserError {
|
||||
fn from(e: IoError) -> Self {
|
||||
ParserError::Io(e)
|
||||
}
|
||||
}
|
||||
|
||||
impl From<ParserError> for Error {
|
||||
fn from(e: ParserError) -> Self {
|
||||
ProtocolError::Parser(e).into()
|
||||
}
|
||||
}
|
||||
|
||||
/// XML parse error wrapper type
|
||||
#[derive(Debug)]
|
||||
pub struct ParseError(pub Cow<'static, str>);
|
||||
|
|
@ -167,7 +128,7 @@ impl fmt::Display for ParseError {
|
|||
#[derive(Debug)]
|
||||
pub enum ProtocolError {
|
||||
/// XML parser error
|
||||
Parser(ParserError),
|
||||
Parser(minidom::Error),
|
||||
/// Error with expected stanza schema
|
||||
Parsers(ParsersError),
|
||||
/// No TLS available
|
||||
|
|
@ -205,12 +166,18 @@ impl fmt::Display for ProtocolError {
|
|||
|
||||
impl StdError for ProtocolError {}
|
||||
|
||||
impl From<ParserError> for ProtocolError {
|
||||
fn from(e: ParserError) -> Self {
|
||||
impl From<minidom::Error> for ProtocolError {
|
||||
fn from(e: minidom::Error) -> Self {
|
||||
ProtocolError::Parser(e)
|
||||
}
|
||||
}
|
||||
|
||||
impl From<minidom::Error> for Error {
|
||||
fn from(e: minidom::Error) -> Self {
|
||||
ProtocolError::Parser(e).into()
|
||||
}
|
||||
}
|
||||
|
||||
impl From<ParsersError> for ProtocolError {
|
||||
fn from(e: ParsersError) -> Self {
|
||||
ProtocolError::Parsers(e)
|
||||
|
|
|
|||
|
|
@ -16,5 +16,5 @@ pub use client::{async_client::Client as AsyncClient, simple_client::Client as S
|
|||
mod component;
|
||||
pub use crate::component::Component;
|
||||
mod error;
|
||||
pub use crate::error::{AuthError, ConnecterError, Error, ParseError, ParserError, ProtocolError};
|
||||
pub use crate::error::{AuthError, ConnecterError, Error, ParseError, ProtocolError};
|
||||
pub use starttls::starttls;
|
||||
|
|
|
|||
|
|
@ -1,23 +1,16 @@
|
|||
//! XML stream parser for XMPP
|
||||
|
||||
use crate::{ParseError, ParserError};
|
||||
use crate::Error;
|
||||
use bytes::{BufMut, BytesMut};
|
||||
use log::{debug, error};
|
||||
use log::debug;
|
||||
use minidom::tree_builder::TreeBuilder;
|
||||
use rxml::{Lexer, PushDriver, RawParser};
|
||||
use std;
|
||||
use std::borrow::Cow;
|
||||
use std::collections::vec_deque::VecDeque;
|
||||
use std::collections::HashMap;
|
||||
use std::default::Default;
|
||||
use std::fmt::Write;
|
||||
use std::io;
|
||||
use std::iter::FromIterator;
|
||||
use std::str::from_utf8;
|
||||
use std::sync::Arc;
|
||||
use std::sync::Mutex;
|
||||
use tokio_util::codec::{Decoder, Encoder};
|
||||
use xml5ever::buffer_queue::BufferQueue;
|
||||
use xml5ever::interface::Attribute;
|
||||
use xml5ever::tokenizer::{Tag, TagKind, Token, TokenSink, XmlTokenizer};
|
||||
use xmpp_parsers::Element;
|
||||
|
||||
/// Anything that can be sent or received on an XMPP/XML stream
|
||||
|
|
@ -33,175 +26,24 @@ pub enum Packet {
|
|||
StreamEnd,
|
||||
}
|
||||
|
||||
type QueueItem = Result<Packet, ParserError>;
|
||||
|
||||
/// Parser state
|
||||
struct ParserSink {
|
||||
// Ready stanzas, shared with XMPPCodec
|
||||
queue: Arc<Mutex<VecDeque<QueueItem>>>,
|
||||
// Parsing stack
|
||||
stack: Vec<Element>,
|
||||
ns_stack: Vec<HashMap<Option<String>, String>>,
|
||||
}
|
||||
|
||||
impl ParserSink {
|
||||
pub fn new(queue: Arc<Mutex<VecDeque<QueueItem>>>) -> Self {
|
||||
ParserSink {
|
||||
queue,
|
||||
stack: vec![],
|
||||
ns_stack: vec![],
|
||||
}
|
||||
}
|
||||
|
||||
fn push_queue(&self, pkt: Packet) {
|
||||
self.queue.lock().unwrap().push_back(Ok(pkt));
|
||||
}
|
||||
|
||||
fn push_queue_error(&self, e: ParserError) {
|
||||
self.queue.lock().unwrap().push_back(Err(e));
|
||||
}
|
||||
|
||||
/// Lookup XML namespace declaration for given prefix (or no prefix)
|
||||
fn lookup_ns(&self, prefix: &Option<String>) -> Option<&str> {
|
||||
for nss in self.ns_stack.iter().rev() {
|
||||
if let Some(ns) = nss.get(prefix) {
|
||||
return Some(ns);
|
||||
}
|
||||
}
|
||||
|
||||
None
|
||||
}
|
||||
|
||||
fn handle_start_tag(&mut self, tag: Tag) {
|
||||
let mut nss = HashMap::new();
|
||||
let is_prefix_xmlns = |attr: &Attribute| {
|
||||
attr.name
|
||||
.prefix
|
||||
.as_ref()
|
||||
.map(|prefix| prefix.eq_str_ignore_ascii_case("xmlns"))
|
||||
.unwrap_or(false)
|
||||
};
|
||||
for attr in &tag.attrs {
|
||||
match attr.name.local.as_ref() {
|
||||
"xmlns" => {
|
||||
nss.insert(None, attr.value.as_ref().to_owned());
|
||||
}
|
||||
prefix if is_prefix_xmlns(attr) => {
|
||||
nss.insert(Some(prefix.to_owned()), attr.value.as_ref().to_owned());
|
||||
}
|
||||
_ => (),
|
||||
}
|
||||
}
|
||||
self.ns_stack.push(nss);
|
||||
|
||||
let el = {
|
||||
let el_ns = self
|
||||
.lookup_ns(&tag.name.prefix.map(|prefix| prefix.as_ref().to_owned()))
|
||||
.unwrap();
|
||||
let mut el_builder = Element::builder(tag.name.local.as_ref(), el_ns);
|
||||
for attr in &tag.attrs {
|
||||
match attr.name.local.as_ref() {
|
||||
"xmlns" => (),
|
||||
_ if is_prefix_xmlns(attr) => (),
|
||||
_ => {
|
||||
let attr_name = if let Some(ref prefix) = attr.name.prefix {
|
||||
Cow::Owned(format!("{}:{}", prefix, attr.name.local))
|
||||
} else {
|
||||
Cow::Borrowed(attr.name.local.as_ref())
|
||||
};
|
||||
el_builder = el_builder.attr(attr_name, attr.value.as_ref());
|
||||
}
|
||||
}
|
||||
}
|
||||
el_builder.build()
|
||||
};
|
||||
|
||||
if self.stack.is_empty() {
|
||||
let attrs = HashMap::from_iter(tag.attrs.iter().map(|attr| {
|
||||
(
|
||||
attr.name.local.as_ref().to_owned(),
|
||||
attr.value.as_ref().to_owned(),
|
||||
)
|
||||
}));
|
||||
self.push_queue(Packet::StreamStart(attrs));
|
||||
}
|
||||
|
||||
self.stack.push(el);
|
||||
}
|
||||
|
||||
fn handle_end_tag(&mut self) {
|
||||
let el = self.stack.pop().unwrap();
|
||||
self.ns_stack.pop();
|
||||
|
||||
match self.stack.len() {
|
||||
// </stream:stream>
|
||||
0 => self.push_queue(Packet::StreamEnd),
|
||||
// </stanza>
|
||||
1 => self.push_queue(Packet::Stanza(el)),
|
||||
len => {
|
||||
let parent = &mut self.stack[len - 1];
|
||||
parent.append_child(el);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl TokenSink for ParserSink {
|
||||
fn process_token(&mut self, token: Token) {
|
||||
match token {
|
||||
Token::TagToken(tag) => match tag.kind {
|
||||
TagKind::StartTag => self.handle_start_tag(tag),
|
||||
TagKind::EndTag => self.handle_end_tag(),
|
||||
TagKind::EmptyTag => {
|
||||
self.handle_start_tag(tag);
|
||||
self.handle_end_tag();
|
||||
}
|
||||
TagKind::ShortTag => self.push_queue_error(ParserError::ShortTag),
|
||||
},
|
||||
Token::CharacterTokens(tendril) => match self.stack.len() {
|
||||
0 | 1 => self.push_queue(Packet::Text(tendril.into())),
|
||||
len => {
|
||||
let el = &mut self.stack[len - 1];
|
||||
el.append_text_node(tendril);
|
||||
}
|
||||
},
|
||||
Token::EOFToken => self.push_queue(Packet::StreamEnd),
|
||||
Token::ParseError(s) => {
|
||||
self.push_queue_error(ParserError::Parse(ParseError(s)));
|
||||
}
|
||||
_ => (),
|
||||
}
|
||||
}
|
||||
|
||||
// fn end(&mut self) {
|
||||
// }
|
||||
}
|
||||
|
||||
/// Stateful encoder/decoder for a bytestream from/to XMPP `Packet`
|
||||
pub struct XMPPCodec {
|
||||
/// Outgoing
|
||||
ns: Option<String>,
|
||||
/// Incoming
|
||||
parser: XmlTokenizer<ParserSink>,
|
||||
/// For handling incoming truncated utf8
|
||||
// TODO: optimize using tendrils?
|
||||
buf: Vec<u8>,
|
||||
/// Shared with ParserSink
|
||||
queue: Arc<Mutex<VecDeque<QueueItem>>>,
|
||||
driver: PushDriver<RawParser>,
|
||||
stanza_builder: TreeBuilder,
|
||||
}
|
||||
|
||||
impl XMPPCodec {
|
||||
/// Constructor
|
||||
pub fn new() -> Self {
|
||||
let queue = Arc::new(Mutex::new(VecDeque::new()));
|
||||
let sink = ParserSink::new(queue.clone());
|
||||
// TODO: configure parser?
|
||||
let parser = XmlTokenizer::new(sink, Default::default());
|
||||
let stanza_builder = TreeBuilder::new();
|
||||
let driver = PushDriver::wrap(Lexer::new(), RawParser::new());
|
||||
XMPPCodec {
|
||||
ns: None,
|
||||
parser,
|
||||
queue,
|
||||
buf: vec![],
|
||||
driver,
|
||||
stanza_builder,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -214,57 +56,53 @@ impl Default for XMPPCodec {
|
|||
|
||||
impl Decoder for XMPPCodec {
|
||||
type Item = Packet;
|
||||
type Error = ParserError;
|
||||
type Error = Error;
|
||||
|
||||
fn decode(&mut self, buf: &mut BytesMut) -> Result<Option<Self::Item>, Self::Error> {
|
||||
let buf1: Box<dyn AsRef<[u8]>> = if !self.buf.is_empty() && !buf.is_empty() {
|
||||
let mut prefix = std::mem::replace(&mut self.buf, vec![]);
|
||||
prefix.extend_from_slice(&buf.split_to(buf.len()));
|
||||
Box::new(prefix)
|
||||
} else {
|
||||
Box::new(buf.split_to(buf.len()))
|
||||
};
|
||||
let buf1 = buf1.as_ref().as_ref();
|
||||
match from_utf8(buf1) {
|
||||
Ok(s) => {
|
||||
debug!("<< {:?}", s);
|
||||
if !s.is_empty() {
|
||||
let mut buffer_queue = BufferQueue::new();
|
||||
let tendril = FromIterator::from_iter(s.chars());
|
||||
buffer_queue.push_back(tendril);
|
||||
self.parser.feed(&mut buffer_queue);
|
||||
loop {
|
||||
let token = match self.driver.parse(buf, false) {
|
||||
Ok(Some(token)) => token,
|
||||
Ok(None) => break,
|
||||
Err(rxml::Error::IO(e)) if e.kind() == std::io::ErrorKind::WouldBlock => break,
|
||||
Err(e) => return Err(minidom::Error::from(e).into()),
|
||||
};
|
||||
|
||||
let had_stream_root = self.stanza_builder.depth() > 0;
|
||||
self.stanza_builder.process_event(token)?;
|
||||
let has_stream_root = self.stanza_builder.depth() > 0;
|
||||
|
||||
if !had_stream_root && has_stream_root {
|
||||
let root = self.stanza_builder.top().unwrap();
|
||||
let attrs =
|
||||
root.attrs()
|
||||
.map(|(name, value)| (name.to_owned(), value.to_owned()))
|
||||
.chain(root.prefixes.declared_prefixes().iter().map(
|
||||
|(prefix, namespace)| {
|
||||
(
|
||||
prefix
|
||||
.as_ref()
|
||||
.map(|prefix| format!("xmlns:{}", prefix))
|
||||
.unwrap_or_else(|| "xmlns".to_owned()),
|
||||
namespace.clone(),
|
||||
)
|
||||
},
|
||||
))
|
||||
.collect();
|
||||
return Ok(Some(Packet::StreamStart(attrs)));
|
||||
} else if self.stanza_builder.depth() == 1 {
|
||||
self.driver.release_temporaries();
|
||||
|
||||
if let Some(stanza) = self.stanza_builder.unshift_child() {
|
||||
return Ok(Some(Packet::Stanza(stanza)));
|
||||
}
|
||||
}
|
||||
// Remedies for truncated utf8
|
||||
Err(e) if e.valid_up_to() >= buf1.len() - 3 => {
|
||||
// Prepare all the valid data
|
||||
let mut b = BytesMut::with_capacity(e.valid_up_to());
|
||||
b.put(&buf1[0..e.valid_up_to()]);
|
||||
} else if let Some(_) = self.stanza_builder.root.take() {
|
||||
self.driver.release_temporaries();
|
||||
|
||||
// Retry
|
||||
let result = self.decode(&mut b);
|
||||
|
||||
// Keep the tail back in
|
||||
self.buf.extend_from_slice(&buf1[e.valid_up_to()..]);
|
||||
|
||||
return result;
|
||||
}
|
||||
Err(e) => {
|
||||
error!(
|
||||
"error {} at {}/{} in {:?}",
|
||||
e,
|
||||
e.valid_up_to(),
|
||||
buf1.len(),
|
||||
buf1
|
||||
);
|
||||
return Err(ParserError::Utf8(e));
|
||||
return Ok(Some(Packet::StreamEnd));
|
||||
}
|
||||
}
|
||||
|
||||
match self.queue.lock().unwrap().pop_front() {
|
||||
None => Ok(None),
|
||||
Some(result) => result.map(|pkt| Some(pkt)),
|
||||
}
|
||||
Ok(None)
|
||||
}
|
||||
|
||||
fn decode_eof(&mut self, buf: &mut BytesMut) -> Result<Option<Self::Item>, Self::Error> {
|
||||
|
|
@ -392,7 +230,6 @@ mod tests {
|
|||
Ok(Some(Packet::StreamStart(_))) => true,
|
||||
_ => false,
|
||||
});
|
||||
b.clear();
|
||||
b.put_slice(b"</stream:stream>");
|
||||
let r = c.decode(&mut b);
|
||||
assert!(match r {
|
||||
|
|
@ -412,7 +249,6 @@ mod tests {
|
|||
_ => false,
|
||||
});
|
||||
|
||||
b.clear();
|
||||
b.put_slice("<test>ß</test".as_bytes());
|
||||
let r = c.decode(&mut b);
|
||||
assert!(match r {
|
||||
|
|
@ -420,7 +256,6 @@ mod tests {
|
|||
_ => false,
|
||||
});
|
||||
|
||||
b.clear();
|
||||
b.put_slice(b">");
|
||||
let r = c.decode(&mut b);
|
||||
assert!(match r {
|
||||
|
|
@ -440,7 +275,6 @@ mod tests {
|
|||
_ => false,
|
||||
});
|
||||
|
||||
b.clear();
|
||||
b.put(&b"<test>\xc3"[..]);
|
||||
let r = c.decode(&mut b);
|
||||
assert!(match r {
|
||||
|
|
@ -448,7 +282,6 @@ mod tests {
|
|||
_ => false,
|
||||
});
|
||||
|
||||
b.clear();
|
||||
b.put(&b"\x9f</test>"[..]);
|
||||
let r = c.decode(&mut b);
|
||||
assert!(match r {
|
||||
|
|
@ -469,7 +302,6 @@ mod tests {
|
|||
_ => false,
|
||||
});
|
||||
|
||||
b.clear();
|
||||
b.put_slice(b"<status xml:lang='en'>Test status</status>");
|
||||
let r = c.decode(&mut b);
|
||||
assert!(match r {
|
||||
|
|
@ -503,8 +335,11 @@ mod tests {
|
|||
block_on(framed.send(Packet::Stanza(stanza))).expect("send");
|
||||
assert_eq!(
|
||||
framed.get_ref().get_ref(),
|
||||
&("<message xmlns=\"jabber:client\"><body>".to_owned() + &text + "</body></message>")
|
||||
.as_bytes()
|
||||
&format!(
|
||||
"<message xmlns=\"jabber:client\"><body>{}</body></message>",
|
||||
text
|
||||
)
|
||||
.as_bytes()
|
||||
);
|
||||
}
|
||||
|
||||
|
|
@ -519,7 +354,6 @@ mod tests {
|
|||
_ => false,
|
||||
});
|
||||
|
||||
b.clear();
|
||||
b.put_slice(b"<message ");
|
||||
b.put_slice(b"type='chat'><body>Foo</body></message>");
|
||||
let r = c.decode(&mut b);
|
||||
|
|
|
|||
|
|
@ -54,7 +54,7 @@ impl<S: AsyncRead + AsyncWrite + Unpin> XMPPStream<S> {
|
|||
}
|
||||
|
||||
/// Send a `<stream:stream>` start tag
|
||||
pub async fn start<'a>(stream: S, jid: Jid, ns: String) -> Result<Self, Error> {
|
||||
pub async fn start(stream: S, jid: Jid, ns: String) -> Result<Self, Error> {
|
||||
let xmpp_stream = Framed::new(stream, XMPPCodec::new());
|
||||
stream_start::start(xmpp_stream, jid, ns).await
|
||||
}
|
||||
|
|
@ -65,7 +65,7 @@ impl<S: AsyncRead + AsyncWrite + Unpin> XMPPStream<S> {
|
|||
}
|
||||
|
||||
/// Re-run `start()`
|
||||
pub async fn restart<'a>(self) -> Result<Self, Error> {
|
||||
pub async fn restart(self) -> Result<Self, Error> {
|
||||
let stream = self.stream.into_inner().unwrap().into_inner();
|
||||
Self::start(stream, self.jid, self.ns).await
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in a new issue