diff --git a/lib/ch/row_binary.ex b/lib/ch/row_binary.ex index 5ba2a7b4..63867645 100644 --- a/lib/ch/row_binary.ex +++ b/lib/ch/row_binary.ex @@ -676,7 +676,8 @@ defmodule Ch.RowBinary do end defp decode_rows!(data, types) do - {rows, remaining_data, state} = decode_rows(types, data, [], [], types) + {rows, remaining_data, state} = + decode_rows(types, data, [], [], types, :row, types, 1, []) case state do nil -> @@ -686,6 +687,16 @@ defmodule Ch.RowBinary do raise ArgumentError, """ incomplete RowBinary data: ran out of bytes while decoding + Expected to decode: #{inspect(types_rest)} + Remaining bytes: #{byte_size(remaining_data)} bytes + Partial row: #{inspect(row)} + Completed rows: #{length(rows)} + """ + + {:cont, types_rest, row, _kind, _schema, _left, _outer} -> + raise ArgumentError, """ + incomplete RowBinary data: ran out of bytes while decoding + Expected to decode: #{inspect(types_rest)} Remaining bytes: #{byte_size(remaining_data)} bytes Partial row: #{inspect(row)} @@ -697,8 +708,14 @@ defmodule Ch.RowBinary do @doc false def decode_rows_continue(<>, types, state) do case state do - {:cont, types_rest, row} -> decode_rows(types_rest, data, row, [], types) - nil -> decode_rows(types, data, [], [], types) + {:cont, types_rest, row} -> + decode_rows(types_rest, data, row, [], types, :row, types, 1, []) + + {:cont, types_rest, row, kind, schema, left, outer} -> + decode_rows(types_rest, data, row, [], types, kind, schema, left, outer) + + nil -> + decode_rows(types, data, [], [], types, :row, types, 1, []) end end @@ -834,7 +851,75 @@ defmodule Ch.RowBinary do end end - @compile inline: [decode_string_decode_rows: 5] + # Container state stays as separate function arguments at runtime. These macros only keep + # decoder clauses from repeating that plumbing for every successful or partial value. + defmacrop decode_next(types_rest, bin, row) do + quote do + decode_rows( + unquote(types_rest), + unquote(bin), + unquote(row), + var!(rows), + var!(types), + var!(kind), + var!(schema), + var!(left), + var!(outer) + ) + end + end + + defmacrop decode_more(bin, types_rest, row) do + quote do + to_be_continued( + var!(rows), + unquote(bin), + unquote(types_rest), + unquote(row), + var!(types), + var!(kind), + var!(schema), + var!(left), + var!(outer) + ) + end + end + + defmacrop dynamic_next(rest, dynamic, types_rest, row) do + quote do + decode_dynamic_continue( + unquote(rest), + unquote(dynamic), + unquote(types_rest), + unquote(row), + var!(rows), + var!(types), + var!(kind), + var!(schema), + var!(left), + var!(outer) + ) + end + end + + defmacrop dynamic_decode(rest, dynamic, types_rest, row) do + quote do + decode_dynamic( + unquote(rest), + unquote(dynamic), + unquote(types_rest), + unquote(row), + var!(rows), + var!(types), + var!(kind), + var!(schema), + var!(left), + var!(outer) + ) + end + end + + @compile inline: [decode_string_decode_rows: 9] for {pattern, size} <- varints do defp decode_string_decode_rows( @@ -842,17 +927,31 @@ defmodule Ch.RowBinary do types_rest, row, rows, - types + types, + kind, + schema, + left, + outer ) do - decode_rows(types_rest, bin, [s | row], rows, types) + decode_next(types_rest, bin, [s | row]) end end - defp decode_string_decode_rows(<>, types_rest, row, rows, _types) do - to_be_continued(rows, bin, [:string | types_rest], row) + defp decode_string_decode_rows( + <>, + types_rest, + row, + rows, + types, + kind, + schema, + left, + outer + ) do + decode_more(bin, [:string | types_rest], row) end - @compile inline: [decode_string_json_decode_rows: 5] + @compile inline: [decode_string_json_decode_rows: 9] for {pattern, size} <- varints do defp decode_string_json_decode_rows( @@ -860,19 +959,44 @@ defmodule Ch.RowBinary do types_rest, row, rows, - types + types, + kind, + schema, + left, + outer ) do - decode_rows(types_rest, bin, [JSON.decode!(s) | row], rows, types) + decode_next(types_rest, bin, [JSON.decode!(s) | row]) end end - defp decode_string_json_decode_rows(<>, types_rest, row, rows, _types) do - to_be_continued(rows, bin, [:json | types_rest], row) + defp decode_string_json_decode_rows( + <>, + types_rest, + row, + rows, + types, + kind, + schema, + left, + outer + ) do + decode_more(bin, [:json | types_rest], row) end - @compile inline: [decode_array_decode_rows: 6] - defp decode_array_decode_rows(<<0, bin::bytes>>, _type, types_rest, row, rows, types) do - decode_rows(types_rest, bin, [[] | row], rows, types) + @compile inline: [decode_array_decode_rows: 10] + defp decode_array_decode_rows( + <<0, bin::bytes>>, + _type, + types_rest, + row, + rows, + types, + kind, + schema, + left, + outer + ) do + decode_next(types_rest, bin, [[] | row]) end for {pattern, size} <- varints do @@ -882,19 +1006,44 @@ defmodule Ch.RowBinary do types_rest, row, rows, - types + types, + kind, + schema, + left, + outer ) do - array_types = List.duplicate(type, unquote(size)) - types_rest = array_types ++ [{:array_over, row} | types_rest] - decode_rows(types_rest, bin, [], rows, types) + array_schema = [type] + + decode_rows( + array_schema, + bin, + [], + rows, + types, + :array, + array_schema, + unquote(size), + [{types_rest, row, kind, schema, left} | outer] + ) end end - defp decode_array_decode_rows(<>, type, types_rest, row, rows, _types) do - to_be_continued(rows, bin, [{:array, type} | types_rest], row) + defp decode_array_decode_rows( + <>, + type, + types_rest, + row, + rows, + types, + kind, + schema, + left, + outer + ) do + decode_more(bin, [{:array, type} | types_rest], row) end - @compile inline: [decode_map_decode_rows: 7] + @compile inline: [decode_map_decode_rows: 11] defp decode_map_decode_rows( <<0, bin::bytes>>, _key_type, @@ -902,9 +1051,13 @@ defmodule Ch.RowBinary do types_rest, row, rows, - types + types, + kind, + schema, + left, + outer ) do - decode_rows(types_rest, bin, [%{} | row], rows, types) + decode_next(types_rest, bin, [%{} | row]) end for {pattern, size} <- varints do @@ -915,25 +1068,44 @@ defmodule Ch.RowBinary do types_rest, row, rows, - types + types, + kind, + schema, + left, + outer ) do - types_rest = - map_types(unquote(size), key_type, value_type) ++ [{:map_over, row} | types_rest] + map_schema = [key_type, value_type] - decode_rows(types_rest, bin, [], rows, types) + decode_rows( + map_schema, + bin, + [], + rows, + types, + {:map, %{}}, + map_schema, + unquote(size), + [{types_rest, row, kind, schema, left} | outer] + ) end end - defp decode_map_decode_rows(<>, key_type, value_type, types_rest, row, rows, _types) do - to_be_continued(rows, bin, [{:map, key_type, value_type} | types_rest], row) - end - - defp map_types(count, key_type, value_type) when count > 0 do - [key_type, value_type | map_types(count - 1, key_type, value_type)] + defp decode_map_decode_rows( + <>, + key_type, + value_type, + types_rest, + row, + rows, + types, + kind, + schema, + left, + outer + ) do + decode_more(bin, [{:map, key_type, value_type} | types_rest], row) end - defp map_types(0, _key_type, _value_types), do: [] - # https://clickhouse.com/docs/sql-reference/data-types/data-types-binary-encoding dynamic_types = [ nothing: 0x00, @@ -969,15 +1141,30 @@ defmodule Ch.RowBinary do types_rest, row, rows, - types + types, + kind, + schema, + left, + outer ) do - decode_dynamic_continue(rest, [unquote(type) | dynamic], types_rest, row, rows, types) + dynamic_next(rest, [unquote(type) | dynamic], types_rest, row) end end # DateTime 0x11 - defp decode_dynamic(<<0x11, rest::bytes>>, dynamic, types_rest, row, rows, types) do - decode_dynamic_continue(rest, [{:datetime, nil} | dynamic], types_rest, row, rows, types) + defp decode_dynamic( + <<0x11, rest::bytes>>, + dynamic, + types_rest, + row, + rows, + types, + kind, + schema, + left, + outer + ) do + dynamic_next(rest, [{:datetime, nil} | dynamic], types_rest, row) end # DateTime(time_zone) 0x12 @@ -988,9 +1175,13 @@ defmodule Ch.RowBinary do types_rest, row, rows, - types + types, + kind, + schema, + left, + outer ) do - decode_dynamic_continue(rest, [{:datetime, tz} | dynamic], types_rest, row, rows, types) + dynamic_next(rest, [{:datetime, tz} | dynamic], types_rest, row) end end @@ -1001,15 +1192,17 @@ defmodule Ch.RowBinary do types_rest, row, rows, - types + types, + kind, + schema, + left, + outer ) do - decode_dynamic_continue( + dynamic_next( rest, [decoding_type({:datetime64, precision}) | dynamic], types_rest, - row, - rows, - types + row ) end @@ -1021,15 +1214,17 @@ defmodule Ch.RowBinary do types_rest, row, rows, - types + types, + kind, + schema, + left, + outer ) do - decode_dynamic_continue( + dynamic_next( rest, [decoding_type({:datetime64, precision, tz}) | dynamic], types_rest, - row, - rows, - types + row ) end end @@ -1042,7 +1237,11 @@ defmodule Ch.RowBinary do types_rest, row, rows, - types + types, + kind, + schema, + left, + outer ) do decode_dynamic_continue( rest, @@ -1050,7 +1249,11 @@ defmodule Ch.RowBinary do types_rest, row, rows, - types + types, + kind, + schema, + left, + outer ) end end @@ -1066,32 +1269,34 @@ defmodule Ch.RowBinary do types_rest, row, rows, - types + types, + kind, + schema, + left, + outer ) do - decode_dynamic_continue( - rest, - [{:decimal, unquote(size), scale} | dynamic], - types_rest, - row, - rows, - types - ) + dynamic_next(rest, [{:decimal, unquote(size), scale} | dynamic], types_rest, row) end end # Array(T) 0x1E - defp decode_dynamic(<<0x1E, rest::bytes>>, dynamic, types_rest, row, rows, types) do - decode_dynamic_continue(rest, [:array | dynamic], types_rest, row, rows, types) - end - - # Nullable(T) 0x23 - defp decode_dynamic(<<0x23, rest::bytes>>, dynamic, types_rest, row, rows, types) do - decode_dynamic_continue(rest, [:nullable | dynamic], types_rest, row, rows, types) - end - + # Nullable(T) 0x23 # LowCardinality(T) 0x26 - defp decode_dynamic(<<0x26, rest::bytes>>, dynamic, types_rest, row, rows, types) do - decode_dynamic_continue(rest, [:low_cardinality | dynamic], types_rest, row, rows, types) + for {code, wrapper} <- [{0x1E, :array}, {0x23, :nullable}, {0x26, :low_cardinality}] do + defp decode_dynamic( + <>, + dynamic, + types_rest, + row, + rows, + types, + kind, + schema, + left, + outer + ) do + dynamic_next(rest, [unquote(wrapper) | dynamic], types_rest, row) + end end # TODO @@ -1130,18 +1335,51 @@ defmodule Ch.RowBinary do } for {type, code} <- unsupported_dynamic_types do - defp decode_dynamic(<>, _dynamic, _types_rest, _row, _rows, _types) do + defp decode_dynamic( + <>, + _dynamic, + _types_rest, + _row, + _rows, + _types, + _kind, + _schema, + _left, + _outer + ) do raise ArgumentError, "unsupported dynamic type #{unquote(type)}" end end - defp decode_dynamic(<>, dynamic, types_rest, row, rows, _types) do - to_be_continued(rows, bin, [{:dynamic, dynamic} | types_rest], row) + defp decode_dynamic( + <>, + dynamic, + types_rest, + row, + rows, + types, + kind, + schema, + left, + outer + ) do + decode_more(bin, [{:dynamic, dynamic} | types_rest], row) end - @compile inline: [decode_dynamic_continue: 6] + @compile inline: [decode_dynamic_continue: 10] - defp decode_dynamic_continue(<>, dynamic, types_rest, row, rows, types) do + defp decode_dynamic_continue( + <>, + dynamic, + types_rest, + row, + rows, + types, + kind, + schema, + left, + outer + ) do continue? = case dynamic do [:array | _] -> true @@ -1151,10 +1389,10 @@ defmodule Ch.RowBinary do end if continue? do - decode_dynamic(rest, dynamic, types_rest, row, rows, types) + dynamic_decode(rest, dynamic, types_rest, row) else type = build_dynamic_type(:lists.reverse(dynamic)) - decode_rows([type | types_rest], rest, row, rows, types) + decode_next([type | types_rest], rest, row) end end @@ -1226,110 +1464,162 @@ defmodule Ch.RowBinary do for {type, clauses} <- simple_types do fun = :"decode_#{type}_decode_rows" - @compile inline: [{fun, 5}] + @compile inline: [{fun, 9}] for %{pattern: pattern, value: value} <- List.wrap(clauses) do - defp unquote(fun)(<>, types_rest, row, rows, types) do - decode_rows(types_rest, rest, [unquote(value) | row], rows, types) + defp unquote(fun)( + <>, + types_rest, + row, + rows, + types, + kind, + schema, + left, + outer + ) do + decode_rows( + types_rest, + rest, + [unquote(value) | row], + rows, + types, + kind, + schema, + left, + outer + ) end end - defp unquote(fun)(<>, types_rest, row, rows, _types) do - to_be_continued(rows, bin, [unquote(type) | types_rest], row) + defp unquote(fun)( + <>, + types_rest, + row, + rows, + types, + kind, + schema, + left, + outer + ) do + decode_more(bin, [unquote(type) | types_rest], row) end end - defp decode_rows([type | types_rest], <>, row, rows, types) do + # The active sequence is carried in kind/schema/left. Only nested parents are stored in outer, + # so decoding another collection item does not allocate a replacement continuation frame. + defp decode_rows( + [type | types_rest], + <>, + row, + rows, + types, + kind, + schema, + left, + outer + ) do case type do :u8 -> - decode_u8_decode_rows(bin, types_rest, row, rows, types) + decode_u8_decode_rows(bin, types_rest, row, rows, types, kind, schema, left, outer) :u16 -> - decode_u16_decode_rows(bin, types_rest, row, rows, types) + decode_u16_decode_rows(bin, types_rest, row, rows, types, kind, schema, left, outer) :u32 -> - decode_u32_decode_rows(bin, types_rest, row, rows, types) + decode_u32_decode_rows(bin, types_rest, row, rows, types, kind, schema, left, outer) :u64 -> - decode_u64_decode_rows(bin, types_rest, row, rows, types) + decode_u64_decode_rows(bin, types_rest, row, rows, types, kind, schema, left, outer) :u128 -> - decode_u128_decode_rows(bin, types_rest, row, rows, types) + decode_u128_decode_rows(bin, types_rest, row, rows, types, kind, schema, left, outer) :u256 -> - decode_u256_decode_rows(bin, types_rest, row, rows, types) + decode_u256_decode_rows(bin, types_rest, row, rows, types, kind, schema, left, outer) :i8 -> - decode_i8_decode_rows(bin, types_rest, row, rows, types) + decode_i8_decode_rows(bin, types_rest, row, rows, types, kind, schema, left, outer) :i16 -> - decode_i16_decode_rows(bin, types_rest, row, rows, types) + decode_i16_decode_rows(bin, types_rest, row, rows, types, kind, schema, left, outer) :i32 -> - decode_i32_decode_rows(bin, types_rest, row, rows, types) + decode_i32_decode_rows(bin, types_rest, row, rows, types, kind, schema, left, outer) :i64 -> - decode_i64_decode_rows(bin, types_rest, row, rows, types) + decode_i64_decode_rows(bin, types_rest, row, rows, types, kind, schema, left, outer) :i128 -> - decode_i128_decode_rows(bin, types_rest, row, rows, types) + decode_i128_decode_rows(bin, types_rest, row, rows, types, kind, schema, left, outer) :i256 -> - decode_i256_decode_rows(bin, types_rest, row, rows, types) + decode_i256_decode_rows(bin, types_rest, row, rows, types, kind, schema, left, outer) :f32 -> - decode_f32_decode_rows(bin, types_rest, row, rows, types) + decode_f32_decode_rows(bin, types_rest, row, rows, types, kind, schema, left, outer) :f64 -> - decode_f64_decode_rows(bin, types_rest, row, rows, types) + decode_f64_decode_rows(bin, types_rest, row, rows, types, kind, schema, left, outer) :string -> - decode_string_decode_rows(bin, types_rest, row, rows, types) + decode_string_decode_rows(bin, types_rest, row, rows, types, kind, schema, left, outer) :json -> # assuming it arrives as text and not "native" binary JSON # i.e. assumes `settings: [output_format_binary_write_json_as_string: 1]` # TODO - decode_string_json_decode_rows(bin, types_rest, row, rows, types) + decode_string_json_decode_rows( + bin, + types_rest, + row, + rows, + types, + kind, + schema, + left, + outer + ) :dynamic -> - decode_dynamic(bin, _dynamic = [], types_rest, row, rows, types) + dynamic_decode(bin, _dynamic = [], types_rest, row) {:dynamic, dynamic} -> - decode_dynamic(bin, dynamic, types_rest, row, rows, types) + dynamic_decode(bin, dynamic, types_rest, row) {:fixed_string, size} -> case bin do <> -> - decode_rows(types_rest, rest, [s | row], rows, types) + decode_next(types_rest, rest, [s | row]) _ -> - to_be_continued(rows, bin, [type | types_rest], row) + decode_more(bin, [type | types_rest], row) end :boolean -> - decode_boolean_decode_rows(bin, types_rest, row, rows, types) + decode_boolean_decode_rows(bin, types_rest, row, rows, types, kind, schema, left, outer) :uuid -> - decode_uuid_decode_rows(bin, types_rest, row, rows, types) + decode_uuid_decode_rows(bin, types_rest, row, rows, types, kind, schema, left, outer) :date -> - decode_date_decode_rows(bin, types_rest, row, rows, types) + decode_date_decode_rows(bin, types_rest, row, rows, types, kind, schema, left, outer) :date32 -> - decode_date32_decode_rows(bin, types_rest, row, rows, types) + decode_date32_decode_rows(bin, types_rest, row, rows, types, kind, schema, left, outer) :time -> - decode_time_decode_rows(bin, types_rest, row, rows, types) + decode_time_decode_rows(bin, types_rest, row, rows, types, kind, schema, left, outer) {:time64, time_unit} -> case bin do <> -> time = time_after_midnight(ticks, time_unit) - decode_rows(types_rest, bin, [time | row], rows, types) + decode_next(types_rest, bin, [time | row]) _ -> - to_be_continued(rows, bin, [type | types_rest], row) + decode_more(bin, [type | types_rest], row) end {:datetime, timezone} -> @@ -1344,10 +1634,10 @@ defmodule Ch.RowBinary do _ -> DateTime.shift_zone!(dt, timezone) end - decode_rows(types_rest, bin, [dt | row], rows, types) + decode_next(types_rest, bin, [dt | row]) _ -> - to_be_continued(rows, bin, [type | types_rest], row) + decode_more(bin, [type | types_rest], row) end {:decimal, size, scale} -> @@ -1355,64 +1645,93 @@ defmodule Ch.RowBinary do <> -> sign = if val < 0, do: -1, else: 1 d = Decimal.new(sign, abs(val), -scale) - decode_rows(types_rest, bin, [d | row], rows, types) + decode_next(types_rest, bin, [d | row]) _ -> - to_be_continued(rows, bin, [type | types_rest], row) + decode_more(bin, [type | types_rest], row) end {:nullable, inner_type} -> case bin do <> -> case b do - 0 -> decode_rows([inner_type | types_rest], bin, row, rows, types) - 1 -> decode_rows(types_rest, bin, [nil | row], rows, types) + 0 -> + decode_next([inner_type | types_rest], bin, row) + + 1 -> + decode_next(types_rest, bin, [nil | row]) end _ -> - to_be_continued(rows, bin, [type | types_rest], row) + decode_more(bin, [type | types_rest], row) end :nothing -> - decode_rows(types_rest, bin, [nil | row], rows, types) + decode_next(types_rest, bin, [nil | row]) {:array, inner_type} -> - decode_array_decode_rows(bin, inner_type, types_rest, row, rows, types) - - {:array_over, original_row} -> - decode_rows(types_rest, bin, [:lists.reverse(row) | original_row], rows, types) + decode_array_decode_rows( + bin, + inner_type, + types_rest, + row, + rows, + types, + kind, + schema, + left, + outer + ) {:map, key_type, value_type} -> - decode_map_decode_rows(bin, key_type, value_type, types_rest, row, rows, types) - - {:map_over, original_row} -> - map = row |> Enum.chunk_every(2) |> Enum.map(fn [v, k] -> {k, v} end) |> Map.new() - decode_rows(types_rest, bin, [map | original_row], rows, types) + decode_map_decode_rows( + bin, + key_type, + value_type, + types_rest, + row, + rows, + types, + kind, + schema, + left, + outer + ) + + {:tuple, []} -> + decode_next(types_rest, bin, [{} | row]) {:tuple, tuple_types} -> - decode_rows(tuple_types ++ [{:tuple_over, row} | types_rest], bin, [], rows, types) - - {:tuple_over, original_row} -> - tuple = row |> :lists.reverse() |> List.to_tuple() - decode_rows(types_rest, bin, [tuple | original_row], rows, types) + decode_rows( + tuple_types, + bin, + [], + rows, + types, + :tuple, + tuple_types, + 1, + [{types_rest, row, kind, schema, left} | outer] + ) {:variant, variant_types} -> case bin do <<255, bin::bytes>> -> # 255 is the variant type index for "nothing" - decode_rows(types_rest, bin, [nil | row], rows, types) + decode_next(types_rest, bin, [nil | row]) # TODO varint? <> when variant_type_index < tuple_size(variant_types) -> variant_type = elem(variant_types, variant_type_index) - decode_rows([variant_type | types_rest], bin, row, rows, types) + + decode_next([variant_type | types_rest], bin, row) <> -> raise ArgumentError, "invalid Variant type index: #{variant_type_index}" _ -> - to_be_continued(rows, bin, [type | types_rest], row) + decode_more(bin, [type | types_rest], row) end {:datetime64, time_unit, timezone} -> @@ -1427,54 +1746,107 @@ defmodule Ch.RowBinary do _ -> DateTime.shift_zone!(dt, timezone) end - decode_rows(types_rest, bin, [dt | row], rows, types) + decode_next(types_rest, bin, [dt | row]) _ -> - to_be_continued(rows, bin, [type | types_rest], row) + decode_more(bin, [type | types_rest], row) end {:enum8, mapping} -> case bin do <> -> - decode_rows(types_rest, bin, [Map.fetch!(mapping, v) | row], rows, types) + decode_next(types_rest, bin, [Map.fetch!(mapping, v) | row]) _ -> - to_be_continued(rows, bin, [type | types_rest], row) + decode_more(bin, [type | types_rest], row) end {:enum16, mapping} -> case bin do <> -> - decode_rows(types_rest, bin, [Map.fetch!(mapping, v) | row], rows, types) + decode_next(types_rest, bin, [Map.fetch!(mapping, v) | row]) _ -> - to_be_continued(rows, bin, [type | types_rest], row) + decode_more(bin, [type | types_rest], row) end :ipv4 -> - decode_ipv4_decode_rows(bin, types_rest, row, rows, types) + decode_ipv4_decode_rows(bin, types_rest, row, rows, types, kind, schema, left, outer) :ipv6 -> - decode_ipv6_decode_rows(bin, types_rest, row, rows, types) + decode_ipv6_decode_rows(bin, types_rest, row, rows, types, kind, schema, left, outer) :point -> - decode_point_decode_rows(bin, types_rest, row, rows, types) + decode_point_decode_rows(bin, types_rest, row, rows, types, kind, schema, left, outer) end end - defp decode_rows([], <<>> = empty, row, rows, _types) do + defp decode_rows([], <>, row, rows, types, :array, schema, left, outer) + when left > 1 do + decode_rows(schema, bin, row, rows, types, :array, schema, left - 1, outer) + end + + defp decode_rows([], <>, row, rows, types, {:map, map}, schema, left, outer) do + [value, key] = row + map = Map.put(map, key, value) + + if left > 1 do + decode_rows(schema, bin, [], rows, types, {:map, map}, schema, left - 1, outer) + else + restore_parent(map, bin, rows, types, outer) + end + end + + defp decode_rows([], <>, row, rows, types, :array, _schema, 1, outer) do + restore_parent(:lists.reverse(row), bin, rows, types, outer) + end + + defp decode_rows([], <>, row, rows, types, :tuple, _schema, 1, outer) do + tuple = row |> :lists.reverse() |> List.to_tuple() + restore_parent(tuple, bin, rows, types, outer) + end + + defp decode_rows([], <<>> = empty, row, rows, _types, :row, _schema, 1, []) do rows = :lists.reverse([:lists.reverse(row) | rows]) {rows, empty, _no_state = nil} end - defp decode_rows([], <>, row, rows, types) do + defp decode_rows([], <>, row, rows, types, :row, schema, 1, []) do row = :lists.reverse(row) - decode_rows(types, bin, [], [row | rows], types) + decode_rows(schema, bin, [], [row | rows], types, :row, schema, 1, []) end - @compile inline: [to_be_continued: 4] - defp to_be_continued(rows, bin, types_rest, row) do - {:lists.reverse(rows), bin, {:cont, types_rest, row}} + @compile inline: [restore_parent: 5] + defp restore_parent( + value, + bin, + rows, + types, + [{parent_types, parent_row, parent_kind, parent_schema, parent_left} | outer] + ) do + decode_rows( + parent_types, + bin, + [value | parent_row], + rows, + types, + parent_kind, + parent_schema, + parent_left, + outer + ) + end + + @compile inline: [to_be_continued: 9] + defp to_be_continued(rows, bin, types_rest, row, _types, kind, schema, left, outer) do + state = + if kind == :row and outer == [] do + {:cont, types_rest, row} + else + {:cont, types_rest, row, kind, schema, left, outer} + end + + {:lists.reverse(rows), bin, state} end @compile inline: [decimal_size: 1]