use crate::{PeerId, ResolveHandle};
use alloc::string::String;
use asimov_id::Handle;
use core::str::FromStr;
use csv_async::{AsyncReader, AsyncReaderBuilder, StringRecord};
use futures_lite::{Stream, StreamExt, pin, stream};
use iroh::EndpointId;
use std::io::{Error, Result};
use tokio::fs::File;
use tokio::io::AsyncReadExt;
#[cfg(not(feature = "std"))]
use alloc::collections::BTreeSet as Set;
#[cfg(feature = "std")]
use std::collections::HashSet as Set;
pub struct CsvHandleResolver(AsyncReader<File>);
impl CsvHandleResolver {
pub async fn open(path: &str) -> std::io::Result<Self> {
let file = File::open(path).await?;
let reader = AsyncReaderBuilder::new()
.has_headers(false)
.create_reader(file);
Ok(Self::from(reader))
}
pub fn handles(&mut self) -> impl Stream<Item = Result<Handle>> + Send {
async_stream::stream! {
let mut handles = Set::new();
let records = self.records();
pin!(records);
while let Some(record) = records.next().await {
let (handle, _) = record?;
if handles.contains(&handle) {
continue; }
handles.insert(handle.clone());
yield Ok(handle);
}
}
}
pub fn records(&mut self) -> impl Stream<Item = Result<(Handle, PeerId)>> + Send {
async_stream::stream! {
self.0.rewind().await?;
let mut record = StringRecord::new();
while let Ok(true) = self.0.read_record(&mut record).await {
let Some(record_handle) = record.get(0) else {
continue; };
let Some(record_endpoint) = record.get(1) else {
continue; };
let Ok(handle) = record_handle.parse::<Handle>() else {
continue; };
let Ok(endpoint) = record_endpoint.parse::<PeerId>() else {
continue; };
yield Ok((handle, endpoint));
}
}
}
}
impl From<AsyncReader<File>> for CsvHandleResolver {
fn from(reader: AsyncReader<File>) -> Self {
Self(reader)
}
}
impl ResolveHandle for CsvHandleResolver {
type Error = std::io::Error;
fn resolve_handle(
&mut self,
handle: impl Into<Handle>,
) -> impl Stream<Item = Result<PeerId>> + Send {
let handle = handle.into();
async_stream::stream! {
let mut endpoints = Set::new();
let records = self.records();
pin!(records);
while let Some(record) = records.next().await {
let (record_handle, record_endpoint) = record?;
if record_handle != handle {
continue; }
if endpoints.contains(&record_endpoint) {
continue; }
endpoints.insert(record_endpoint.clone());
yield Ok(record_endpoint);
}
}
}
}