mod mock;
use faktory::*;
use std::{io, sync::Arc, time::Duration};
use tokio::io::BufStream;
use tokio::{spawn, sync::Mutex, time::sleep};
use tokio_util::sync::CancellationToken;
#[tokio::test(flavor = "multi_thread")]
async fn hello() {
let mut s = mock::Stream::default();
let w: Worker<io::Error> = WorkerBuilder::default()
.hostname("host".to_string())
.wid(WorkerId::new("wid"))
.labels([
"will".to_string(),
"be!".to_string(),
"overwritten".to_string(),
])
.labels(["foo".to_string(), "bar".to_string()])
.add_to_labels(["will".to_string()])
.add_to_labels(["be".to_string(), "added".to_string()])
.register_fn("never_called", |_j: Job| async move { unreachable!() })
.connect_with(BufStream::new(s.clone()), None)
.await
.unwrap();
let written = s.pop_bytes_written(0);
assert!(written.starts_with(b"HELLO {"));
let written: serde_json::Value = serde_json::from_slice(&written[b"HELLO ".len()..]).unwrap();
let written = written.as_object().unwrap();
assert_eq!(
written.get("hostname").and_then(|h| h.as_str()),
Some("host")
);
assert_eq!(written.get("wid").and_then(|h| h.as_str()), Some("wid"));
assert_eq!(written.get("pid").map(|h| h.is_number()), Some(true));
assert_eq!(written.get("v").and_then(|h| h.as_i64()), Some(2));
let labels = written["labels"].as_array().unwrap();
assert_eq!(labels, &["foo", "bar", "will", "be", "added"]);
drop(w);
let written = s.pop_bytes_written(0);
assert_eq!(written, b"END\r\n");
}
#[tokio::test(flavor = "multi_thread")]
async fn hello_pwd() {
let mut s = mock::Stream::with_salt(1545, "55104dc76695721d");
let w: Worker<io::Error> = WorkerBuilder::default()
.register_fn("never_called", |_j: Job| async move { unreachable!() })
.connect_with(BufStream::new(s.clone()), Some("foobar".to_string()))
.await
.unwrap();
let written = s.pop_bytes_written(0);
assert!(written.starts_with(b"HELLO {"));
let written: serde_json::Value = serde_json::from_slice(&written[b"HELLO ".len()..]).unwrap();
let written = written.as_object().unwrap();
assert_eq!(
written.get("pwdhash").and_then(|h| h.as_str()),
Some("6d877f8e5544b1f2598768f817413ab8a357afffa924dedae99eb91472d4ec30")
);
drop(w);
}
#[tokio::test(flavor = "multi_thread")]
async fn dequeue() {
let mut s = mock::Stream::default();
let mut w = WorkerBuilder::default()
.register_fn("foobar", |job: Job| async move {
assert_eq!(job.args(), &["z"]);
Ok::<(), io::Error>(())
})
.connect_with(BufStream::new(s.clone()), None)
.await
.unwrap();
s.ignore(0);
s.push_bytes_to_read(
0,
b"$188\r\n\
{\
\"jid\":\"foojid\",\
\"queue\":\"default\",\
\"jobtype\":\"foobar\",\
\"args\":[\"z\"],\
\"created_at\":\"2017-11-01T21:02:35.772981326Z\",\
\"enqueued_at\":\"2017-11-01T21:02:35.773318394Z\",\
\"reserve_for\":600,\
\"retry\":25\
}\r\n",
);
s.ok(0); if let Err(e) = w.run_one(0, &["default"]).await {
println!("{:?}", e);
unreachable!();
}
let written = s.pop_bytes_written(0);
assert_eq!(
written,
&b"FETCH default\r\n\
ACK {\"jid\":\"foojid\"}\r\n"[..]
);
}
#[tokio::test(flavor = "multi_thread")]
async fn dequeue_first_empty() {
let mut s = mock::Stream::default();
let mut w = WorkerBuilder::default()
.register_fn("foobar", |job: Job| async move {
assert_eq!(job.args(), &["z"]);
Ok::<(), io::Error>(())
})
.connect_with(BufStream::new(s.clone()), None)
.await
.unwrap();
s.ignore(0);
s.push_bytes_to_read(
0,
b"$0\r\n\r\n$188\r\n\
{\
\"jid\":\"foojid\",\
\"queue\":\"default\",\
\"jobtype\":\"foobar\",\
\"args\":[\"z\"],\
\"created_at\":\"2017-11-01T21:02:35.772981326Z\",\
\"enqueued_at\":\"2017-11-01T21:02:35.773318394Z\",\
\"reserve_for\":600,\
\"retry\":25\
}\r\n",
);
s.ok(0);
match w.run_one(0, &["default"]).await {
Ok(did_work) => assert!(!did_work),
Err(e) => {
println!("{:?}", e);
unreachable!();
}
}
match w.run_one(0, &["default"]).await {
Ok(did_work) => assert!(did_work),
Err(e) => {
println!("{:?}", e);
unreachable!();
}
}
let written = s.pop_bytes_written(0);
assert_eq!(
written,
&b"\
FETCH default\r\n\
FETCH default\r\n\
ACK {\"jid\":\"foojid\"}\r\n\
"[..]
);
}
#[tokio::test(flavor = "multi_thread")]
async fn well_behaved() {
let mut s = mock::Stream::new(2); let mut w = WorkerBuilder::default()
.wid(WorkerId::new("wid"))
.register_fn("foobar", |_| async move {
sleep(Duration::from_secs(7)).await;
Ok::<(), io::Error>(())
})
.connect_with(BufStream::new(s.clone()), None)
.await
.unwrap();
s.ignore(0);
s.push_bytes_to_read(
1,
b"$182\r\n\
{\
\"jid\":\"jid\",\
\"queue\":\"default\",\
\"jobtype\":\"foobar\",\
\"args\":[],\
\"created_at\":\"2017-11-01T21:02:35.772981326Z\",\
\"enqueued_at\":\"2017-11-01T21:02:35.773318394Z\",\
\"reserve_for\":600,\
\"retry\":25\
}\r\n",
);
let jh = spawn(async move { w.run(&["default"]).await });
s.push_bytes_to_read(0, b"+{\"state\":\"quiet\"}\r\n");
s.ok(1);
s.push_bytes_to_read(0, b"+{\"state\":\"terminate\"}\r\n");
let details = jh.await.unwrap().unwrap();
assert_eq!(details.reason, StopReason::ServerInstruction);
assert_eq!(details.workers_still_running, 0);
let written = s.pop_bytes_written(0);
let msgs = "\
BEAT {\"wid\":\"wid\"}\r\n\
BEAT {\"wid\":\"wid\"}\r\n\
END\r\n";
assert_eq!(std::str::from_utf8(&written[..]).unwrap(), msgs);
let written = s.pop_bytes_written(1);
let msgs = "\r\n\
FETCH default\r\n\
ACK {\"jid\":\"jid\"}\r\n\
END\r\n";
assert_eq!(
std::str::from_utf8(&written[(written.len() - msgs.len())..]).unwrap(),
msgs
);
}
#[tokio::test(flavor = "multi_thread")]
async fn no_first_job() {
let mut s = mock::Stream::new(2); let mut w = WorkerBuilder::default()
.wid(WorkerId::new("wid"))
.register_fn("foobar", |_| async move {
sleep(Duration::from_secs(7)).await;
Ok::<(), io::Error>(())
})
.connect_with(BufStream::new(s.clone()), None)
.await
.unwrap();
s.ignore(0);
s.push_bytes_to_read(
1,
b"$0\r\n\r\n$182\r\n\
{\
\"jid\":\"jid\",\
\"queue\":\"default\",\
\"jobtype\":\"foobar\",\
\"args\":[],\
\"created_at\":\"2017-11-01T21:02:35.772981326Z\",\
\"enqueued_at\":\"2017-11-01T21:02:35.773318394Z\",\
\"reserve_for\":600,\
\"retry\":25\
}\r\n",
);
let jh = spawn(async move { w.run(&["default"]).await });
s.push_bytes_to_read(0, b"+{\"state\":\"quiet\"}\r\n");
s.ok(1);
s.push_bytes_to_read(0, b"+{\"state\":\"terminate\"}\r\n");
let details = jh.await.unwrap().unwrap();
assert_eq!(details.reason, StopReason::ServerInstruction);
assert_eq!(details.workers_still_running, 0);
let written = s.pop_bytes_written(0);
let msgs = "\
BEAT {\"wid\":\"wid\"}\r\n\
BEAT {\"wid\":\"wid\"}\r\n\
END\r\n";
assert_eq!(std::str::from_utf8(&written[..]).unwrap(), msgs);
let written = s.pop_bytes_written(1);
let msgs = "\r\n\
FETCH default\r\n\
FETCH default\r\n\
ACK {\"jid\":\"jid\"}\r\n\
END\r\n";
assert_eq!(
std::str::from_utf8(&written[(written.len() - msgs.len())..]).unwrap(),
msgs
);
}
#[tokio::test(flavor = "multi_thread")]
async fn well_behaved_many() {
let mut s = mock::Stream::new(3); let mut w = WorkerBuilder::default()
.workers(2)
.wid(WorkerId::new("wid"))
.register_fn("foobar", |_| async move {
sleep(Duration::from_secs(7)).await;
Ok::<(), io::Error>(())
})
.connect_with(BufStream::new(s.clone()), None)
.await
.unwrap();
s.ignore(0);
for i in 0..2 {
s.push_bytes_to_read(
i + 1,
b"$182\r\n\
{\
\"jid\":\"jid\",\
\"queue\":\"default\",\
\"jobtype\":\"foobar\",\
\"args\":[],\
\"created_at\":\"2017-11-01T21:02:35.772981326Z\",\
\"enqueued_at\":\"2017-11-01T21:02:35.773318394Z\",\
\"reserve_for\":600,\
\"retry\":25\
}\r\n",
);
}
let jh = spawn(async move { w.run(&["default"]).await });
s.push_bytes_to_read(0, b"+{\"state\":\"quiet\"}\r\n");
s.ok(1);
s.ok(2);
s.push_bytes_to_read(0, b"+{\"state\":\"terminate\"}\r\n");
let details = jh.await.unwrap().unwrap();
assert_eq!(details.reason, StopReason::ServerInstruction);
assert_eq!(details.workers_still_running, 0);
let written = s.pop_bytes_written(0);
let msgs = "\
BEAT {\"wid\":\"wid\"}\r\n\
BEAT {\"wid\":\"wid\"}\r\n\
END\r\n";
assert_eq!(std::str::from_utf8(&written[..]).unwrap(), msgs);
for i in 0..2 {
let written = s.pop_bytes_written(i + 1);
let msgs = "\r\n\
FETCH default\r\n\
ACK {\"jid\":\"jid\"}\r\n\
END\r\n";
assert_eq!(
std::str::from_utf8(&written[(written.len() - msgs.len())..]).unwrap(),
msgs
);
}
}
#[tokio::test(flavor = "multi_thread")]
async fn terminate() {
let mut s = mock::Stream::new(2);
let mut w: Worker<io::Error> = WorkerBuilder::default()
.hostname("machine".into())
.wid(WorkerId::new("wid"))
.register_fn("foobar", |_| async move {
loop {
sleep(Duration::from_secs(5)).await;
}
})
.connect_with(BufStream::new(s.clone()), None)
.await
.unwrap();
s.ignore(0);
s.push_bytes_to_read(
1,
b"$186\r\n\
{\
\"jid\":\"forever\",\
\"queue\":\"default\",\
\"jobtype\":\"foobar\",\
\"args\":[],\
\"created_at\":\"2017-11-01T21:02:35.772981326Z\",\
\"enqueued_at\":\"2017-11-01T21:02:35.773318394Z\",\
\"reserve_for\":600,\
\"retry\":25\
}\r\n",
);
let jh = spawn(async move {
w.run(&["default"]).await
});
s.push_bytes_to_read(0, b"+{\"state\":\"terminate\"}\r\n");
let details = jh.await.unwrap().unwrap();
assert_eq!(details.reason, StopReason::ServerInstruction);
assert_eq!(details.workers_still_running, 1);
let written = s.pop_bytes_written(0);
let beat = b"BEAT {\"wid\":\"wid\"}\r\nFAIL ";
assert_eq!(&written[0..beat.len()], &beat[..]);
assert!(written.ends_with(b"\r\nEND\r\n"));
let written: serde_json::Value =
serde_json::from_slice(&written[beat.len()..(written.len() - b"\r\nEND\r\n".len())])
.unwrap();
assert_eq!(
written
.as_object()
.and_then(|o| o.get("jid"))
.and_then(|v| v.as_str()),
Some("forever")
);
assert_eq!(written.get("errtype").unwrap().as_str(), Some("unknown"));
assert_eq!(written.get("message").unwrap().as_str(), Some("terminated"));
sleep(Duration::from_millis(500)).await;
let written = s.pop_bytes_written(1);
assert!(written.starts_with(b"HELLO {\"hostname\":\"machine\",\"wid\":\"wid\""));
assert!(written.ends_with(b"\r\nFETCH default\r\nEND\r\n"));
}
#[tokio::test(flavor = "multi_thread")]
async fn heart_broken() {
let mut s = mock::Stream::new_unchecked(3);
let token = CancellationToken::new();
let child_token = token.child_token();
let signal = async move { child_token.cancelled().await };
let w: Worker<io::Error> = Worker::builder()
.with_graceful_shutdown(signal)
.shutdown_timeout(Duration::from_millis(500))
.register_fn("foobar", |_j| async move {
println!("{:?}", _j);
sleep(Duration::from_secs(7)).await;
Ok(())
})
.connect_with(BufStream::new(s.clone()), None)
.await
.unwrap();
s.push_bytes_to_read(
1,
b"$186\r\n\
{\
\"jid\":\"forever\",\
\"queue\":\"default\",\
\"jobtype\":\"foobar\",\
\"args\":[],\
\"created_at\":\"2024-07-18T17:41:35.772981326Z\",\
\"enqueued_at\":\"2024-07-18T17:41:35.773318394Z\",\
\"reserve_for\":600,\
\"retry\":25\
}\r\n",
);
s.ignore(0);
let w = Arc::new(Mutex::new(w));
let w_clone = w.clone();
let jh = spawn(async move { w_clone.lock().await.run(&["default"]).await });
s.ok(1);
s.push_bytes_to_read(0, b"+{\"state\":\"heartbroken response\"}\r\n");
let error = jh.await.expect("joined ok").unwrap_err();
match error {
Error::Protocol(error::Protocol::BadType { expected, received }) => {
assert_eq!(expected, "heartbeat response");
assert_eq!(received, "{\"state\":\"heartbroken response\"}")
}
e => unreachable!("{:?}", e),
}
s.pop_bytes_written(0);
s.push_bytes_to_read(
2,
b"$186\r\n\
{\
\"jid\":\"ornever\",\
\"queue\":\"default\",\
\"jobtype\":\"foobar\",\
\"args\":[],\
\"created_at\":\"2024-07-18T17:41:36.772981326Z\",\
\"enqueued_at\":\"2024-07-18T17:41:36.773318394Z\",\
\"reserve_for\":600,\
\"retry\":25\
}\r\n",
);
let jh = spawn(async move { w.lock().await.run(&["default"]).await });
tokio::time::sleep(Duration::from_secs(1)).await;
token.cancel();
let stop_details = jh
.await
.expect("joined ok")
.expect("stop details rather than error");
assert_eq!(stop_details.reason, StopReason::GracefulShutdown);
assert_eq!(stop_details.workers_still_running, 1);
}