From 0bc563255b35f927da4ddd466c11e88b778ea7df Mon Sep 17 00:00:00 2001 From: Paul Masurel Date: Tue, 4 Aug 2026 17:06:10 +0200 Subject: [PATCH] Add columnar V3 payload encoding tags --- .../v1_string_bytes.columnar | Bin 0 -> 595 bytes .../v2_string_bytes.columnar | Bin 0 -> 633 bytes columnar/compat_tests_data/v3.columnar | Bin 0 -> 42007 bytes columnar/src/column/serialize.rs | 65 +++++++- .../src/column_index/multivalued_index.rs | 2 +- columnar/src/columnar/format_version.rs | 9 +- .../src/columnar/merge/merge_dict_column.rs | 10 +- columnar/src/columnar/writer/mod.rs | 6 +- columnar/src/compat_tests.rs | 156 +++++++++++++++++- columnar/src/payload_encoding.rs | 36 ++++ columnar/src/tests.rs | 4 +- plain_string_column_plan.md | 52 +++++- 12 files changed, 318 insertions(+), 22 deletions(-) create mode 100644 columnar/compat_tests_data/v1_string_bytes.columnar create mode 100644 columnar/compat_tests_data/v2_string_bytes.columnar create mode 100644 columnar/compat_tests_data/v3.columnar diff --git a/columnar/compat_tests_data/v1_string_bytes.columnar b/columnar/compat_tests_data/v1_string_bytes.columnar new file mode 100644 index 0000000000000000000000000000000000000000..07d72faf1f648661f76bc4cd27a0553cb0d5ecef GIT binary patch literal 595 zcmWe+00ILBhW`ePK+FWh%nbiQVlW^HWw1bLW-#5**xd4j5y&(JVz3ee1||j}pebOI z2~0ABNT@b2r3mG4KxvR!ObiW8Ev+4H3=NI#9UaGhf#hs~m=mPNfPwL!0Tbgt5E}*r zp$sM{4KjqW1!NGA211}gjZI8EK%OoT3js01g5r`Q-L%r299;uRLrEh^V~`*WNI@Am zJZJ`#Rt71#05T-EG^Zp}*FZN>CdtYWOd5enV-N{72ux{1IfT6G4pawr;{%ZK`2{7J z`FV*zcgh+en+9_uVLyYM3o?^|oxvs4-8DYoKvHE%YH>Ws_YBNsVD~dHFN0aczzhl- nh%8G9M3!Xb7?p?1RFdV7W%>G9L%55&NM{bz#A!7l9Hb&Z zizhCHyc8?u%=4jM!pNdlQQ@iYd-iOQrq;#vIEfo{7-U7jFEejrdV@;dh~#pR`i^zB zwc6gPvrX7G$O;j=O$q)neG{*8+wb&Z_u6;5S1oLPE&2ICX#4Lq&;-N@^RR#3Jm?+2 zwD0|C|AZC-=9WhS^Bf)mL+hGlBR|VVzD@juO!=$AGumXrU=$zNY3|8KRSq+YyeqKHw% zDH0R`g;FFbQWR;5424l-DRLBfijYEzxTA5u=DxBq#z3rAShwDAE)e3ZuwU|GB1w^=NK<4ej3P^sqsUW)6jIC`#oSTM9mU*H z%pDXlia14rBA`%;Bt?oMO_8B6iY!HrB2N)g$Q16FG8p})7f<1iDcmuIJ1Allaf$>* zK%o>#iWEhfB12&mS&AG*o+6}>soXJ@JEn5SRPLC{9TYK&I7NaYpiqh=MT#O#k)be( zEJcnYPZ3hcH13$j9n-jD8h1?N4vH8>oFYLHP$)%`B1Ms=$WRzXmLf-yrwA#eggZ*O zql7z3xWj*`tcfB<5vNE{1QbeIM6fuf8MS>!rP>Li)iXu&sp)iUpMUEm*5mLxZ?wH9PGr40X zcg*AtiWo(lB0&*QC`FPYMUke+P#8s)B1e&@2q|P1cg*6BS==#;J7#eQMT{a&k)Q}D zlp;xyqDWI@D2yUYk)y~{gcLHHJ7#moZ0?xN9kaQEB1RFXNKgb6N|B^UQKTs{6h@Jy z$Wi1eLJIz;{L-0a+)>6IW!zE59sWy8V-#_U1Vuoh6iJE{MVcZ*VH8=497UcYq>wrM z=P`#n=5WUx|2M3k!yObcia14rBA`%;Bt?oMO_8B6iY!HrB2N)g$XxE2%N=vMV=i~h zak|ITsrpQnjMV2B*k*5eLWFB|SUQ<=j!u9TYK&I7NaYpiqh=MT#O#k)be(EJcnY zPZ3hceD0Xf9rL+kK6lLL4vH8>oFYLHP$)%`B1Ms=$WRzXmLf-yrwA!z0e39mjs@JY zfIAj&2Stn`PLZGpD3l^ek)lXbWGIXxOOd0$^iWo(lB0&*QC`FPYMUke+P#8s)B1e&@2q|O< zcP!zKCET%uJC<+3HRoqd<9aY>x5u=DxBq#z3rAShw zDAE)e3ZuwUNCWC}I?GiUdVKp%h7q6h)dMLtzwIiX26rBBYSz z8=~Lz%;ns%oI93t$8zqVh*88T5)=W2QY0x-6lsbKg;8WFauj)rkV2}tqnbOaxucpp zs=0$AMiHk-Py`f8k)%jbq$x5KMvq*USjinLxq~7`5vNE{1QbeLNwcJt59ktv+5u=DxBq#z3rAShwDAE)e3ZuwU!rP>Li)iXu&sp)iUpMUEm*5mLxn?pVtmYq?`BcdX?OiWo(lB0&*Q zC`FPYMUke+P#8s)B1e&@2r1+o?l^}#&f$)8xZ@n|pome#DH0R`g;FFbQWR;5424l- zDRLBfijYF;xTB6c>bRqhJL{f zbGhSO?l_k_C}I?GiUdVKp%h7q6h)dMLtzwIiX26rBBYRY+_8>3)^W!=?pVhi6fuf8 zMS>!rP>Li)iXu&sp)iUpMUEm*5mLx{?pV(q>!Ul0SFIOWx4uXwiT`o?UF)7%|IEq1 z8q!n0HY9`5yE%HdMDI}aZjIh;(Yrl*hog5#^zMw_UC}!dy_x9U9ld*^*P?fC^zMt^ z{n2|MdJjhLq3Asvy+@*VR30DKjGFO3{nwqhZ`^d}_}~9{a`S)a(o^Ti*vUtR8Ns4J8B%YjygxZqruVWXmT_=S{$v8Hb=Xo z!_n#Ja&$X-9KDV{N8spp3^?Qt|9l)Jj#5XNquf#9sB~00svR|sT1TCu-qGM_bTm1d z9W9PlN1LPF(c$QHbUC^mJ&s;SpCfSeI|dwb>c2+gU$|mNiKEm}<|ubmI4T`gj%r7Z zqt;R9sCP6t8XZlJW=D&o)zRi?cXT*99bJxYM~|b|(dP&p{f>b<<*mgcMdM$dN#j3i zUX1<_aWx@DhNOrcQe;Rbu|tXs$z*m&ks&E&hZGr-DeRCULo$^eQe;S`u|tXsNeMfo z$dF8DhZGr-8SIcELvk8Bq{xtzvO|gt$xL=gks+DI4kNA*o}B6d96p*&#)SWF0%C$dIgOhZGr-dUi;WAvup7Qe;TZXNMFS zl9#YUiVR5uJEX{vyp$bMWJoSxhZGr-m$5^N3`rw9q{xuGoE=hRNG@cD6d96>*daxR zq=_9;WJq4Y4k*&#)SDKaGO?2sZuaydJs$dFvY4k7MhUBg6kRn6U#SSSlByVGf6d97M*daxRm zhD6yRMTX>dc1V#Sxq}^2WJvC0hZGr-L3T)yA-RhkQe;Tp#||knBzLnziVVqSc1V#S zc|SX($dG)19a3aSKFAI!G9+8rAw`BH$qp$pB=@jGiVVqz*daxRWQZM7WJo^D4kAw`B{D?6mfkbIOKQe;T(V}}$Ol8>=NiVVp%c1V#S`8YeI$dIJiAw`Dd z6YP*8L$aM6Qe;Rz$qp$pB=@sJiVVr8*daxRWSAXNWJo^E4k z4kFtB)?&Y6d97=vO|gt$q{x)ks*1C9a3aSo@R#>8Is?z zLy8Q^C_ALcko=wFtB>!QD z6d4kqz1W1*waH$T@qeQtDbhPm{cXpm{;uj_mFeSuU-G0=e|zc4Kiz)vZ%f{|>F)N+ zC-le93w`19o4_Yf`2I}b!y8WvpI_ngo5&|n`2I}f(<}V@(C2aKNI=% z3O|1e-=D(wXCj|K;rla@Pp|Ovr||tLe19hL2^79R6Z!NCKYt3}pThTNBA-Cv`!kVG zukiDy@ck)#eL_u{ zaFjY`Im#S!9p#P%jta+ON2Ozlqsp<&QSDgisBx@u)H>EW>KyAF^^Wr$4UP*OjgAW) zO^ywYX2->j7RM!yR>!4|HpgX-cE=Tt4oBS4>A2F-<+#ex?YPF#$u+0=eWTU zIBs_IJ8pFhI5s-uPXGDiD0WPBlsINMN*%KtWsbRya>oKkg=4X!(y_!*}+~5ctH#_k#KapA_>?v{;r7CllTvSj=pY*cjW zA9d6%`=qvLs3nv!trG>OfWf(rrsmEbUshcUjHyEz6HBZ?E26UAaP6jIL-| zxoc(Fs$kWjRgE>nHN~rYR_|Y3w`S{_<7+z4va_mdH`gAmZ96-2cE#F(wMW)ApR@Cv z(z?F7gLMt(Za-Jnb+6mEu6F&<`fz)cA_wR}^pP z*|2{@-7B}g^7tz|o2|L}RhwUR^i^#aXD+UI^}wr-yt=t%XG`g8`d)MJH4T?+zeHZ! z{n~x6t!*7@4O=^6dtz0W4qiHTY3u7oURU1M-*&jI>GeBaUvgRRWd|;+f5WynoOnZ5 z``-4N%eP#9?DF<2c3)BXMt$Sx8(TVdb(Fm+c+;UbHO7bI#c%F;^Zqy2b#CoE-r4yU zdrS3|o3A{2W!qabZ>{JW=sMEX{I;EME4`}ks)JWGynXxIrMtU(Uw7@*Lsy4acU-gQ znySQLVl2_xGtyIjZU420uWh<+$8{y|=zYh5chvW8>pjui_0GNTths*6^~bJnf7kAJ zRrYD$XkW{_cfGsphTw)nH#FWjd}DFY6YLM_ZrXa&@tZntwwtSO*?h~+j!o|B3s%KDGB#HN#tmj}5nfdiSR*AJ7LzA87f^uFsV12zDIW(fHZn&lW$}^Wgpm z>(X1($J3oVZD;l8Hh=Eu=h{A>`FzC}2EK6Q3(dQB?kfFa-xm*lvEfVGza(Go{_?&r z*NzO0gd-ha+4GgEhXx-Sd#LrRBVR4g^k)udn!dKB~?YMP**S=f* z=;lX{KHB!Z%=aqx4(vU$xB2@!zhC-T-(v?KYxu$TAIQG$ef##+{&47r;SW20wC6`v zKOX$?*pFNHkL)kc_Gb@gn|`w6Cnb;fK7Qcw`UBezoH)?+)4e~fd1A{G$DU~a+3ufJ z9@K-Q2U~u=>*r-p22UP(vhf$gzbHP`b7=pex?gVn8g; z)BnujXPWXm@+D)vV+Y3S|FG>3C;rg&$Gv~7`O}s^9s5)J(cMQY|Ezx={d3D-cKxO7 zufbmr{k8Gf@Uh}&d!F6@Y~6EPpF93s=ilsa)#2vwXxR4m%-<`XA9()A^UeR*`H#}$ zea8I{U7&yH5VY;Qs6XmwNX2U#^)VQhV~(r^vIff1%_bPZj-iMql^U*WY-xNY^zt z_4bPNz4NBt;93#AvF~j+oc#6w{-Zwni}C+U+Eaf_9{+P&( column_index: SerializableColumnIndex<'_>, @@ -106,8 +106,43 @@ pub fn open_column_u128_as_compact_u64( } pub fn open_column_bytes(data: OwnedBytes, format_version: Version) -> io::Result { + let data = match format_version { + Version::V1 | Version::V2 => data, + Version::V3 => { + if data.is_empty() { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "missing string/byte payload encoding tag", + )); + } + let (encoding_bytes, payload) = data.split(1); + let encoding = PayloadEncoding::try_from_code(encoding_bytes.as_slice()[0]) + .map_err(io::Error::from)?; + match encoding { + PayloadEncoding::Dictionary => payload, + PayloadEncoding::Plain => { + return Err(io::Error::new( + io::ErrorKind::Unsupported, + "plain string/byte column decoding is not implemented yet", + )); + } + } + } + }; + if data.len() < 4 { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "truncated dictionary string/byte column payload", + )); + } let (body, dictionary_len_bytes) = data.rsplit(4); let dictionary_len = u32::from_le_bytes(dictionary_len_bytes.as_slice().try_into().unwrap()); + if dictionary_len as usize > body.len() { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + "dictionary length exceeds string/byte column payload", + )); + } let (dictionary_bytes, column_bytes) = body.split(dictionary_len as usize); let dictionary = Arc::new(Dictionary::from_bytes(dictionary_bytes)?); let term_ord_column = crate::column::open_column_u64::(column_bytes, format_version)?; @@ -125,3 +160,31 @@ pub fn open_column_str(data: OwnedBytes, format_version: Version) -> io::Result< }; Ok(DictionaryEncodedStrColumn::wrap(bytes_column).into()) } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_v3_payload_encoding_tag_errors() { + let error = open_column_bytes(OwnedBytes::new(Vec::new()), Version::V3).unwrap_err(); + assert_eq!(error.kind(), io::ErrorKind::InvalidData); + + let error = open_column_bytes(OwnedBytes::new(vec![u8::MAX]), Version::V3).unwrap_err(); + assert_eq!(error.kind(), io::ErrorKind::InvalidData); + + let error = open_column_bytes( + OwnedBytes::new(vec![PayloadEncoding::Plain.to_code()]), + Version::V3, + ) + .unwrap_err(); + assert_eq!(error.kind(), io::ErrorKind::Unsupported); + + let error = open_column_bytes( + OwnedBytes::new(vec![PayloadEncoding::Dictionary.to_code()]), + Version::V3, + ) + .unwrap_err(); + assert_eq!(error.kind(), io::ErrorKind::InvalidData); + } +} diff --git a/columnar/src/column_index/multivalued_index.rs b/columnar/src/column_index/multivalued_index.rs index ad7efd363..52393f552 100644 --- a/columnar/src/column_index/multivalued_index.rs +++ b/columnar/src/column_index/multivalued_index.rs @@ -55,7 +55,7 @@ pub fn open_multivalued_index( start_index_column, })) } - Version::V2 => { + Version::V2 | Version::V3 => { let (body_bytes, optional_index_len) = bytes.rsplit(4); let optional_index_len = u32::from_le_bytes(optional_index_len.as_slice().try_into().unwrap()); diff --git a/columnar/src/columnar/format_version.rs b/columnar/src/columnar/format_version.rs index 1314fb754..52aaaf602 100644 --- a/columnar/src/columnar/format_version.rs +++ b/columnar/src/columnar/format_version.rs @@ -23,13 +23,14 @@ pub fn parse_footer(footer_bytes: [u8; VERSION_FOOTER_NUM_BYTES]) -> Result write!(f, "v1"), Version::V2 => write!(f, "v2"), + Version::V3 => write!(f, "v3"), } } } @@ -51,6 +53,7 @@ impl Version { match code { 1u32 => Ok(Version::V1), 2u32 => Ok(Version::V2), + 3u32 => Ok(Version::V3), _ => Err(InvalidData), } } @@ -65,7 +68,7 @@ mod tests { #[test] fn test_footer_deserialization() { let parsed_version: Version = parse_footer(footer()).unwrap(); - assert_eq!(Version::V2, parsed_version); + assert_eq!(Version::V3, parsed_version); } #[test] @@ -83,6 +86,6 @@ mod tests { valid_versions.insert(i); } } - assert_eq!(valid_versions.len(), 2); + assert_eq!(valid_versions.len(), 3); } } diff --git a/columnar/src/columnar/merge/merge_dict_column.rs b/columnar/src/columnar/merge/merge_dict_column.rs index 336cd9b9f..dae870ec5 100644 --- a/columnar/src/columnar/merge/merge_dict_column.rs +++ b/columnar/src/columnar/merge/merge_dict_column.rs @@ -7,9 +7,11 @@ use super::term_merger::{TermMerger, TermsWithSegmentOrd}; use crate::column::serialize_column_mappable_to_u64; use crate::column_index::SerializableColumnIndex; use crate::iterable::Iterable; -use crate::{BytesColumn, DictionaryEncodedBytesColumn, MergeRowOrder, ShuffleMergeOrder}; +use crate::{ + BytesColumn, DictionaryEncodedBytesColumn, MergeRowOrder, PayloadEncoding, ShuffleMergeOrder, +}; -// Serialize [Dictionary, Column, dictionary num bytes U32::LE] +// V3 serialize [PayloadEncoding, Dictionary, Column, dictionary num bytes U32::LE] // Column: [Column Index, Column Values, column index num bytes U32::LE] pub fn merge_bytes_or_str_column( column_index: SerializableColumnIndex<'_>, @@ -17,7 +19,9 @@ pub fn merge_bytes_or_str_column( merge_row_order: &MergeRowOrder, output: &mut impl Write, ) -> io::Result<()> { - // Serialize dict and generate mapping for values + output.write_all(&[PayloadEncoding::Dictionary.to_code()])?; + // Serialize dict and generate mapping for values. + // The encoding tag is intentionally excluded from `dictionary_num_bytes`. let mut output = CountingWriter::wrap(output); // TODO !!! Remove useless terms. let term_ord_mapping = serialize_merged_dict(bytes_columns, merge_row_order, &mut output)?; diff --git a/columnar/src/columnar/writer/mod.rs b/columnar/src/columnar/writer/mod.rs index 999ccd058..abe8d15ca 100644 --- a/columnar/src/columnar/writer/mod.rs +++ b/columnar/src/columnar/writer/mod.rs @@ -22,7 +22,7 @@ use crate::columnar::writer::column_writers::{ use crate::columnar::writer::value_index::{IndexBuilder, PreallocatedIndexBuilders}; use crate::dictionary::{DictionaryBuilder, TermIdMapping, UnorderedId}; use crate::value::{Coerce, NumericalType, NumericalValue}; -use crate::{Cardinality, RowId}; +use crate::{Cardinality, PayloadEncoding, RowId}; /// This is a set of buffers that are used to temporarily write the values into before passing them /// to the fast field codecs. @@ -532,7 +532,7 @@ fn collect_sort_order_from_ops( .collect() } -// Serialize [Dictionary, Column, dictionary num bytes U32::LE] +// V3 serialize [PayloadEncoding, Dictionary, Column, dictionary num bytes U32::LE] // Column: [Column Index, Column Values, column index num bytes U32::LE] #[expect(clippy::too_many_arguments)] fn serialize_bytes_or_str_column( @@ -550,6 +550,8 @@ fn serialize_bytes_or_str_column( u64_values, .. } = buffers; + let mut wrt = wrt; + wrt.write_all(&[PayloadEncoding::Dictionary.to_code()])?; let mut counting_writer = CountingWriter::wrap(wrt); let term_id_mapping: TermIdMapping = dictionary_builder.serialize(arena, &mut counting_writer)?; diff --git a/columnar/src/compat_tests.rs b/columnar/src/compat_tests.rs index 64e615973..61130a22d 100644 --- a/columnar/src/compat_tests.rs +++ b/columnar/src/compat_tests.rs @@ -3,11 +3,12 @@ use std::path::PathBuf; use itertools::Itertools; use crate::{ - CURRENT_VERSION, Cardinality, Column, ColumnarReader, DynamicColumn, StackMergeOrder, - merge_columnar, + CURRENT_VERSION, Cardinality, Column, ColumnarReader, DictionaryEncodedBytesColumn, + DictionaryEncodedStrColumn, DynamicColumn, PayloadEncoding, StackMergeOrder, merge_columnar, }; const NUM_DOCS: u32 = u16::MAX as u32; +const STRING_BYTES_NUM_DOCS: u32 = 4; fn generate_columnar(num_docs: u32, value_offset: u64) -> Vec { use crate::ColumnarWriter; @@ -32,6 +33,32 @@ fn generate_columnar(num_docs: u32, value_offset: u64) -> Vec { wrt } +fn generate_string_bytes_columnar() -> Vec { + use crate::ColumnarWriter; + + let mut columnar_writer = ColumnarWriter::default(); + for doc in 0..STRING_BYTES_NUM_DOCS { + columnar_writer.record_str(doc, "str_full", &format!("str-full-{doc}")); + columnar_writer.record_bytes(doc, "bytes_full", &[doc as u8, 0, 255]); + + if doc.is_multiple_of(2) { + columnar_writer.record_str(doc, "str_optional", &format!("str-optional-{doc}")); + columnar_writer.record_bytes(doc, "bytes_optional", &[doc as u8, 1, 254]); + } + + columnar_writer.record_str(doc, "str_multi", &format!("str-multi-{doc}-a")); + columnar_writer.record_str(doc, "str_multi", &format!("str-multi-{doc}-b")); + columnar_writer.record_bytes(doc, "bytes_multi", &[doc as u8, 2, 0]); + columnar_writer.record_bytes(doc, "bytes_multi", &[doc as u8, 2, 255]); + } + + let mut output = Vec::new(); + columnar_writer + .serialize(STRING_BYTES_NUM_DOCS, None, &mut output) + .unwrap(); + output +} + #[test] /// Writes a columnar for the CURRENT_VERSION to disk. fn create_format() { @@ -63,6 +90,131 @@ fn test_format_v2() { test_format(&path); } +#[test] +fn test_format_v3() { + let path = path_for_version("v3"); + test_format(&path); +} + +#[test] +fn test_string_bytes_format_v1() { + test_string_bytes_format("v1"); +} + +#[test] +fn test_string_bytes_format_v2() { + test_string_bytes_format("v2"); +} + +fn test_string_bytes_format(version: &str) { + let fixture_path = format!("./compat_tests_data/{version}_string_bytes.columnar"); + let fixture_reader = ColumnarReader::open(std::fs::read(fixture_path).unwrap()).unwrap(); + check_string_bytes_columns(&fixture_reader, 1); + + let current_reader = ColumnarReader::open(generate_string_bytes_columnar()).unwrap(); + check_string_bytes_columns(¤t_reader, 1); + + let readers = [&fixture_reader, ¤t_reader]; + let merge_row_order = StackMergeOrder::stack(&readers); + let mut output = Vec::new(); + merge_columnar(&readers, &[], merge_row_order.into(), &mut output).unwrap(); + let merged_reader = ColumnarReader::open(output).unwrap(); + check_string_bytes_columns(&merged_reader, 2); +} + +fn check_string_bytes_columns(reader: &ColumnarReader, repetitions: u32) { + let num_docs = STRING_BYTES_NUM_DOCS * repetitions; + + let str_full = open_str_column(reader, "str_full"); + assert_eq!(str_full.get_cardinality(), Cardinality::Full); + let str_optional = open_str_column(reader, "str_optional"); + assert_eq!(str_optional.get_cardinality(), Cardinality::Optional); + let str_multi = open_str_column(reader, "str_multi"); + assert_eq!(str_multi.get_cardinality(), Cardinality::Multivalued); + + let bytes_full = open_bytes_column(reader, "bytes_full"); + assert_eq!(bytes_full.get_cardinality(), Cardinality::Full); + let bytes_optional = open_bytes_column(reader, "bytes_optional"); + assert_eq!(bytes_optional.get_cardinality(), Cardinality::Optional); + let bytes_multi = open_bytes_column(reader, "bytes_multi"); + assert_eq!(bytes_multi.get_cardinality(), Cardinality::Multivalued); + + for row_id in 0..num_docs { + let doc = row_id % STRING_BYTES_NUM_DOCS; + assert_eq!( + str_values(&str_full, row_id), + vec![format!("str-full-{doc}")] + ); + assert_eq!( + str_values(&str_optional, row_id), + if doc.is_multiple_of(2) { + vec![format!("str-optional-{doc}")] + } else { + Vec::new() + } + ); + assert_eq!( + str_values(&str_multi, row_id), + vec![format!("str-multi-{doc}-a"), format!("str-multi-{doc}-b")] + ); + + assert_eq!( + bytes_values(&bytes_full, row_id), + vec![vec![doc as u8, 0, 255]] + ); + assert_eq!( + bytes_values(&bytes_optional, row_id), + if doc.is_multiple_of(2) { + vec![vec![doc as u8, 1, 254]] + } else { + Vec::new() + } + ); + assert_eq!( + bytes_values(&bytes_multi, row_id), + vec![vec![doc as u8, 2, 0], vec![doc as u8, 2, 255]] + ); + } +} + +fn open_str_column(reader: &ColumnarReader, name: &str) -> DictionaryEncodedStrColumn { + let DynamicColumn::Str(column) = reader.read_columns(name).unwrap()[0].open().unwrap() else { + panic!("expected a string column") + }; + assert_eq!(column.payload_encoding(), PayloadEncoding::Dictionary); + column.as_dictionary_encoded().unwrap().clone() +} + +fn open_bytes_column(reader: &ColumnarReader, name: &str) -> DictionaryEncodedBytesColumn { + let DynamicColumn::Bytes(column) = reader.read_columns(name).unwrap()[0].open().unwrap() else { + panic!("expected a byte column") + }; + assert_eq!(column.payload_encoding(), PayloadEncoding::Dictionary); + column.as_dictionary_encoded().unwrap().clone() +} + +fn str_values(column: &DictionaryEncodedStrColumn, row_id: u32) -> Vec { + column + .term_ords(row_id) + .map(|term_ord| { + let mut output = String::new(); + assert!(column.ord_to_str(term_ord, &mut output).unwrap()); + output + }) + .collect() +} + +fn bytes_values(column: &DictionaryEncodedBytesColumn, row_id: u32) -> Vec> { + column + .term_ords(row_id) + .map(|term_ord| { + let mut output = Vec::new(); + assert!(column.ord_to_bytes(term_ord, &mut output).unwrap()); + output + }) + .collect() +} + fn test_format(path: &str) { let file_content = std::fs::read(path).unwrap(); let reader = ColumnarReader::open(file_content).unwrap(); diff --git a/columnar/src/payload_encoding.rs b/columnar/src/payload_encoding.rs index c7569b0bf..eba4cf737 100644 --- a/columnar/src/payload_encoding.rs +++ b/columnar/src/payload_encoding.rs @@ -1,5 +1,7 @@ use serde::{Deserialize, Serialize}; +use crate::InvalidData; + /// Encoding used to store string and byte column payloads. #[derive(Clone, Copy, Debug, Default, Eq, PartialEq, Serialize, Deserialize)] #[serde(rename_all = "lowercase")] @@ -10,3 +12,37 @@ pub enum PayloadEncoding { /// Store values directly without assigning them dictionary ordinals. Plain, } + +impl PayloadEncoding { + pub(crate) fn to_code(self) -> u8 { + match self { + PayloadEncoding::Dictionary => 0, + PayloadEncoding::Plain => 1, + } + } + + pub(crate) fn try_from_code(code: u8) -> Result { + match code { + 0 => Ok(PayloadEncoding::Dictionary), + 1 => Ok(PayloadEncoding::Plain), + _ => Err(InvalidData), + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_payload_encoding_codes() { + let mut valid_codes = Vec::new(); + for code in u8::MIN..=u8::MAX { + if let Ok(encoding) = PayloadEncoding::try_from_code(code) { + assert_eq!(encoding.to_code(), code); + valid_codes.push(code); + } + } + assert_eq!(valid_codes, [0, 1]); + } +} diff --git a/columnar/src/tests.rs b/columnar/src/tests.rs index 22e6801e5..83e45612b 100644 --- a/columnar/src/tests.rs +++ b/columnar/src/tests.rs @@ -26,7 +26,7 @@ fn test_dataframe_writer_str() { assert_eq!(columnar.num_columns(), 1); let cols: Vec = columnar.read_columns("my_string").unwrap(); assert_eq!(cols.len(), 1); - assert_eq!(cols[0].num_bytes(), 73); + assert_eq!(cols[0].num_bytes(), 74); } #[test] @@ -40,7 +40,7 @@ fn test_dataframe_writer_bytes() { assert_eq!(columnar.num_columns(), 1); let cols: Vec = columnar.read_columns("my_string").unwrap(); assert_eq!(cols.len(), 1); - assert_eq!(cols[0].num_bytes(), 73); + assert_eq!(cols[0].num_bytes(), 74); } #[test] diff --git a/plain_string_column_plan.md b/plain_string_column_plan.md index 345519eb6..92cc155bf 100644 --- a/plain_string_column_plan.md +++ b/plain_string_column_plan.md @@ -6,8 +6,8 @@ Add a user-selectable payload encoding for string and byte fast fields: ```rust pub enum PayloadEncoding { - Plain, Dictionary, + Plain, } ``` @@ -50,9 +50,9 @@ and on-disk format all need it. Re-export it from `tantivy::schema` for normal T ```rust #[derive(Clone, Copy, Debug, Default, Eq, PartialEq, Serialize, Deserialize)] pub enum PayloadEncoding { - Plain, #[default] Dictionary, + Plain, } ``` @@ -246,12 +246,47 @@ V3 string and byte payloads start with a stable encoding discriminant: encoding tag | encoding-specific payload ``` -For `Dictionary`, the bytes after the tag use the current dictionary payload layout unchanged. -For `Plain`, define a versioned layout containing the column index, OnPair16 model, compressed -payload, offsets, and explicit region lengths in a footer. +The one-byte encoding tags are fixed independently of the Rust enum declaration order: -The exact ordering should allow `OwnedBytes::split` operations without copying. All serialized -lengths must be checked before slicing. +```text +0 = Dictionary +1 = Plain +2..=255 = reserved; readers reject them +``` + +For `Dictionary`, the bytes after the tag use the V2 dictionary payload layout unchanged: + +```text +0u8 | dictionary | column index and term ordinals | dictionary_num_bytes:u32 LE +``` + +`dictionary_num_bytes` counts only `dictionary`, excluding the encoding tag. For `Plain`, the V3 +layout is: + +```text +1u8 +| column_index +| onpair16_model +| compressed_values +| value_offsets +| column_index_num_bytes:u32 LE +| model_num_bytes:u32 LE +| compressed_values_num_bytes:u64 LE +| value_offsets_num_bytes:u32 LE +| num_values:u32 LE +``` + +`value_offsets` is a serialized monotonic `u64` column with exactly `num_values + 1` entries. Its +first entry is zero, its last entry equals `compressed_values_num_bytes`, and every adjacent pair +delimits one independently compressed value. `onpair16_model` uses the codec's canonical model +serialization, which must be frozen alongside the standalone plain-column implementation. + +The fixed 24-byte footer is read from the end first. The four regions are then split from left to +right without copying. Checked conversion to `usize`, checked length sums, and region bounds are +required before any `OwnedBytes::split` call; unrecognized tags and trailing bytes are invalid. + +V1 and V2 never consume a tag. V3 always consumes exactly one tag byte for string and byte +payloads, including dictionary payloads. Other column types do not need an encoding tag and retain their existing V3 representation unless the version implementation requires a uniform envelope. @@ -259,7 +294,8 @@ the version implementation requires a uniform envelope. ### Compatibility tests Expand `columnar/src/compat_tests.rs` fixtures so V1 and V2 include dictionary-encoded string and -byte columns for all supported cardinalities. Tests must: +byte columns for all supported cardinalities. The historical fixtures are +`v1_string_bytes.columnar` and `v2_string_bytes.columnar`. Tests must: - Open and read old values through the new enum variants. - Assert that old columns become `DictionaryEncoded`.