1use std::io::{Read, Write};
10
11use ytsaurus_format::DataFormat;
12use ytsaurus_skiff::Value;
13
14use crate::{
15 Event, JobError, JobReader, JobWriter, Result, SkiffJobReader, SkiffJobWriter, SkiffRow,
16 TableId,
17};
18
19#[derive(Debug)]
21pub enum WorkerReader<R> {
22 Yson(JobReader<R>),
24 Skiff(SkiffJobReader<R>),
26}
27
28#[derive(Debug)]
30pub enum WorkerEvent<'input> {
31 Yson(Event<'input>),
33 Skiff(SkiffRow),
35}
36
37impl WorkerReader<std::io::BufReader<std::io::Stdin>> {
38 pub fn from_stdin(format: DataFormat) -> Result<Self> {
50 Self::new(
51 std::io::BufReader::with_capacity(crate::skiff::STDIN_BUFFER_BYTES, std::io::stdin()),
52 format,
53 )
54 }
55}
56
57impl<R: Read> WorkerReader<R> {
58 pub fn new(input: R, format: DataFormat) -> Result<Self> {
65 match format {
66 DataFormat::Yson(format) => Ok(Self::Yson(JobReader::with_format(input, format))),
67 DataFormat::Skiff(format) => Ok(Self::Skiff(SkiffJobReader::new(input, format)?)),
68 _ => Err(JobError::UnsupportedDataFormat),
69 }
70 }
71
72 pub fn next_event(&mut self) -> Result<Option<WorkerEvent<'_>>> {
77 match self {
78 Self::Yson(reader) => reader
79 .next_event()
80 .map(|event| event.map(WorkerEvent::Yson)),
81 Self::Skiff(reader) => reader.next_row().map(|row| row.map(WorkerEvent::Skiff)),
82 }
83 }
84}
85
86pub enum WorkerRow<'row> {
88 YsonRaw(&'row [u8]),
91 Skiff(&'row Value),
93}
94
95pub enum WorkerWriter {
100 Yson(JobWriter),
102 Skiff(SkiffJobWriter),
104}
105
106impl std::fmt::Debug for WorkerWriter {
107 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
108 match self {
109 Self::Yson(writer) => formatter
110 .debug_tuple("WorkerWriter::Yson")
111 .field(writer)
112 .finish(),
113 Self::Skiff(writer) => formatter
114 .debug_tuple("WorkerWriter::Skiff")
115 .field(writer)
116 .finish(),
117 }
118 }
119}
120
121impl WorkerWriter {
122 #[cfg(unix)]
132 pub fn descriptors(format: DataFormat, table_count: usize) -> Result<Self> {
133 match format {
134 DataFormat::Yson(format) => {
135 JobWriter::descriptors_with_format(table_count, format).map(Self::Yson)
136 }
137 DataFormat::Skiff(format) => {
138 let schemas = format.table_schemas().len();
139 if schemas != table_count {
140 return Err(JobError::SkiffOutputSchemaCount {
141 sinks: table_count,
142 schemas,
143 });
144 }
145 SkiffJobWriter::descriptors(format).map(Self::Skiff)
146 }
147 _ => Err(JobError::UnsupportedDataFormat),
148 }
149 }
150
151 pub fn from_writers(tables: Vec<Box<dyn Write>>, format: DataFormat) -> Result<Self> {
161 match format {
162 DataFormat::Yson(format) => Ok(Self::Yson(JobWriter::from_writers(tables, format))),
163 DataFormat::Skiff(format) => {
164 SkiffJobWriter::from_writers(tables, format).map(Self::Skiff)
165 }
166 _ => Err(JobError::UnsupportedDataFormat),
167 }
168 }
169
170 #[must_use]
172 pub fn table_count(&self) -> usize {
173 match self {
174 Self::Yson(writer) => writer.table_count(),
175 Self::Skiff(writer) => writer.table_count(),
176 }
177 }
178
179 pub fn write(&mut self, table: impl Into<TableId>, row: WorkerRow<'_>) -> Result<()> {
189 match (self, row) {
190 (Self::Yson(writer), WorkerRow::YsonRaw(row)) => writer.write_raw(table, row),
191 (Self::Skiff(writer), WorkerRow::Skiff(row)) => writer.write(table, row),
192 (Self::Yson(_), WorkerRow::Skiff(_)) => Err(JobError::WorkerRowFormatMismatch {
193 writer: "YSON",
194 row: "Skiff",
195 }),
196 (Self::Skiff(_), WorkerRow::YsonRaw(_)) => Err(JobError::WorkerRowFormatMismatch {
197 writer: "Skiff",
198 row: "YSON",
199 }),
200 }
201 }
202
203 pub fn flush(&mut self) -> Result<()> {
205 match self {
206 Self::Yson(writer) => writer.flush(),
207 Self::Skiff(writer) => writer.flush(),
208 }
209 }
210
211 pub fn finish(&mut self) -> Result<()> {
213 match self {
214 Self::Yson(writer) => writer.finish(),
215 Self::Skiff(writer) => writer.finish(),
216 }
217 }
218}
219
220#[cfg(test)]
221mod tests {
222 use std::io::Cursor;
223
224 use ytsaurus_format::SkiffFormat;
225 use ytsaurus_skiff::{Encoder, Schema, SchemaRef, Value, WireType};
226
227 use super::*;
228
229 fn skiff_format() -> SkiffFormat {
230 SkiffFormat::new(vec![SchemaRef::Inline(Schema::tuple([Schema::named(
231 "value",
232 WireType::String32,
233 )]))])
234 .unwrap()
235 }
236
237 #[test]
238 fn reader_selects_yson_and_skiff_from_the_same_enum() {
239 let mut yson =
240 WorkerReader::new(Cursor::new(b"{value=one};"), DataFormat::text_yson()).unwrap();
241 assert!(matches!(
242 yson.next_event().unwrap(),
243 Some(WorkerEvent::Yson(_))
244 ));
245
246 let schema = skiff_format().table_schema(0).unwrap().clone();
247 let mut encoder = Encoder::new(Vec::new(), schema).unwrap();
248 encoder
249 .write(&Value::Tuple(vec![Value::Bytes(b"one".to_vec())]))
250 .unwrap();
251 let stream = encoder.into_inner().unwrap();
252 let mut skiff =
253 WorkerReader::new(Cursor::new(stream), DataFormat::skiff(skiff_format())).unwrap();
254 assert!(matches!(
255 skiff.next_event().unwrap(),
256 Some(WorkerEvent::Skiff(_))
257 ));
258 }
259
260 #[test]
261 fn writer_rejects_a_row_from_the_other_format() {
262 let mut writer =
263 WorkerWriter::from_writers(vec![Box::new(Vec::new())], DataFormat::binary_yson())
264 .unwrap();
265 let error = writer
266 .write(0, WorkerRow::Skiff(&Value::Tuple(Vec::new())))
267 .unwrap_err();
268 assert_eq!(error.kind(), "worker_row_format_mismatch");
269 }
270}