+ async source
(reader and converter)
This commit is contained in:
+216
-91
@@ -1,6 +1,4 @@
|
||||
use serde_derive::{Deserialize, Serialize};
|
||||
use std::io::{Read, Seek};
|
||||
use zip::ZipArchive;
|
||||
|
||||
#[derive(Debug, Default, Clone, Serialize, Deserialize)]
|
||||
pub struct SourceQuestion {
|
||||
@@ -105,102 +103,229 @@ pub struct SourceQuestionsBatch {
|
||||
pub questions: Vec<SourceQuestion>,
|
||||
}
|
||||
|
||||
pub struct SourceQuestionsZipReader<R>
|
||||
where
|
||||
R: Read + Seek,
|
||||
{
|
||||
zipfile: ZipArchive<R>,
|
||||
index: Option<usize>,
|
||||
}
|
||||
#[cfg(feature = "source")]
|
||||
pub mod reader_sync {
|
||||
use std::io::{Read, Seek};
|
||||
use zip::ZipArchive;
|
||||
|
||||
impl<R> SourceQuestionsZipReader<R>
|
||||
where
|
||||
R: Read + Seek,
|
||||
{
|
||||
fn new(zipfile: ZipArchive<R>) -> Self {
|
||||
SourceQuestionsZipReader {
|
||||
zipfile,
|
||||
index: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
use super::SourceQuestionsBatch;
|
||||
|
||||
impl<R> Iterator for SourceQuestionsZipReader<R>
|
||||
where
|
||||
R: Read + Seek,
|
||||
{
|
||||
type Item = (String, Result<SourceQuestionsBatch, serde_json::Error>);
|
||||
|
||||
fn next(&mut self) -> Option<Self::Item> {
|
||||
if self.index.is_none() && !self.zipfile.is_empty() {
|
||||
self.index = Some(0);
|
||||
}
|
||||
|
||||
match self.index {
|
||||
Some(i) if i < self.zipfile.len() => {
|
||||
self.index = Some(i + 1);
|
||||
|
||||
self.nth(i)
|
||||
}
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
|
||||
fn nth(&mut self, n: usize) -> Option<Self::Item> {
|
||||
if self.zipfile.len() <= n {
|
||||
return None;
|
||||
}
|
||||
self.index = Some(n + 1);
|
||||
|
||||
let file = self.zipfile.by_index(n).unwrap();
|
||||
let name = file.mangled_name();
|
||||
let name_str = name.to_str().unwrap();
|
||||
|
||||
let data: Result<SourceQuestionsBatch, _> = serde_json::from_reader(file);
|
||||
|
||||
Some((String::from(name_str), data))
|
||||
}
|
||||
|
||||
fn size_hint(&self) -> (usize, Option<usize>) {
|
||||
let len = self.zipfile.len();
|
||||
let index = self.index.unwrap_or(0);
|
||||
let rem = if len > index + 1 {
|
||||
len - (index + 1)
|
||||
} else {
|
||||
0
|
||||
};
|
||||
(rem, Some(rem))
|
||||
}
|
||||
|
||||
fn count(self) -> usize
|
||||
pub struct SourceQuestionsZipReader<R>
|
||||
where
|
||||
Self: Sized,
|
||||
R: Read + Seek,
|
||||
{
|
||||
self.zipfile.len()
|
||||
zipfile: ZipArchive<R>,
|
||||
index: Option<usize>,
|
||||
}
|
||||
|
||||
impl<R> SourceQuestionsZipReader<R>
|
||||
where
|
||||
R: Read + Seek,
|
||||
{
|
||||
fn new(zipfile: ZipArchive<R>) -> Self {
|
||||
SourceQuestionsZipReader {
|
||||
zipfile,
|
||||
index: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl<R> Iterator for SourceQuestionsZipReader<R>
|
||||
where
|
||||
R: Read + Seek,
|
||||
{
|
||||
type Item = (String, Result<SourceQuestionsBatch, serde_json::Error>);
|
||||
|
||||
fn next(&mut self) -> Option<Self::Item> {
|
||||
if self.index.is_none() && !self.zipfile.is_empty() {
|
||||
self.index = Some(0);
|
||||
}
|
||||
|
||||
match self.index {
|
||||
Some(i) if i < self.zipfile.len() => {
|
||||
self.index = Some(i + 1);
|
||||
|
||||
self.nth(i)
|
||||
}
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
|
||||
fn nth(&mut self, n: usize) -> Option<Self::Item> {
|
||||
if self.zipfile.len() <= n {
|
||||
return None;
|
||||
}
|
||||
self.index = Some(n + 1);
|
||||
|
||||
let file = self.zipfile.by_index(n).unwrap();
|
||||
let name = file.mangled_name();
|
||||
let name_str = name.to_str().unwrap();
|
||||
|
||||
let data: Result<SourceQuestionsBatch, _> = serde_json::from_reader(file);
|
||||
|
||||
Some((String::from(name_str), data))
|
||||
}
|
||||
|
||||
fn size_hint(&self) -> (usize, Option<usize>) {
|
||||
let len = self.zipfile.len();
|
||||
let index = self.index.unwrap_or(0);
|
||||
let rem = if len > index + 1 {
|
||||
len - (index + 1)
|
||||
} else {
|
||||
0
|
||||
};
|
||||
(rem, Some(rem))
|
||||
}
|
||||
|
||||
fn count(self) -> usize
|
||||
where
|
||||
Self: Sized,
|
||||
{
|
||||
self.zipfile.len()
|
||||
}
|
||||
}
|
||||
|
||||
impl<R> ExactSizeIterator for SourceQuestionsZipReader<R>
|
||||
where
|
||||
R: Read + Seek,
|
||||
{
|
||||
fn len(&self) -> usize {
|
||||
self.zipfile.len()
|
||||
}
|
||||
}
|
||||
|
||||
pub trait ReadSourceQuestionsBatches<R>
|
||||
where
|
||||
R: Read + Seek,
|
||||
{
|
||||
fn source_questions(self) -> SourceQuestionsZipReader<R>;
|
||||
}
|
||||
|
||||
impl<R> ReadSourceQuestionsBatches<R> for ZipArchive<R>
|
||||
where
|
||||
R: Read + Seek,
|
||||
{
|
||||
fn source_questions(self) -> SourceQuestionsZipReader<R> {
|
||||
SourceQuestionsZipReader::new(self)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl<R> ExactSizeIterator for SourceQuestionsZipReader<R>
|
||||
where
|
||||
R: Read + Seek,
|
||||
{
|
||||
fn len(&self) -> usize {
|
||||
self.zipfile.len()
|
||||
#[cfg(feature = "source")]
|
||||
pub use reader_sync::{ReadSourceQuestionsBatches, SourceQuestionsZipReader};
|
||||
|
||||
#[cfg(feature = "source_async")]
|
||||
pub mod reader_async {
|
||||
use crate::util::ErrorToString;
|
||||
|
||||
use async_zip::tokio::read::seek::ZipFileReader;
|
||||
use futures_core::stream::Stream;
|
||||
use futures_core::Future;
|
||||
use futures_util::{pin_mut, AsyncReadExt};
|
||||
|
||||
use std::pin::Pin;
|
||||
use std::task::{Context, Poll};
|
||||
|
||||
use tokio::io::{AsyncRead, AsyncSeek};
|
||||
|
||||
use super::SourceQuestionsBatch;
|
||||
|
||||
pub struct SourceQuestionsZipReaderAsync<R>
|
||||
where
|
||||
R: AsyncRead + AsyncSeek + Unpin,
|
||||
{
|
||||
zipfile: ZipFileReader<R>,
|
||||
index: Option<usize>,
|
||||
}
|
||||
}
|
||||
|
||||
pub trait ReadSourceQuestionsBatches<R>
|
||||
where
|
||||
R: Read + Seek,
|
||||
{
|
||||
fn source_questions(self) -> SourceQuestionsZipReader<R>;
|
||||
}
|
||||
|
||||
impl<R> ReadSourceQuestionsBatches<R> for ZipArchive<R>
|
||||
where
|
||||
R: Read + Seek,
|
||||
{
|
||||
fn source_questions(self) -> SourceQuestionsZipReader<R> {
|
||||
SourceQuestionsZipReader::new(self)
|
||||
|
||||
impl<R> SourceQuestionsZipReaderAsync<R>
|
||||
where
|
||||
R: AsyncRead + AsyncSeek + Unpin,
|
||||
{
|
||||
fn new(zipfile: ZipFileReader<R>) -> Self {
|
||||
SourceQuestionsZipReaderAsync {
|
||||
zipfile,
|
||||
index: None,
|
||||
}
|
||||
}
|
||||
async fn parse_zip_entry(
|
||||
&mut self,
|
||||
) -> Result<(String, Result<SourceQuestionsBatch, serde_json::Error>), String>
|
||||
where
|
||||
R: AsyncRead + AsyncSeek + Unpin,
|
||||
{
|
||||
let mut reader = self
|
||||
.zipfile
|
||||
.reader_with_entry(self.index.unwrap())
|
||||
.await
|
||||
.str_err()?;
|
||||
let filename = reader.entry().filename().clone().into_string().str_err()?;
|
||||
let mut data: Vec<u8> = Vec::new();
|
||||
reader.read_to_end(&mut data).await.str_err()?;
|
||||
let parsed: Result<SourceQuestionsBatch, _> = serde_json::from_slice(&data);
|
||||
self.index = Some(self.index.unwrap() + 1);
|
||||
|
||||
Ok((filename, parsed))
|
||||
}
|
||||
}
|
||||
|
||||
impl<R> Stream for SourceQuestionsZipReaderAsync<R>
|
||||
where
|
||||
R: AsyncRead + AsyncSeek + Unpin,
|
||||
{
|
||||
type Item = (String, Result<SourceQuestionsBatch, serde_json::Error>);
|
||||
|
||||
fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
|
||||
let len = self.zipfile.file().entries().len();
|
||||
if self.index.is_none() && len > 0 {
|
||||
self.index = Some(0);
|
||||
}
|
||||
|
||||
let index = &mut self.index;
|
||||
if index.unwrap() == len {
|
||||
return Poll::Ready(None);
|
||||
}
|
||||
|
||||
let future = self.parse_zip_entry();
|
||||
pin_mut!(future);
|
||||
match Pin::new(&mut future).poll(cx) {
|
||||
Poll::Ready(Ok(item)) => Poll::Ready(Some(item)),
|
||||
Poll::Ready(Err(_)) => Poll::Ready(None),
|
||||
Poll::Pending => Poll::Pending,
|
||||
}
|
||||
}
|
||||
|
||||
fn size_hint(&self) -> (usize, Option<usize>) {
|
||||
let len = self.zipfile.file().entries().len();
|
||||
if self.index.is_none() {
|
||||
return (len, Some(len));
|
||||
}
|
||||
|
||||
let index = self.index.unwrap();
|
||||
let rem = if len > index + 1 {
|
||||
len - (index + 1)
|
||||
} else {
|
||||
0
|
||||
};
|
||||
(rem, Some(rem))
|
||||
}
|
||||
}
|
||||
|
||||
pub trait ReadSourceQuestionsBatchesAsync<R>
|
||||
where
|
||||
R: AsyncRead + AsyncSeek + Unpin,
|
||||
{
|
||||
fn source_questions(self) -> SourceQuestionsZipReaderAsync<R>;
|
||||
}
|
||||
|
||||
impl<R> ReadSourceQuestionsBatchesAsync<R> for ZipFileReader<R>
|
||||
where
|
||||
R: AsyncRead + AsyncSeek + Unpin,
|
||||
{
|
||||
fn source_questions(self) -> SourceQuestionsZipReaderAsync<R> {
|
||||
SourceQuestionsZipReaderAsync::new(self)
|
||||
}
|
||||
}
|
||||
}
|
||||
#[cfg(feature = "source_async")]
|
||||
pub use reader_async::{ReadSourceQuestionsBatchesAsync, SourceQuestionsZipReaderAsync};
|
||||
|
||||
Reference in New Issue
Block a user