use bytes::{Buf, Bytes};
use tracing::{debug, warn};
use super::attrs::{decode_getattr_response, standard_getattr_bitmap};
use super::callback::RecallNotification;
use super::compound::OpenArgs;
use super::mount::{Mount41, decode_fh, extract_open_delegation, extract_stateid};
use super::pnfs_io::PnfsWriteOutcome;
use super::setattr::encode_setattr;
use super::state::{AccessMode, StateId};
use crate::error::{NfsError, Result};
use crate::mount;
impl Mount41 {
pub(crate) async fn write(&self, fh: Bytes, offset: u64, data: Bytes) -> Result<u32> {
self.refresh_layout_for_write(&fh, offset).await?;
let _io_guard = self.layout_manager.read_file_io(&fh).await;
match self.pnfs_write(&fh, offset, data.clone()).await {
PnfsWriteOutcome::NotAttempted => {}
PnfsWriteOutcome::Attempted(result) => return result,
}
self.mds_write(&fh, offset, data).await
}
async fn mds_write(&self, fh: &Bytes, offset: u64, data: Bytes) -> Result<u32> {
let sid = self
.state
.has_open(fh, AccessMode::Write)
.await
.unwrap_or_else(StateId::anonymous);
let stateid = sid.raw;
let data_len = data.len() as u32;
let fh_ref = fh.clone();
let resp = self
.compound_write("write", data, |b| {
b.require_generation(sid.generation)
.putfh(&fh_ref)
.write_header(&stateid, offset, 2 , data_len)
})
.await?;
resp.op_ok(1)?; let write_op = resp.op_ok(2)?;
let mut d = write_op.data.clone();
if d.remaining() < 16 {
return Err(NfsError::Xdr("WRITE result too short".to_string()));
}
let count = d.get_u32();
let committed = d.get_u32();
d.advance(8); if committed != 2
{
self.commit(fh.clone(), offset, count).await?;
}
Ok(count)
}
pub(crate) async fn write_path(&self, path: &str, offset: u64, data: Bytes) -> Result<u32> {
let obj = self.lookup_path(path).await?;
self.write(obj.fh, offset, data).await
}
pub(crate) async fn open(
&self,
dir_fh: Bytes,
filename: &str,
access: u32,
) -> Result<mount::ObjRes> {
let access_mode = match access {
crate::OPEN_READ => AccessMode::Read,
crate::OPEN_WRITE => AccessMode::Write,
_ => AccessMode::Both,
};
let bitmap = standard_getattr_bitmap();
let open_args = OpenArgs {
seqid: 0,
share_access: access_mode.share_access(),
share_deny: 0, client_id: self.session_holder.get().await.client_id(),
owner: Bytes::from_static(b"nfs-rs"),
create: false,
create_attrs_mask: vec![],
create_attrs_vals: vec![],
claim_file: filename.to_string(),
want_no_delegation: !self.retain_delegations,
};
let resp = self
.compound("open", |b| {
b.putfh(&dir_fh).open(&open_args).getfh().getattr(&bitmap)
})
.await?;
resp.op_ok(1)?; let open_op = resp.op_ok(2)?; let mut open_data = open_op.data.clone();
let stateid = extract_stateid(&mut open_data)?;
let delegation = extract_open_delegation(&mut open_data);
let getfh = resp.op_ok(3)?; let mut fh_data = getfh.data.clone();
let fh = decode_fh(&mut fh_data)?;
let getattr = resp.op_ok(4)?; let mut attr_data = getattr.data.clone();
let attr = decode_getattr_response(&mut attr_data)?;
if !self.retain_delegations {
self.return_unsolicited_delegation(delegation, &fh);
}
self.state
.register_open(
&fh,
StateId::from_bytes_at(&stateid, resp.session_generation),
access_mode,
)
.await?;
Ok(mount::ObjRes {
fh,
attr: Some(attr),
})
}
fn return_unsolicited_delegation(&self, delegation: Option<[u8; 16]>, fh: &Bytes) {
if let Some(deleg_sid) = delegation {
let notification = RecallNotification::Delegation {
stateid: deleg_sid,
truncate: false,
fh: fh.clone(),
};
if let Err(e) = self.recall_tx.try_send(notification) {
debug!(error = %e, "proactive DELEGRETURN enqueue failed, deferring to CB_RECALL");
}
}
}
pub(crate) async fn open_path(&self, path: &str, access: u32) -> Result<mount::ObjRes> {
let (dir, name) = crate::split_path(path)?;
let dir_obj = self.lookup_path(&dir).await?;
self.open(dir_obj.fh, &name, access).await
}
pub(crate) async fn close_file(&self, fh: Bytes) -> Result<()> {
let _io_guard = self.layout_manager.write_file_io(&fh).await;
self.flush_layoutcommit(&fh).await?;
if let Some(layout) = self.layout_manager.get_layout(&fh).await
&& layout.return_on_close
{
self.layoutreturn_file(&fh).await?;
}
if let Some(sid) = self.state.release(&fh).await {
let _ = self
.compound("close", |b| {
b.require_generation(sid.generation)
.putfh(&fh)
.close(0, &sid.raw)
})
.await;
}
Ok(())
}
pub(crate) async fn create(
&self,
dir_fh: Bytes,
filename: &str,
mode: Option<u32>,
) -> Result<mount::ObjRes> {
let bitmap = standard_getattr_bitmap();
let open_args = OpenArgs {
seqid: 0,
share_access: 0x00000003, share_deny: 0, client_id: self.session_holder.get().await.client_id(),
owner: Bytes::from_static(b"nfs-rs-create"),
create: true,
create_attrs_mask: vec![],
create_attrs_vals: vec![],
claim_file: filename.to_string(),
want_no_delegation: !self.retain_delegations,
};
let resp = self
.compound("create", |b| {
b.putfh(&dir_fh).open(&open_args).getfh().getattr(&bitmap)
})
.await?;
resp.op_ok(1)?; let open_op = resp.op_ok(2)?; let mut open_data = open_op.data.clone();
let stateid = extract_stateid(&mut open_data)?;
let delegation = extract_open_delegation(&mut open_data);
let getfh = resp.op_ok(3)?;
let mut fh_data = getfh.data.clone();
let fh = decode_fh(&mut fh_data)?;
if !self.retain_delegations {
self.return_unsolicited_delegation(delegation, &fh);
}
let getattr = resp.op_ok(4)?;
let mut attr_data = getattr.data.clone();
let mut attr = decode_getattr_response(&mut attr_data)?;
if let Some(m) = mode {
let (attrmask, attr_vals) = encode_setattr(Some(m), None, None, None, None, None);
match self
.compound("setattr", |b| {
b.putfh(&fh).setattr(&stateid, &attrmask, &attr_vals)
})
.await
{
Err(e) => {
warn!(error = %e, mode = m, "create: SETATTR mode failed after file creation");
}
_ => {
attr.file_mode = m;
}
}
}
self.state
.register_open(
&fh,
StateId::from_bytes_at(&stateid, resp.session_generation),
AccessMode::Both,
)
.await?;
Ok(mount::ObjRes {
fh,
attr: Some(attr),
})
}
pub(crate) async fn create_path(&self, path: &str, mode: Option<u32>) -> Result<mount::ObjRes> {
let (dir, name) = crate::split_path(path)?;
let dir_obj = self.lookup_path(&dir).await?;
self.create(dir_obj.fh, &name, mode).await
}
pub(crate) async fn mkdir(
&self,
dir_fh: Bytes,
dirname: &str,
mode: u32,
) -> Result<mount::ObjRes> {
let bitmap = standard_getattr_bitmap();
let resp = self
.compound("mkdir", |b| {
b.putfh(&dir_fh)
.create(2 , dirname, &[], &[])
.getfh()
.getattr(&bitmap)
})
.await?;
resp.op_ok(1)?; resp.op_ok(2)?; if let Some(create_op) = resp.results.get(2) {
let mut cdata = create_op.data.clone();
if cdata.remaining() >= 20 {
let atomic = cdata.get_u32() != 0;
let before = cdata.get_u64();
let after = cdata.get_u64();
debug!(atomic, before, after, name = dirname, "CREATE change_info");
}
}
let getfh = resp.op_ok(3)?;
let mut fh_data = getfh.data.clone();
let fh = decode_fh(&mut fh_data)?;
let getattr = resp.op_ok(4)?;
let mut attr_data = getattr.data.clone();
let mut attr = decode_getattr_response(&mut attr_data)?;
if let Err(e) = self
.setattr(fh.clone(), None, Some(mode), None, None, None, None, None)
.await
{
warn!(error = %e, mode, "mkdir: SETATTR mode failed after dir creation");
} else {
attr.file_mode = mode;
}
Ok(mount::ObjRes {
fh,
attr: Some(attr),
})
}
pub(crate) async fn mkdir_path(&self, path: &str, mode: u32) -> Result<mount::ObjRes> {
let (dir, name) = crate::split_path(path)?;
let dir_obj = self.lookup_path(&dir).await?;
self.mkdir(dir_obj.fh, &name, mode).await
}
}