use std::io::BufReader;
use std::os::unix::net::UnixStream;
use super::*;
use crate::chat::{UiIoMeter, UiWriter};
#[test]
fn paste_upload_version_gate() {
assert!(!supported(None));
assert!(!supported(Some(tau_proto::ProtocolVersion::new(8, 0))));
assert!(supported(Some(tau_proto::ProtocolVersion::new(8, 1))));
}
#[test]
fn paste_upload_correlates_retries_and_cancel_without_closing_ui() {
let (ui, peer) = UnixStream::pair().expect("socket pair");
peer.set_read_timeout(Some(Duration::from_secs(2)))
.expect("read timeout");
let writer = Arc::new(Mutex::new(UiWriter::new(ui, UiIoMeter::default())));
let (control, worker) = PasteUpload::spawn(writer.clone());
let mut reader = tau_proto::HarnessInputReader::new(BufReader::new(peer));
let (_term, handle, _input) = tau_cli_term_raw::Term::new_virtual(
80,
24,
"> ",
Box::new(std::io::sink()),
tau_cli_term_raw::CursorShape::Bar,
);
let session = "s1".parse().expect("session");
let text: Arc<str> = "x".repeat(8192).into();
control.start(session, 1, text, handle);
let available = read_request(&mut reader);
assert!(matches!(available.op, ArtifactOp::Available));
control.deliver(ArtifactResult {
request_id: "unrelated".parse().expect("request id"),
result: Ok(ArtifactValue::Done),
});
assert_eq!(
*control.expected.lock().expect("correlation lock"),
Some(available.request_id.clone())
);
control.deliver(ArtifactResult {
request_id: available.request_id,
result: Ok(ArtifactValue::Done),
});
let begin = read_request(&mut reader);
assert!(matches!(begin.op, ArtifactOp::Begin { .. }));
let upload_id: tau_proto::ArtifactUploadId = "upload-1".parse().expect("upload id");
control.deliver(ArtifactResult {
request_id: begin.request_id,
result: Ok(ArtifactValue::Upload {
upload: upload_id.clone(),
}),
});
let write = read_request(&mut reader);
assert!(matches!(write.op, ArtifactOp::Write { offset: 0, .. }));
control.cancel(1);
let abort = read_request(&mut reader);
assert!(matches!(abort.op, ArtifactOp::Abort { upload } if upload == upload_id));
control.deliver(ArtifactResult {
request_id: write.request_id,
result: Ok(ArtifactValue::Written { next_offset: 8192 }),
});
send_frame(
&writer,
&HarnessInputMessage::GetCurrentSession(tau_proto::GetCurrentSession {
request_id: "still-live".to_owned(),
}),
)
.expect("ordinary UI request");
assert!(matches!(
reader.read_message().expect("ordinary UI frame"),
Some(HarnessInputMessage::GetCurrentSession(_))
));
control.stop();
worker.join().expect("worker exit");
}
#[test]
fn paste_upload_retry_preserves_identity_and_inserts_only_verified_reference() {
use crossterm::event::{KeyCode, KeyEvent, KeyModifiers};
use tau_cli_term_raw::{Event, RawEvent};
let (ui, peer) = UnixStream::pair().expect("socket pair");
peer.set_read_timeout(Some(Duration::from_secs(2)))
.expect("read timeout");
let writer = Arc::new(Mutex::new(UiWriter::new(ui, UiIoMeter::default())));
let (control, worker) = PasteUpload::spawn(writer);
let mut reader = tau_proto::HarnessInputReader::new(BufReader::new(peer));
let (term, handle, input) = tau_cli_term_raw::Term::new_virtual(
80,
24,
"> ",
Box::new(std::io::sink()),
tau_cli_term_raw::CursorShape::Bar,
);
handle.enable_paste_uploads(8192);
handle.set_buffer("draft: ".to_owned(), 7);
let text = "x".repeat(8192);
input
.send(RawEvent::Paste(text.clone()))
.expect("paste input");
let Event::PasteUpload { id, text: source } = term.get_next_event().expect("upload event")
else {
panic!("paste")
};
control.start("s1".parse().expect("session"), id, source, handle.clone());
reply(&control, read_request(&mut reader), Ok(ArtifactValue::Done));
reply(
&control,
read_request(&mut reader),
Ok(ArtifactValue::Upload {
upload: "retry-upload".parse().expect("upload id"),
}),
);
let first_write = read_request(&mut reader);
reply(
&control,
first_write.clone(),
Err(tau_proto::ArtifactError::Busy),
);
assert!(matches!(
term.get_next_event().expect("failure notice"),
Event::Notice(_)
));
assert_eq!(handle.get_buffer(), "draft: ");
input
.send(RawEvent::Key(KeyEvent::new(
KeyCode::Enter,
KeyModifiers::NONE,
)))
.expect("retry key");
let Event::PasteUpload { id, text: source } = term.get_next_event().expect("retry event")
else {
panic!("retry")
};
control.start("s1".parse().expect("session"), id, source, handle.clone());
let retry = read_request(&mut reader);
assert_eq!(first_write.op, retry.op);
assert_ne!(first_write.request_id, retry.request_id);
reply(
&control,
retry,
Ok(ArtifactValue::Written {
next_offset: text.len() as u64,
}),
);
let finalize = read_request(&mut reader);
assert!(matches!(finalize.op, ArtifactOp::Finalize { .. }));
let key = format!("blake3:{}", blake3::hash(text.as_bytes()).to_hex())
.parse()
.expect("digest key");
let reference = tau_proto::artifact_reference(&key);
reply(
&control,
finalize,
Ok(ArtifactValue::Descriptor(tau_proto::ArtifactDescriptor {
key,
size: tau_proto::ArtifactSize::new(text.len() as u64).expect("object size"),
})),
);
assert!(matches!(
term.get_next_event().expect("completion event"),
Event::BufferChanged
));
assert_eq!(handle.get_buffer(), format!("draft: {reference}"));
control.stop();
worker.join().expect("worker exit");
}
fn reply(
control: &PasteUpload,
request: ArtifactRequest,
result: Result<ArtifactValue, tau_proto::ArtifactError>,
) {
control.deliver(ArtifactResult {
request_id: request.request_id,
result,
});
}
fn read_request(
reader: &mut tau_proto::HarnessInputReader<BufReader<UnixStream>>,
) -> ArtifactRequest {
let Some(HarnessInputMessage::ArtifactRequest(request)) =
reader.read_message().expect("artifact request")
else {
panic!("expected artifact request")
};
request
}