-
Notifications
You must be signed in to change notification settings - Fork 0
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Merge pull request #22 from RDMA-Rust/dev/post-send
feat: introduce PostSendGuard for extended QP with basic support for polling extended CQ
- Loading branch information
Showing
11 changed files
with
800 additions
and
6 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,113 @@ | ||
use core::time; | ||
use std::thread; | ||
|
||
use sideway::verbs::{ | ||
address::AddressHandleAttribute, | ||
device, | ||
device_context::Mtu, | ||
queue_pair::{PostSendGuard, QueuePair, QueuePairAttribute, QueuePairState, SetInlineData, WorkRequestFlags}, | ||
AccessFlags, | ||
}; | ||
|
||
fn main() -> Result<(), Box<dyn std::error::Error>> { | ||
let device_list = device::DeviceList::new()?; | ||
for device in &device_list { | ||
let ctx = device.open().unwrap(); | ||
|
||
let pd = ctx.alloc_pd().unwrap(); | ||
let mr = pd.reg_managed_mr(64).unwrap(); | ||
|
||
let _comp_channel = ctx.create_comp_channel().unwrap(); | ||
let mut cq_builder = ctx.create_cq_builder(); | ||
let sq = cq_builder.setup_cqe(128).build_ex().unwrap(); | ||
let rq = cq_builder.setup_cqe(128).build_ex().unwrap(); | ||
|
||
let mut builder = pd.create_qp_builder(); | ||
|
||
let mut qp = builder | ||
.setup_max_inline_data(128) | ||
.setup_send_cq(&sq) | ||
.setup_recv_cq(&rq) | ||
.build_ex() | ||
.unwrap(); | ||
|
||
println!("qp pointer is {:?}", qp); | ||
// modify QP to INIT state | ||
let mut attr = QueuePairAttribute::new(); | ||
attr.setup_state(QueuePairState::Init) | ||
.setup_pkey_index(0) | ||
.setup_port(1) | ||
.setup_access_flags(AccessFlags::LocalWrite | AccessFlags::RemoteWrite); | ||
qp.modify(&attr).unwrap(); | ||
|
||
assert_eq!(QueuePairState::Init, qp.state()); | ||
|
||
// modify QP to RTR state, set dest qp as itself | ||
let mut attr = QueuePairAttribute::new(); | ||
attr.setup_state(QueuePairState::ReadyToReceive) | ||
.setup_path_mtu(Mtu::Mtu1024) | ||
.setup_dest_qp_num(qp.qp_number()) | ||
.setup_rq_psn(1) | ||
.setup_max_dest_read_atomic(0) | ||
.setup_min_rnr_timer(0); | ||
// setup address vector | ||
let mut ah_attr = AddressHandleAttribute::new(); | ||
let gid_entries = ctx.query_gid_table().unwrap(); | ||
|
||
ah_attr | ||
.setup_dest_lid(1) | ||
.setup_port(1) | ||
.setup_service_level(1) | ||
.setup_grh_src_gid_index(gid_entries[0].gid_index().try_into().unwrap()) | ||
.setup_grh_dest_gid(&gid_entries[0].gid()) | ||
.setup_grh_hop_limit(64); | ||
attr.setup_address_vector(&ah_attr); | ||
qp.modify(&attr).unwrap(); | ||
|
||
assert_eq!(QueuePairState::ReadyToReceive, qp.state()); | ||
|
||
// modify QP to RTS state | ||
let mut attr = QueuePairAttribute::new(); | ||
attr.setup_state(QueuePairState::ReadyToSend) | ||
.setup_sq_psn(1) | ||
.setup_timeout(12) | ||
.setup_retry_cnt(7) | ||
.setup_rnr_retry(7) | ||
.setup_max_read_atomic(0); | ||
|
||
qp.modify(&attr).unwrap(); | ||
|
||
assert_eq!(QueuePairState::ReadyToSend, qp.state()); | ||
|
||
let mut guard = qp.start_post_send(); | ||
let buf = vec![0, 1, 2, 3]; | ||
|
||
let write_handle = guard | ||
.construct_wr(233, WorkRequestFlags::Signaled) | ||
.setup_write(mr.rkey(), mr.buf.data.as_ptr() as _); | ||
|
||
write_handle.setup_inline_data(&buf); | ||
|
||
let _err = guard.post().unwrap(); | ||
|
||
thread::sleep(time::Duration::from_millis(10)); | ||
|
||
// poll for the completion | ||
{ | ||
let mut poller = sq.start_poll().unwrap(); | ||
let mut wc = poller.iter_mut(); | ||
println!("wr_id {}, status: {}, opcode: {}", wc.wr_id(), wc.status(), wc.opcode()); | ||
assert_eq!(wc.wr_id(), 233); | ||
while let Some(wc) = wc.next() { | ||
println!("wr_id {}, status: {}, opcode: {}", wc.wr_id(), wc.status(), wc.opcode()) | ||
} | ||
} | ||
|
||
unsafe { | ||
let slice = std::slice::from_raw_parts(mr.buf.data.as_ptr(), mr.buf.len); | ||
println!("Buffer contents: {:?}", slice); | ||
} | ||
} | ||
|
||
Ok(()) | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.