use crate::envelope;
use crate::error::{Error, Result};
use crate::session::ReadWrite;
use crate::session::Session;
use crate::transport::Transport;
use crate::wire::{
cmd, ui, Bank, Dependency, Location, Message, ObjectClass, Partition, ProgramInfo, Service,
Status,
};
pub async fn status<T: Transport, C>(session: &mut Session<'_, T, C>) -> Result<Status> {
let class = session.class();
let resp = session
.request(
Service::Program,
10,
cmd::STATUS,
&class.to_raw().to_be_bytes(),
)
.await?;
Status::decode(class, &resp)
}
pub async fn inventory<T: Transport>(transport: &mut T) -> Result<Vec<Status>> {
let mut out = Vec::new();
for class in ObjectClass::INVENTORY {
let mut session = match Session::open(transport, class).await {
Ok(s) => s,
Err(_) => continue,
};
match status(&mut session).await {
Ok(s) => {
session.commit().await?;
out.push(s);
}
Err(_) => {
session.commit().await?;
}
}
}
Ok(out)
}
pub async fn info<T: Transport, C>(
session: &mut Session<'_, T, C>,
at: Location,
) -> Result<ProgramInfo> {
let mut args = Vec::new();
at.write_to(&mut args);
let resp = session
.request(Service::Program, 10, cmd::INFO, &args)
.await?;
let info = ProgramInfo::decode(&resp)?;
if info.location != at {
return Err(Error::UnexpectedLocation {
requested: at,
reported: info.location,
});
}
Ok(info)
}
pub async fn read_program<T: Transport, C>(
session: &mut Session<'_, T, C>,
at: Location,
) -> Result<Vec<u8>> {
let (meta, body) = transfer_out(session, at).await?;
let file = envelope::wrap(&meta.format, at, meta.version, &body)?;
if let Some(expected) = meta.crc32 {
let actual = envelope::crc32(&body);
if expected != actual {
return Err(Error::Envelope(format!(
"body checksum mismatch: device reported {expected:08x}, received {actual:08x}"
)));
}
}
Ok(file)
}
pub async fn read_body<T: Transport, C>(
session: &mut Session<'_, T, C>,
at: Location,
) -> Result<Vec<u8>> {
Ok(transfer_out(session, at).await?.1)
}
const READ_CHUNK: u32 = 32720;
const WRITE_CHUNK: usize = 32720;
#[cfg(any(feature = "fault-injection", test))]
fn parse_chunk(name: &str, value: Option<&str>, default: u64) -> Result<u64> {
let Some(value) = value else {
return Ok(default);
};
value
.parse()
.ok()
.filter(|&size| size > 0)
.ok_or_else(|| Error::InvalidArgument(format!("{name} must be a positive integer")))
}
#[cfg(feature = "fault-injection")]
fn chunk_override(name: &str, default: u64) -> Result<u64> {
match std::env::var(name) {
Ok(value) => parse_chunk(name, Some(&value), default),
Err(std::env::VarError::NotPresent) => Ok(default),
Err(std::env::VarError::NotUnicode(_)) => {
Err(Error::InvalidArgument(format!("{name} must be UTF-8")))
}
}
}
#[cfg(feature = "fault-injection")]
fn read_chunk() -> Result<u32> {
let size = chunk_override("NORD_READ_CHUNK", READ_CHUNK.into())?;
u32::try_from(size).map_err(|_| Error::InvalidArgument("NORD_READ_CHUNK exceeds u32".into()))
}
#[cfg(not(feature = "fault-injection"))]
fn read_chunk() -> Result<u32> {
Ok(READ_CHUNK)
}
#[cfg(feature = "fault-injection")]
fn write_chunk() -> Result<usize> {
let size = chunk_override("NORD_WRITE_CHUNK", WRITE_CHUNK as u64)?;
usize::try_from(size)
.map_err(|_| Error::InvalidArgument("NORD_WRITE_CHUNK exceeds usize".into()))
}
#[cfg(not(feature = "fault-injection"))]
fn write_chunk() -> Result<usize> {
Ok(WRITE_CHUNK)
}
fn library_block(class: ObjectClass) -> usize {
match class {
ObjectClass::Sample => 131_072,
_ => 262_144,
}
}
async fn transfer_out<T: Transport, C>(
session: &mut Session<'_, T, C>,
at: Location,
) -> Result<(ProgramInfo, Vec<u8>)> {
let chunk_size = read_chunk()?;
let meta = info(session, at).await?;
session.notify(&ui::label("Uploading...")?).await?;
let mut args = Vec::new();
at.write_to(&mut args);
session
.request(Service::Program, 10, cmd::BEGIN_READ, &args)
.await?;
let mut body = Vec::with_capacity((meta.body_len as usize).min(1 << 20));
let mut painted = None;
while (body.len() as u32) < meta.body_len {
let offset = body.len() as u32;
let want = chunk_size.min(meta.body_len - offset);
let mut req = args.clone();
req.extend_from_slice(&offset.to_be_bytes());
req.extend_from_slice(&want.to_be_bytes());
let resp = session
.request(Service::Program, 10, cmd::READ, &req)
.await?;
let chunk = read_payload(resp.payload(), at, offset, want)?;
body.extend_from_slice(chunk);
let pct = (body.len() as u64 * 100 / (meta.body_len.max(1)) as u64) as u16;
if painted != Some(pct) {
session.notify(&ui::percent(pct)).await?;
painted = Some(pct);
}
}
if painted != Some(100) {
session.notify(&ui::percent(100)).await?;
}
session
.request(Service::Program, 10, cmd::END_TRANSFER, &args)
.await?;
Ok((meta, body))
}
fn read_payload(payload: &[u8], at: Location, offset: u32, length: u32) -> Result<&[u8]> {
if payload.len() < 16 {
return Err(Error::Truncated {
got: payload.len(),
need: 16,
});
}
let word = |start| u32::from_be_bytes(payload[start..start + 4].try_into().unwrap());
let echoed = (word(0), word(4), word(8), word(12));
let expected = (at.bank, at.slot, offset, length);
if echoed != expected {
return Err(Error::Transport(format!(
"READ response echoed {echoed:?}, expected {expected:?}"
)));
}
let body = &payload[16..];
if body.len() != length as usize {
return Err(Error::Transport(format!(
"asked for {length} bytes at offset {offset} but the device sent {}",
body.len()
)));
}
Ok(body)
}
const CLEANING_POLLS: u32 = 120;
const CLEANING_POLL_SPACING: std::time::Duration = std::time::Duration::from_millis(250);
async fn clean_library<T: Transport>(
session: &mut Session<'_, T, ReadWrite>,
blocks: u32,
) -> Result<()> {
session.notify(&ui::label("Cleaning...")?).await?;
session.notify(&ui::percent(0)).await?;
session
.request(
Service::Program,
10,
cmd::WRITE_PREPARE,
&blocks.to_be_bytes(),
)
.await?;
let mut painted = Some(0);
for polls in 0..CLEANING_POLLS {
if polls > 0 {
crate::sleep::sleep(CLEANING_POLL_SPACING).await;
}
let resp = session
.request(Service::Program, 10, cmd::WRITE_PREPARE_2, &[])
.await?;
let p = resp.payload();
if p.len() >= 12 {
let requested = u32::from_be_bytes(p[0..4].try_into().unwrap());
let done = u32::from_be_bytes(p[4..8].try_into().unwrap());
let running = u32::from_be_bytes(p[8..12].try_into().unwrap());
if running == 0 {
if painted != Some(100) {
session.notify(&ui::percent(100)).await?;
}
return Ok(());
}
let pct = (done as u64 * 100 / requested.max(1) as u64).min(99) as u16;
if painted != Some(pct) {
session.notify(&ui::percent(pct)).await?;
painted = Some(pct);
}
}
}
Err(Error::Transport(format!(
"the library's cleaning pass did not report ready within {} polls",
CLEANING_POLLS
)))
}
pub async fn write<T: Transport>(
session: &mut Session<'_, T, ReadWrite>,
at: Location,
file: &[u8],
name: &str,
timestamp: u32,
) -> Result<()> {
let file = envelope::unwrap(file)?;
let body = &file.body.0;
let chunk_size = write_chunk()?;
let body_len = u32::try_from(body.len())
.map_err(|_| Error::InvalidArgument("the body is larger than the wire format".into()))?;
let name_len = u32::try_from(name.len())
.map_err(|_| Error::InvalidArgument("the name is larger than the wire format".into()))?;
if matches!(session.class(), ObjectClass::Piano | ObjectClass::Sample) {
let needed = body.len().div_ceil(library_block(session.class())) as u32;
let free = status(session).await?.free;
if needed > free {
clean_library(session, needed - free).await?;
}
}
session.notify(&ui::label("Downloading...")?).await?;
let mut begin = Vec::new();
at.write_to(&mut begin);
begin.extend_from_slice(&body_len.to_be_bytes());
begin.extend_from_slice(&file.header.tag);
begin.extend_from_slice(×tamp.to_be_bytes());
begin.extend_from_slice(&u32::MAX.to_be_bytes());
begin.extend_from_slice(&name_len.to_be_bytes());
begin.extend_from_slice(name.as_bytes());
session
.request(Service::Program, 10, cmd::BEGIN_WRITE, &begin)
.await?;
let mut offset = 0usize;
let mut painted = None;
while offset < body.len() {
let end = offset.saturating_add(chunk_size).min(body.len());
let chunk = &body[offset..end];
let mut data = Vec::new();
at.write_to(&mut data);
data.extend_from_slice(&(offset as u32).to_be_bytes());
data.extend_from_slice(&(chunk.len() as u32).to_be_bytes());
data.extend_from_slice(chunk);
if end == body.len() {
session
.request(Service::Program, 10, cmd::WRITE_DATA, &data)
.await?;
} else {
let msg = Message::new(Service::Program, 10, cmd::WRITE_DATA, data);
session.notify(&msg).await?;
}
offset = end;
let pct = (offset as u64 * 100 / (body.len().max(1)) as u64) as u16;
if painted != Some(pct) {
session.notify(&ui::percent(pct)).await?;
painted = Some(pct);
}
}
if painted != Some(100) {
session.notify(&ui::percent(100)).await?;
}
let mut args = Vec::new();
at.write_to(&mut args);
session
.request(Service::Program, 10, cmd::END_TRANSFER, &args)
.await?;
Ok(())
}
pub async fn select<T: Transport, C>(session: &mut Session<'_, T, C>, at: Location) -> Result<()> {
let mut args = Vec::new();
at.write_to(&mut args);
session
.request(Service::Program, 10, cmd::SELECT, &args)
.await?;
Ok(())
}
async fn drain<T: Transport>(transport: &mut T) -> Result<()> {
for _ in 0..DRAIN_CAP {
match transport
.read_timeout(crate::transport::READ_BUFFER, DRAIN_LIMIT)
.await?
{
Some(_) => continue,
None => break,
}
}
Ok(())
}
const DRAIN_LIMIT: std::time::Duration = std::time::Duration::from_millis(300);
const DRAIN_CAP: usize = 16;
pub async fn recover<T: Transport>(transport: &mut T) -> Result<()> {
drain(transport).await?;
let goodbye = Message::new(Service::Ui, ui::SUBSYSTEM, ui::GOODBYE, Vec::new());
transport.write(&goodbye.encode()).await?;
let _ = transport
.read_timeout(crate::transport::READ_BUFFER, DRAIN_LIMIT)
.await?;
let close = Message::new(Service::Program, 10, cmd::SESSION_CLOSE, Vec::new());
transport.write(&close.encode()).await?;
let _ = transport
.read_timeout(crate::transport::READ_BUFFER, DRAIN_LIMIT)
.await?;
Ok(())
}
pub async fn partitions<T: Transport, C>(
session: &mut Session<'_, T, C>,
) -> Result<Vec<Partition>> {
let resp = session
.request(Service::Program, 10, cmd::PARTITIONS, &[])
.await?;
Partition::decode_all(&resp)
}
pub async fn banks<T: Transport, C>(
session: &mut Session<'_, T, C>,
partition: u32,
) -> Result<Vec<Bank>> {
let resp = session
.request(Service::Program, 10, cmd::BANKS, &partition.to_be_bytes())
.await?;
Bank::decode_all(&resp)
}
pub async fn check_address<T: Transport, C>(
session: &mut Session<'_, T, C>,
at: Location,
) -> Result<Option<String>> {
let banks = banks(session, session.class().to_raw()).await?;
let Some(bank) = banks.get(at.bank as usize) else {
let names: Vec<&str> = banks.iter().map(|b| b.name.as_str()).collect();
return Ok(Some(format!(
"bank {} does not exist; this class has {} ({})",
at.user_bank(),
banks.len(),
names.join(", ")
)));
};
if bank.is_bounded() && at.slot >= bank.slots {
return Ok(Some(format!(
"\"{}\" holds {} slots, so slot {} is out of range",
bank.name,
bank.slots,
at.user_slot()
)));
}
Ok(None)
}
pub async fn focus<T: Transport, C>(session: &mut Session<'_, T, C>) -> Result<Location> {
let resp = session
.request(Service::Program, 10, cmd::FOCUS, &[])
.await?;
let p = resp.payload();
if p.len() < 8 {
return Err(Error::Truncated {
got: p.len(),
need: 8,
});
}
Ok(Location {
bank: u32::from_be_bytes(p[0..4].try_into().unwrap()),
slot: u32::from_be_bytes(p[4..8].try_into().unwrap()),
})
}
pub const ENUMERATION_DISABLED: u32 = 0x11;
pub const SLOT_BOUNDARY: u32 = 0xffff_ffff;
pub async fn next_occupied<T: Transport, C>(
session: &mut Session<'_, T, C>,
at: Location,
) -> Result<Option<Location>> {
let mut args = Vec::new();
at.write_to(&mut args);
args.extend_from_slice(&0u32.to_be_bytes());
match session
.request(Service::Program, 10, cmd::NEXT_SLOT, &args)
.await
{
Ok(resp) => {
let p = resp.payload();
if p.len() < 8 {
return Err(Error::Truncated {
got: p.len(),
need: 8,
});
}
Ok(Some(Location {
bank: u32::from_be_bytes(p[0..4].try_into().unwrap()),
slot: u32::from_be_bytes(p[4..8].try_into().unwrap()),
}))
}
Err(Error::DeviceStatus(1)) => Ok(None),
Err(e) => Err(e),
}
}
pub async fn occupied_slots<T: Transport, C>(
session: &mut Session<'_, T, C>,
cap: usize,
) -> Result<Vec<Location>> {
let mut found: Vec<Location> = Vec::new();
let mut empty_banks = 0;
for bank in 0..MAX_BANKS {
if found.len() >= cap {
break;
}
match info(session, Location { bank, slot: 0 }).await {
Ok(_) | Err(Error::DeviceStatus(1)) => {}
Err(Error::DeviceStatus(3)) => break,
Err(e) => return Err(e),
}
let before = found.len();
let mut at = Location {
bank,
slot: SLOT_BOUNDARY,
};
while found.len() < cap {
match next_occupied(session, at).await? {
Some(next)
if next.bank == bank && (at.slot == SLOT_BOUNDARY || next.slot > at.slot) =>
{
found.push(next);
at = next;
}
_ => break,
}
}
if found.len() == before {
empty_banks += 1;
if empty_banks >= EMPTY_BANKS_BEFORE_STOP {
break;
}
} else {
empty_banks = 0;
}
}
Ok(found)
}
const MAX_BANKS: u32 = 64;
const EMPTY_BANKS_BEFORE_STOP: u32 = 2;
pub async fn required_dependencies<T: Transport, C>(
session: &mut Session<'_, T, C>,
at: Location,
) -> Result<Vec<Dependency>> {
Ok(dependencies(session, at)
.await?
.into_iter()
.filter(Dependency::is_required)
.collect())
}
pub async fn dependencies<T: Transport, C>(
session: &mut Session<'_, T, C>,
at: Location,
) -> Result<Vec<Dependency>> {
let mut args = Vec::new();
at.write_to(&mut args);
let resp = session
.request(Service::Program, 10, cmd::DEPENDENCIES, &args)
.await?;
Dependency::decode_all(&resp)
}
pub async fn move_object<T: Transport>(
session: &mut Session<'_, T, ReadWrite>,
from: Location,
to: Location,
) -> Result<()> {
let mut args = Vec::new();
from.write_to(&mut args);
to.write_to(&mut args);
session
.request(Service::Program, 10, cmd::MOVE, &args)
.await?;
Ok(())
}
pub async fn delete<T: Transport>(
session: &mut Session<'_, T, ReadWrite>,
at: Location,
) -> Result<()> {
session.notify(&ui::label("Deleting...")?).await?;
let mut args = Vec::new();
at.write_to(&mut args);
session
.request(Service::Program, 10, cmd::DELETE, &args)
.await?;
Ok(())
}
pub async fn rename<T: Transport>(
session: &mut Session<'_, T, ReadWrite>,
at: Location,
name: &str,
) -> Result<()> {
let mut args = Vec::new();
at.write_to(&mut args);
let name_len = u32::try_from(name.len())
.map_err(|_| Error::InvalidArgument("the name is larger than the wire format".into()))?;
args.extend_from_slice(&name_len.to_be_bytes());
args.extend_from_slice(name.as_bytes());
session
.request(Service::Program, 10, cmd::RENAME, &args)
.await?;
Ok(())
}
pub async fn duplicate<T: Transport>(
session: &mut Session<'_, T, ReadWrite>,
from: Location,
to: Location,
) -> Result<()> {
let mut args = Vec::new();
from.write_to(&mut args);
to.write_to(&mut args);
session
.request(Service::Program, 10, cmd::COPY, &args)
.await?;
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_read_chunk_must_echo_its_request() {
let at = Location { bank: 2, slot: 3 };
let mut payload = Vec::new();
for word in [at.bank, at.slot, 40, 3] {
payload.extend_from_slice(&word.to_be_bytes());
}
payload.extend_from_slice(&[1, 2, 3]);
assert_eq!(read_payload(&payload, at, 40, 3).unwrap(), [1, 2, 3]);
payload[3] ^= 1;
assert!(read_payload(&payload, at, 40, 3).is_err());
}
#[test]
fn invalid_chunk_overrides_are_refused() {
assert!(parse_chunk("NORD_READ_CHUNK", Some("0"), READ_CHUNK.into()).is_err());
assert!(parse_chunk("NORD_WRITE_CHUNK", Some("bad"), WRITE_CHUNK as u64).is_err());
assert_eq!(
parse_chunk("NORD_READ_CHUNK", None, READ_CHUNK.into()).unwrap(),
READ_CHUNK.into()
);
}
}