• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

plausible / ch / 5b5d3345939fba904c70b91173e31f75c805ff9f-PR-407

03 Aug 2026 12:51PM UTC coverage: 97.767% (-0.3%) from 98.062%
5b5d3345939fba904c70b91173e31f75c805ff9f-PR-407

Pull #407

github

ruslandoga
Add reusable RowBinary schemas
Pull Request #407: Add reusable RowBinary schemas

29 of 32 new or added lines in 1 file covered. (90.63%)

788 of 806 relevant lines covered (97.77%)

14832.59 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

99.3
/lib/ch/row_binary.ex
1
defmodule Ch.RowBinary do
2
  @moduledoc "Helpers for working with ClickHouse [RowBinary](https://clickhouse.com/docs/en/interfaces/formats/RowBinary) format."
3

4
  defmodule Schema do
5
    @moduledoc false
6

7
    @enforce_keys [:types, :encoding_types]
8
    defstruct [:names, :types, :encoding_types, :names_and_types]
9

10
    @type t :: %__MODULE__{
11
            names: [String.t()] | nil,
12
            types: [String.t()],
13
            encoding_types: [term()],
14
            names_and_types: iodata() | nil
15
          }
16
  end
17

18
  @opaque schema :: Schema.t()
19

20
  # @compile {:bin_opt_info, true}
21
  @dialyzer :no_improper_lists
22

23
  import Bitwise
24

25
  @epoch_gregorian_seconds 62_167_219_200
26
  @epoch_gregorian_days 719_528
27

28
  @doc """
29
  Builds a reusable RowBinary schema.
30

31
  ClickHouse type names and atoms for simple types are accepted. Types are
32
  parsed once when the schema is built and reused when rows are encoded.
33

34
  Pass `:names` to also build a `RowBinaryWithNamesAndTypes` header.
35

36
  Examples:
37

38
      iex> schema = schema(types: ["UInt64", :string])
39
      iex> encode_rows(schema, [[1, "one"], [2, "two"]])
40
      [<<1, 0, 0, 0, 0, 0, 0, 0>>, [3 | "one"], <<2, 0, 0, 0, 0, 0, 0, 0>>, [3 | "two"]]
41

42
      iex> schema = schema(names: ["id", "text"], types: ["UInt64", :string])
43
      iex> IO.iodata_to_binary(encode_names_and_types(schema))
44
      <<2, 2, "id", 4, "text", 6, "UInt64", 6, "String">>
45

46
  """
47
  @spec schema(keyword()) :: schema()
48
  def schema(options) when is_list(options) do
49
    unless Keyword.keyword?(options) do
11✔
NEW
50
      raise ArgumentError, "expected schema options to be a keyword list"
×
51
    end
52

53
    case Keyword.keys(options) -- [:names, :types] do
11✔
54
      [] -> :ok
10✔
55
      unknown -> raise ArgumentError, "unknown schema options: #{inspect(unknown)}"
1✔
56
    end
57

58
    types =
10✔
59
      case Keyword.fetch(options, :types) do
60
        {:ok, types} when is_list(types) ->
61
          types
8✔
62

63
        {:ok, types} ->
64
          raise ArgumentError, "expected schema :types to be a list, got: #{inspect(types)}"
1✔
65

66
        :error ->
67
          raise ArgumentError, "missing required schema option :types"
1✔
68
      end
69

70
    {types, encoding_types} =
8✔
71
      Enum.map(types, &schema_type/1)
72
      |> Enum.unzip()
73

74
    names = Keyword.get(options, :names)
7✔
75
    validate_schema_names!(names, length(types))
7✔
76

77
    names_and_types =
5✔
78
      if names do
3✔
79
        encode_names_and_types(names, types)
2✔
80
      end
81

82
    %Schema{
5✔
83
      names: names,
84
      types: types,
85
      encoding_types: encoding_types,
86
      names_and_types: names_and_types
87
    }
88
  end
89

90
  def schema(options) do
NEW
91
    raise ArgumentError, "expected schema options to be a keyword list, got: #{inspect(options)}"
×
92
  end
93

94
  defp schema_type(type) when is_binary(type), do: {type, encoding_type(type)}
10✔
95

96
  defp schema_type(type) when is_atom(type) do
4✔
97
    encoding_type = encoding_type(type)
4✔
98
    {IO.iodata_to_binary(Ch.Types.encode(type)), encoding_type}
99
  end
100

101
  defp schema_type(type) do
102
    raise ArgumentError,
1✔
103
          "expected schema types to be ClickHouse type strings or atoms, got: #{inspect(type)}"
104
  end
105

106
  defp validate_schema_names!(nil, _type_count), do: :ok
3✔
107

108
  defp validate_schema_names!(names, type_count) when is_list(names) do
109
    unless Enum.all?(names, &is_binary/1) do
4✔
110
      raise ArgumentError,
1✔
111
            "expected schema :names to contain only strings, got: #{inspect(names)}"
112
    end
113

114
    if length(names) != type_count do
3✔
115
      raise ArgumentError,
1✔
116
            "schema names and types must have the same length, got #{length(names)} names and #{type_count} types"
1✔
117
    end
118
  end
119

120
  defp validate_schema_names!(names, _type_count) do
NEW
121
    raise ArgumentError, "expected schema :names to be a list of strings, got: #{inspect(names)}"
×
122
  end
123

124
  @doc """
125
  Returns the encoded names and types header for a schema.
126

127
  The schema must have been built with the `:names` option.
128
  """
129
  @spec encode_names_and_types(schema()) :: iodata()
130
  def encode_names_and_types(%Schema{names_and_types: nil}) do
131
    raise ArgumentError, "can't encode names and types for a schema without names"
1✔
132
  end
133

134
  def encode_names_and_types(%Schema{names_and_types: names_and_types}), do: names_and_types
2✔
135

136
  @doc false
137
  def encode_names_and_types(names, types) do
7✔
138
    [encode(:varint, length(names)), encode_many(names, :string), encode_types(types)]
139
  end
140

141
  defp encode_types([type | types]) do
16✔
142
    encoded =
16✔
143
      case type do
144
        _ when is_binary(type) -> type
15✔
145
        _ -> Ch.Types.encode(type)
1✔
146
      end
147

148
    [encode(:string, encoded) | encode_types(types)]
149
  end
150

151
  defp encode_types([] = done), do: done
7✔
152

153
  @doc """
154
  Encodes a single row to [RowBinary](https://clickhouse.com/docs/en/interfaces/formats/RowBinary) as iodata.
155

156
  Examples:
157

158
      iex> encode_row([], [])
159
      []
160

161
      iex> encode_row([1], ["UInt8"])
162
      [1]
163

164
      iex> encode_row([3, "hello"], ["UInt8", "String"])
165
      [3, [5 | "hello"]]
166

167
  """
168
  @spec encode_row(schema(), list()) :: iodata()
169
  @spec encode_row(list(), [term()]) :: iodata()
170
  def encode_row(%Schema{encoding_types: types}, row) do
171
    _encode_row(row, types)
1✔
172
  end
173

174
  def encode_row(row, types) do
175
    _encode_row(row, encoding_types(types))
23✔
176
  end
177

178
  defp _encode_row([el | els], [type | types]), do: [encode(type, el) | _encode_row(els, types)]
112✔
179
  defp _encode_row([] = done, []), do: done
24✔
180

181
  @doc """
182
  Encodes multiple rows to [RowBinary](https://clickhouse.com/docs/en/interfaces/formats/RowBinary) as iodata.
183

184
  Examples:
185

186
      iex> encode_rows([], [])
187
      []
188

189
      iex> encode_rows([[1]], ["UInt8"])
190
      [1]
191

192
      iex> encode_rows([[3, "hello"], [4, "hi"]], ["UInt8", "String"])
193
      [3, [5 | "hello"], 4, [2 | "hi"]]
194

195
  """
196
  @spec encode_rows(schema(), [list()]) :: iodata()
197
  @spec encode_rows([list()], [term()]) :: iodata()
198
  def encode_rows(%Schema{encoding_types: types}, rows) do
199
    _encode_rows(rows, types)
2✔
200
  end
201

202
  def encode_rows(rows, types) do
203
    _encode_rows(rows, encoding_types(types))
855✔
204
  end
205

206
  @doc false
207
  def _encode_rows([row | rows], types), do: _encode_rows(row, types, rows, types)
4,737✔
208
  def _encode_rows([] = done, _types), do: done
857✔
209

210
  defp _encode_rows([el | els], [t | ts], rows, types) do
12,204✔
211
    [encode(t, el) | _encode_rows(els, ts, rows, types)]
212
  end
213

214
  defp _encode_rows([], [], rows, types), do: _encode_rows(rows, types)
4,737✔
215

216
  @doc false
217
  def encoding_types([type | types]) do
2,268✔
218
    [encoding_type(type) | encoding_types(types)]
219
  end
220

221
  def encoding_types([] = done), do: done
886✔
222

223
  defp encoding_type(type) when is_binary(type) do
224
    encoding_type(Ch.Types.decode(type))
2,178✔
225
  end
226

227
  defp encoding_type(t)
228
       when t in [
229
              :string,
230
              :json,
231
              :dynamic,
232
              :boolean,
233
              :uuid,
234
              :date,
235
              :datetime,
236
              :date32,
237
              :time,
238
              :ipv4,
239
              :ipv6,
240
              :point,
241
              :nothing
242
            ],
243
       do: t
929✔
244

245
  defp encoding_type({:datetime = d, "UTC"}), do: d
2✔
246

247
  defp encoding_type({:datetime, tz}) do
248
    raise ArgumentError, "can't encode DateTime with non-UTC timezone: #{inspect(tz)}"
1✔
249
  end
250

251
  defp encoding_type({:fixed_string, _len} = t), do: t
311✔
252

253
  for size <- [8, 16, 32, 64, 128, 256] do
254
    defp encoding_type(unquote(:"u#{size}") = u), do: u
667✔
255
    defp encoding_type(unquote(:"i#{size}") = i), do: i
40✔
256
  end
257

258
  for size <- [32, 64] do
259
    defp encoding_type(unquote(:"f#{size}") = f), do: f
222✔
260
  end
261

262
  defp encoding_type({:array = a, t}), do: {a, encoding_type(t)}
551✔
263

264
  defp encoding_type({:tuple = t, ts}) do
5✔
265
    {t, Enum.map(ts, &encoding_type/1)}
266
  end
267

268
  defp encoding_type({:variant = v, ts}) do
3✔
269
    {v, Enum.map(ts, &encoding_type/1)}
270
  end
271

272
  defp encoding_type({:map = m, kt, vt}) do
273
    {m, encoding_type(kt), encoding_type(vt)}
8✔
274
  end
275

276
  defp encoding_type({:nullable = n, t}), do: {n, encoding_type(t)}
113✔
277
  defp encoding_type({:low_cardinality, t}), do: encoding_type(t)
202✔
278

279
  defp encoding_type({:decimal, p, s}) do
280
    case decimal_size(p) do
5✔
281
      32 -> {:decimal32, s}
1✔
282
      64 -> {:decimal64, s}
2✔
283
      128 -> {:decimal128, s}
1✔
284
      256 -> {:decimal256, s}
1✔
285
    end
286
  end
287

288
  defp encoding_type({d, _scale} = t)
289
       when d in [:decimal32, :decimal64, :decimal128, :decimal256],
290
       do: t
5✔
291

292
  defp encoding_type({:datetime64 = t, p}), do: {t, time_unit(p)}
1✔
293

294
  defp encoding_type({:datetime64 = t, p, "UTC"}), do: {t, time_unit(p)}
2✔
295

296
  defp encoding_type({:datetime64, _, tz}) do
297
    raise ArgumentError, "can't encode DateTime64 with non-UTC timezone: #{inspect(tz)}"
1✔
298
  end
299

300
  defp encoding_type({:time64 = t, p}), do: {t, time_unit(p)}
105✔
301

302
  defp encoding_type({e, mappings}) when e in [:enum8, :enum16] do
4✔
303
    {e, Map.new(mappings)}
304
  end
305

306
  defp encoding_type({:simple_aggregate_function, _f, t}), do: encoding_type(t)
1✔
307

308
  defp encoding_type(:ring), do: {:array, :point}
1✔
309
  defp encoding_type(:polygon), do: {:array, {:array, :point}}
1✔
310
  defp encoding_type(:multipolygon), do: {:array, {:array, {:array, :point}}}
1✔
311

312
  defp encoding_type(type) do
313
    raise ArgumentError, "unsupported type for encoding: #{inspect(type)}"
1✔
314
  end
315

316
  @doc false
317
  def encode(type, value)
318

319
  def encode(:varint, i) when is_integer(i) and i >= 0 and i < 128, do: i
7,882✔
320
  def encode(:varint, i) when is_integer(i) and i >= 0, do: encode_varint_cont(i)
12✔
321

322
  def encode(:varint, i) when is_integer(i) do
323
    raise ArgumentError, "invalid varint: #{inspect(i)}"
1✔
324
  end
325

326
  def encode(:string, str) do
327
    case str do
5,907✔
328
      _ when is_binary(str) -> [encode(:varint, byte_size(str)) | str]
5,873✔
329
      _ when is_list(str) -> [encode(:varint, IO.iodata_length(str)) | str]
25✔
330
      nil -> 0
3✔
331
    end
332
  end
333

334
  def encode(:json, json) do
335
    # assuming it can be sent as text and not "native" binary JSON
336
    # i.e. assumes `settings: [input_format_binary_read_json_as_string: 1]`
337
    # TODO
338
    encode(:string, JSON.encode_to_iodata!(json))
5✔
339
  end
340

341
  def encode({:fixed_string, size}, str) when byte_size(str) == size do
342
    str
829✔
343
  end
344

345
  def encode({:fixed_string, size}, str) when byte_size(str) < size do
3,444✔
346
    to_pad = size - byte_size(str)
3,444✔
347
    [str | <<0::size(to_pad * 8)>>]
348
  end
349

350
  def encode({:fixed_string, size}, nil), do: <<0::size(size * 8)>>
2✔
351

352
  # UInt8 — [0 : 255]
353
  def encode(:u8, u) when is_integer(u) and u >= 0 and u <= 255, do: u
4,490✔
354
  def encode(:u8, nil), do: 0
4✔
355

356
  def encode(:u8, term) do
357
    raise ArgumentError, "invalid UInt8: #{inspect(term)}"
7✔
358
  end
359

360
  # Int8 — [-128 : 127]
361
  def encode(:i8, i) when is_integer(i) and i >= 0 and i <= 127, do: i
22✔
362
  def encode(:i8, i) when is_integer(i) and i < 0 and i >= -128, do: <<i::signed>>
15✔
363
  def encode(:i8, nil), do: 0
1✔
364

365
  def encode(:i8, term) do
366
    raise ArgumentError, "invalid Int8: #{inspect(term)}"
6✔
367
  end
368

369
  for size <- [16, 32, 64, 128, 256] do
370
    unsigned_max = (1 <<< size) - 1
371
    signed_min = -(1 <<< (size - 1))
372
    signed_max = (1 <<< (size - 1)) - 1
373
    uint = :"u#{size}"
374
    int = :"i#{size}"
375

376
    def encode(unquote(uint), u) when is_integer(u) and u >= 0 and u <= unquote(unsigned_max) do
377
      <<u::unquote(size)-little>>
151✔
378
    end
379

380
    def encode(unquote(int), i)
381
        when is_integer(i) and i >= unquote(signed_min) and i <= unquote(signed_max) do
382
      <<i::unquote(size)-little-signed>>
149✔
383
    end
384

385
    def encode(unquote(uint), nil), do: <<0::unquote(size)>>
3✔
386
    def encode(unquote(int), nil), do: <<0::unquote(size)>>
3✔
387

388
    def encode(unquote(uint), term) do
389
      raise ArgumentError, "invalid UInt#{unquote(size)}: #{inspect(term)}"
15✔
390
    end
391

392
    def encode(unquote(int), term) do
393
      raise ArgumentError, "invalid Int#{unquote(size)}: #{inspect(term)}"
15✔
394
    end
395
  end
396

397
  for size <- [32, 64] do
398
    type = :"f#{size}"
399

400
    def encode(unquote(type), f) when is_number(f) do
401
      <<f::unquote(size)-little-signed-float>>
1,220✔
402
    end
403

404
    def encode(unquote(type), nil), do: <<0::unquote(size)>>
4✔
405
  end
406

407
  def encode({:decimal, precision, scale}, decimal) do
408
    type =
4✔
409
      case decimal_size(precision) do
410
        32 -> :decimal32
1✔
411
        64 -> :decimal64
1✔
412
        128 -> :decimal128
1✔
413
        256 -> :decimal256
1✔
414
      end
415

416
    encode({type, scale}, decimal)
4✔
417
  end
418

419
  for size <- [32, 64, 128, 256] do
420
    type = :"decimal#{size}"
421

422
    def encode({unquote(type), scale} = t, %Decimal{sign: sign, coef: coef, exp: exp} = d) do
423
      cond do
29✔
424
        scale == -exp ->
425
          i = sign * coef
20✔
426
          <<i::unquote(size)-little>>
20✔
427

428
        exp >= 0 ->
9✔
429
          i = sign * coef * Integer.pow(10, exp + scale)
1✔
430
          <<i::unquote(size)-little>>
1✔
431

432
        true ->
8✔
433
          encode(t, Decimal.round(d, scale))
8✔
434
      end
435
    end
436

437
    def encode({unquote(type), _scale}, nil), do: <<0::unquote(size)>>
4✔
438
  end
439

440
  def encode(:boolean, true), do: 1
937✔
441
  def encode(:boolean, false), do: 0
988✔
442
  def encode(:boolean, nil), do: 0
1✔
443

444
  def encode({:array, type}, [_ | _] = l) do
1,973✔
445
    [encode(:varint, length(l)) | encode_many(l, type)]
446
  end
447

448
  def encode({:array, _type}, []), do: 0
269✔
449
  def encode({:array, _type}, nil), do: 0
4✔
450

451
  def encode({:map, k, v}, [_ | _] = m) do
1✔
452
    [encode(:varint, length(m)) | encode_many_kv(m, k, v)]
453
  end
454

455
  def encode({:map, k, v}, m) when is_map(m) do
14✔
456
    [
457
      encode(:varint, map_size(m))
458
      | :maps.fold(fn key, value, acc -> [encode(k, key), encode(v, value) | acc] end, [], m)
15✔
459
    ]
460
  end
461

462
  def encode({:map, _k, _v}, []), do: 0
1✔
463
  def encode({:map, _k, _v}, nil), do: 0
1✔
464

465
  def encode({:tuple, _types} = t, v) when is_tuple(v) do
466
    encode(t, Tuple.to_list(v))
11✔
467
  end
468

469
  def encode({:tuple, types}, values) when is_list(types) and is_list(values) do
470
    encode_row(values, types)
11✔
471
  end
472

473
  def encode({:tuple, types}, nil) when is_list(types) do
474
    Enum.map(types, fn type -> encode(type, nil) end)
1✔
475
  end
476

477
  def encode({:variant, _types}, nil), do: 255
3✔
478

479
  def encode({:variant, types}, value) do
480
    try_encode_variant(types, 0, value)
8✔
481
  end
482

483
  def encode(:datetime, %NaiveDateTime{} = datetime) do
484
    {seconds, _micros} = NaiveDateTime.to_gregorian_seconds(datetime)
17✔
485
    <<seconds - @epoch_gregorian_seconds::32-little>>
17✔
486
  end
487

488
  def encode(:datetime, %DateTime{} = datetime) do
489
    <<DateTime.to_unix(datetime, :second)::32-little>>
5✔
490
  end
491

492
  def encode(:datetime, nil), do: <<0::32>>
1✔
493

494
  def encode({:datetime64, time_unit}, %NaiveDateTime{} = datetime) do
495
    {seconds, micros} = NaiveDateTime.to_gregorian_seconds(datetime)
4✔
496

497
    <<(seconds - @epoch_gregorian_seconds) * time_unit + div(micros * time_unit, 1_000_000)::64-little-signed>>
4✔
498
  end
499

500
  def encode({:datetime64, time_unit}, %DateTime{} = datetime) do
501
    <<DateTime.to_unix(datetime, time_unit)::64-little-signed>>
5✔
502
  end
503

504
  def encode({:datetime64, _time_unit}, nil), do: <<0::64>>
1✔
505

506
  def encode(:date, %Date{} = date) do
507
    <<Date.to_gregorian_days(date) - @epoch_gregorian_days::16-little>>
14✔
508
  end
509

510
  def encode(:date, nil), do: <<0::16>>
1✔
511

512
  def encode(:date32, %Date{} = date) do
513
    <<Date.to_gregorian_days(date) - @epoch_gregorian_days::32-little-signed>>
8✔
514
  end
515

516
  def encode(:date32, nil), do: <<0::32>>
1✔
517

518
  def encode(:time, %Time{} = time) do
519
    {s, _micros} = Time.to_seconds_after_midnight(time)
107✔
520
    <<s::32-little-signed>>
107✔
521
  end
522

523
  def encode(:time, nil), do: <<0::32>>
1✔
524

525
  def encode({:time64, time_unit}, %Time{} = time) do
526
    {s, micros} = Time.to_seconds_after_midnight(time)
117✔
527

528
    micros_as_ticks =
117✔
529
      cond do
530
        time_unit < 1_000_000 -> div(micros, div(1_000_000, time_unit))
72✔
531
        time_unit == 1_000_000 -> micros
45✔
532
        true -> micros * div(time_unit, 1_000_000)
34✔
533
      end
534

535
    ticks = s * time_unit + micros_as_ticks
117✔
536
    <<ticks::64-little-signed>>
117✔
537
  end
538

539
  def encode({:time64, _time_unit}, nil), do: <<0::64>>
1✔
540

541
  def encode(:uuid, <<u1::64, u2::64>>), do: <<u1::64-little, u2::64-little>>
12✔
542

543
  def encode(
544
        :uuid,
545
        <<a1, a2, a3, a4, a5, a6, a7, a8, ?-, b1, b2, b3, b4, ?-, c1, c2, c3, c4, ?-, d1, d2, d3,
546
          d4, ?-, e1, e2, e3, e4, e5, e6, e7, e8, e9, e10, e11, e12>>
547
      ) do
548
    raw =
2✔
549
      <<d(a1)::4, d(a2)::4, d(a3)::4, d(a4)::4, d(a5)::4, d(a6)::4, d(a7)::4, d(a8)::4, d(b1)::4,
550
        d(b2)::4, d(b3)::4, d(b4)::4, d(c1)::4, d(c2)::4, d(c3)::4, d(c4)::4, d(d1)::4, d(d2)::4,
551
        d(d3)::4, d(d4)::4, d(e1)::4, d(e2)::4, d(e3)::4, d(e4)::4, d(e5)::4, d(e6)::4, d(e7)::4,
552
        d(e8)::4, d(e9)::4, d(e10)::4, d(e11)::4, d(e12)::4>>
553

554
    encode(:uuid, raw)
2✔
555
  end
556

557
  def encode(:uuid, nil), do: <<0::128>>
1✔
558

559
  def encode(:ipv4, {a, b, c, d}), do: [d, c, b, a]
6✔
560
  def encode(:ipv4, nil), do: <<0::32>>
1✔
561

562
  def encode(:ipv6, {b1, b2, b3, b4, b5, b6, b7, b8}) do
563
    <<b1::16, b2::16, b3::16, b4::16, b5::16, b6::16, b7::16, b8::16>>
6✔
564
  end
565

566
  def encode(:ipv6, <<_::128>> = encoded), do: encoded
1✔
567
  def encode(:ipv6, nil), do: <<0::128>>
1✔
568

569
  def encode(:point, {x, y}), do: [encode(:f64, x) | encode(:f64, y)]
22✔
570
  def encode(:point, nil), do: <<0::128>>
1✔
571
  def encode(:ring, points), do: encode({:array, :point}, points)
1✔
572
  def encode(:polygon, rings), do: encode({:array, :ring}, rings)
1✔
573
  def encode(:multipolygon, polygons), do: encode({:array, :polygon}, polygons)
1✔
574

575
  # TODO
576
  def encode(:dynamic, value) do
577
    case value do
12✔
578
      _ when is_binary(value) -> [0x15 | encode(:string, value)]
2✔
579
      _ when is_integer(value) and value >= 0 -> [0x04 | encode(:u64, value)]
3✔
580
      _ when is_integer(value) -> [0x0A | encode(:i64, value)]
1✔
581
      _ when is_float(value) -> [0x0E | encode(:f64, value)]
2✔
582
      %Date{} -> [0x0F | encode(:date, value)]
2✔
583
      %NaiveDateTime{} -> [0x11 | encode(:datetime, value)]
1✔
584
      [] -> [0x1E, 0x00]
1✔
585
    end
586
  end
587

588
  # TODO enum8 and enum16 nil
589
  for size <- [8, 16] do
590
    enum_t = :"enum#{size}"
591
    int_t = :"i#{size}"
592

593
    def encode({unquote(enum_t), mapping}, e) do
594
      i =
12✔
595
        case e do
596
          _ when is_integer(e) ->
597
            e
2✔
598

599
          _ when is_binary(e) ->
600
            case Map.fetch(mapping, e) do
10✔
601
              {:ok, res} ->
602
                res
9✔
603

604
              :error ->
605
                raise ArgumentError,
1✔
606
                      "enum value #{inspect(e)} not found in mapping: #{inspect(mapping)}"
607
            end
608
        end
609

610
      encode(unquote(int_t), i)
11✔
611
    end
612
  end
613

614
  def encode({:nullable, _type}, nil), do: 1
843✔
615

616
  def encode({:nullable, type}, value) do
617
    case encode(type, value) do
837✔
618
      e when is_list(e) or is_binary(e) -> [0 | e]
836✔
619
      e -> [0, e]
1✔
620
    end
621
  end
622

623
  defp encode_varint_cont(i) when i < 128, do: <<i>>
12✔
624

625
  defp encode_varint_cont(i) do
17✔
626
    [(i &&& 0b0111_1111) ||| 0b1000_0000 | encode_varint_cont(i >>> 7)]
627
  end
628

629
  defp encode_many([el | rest], type), do: [encode(type, el) | encode_many(rest, type)]
8,784✔
630
  defp encode_many([] = done, _type), do: done
1,980✔
631

632
  defp encode_many_kv([{key, value} | rest], key_type, value_type) do
1✔
633
    [
634
      encode(key_type, key),
635
      encode(value_type, value)
636
      | encode_many_kv(rest, key_type, value_type)
637
    ]
638
  end
639

640
  defp encode_many_kv([] = done, _key_type, _value_type), do: done
1✔
641

642
  # TODO find a better way than try/rescue
643
  defp try_encode_variant([type | types], idx, value) do
644
    try do
13✔
645
      encode(type, value)
13✔
646
    else
647
      encoded -> [idx | encoded]
7✔
648
    rescue
649
      _e -> try_encode_variant(types, idx + 1, value)
6✔
650
    end
651
  end
652

653
  defp try_encode_variant([], _idx, value) do
654
    raise ArgumentError, "no matching type found for encoding #{inspect(value)} as Variant"
1✔
655
  end
656

657
  @compile {:inline, d: 1}
658

659
  defp d(?0), do: 0
1✔
660
  defp d(?1), do: 1
1✔
661
  defp d(?2), do: 2
3✔
662
  defp d(?3), do: 3
1✔
663
  defp d(?4), do: 4
1✔
664
  defp d(?5), do: 5
3✔
665
  defp d(?6), do: 6
1✔
666
  defp d(?7), do: 7
1✔
667
  defp d(?8), do: 8
1✔
668
  defp d(?9), do: 9
1✔
669
  defp d(?A), do: 10
1✔
670
  defp d(?B), do: 11
1✔
671
  defp d(?C), do: 12
1✔
672
  defp d(?D), do: 13
1✔
673
  defp d(?E), do: 14
1✔
674
  defp d(?F), do: 15
1✔
675
  defp d(?a), do: 10
1✔
676
  defp d(?b), do: 11
1✔
677
  defp d(?c), do: 12
1✔
678
  defp d(?d), do: 13
1✔
679
  defp d(?e), do: 14
1✔
680
  defp d(?f), do: 15
1✔
681

682
  varints = [
683
    {_pattern = quote(do: <<0::1, v1::7>>), _value = quote(do: v1)},
684
    {quote(do: <<1::1, v1::7, 0::1, v2::7>>), quote(do: (v2 <<< 7) + v1)},
685
    {quote(do: <<1::1, v1::7, 1::1, v2::7, 0::1, v3::7>>),
686
     quote(do: (v3 <<< 14) + (v2 <<< 7) + v1)},
687
    {quote(do: <<1::1, v1::7, 1::1, v2::7, 1::1, v3::7, 0::1, v4::7>>),
688
     quote(do: (v4 <<< 21) + (v3 <<< 14) + (v2 <<< 7) + v1)},
689
    {quote(do: <<1::1, v1::7, 1::1, v2::7, 1::1, v3::7, 1::1, v4::7, 0::1, v5::7>>),
690
     quote(do: (v5 <<< 28) + (v4 <<< 21) + (v3 <<< 14) + (v2 <<< 7) + v1)},
691
    {quote(do: <<1::1, v1::7, 1::1, v2::7, 1::1, v3::7, 1::1, v4::7, 1::1, v5::7, 0::1, v6::7>>),
692
     quote(do: (v6 <<< 35) + (v5 <<< 28) + (v4 <<< 21) + (v3 <<< 14) + (v2 <<< 7) + v1)},
693
    {quote do
694
       <<1::1, v1::7, 1::1, v2::7, 1::1, v3::7, 1::1, v4::7, 1::1, v5::7, 1::1, v6::7, 0::1,
695
         v7::7>>
696
     end,
697
     quote do
698
       (v7 <<< 42) + (v6 <<< 35) + (v5 <<< 28) + (v4 <<< 21) + (v3 <<< 14) + (v2 <<< 7) + v1
699
     end},
700
    {quote do
701
       <<1::1, v1::7, 1::1, v2::7, 1::1, v3::7, 1::1, v4::7, 1::1, v5::7, 1::1, v6::7, 1::1,
702
         v7::7, 0::1, v8::7>>
703
     end,
704
     quote do
705
       (v8 <<< 49) + (v7 <<< 42) + (v6 <<< 35) + (v5 <<< 28) + (v4 <<< 21) + (v3 <<< 14) +
706
         (v2 <<< 7) + v1
707
     end}
708
  ]
709

710
  @doc false
711
  @spec decode_header(binary()) ::
712
          {:ok, names :: [String.t()], types :: [term], rest :: binary} | :more
713
  def decode_header(row_binary_with_names_and_types)
714

715
  for {pattern, value} <- varints do
716
    def decode_header(<<unquote(pattern), rest::bytes>>) do
717
      decode_header_names(rest, unquote(value), unquote(value), _acc = [])
36✔
718
    end
719
  end
720

721
  def decode_header(<<_bin::bytes>>) do
1✔
722
    :more
723
  end
724

725
  defp decode_header_names(<<rest::bytes>>, 0, count, names) do
726
    decode_header_types(rest, count, _acc = [], :lists.reverse(names))
21✔
727
  end
728

729
  for {pattern, value} <- varints do
730
    defp decode_header_names(
731
           <<unquote(pattern), name::size(unquote(value))-bytes, rest::bytes>>,
732
           left,
733
           count,
734
           acc
735
         ) do
736
      decode_header_names(rest, left - 1, count, [name | acc])
78✔
737
    end
738
  end
739

740
  defp decode_header_names(<<_bin::bytes>>, _left, _count, _acc) do
15✔
741
    :more
742
  end
743

744
  defp decode_header_types(<<rest::bytes>>, 0, types, names) do
745
    {:ok, names, decoding_types_reverse(types), rest}
1✔
746
  end
747

748
  for {pattern, value} <- varints do
749
    defp decode_header_types(
750
           <<unquote(pattern), type::size(unquote(value))-bytes, rest::bytes>>,
751
           count,
752
           acc,
753
           names
754
         ) do
755
      decode_header_types(rest, count - 1, [type | acc], names)
24✔
756
    end
757
  end
758

759
  defp decode_header_types(<<_bin::bytes>>, _count, _acc, _names) do
20✔
760
    :more
761
  end
762

763
  @doc """
764
  Decodes [RowBinaryWithNamesAndTypes](https://clickhouse.com/docs/en/interfaces/formats/RowBinaryWithNamesAndTypes) into rows.
765

766
  Example:
767

768
      iex> decode_rows(<<1, 3, "1+1"::bytes, 5, "UInt8"::bytes, 2>>)
769
      [[2]]
770

771
  """
772
  def decode_rows(row_binary_with_names_and_types)
773
  def decode_rows(<<>>), do: []
1✔
774

775
  for {pattern, value} <- varints do
776
    def decode_rows(<<unquote(pattern), rest::bytes>>) do
777
      skip_names(rest, unquote(value), unquote(value))
5✔
778
    end
779
  end
780

781
  @doc """
782
  Same as `decode_rows/1` but the first element is a list of column names.
783

784
  Example:
785

786
      iex> decode_names_and_rows(<<1, 3, "1+1"::bytes, 5, "UInt8"::bytes, 2>>)
787
      [["1+1"], [2]]
788

789
  """
790
  def decode_names_and_rows(row_binary_with_names_and_types)
791

792
  for {pattern, value} <- varints do
793
    def decode_names_and_rows(<<unquote(pattern), rest::bytes>>) do
794
      decode_names(rest, unquote(value), unquote(value), _acc = [])
3,030✔
795
    end
796
  end
797

798
  @doc """
799
  Decodes [RowBinary](https://clickhouse.com/docs/en/interfaces/formats/RowBinary) into rows.
800

801
  Example:
802

803
      iex> decode_rows(<<1>>, ["UInt8"])
804
      [[1]]
805

806
  """
807
  def decode_rows(row_binary, types)
808
  def decode_rows(<<>>, _types), do: []
1✔
809

810
  def decode_rows(<<data::bytes>>, types) do
811
    decode_rows!(data, decoding_types(types))
436✔
812
  end
813

814
  defp decode_rows!(data, types) do
815
    {rows, remaining_data, state} = decode_rows(types, data, [], [], types)
3,441✔
816

817
    case state do
3,423✔
818
      nil ->
819
        rows
3,421✔
820

821
      {:cont, types_rest, row} ->
822
        raise ArgumentError, """
2✔
823
        incomplete RowBinary data: ran out of bytes while decoding
824

825
        Expected to decode: #{inspect(types_rest)}
826
        Remaining bytes: #{byte_size(remaining_data)} bytes
2✔
827
        Partial row: #{inspect(row)}
828
        Completed rows: #{length(rows)}
2✔
829
        """
830
    end
831
  end
832

833
  @doc false
834
  def decode_rows_continue(<<data::bytes>>, types, state) do
835
    case state do
201,104✔
836
      {:cont, types_rest, row} -> decode_rows(types_rest, data, row, [], types)
201,046✔
837
      nil -> decode_rows(types, data, [], [], types)
58✔
838
    end
839
  end
840

841
  @doc false
842
  def decoding_types([type | types]) do
553✔
843
    [decoding_type(type) | decoding_types(types)]
844
  end
845

846
  def decoding_types([] = done), do: done
508✔
847

848
  defp decoding_types_reverse(types), do: decoding_types_reverse(types, [])
3,007✔
849

850
  defp decoding_types_reverse([type | types], acc) do
851
    decoding_types_reverse(types, [decoding_type(type) | acc])
16,447✔
852
  end
853

854
  defp decoding_types_reverse([], acc), do: acc
3,007✔
855

856
  defp decoding_type(t) when is_binary(t) do
857
    decoding_type(Ch.Types.decode(t))
16,764✔
858
  end
859

860
  defp decoding_type(t)
861
       when t in [
862
              :string,
863
              :json,
864
              :dynamic,
865
              :boolean,
866
              :uuid,
867
              :date,
868
              :date32,
869
              :time,
870
              :time64,
871
              :ipv4,
872
              :ipv6,
873
              :point,
874
              :nothing
875
            ],
876
       do: t
3,218✔
877

878
  defp decoding_type({:datetime, _tz} = t), do: t
19✔
879
  defp decoding_type({:fixed_string, _len} = t), do: t
427✔
880

881
  for size <- [8, 16, 32, 64, 128, 256] do
882
    defp decoding_type(unquote(:"u#{size}") = u), do: u
12,176✔
883
    defp decoding_type(unquote(:"i#{size}") = i), do: i
416✔
884
  end
885

886
  for size <- [32, 64] do
887
    defp decoding_type(unquote(:"f#{size}") = f), do: f
449✔
888
  end
889

890
  defp decoding_type(:datetime = t), do: {t, _tz = nil}
13✔
891

892
  defp decoding_type({:array = a, t}), do: {a, decoding_type(t)}
1,336✔
893

894
  defp decoding_type({:tuple = t, ts}) do
320✔
895
    {t, Enum.map(ts, &decoding_type/1)}
896
  end
897

898
  defp decoding_type({:variant = v, ts}) do
17✔
899
    {v, ts |> Enum.map(&decoding_type/1) |> List.to_tuple()}
900
  end
901

902
  defp decoding_type({:map = m, kt, vt}) do
903
    {m, decoding_type(kt), decoding_type(vt)}
321✔
904
  end
905

906
  defp decoding_type({:nullable = n, t}), do: {n, decoding_type(t)}
548✔
907
  defp decoding_type({:low_cardinality, t}), do: decoding_type(t)
277✔
908

909
  defp decoding_type({:decimal = t, p, s}), do: {t, decimal_size(p), s}
355✔
910
  defp decoding_type({:decimal32, s}), do: {:decimal, 32, s}
1✔
911
  defp decoding_type({:decimal64, s}), do: {:decimal, 64, s}
1✔
912
  defp decoding_type({:decimal128, s}), do: {:decimal, 128, s}
1✔
913
  defp decoding_type({:decimal256, s}), do: {:decimal, 256, s}
1✔
914

915
  defp decoding_type({:datetime64 = t, p}), do: {t, time_unit(p), _tz = nil}
6✔
916
  defp decoding_type({:datetime64 = t, p, tz}), do: {t, time_unit(p), tz}
316✔
917

918
  defp decoding_type({:time64 = t, p}), do: {t, time_unit(p)}
262✔
919

920
  defp decoding_type({e, mappings}) when e in [:enum8, :enum16] do
14✔
921
    {e, Map.new(mappings, fn {k, v} -> {v, k} end)}
28✔
922
  end
923

924
  defp decoding_type({:simple_aggregate_function, _f, t}), do: decoding_type(t)
6✔
925

926
  defp decoding_type(:ring), do: {:array, :point}
1✔
927
  defp decoding_type(:polygon), do: {:array, {:array, :point}}
1✔
928
  defp decoding_type(:multipolygon), do: {:array, {:array, {:array, :point}}}
1✔
929

930
  defp decoding_type(type) do
931
    raise ArgumentError, "unsupported type for decoding: #{inspect(type)}"
1✔
932
  end
933

934
  defp skip_names(<<rest::bytes>>, 0, count), do: decode_types(rest, count, _acc = [])
5✔
935

936
  for {pattern, value} <- varints do
937
    defp skip_names(<<unquote(pattern), _::size(unquote(value))-bytes, rest::bytes>>, left, count) do
938
      skip_names(rest, left - 1, count)
75✔
939
    end
940
  end
941

942
  defp decode_names(<<rest::bytes>>, 0, count, names) do
3,030✔
943
    [:lists.reverse(names) | decode_types(rest, count, _acc = [])]
944
  end
945

946
  for {pattern, value} <- varints do
947
    defp decode_names(
948
           <<unquote(pattern), name::size(unquote(value))-bytes, rest::bytes>>,
949
           left,
950
           count,
951
           acc
952
         ) do
953
      decode_names(rest, left - 1, count, [name | acc])
16,447✔
954
    end
955
  end
956

957
  defp decode_types(<<>>, 0, _types), do: []
29✔
958

959
  defp decode_types(<<rest::bytes>>, 0, types) do
960
    decode_rows!(rest, decoding_types_reverse(types))
3,006✔
961
  end
962

963
  for {pattern, value} <- varints do
964
    defp decode_types(
965
           <<unquote(pattern), type::size(unquote(value))-bytes, rest::bytes>>,
966
           count,
967
           acc
968
         ) do
969
      decode_types(rest, count - 1, [type | acc])
16,522✔
970
    end
971
  end
972

973
  @compile inline: [decode_string_decode_rows: 5]
974

975
  for {pattern, size} <- varints do
976
    defp decode_string_decode_rows(
977
           <<unquote(pattern), s::size(unquote(size))-bytes, bin::bytes>>,
978
           types_rest,
979
           row,
980
           rows,
981
           types
982
         ) do
983
      decode_rows(types_rest, bin, [s | row], rows, types)
9,368✔
984
    end
985
  end
986

987
  defp decode_string_decode_rows(<<bin::bytes>>, types_rest, row, rows, _types) do
988
    to_be_continued(rows, bin, [:string | types_rest], row)
200,144✔
989
  end
990

991
  @compile inline: [decode_string_json_decode_rows: 5]
992

993
  for {pattern, size} <- varints do
994
    defp decode_string_json_decode_rows(
995
           <<unquote(pattern), s::size(unquote(size))-bytes, bin::bytes>>,
996
           types_rest,
997
           row,
998
           rows,
999
           types
1000
         ) do
1001
      decode_rows(types_rest, bin, [JSON.decode!(s) | row], rows, types)
43✔
1002
    end
1003
  end
1004

1005
  defp decode_string_json_decode_rows(<<bin::bytes>>, types_rest, row, rows, _types) do
1006
    to_be_continued(rows, bin, [:json | types_rest], row)
46✔
1007
  end
1008

1009
  @compile inline: [decode_array_decode_rows: 6]
1010
  defp decode_array_decode_rows(<<0, bin::bytes>>, _type, types_rest, row, rows, types) do
1011
    decode_rows(types_rest, bin, [[] | row], rows, types)
401✔
1012
  end
1013

1014
  for {pattern, size} <- varints do
1015
    defp decode_array_decode_rows(
1016
           <<unquote(pattern), bin::bytes>>,
1017
           type,
1018
           types_rest,
1019
           row,
1020
           rows,
1021
           types
1022
         ) do
1023
      array_types = List.duplicate(type, unquote(size))
2,687✔
1024
      types_rest = array_types ++ [{:array_over, row} | types_rest]
2,687✔
1025
      decode_rows(types_rest, bin, [], rows, types)
2,687✔
1026
    end
1027
  end
1028

1029
  defp decode_array_decode_rows(<<bin::bytes>>, type, types_rest, row, rows, _types) do
1030
    to_be_continued(rows, bin, [{:array, type} | types_rest], row)
12✔
1031
  end
1032

1033
  @compile inline: [decode_map_decode_rows: 7]
1034
  defp decode_map_decode_rows(
1035
         <<0, bin::bytes>>,
1036
         _key_type,
1037
         _value_type,
1038
         types_rest,
1039
         row,
1040
         rows,
1041
         types
1042
       ) do
1043
    decode_rows(types_rest, bin, [%{} | row], rows, types)
36✔
1044
  end
1045

1046
  for {pattern, size} <- varints do
1047
    defp decode_map_decode_rows(
1048
           <<unquote(pattern), bin::bytes>>,
1049
           key_type,
1050
           value_type,
1051
           types_rest,
1052
           row,
1053
           rows,
1054
           types
1055
         ) do
1056
      types_rest =
291✔
1057
        map_types(unquote(size), key_type, value_type) ++ [{:map_over, row} | types_rest]
1058

1059
      decode_rows(types_rest, bin, [], rows, types)
291✔
1060
    end
1061
  end
1062

1063
  defp decode_map_decode_rows(<<bin::bytes>>, key_type, value_type, types_rest, row, rows, _types) do
1064
    to_be_continued(rows, bin, [{:map, key_type, value_type} | types_rest], row)
6✔
1065
  end
1066

1067
  defp map_types(count, key_type, value_type) when count > 0 do
1,270✔
1068
    [key_type, value_type | map_types(count - 1, key_type, value_type)]
1069
  end
1070

1071
  defp map_types(0, _key_type, _value_types), do: []
291✔
1072

1073
  # https://clickhouse.com/docs/sql-reference/data-types/data-types-binary-encoding
1074
  dynamic_types = [
1075
    nothing: 0x00,
1076
    u8: 0x01,
1077
    u16: 0x02,
1078
    u32: 0x03,
1079
    u64: 0x04,
1080
    u128: 0x05,
1081
    u256: 0x06,
1082
    i8: 0x07,
1083
    i16: 0x08,
1084
    i32: 0x09,
1085
    i64: 0x0A,
1086
    i128: 0x0B,
1087
    i256: 0x0C,
1088
    f32: 0x0D,
1089
    f64: 0x0E,
1090
    date: 0x0F,
1091
    date32: 0x10,
1092
    string: 0x15,
1093
    uuid: 0x1D,
1094
    ipv4: 0x28,
1095
    ipv6: 0x29,
1096
    boolean: 0x2D
1097
  ]
1098

1099
  # TODO compile inline?
1100

1101
  for {type, code} <- dynamic_types do
1102
    defp decode_dynamic(
1103
           <<unquote(code), rest::bytes>>,
1104
           dynamic,
1105
           types_rest,
1106
           row,
1107
           rows,
1108
           types
1109
         ) do
1110
      decode_dynamic_continue(rest, [unquote(type) | dynamic], types_rest, row, rows, types)
120✔
1111
    end
1112
  end
1113

1114
  # DateTime 0x11
1115
  defp decode_dynamic(<<0x11, rest::bytes>>, dynamic, types_rest, row, rows, types) do
1116
    decode_dynamic_continue(rest, [{:datetime, nil} | dynamic], types_rest, row, rows, types)
2✔
1117
  end
1118

1119
  # DateTime(time_zone) 0x12 <var_uint_time_zone_name_size><time_zone_name_data>
1120
  for {pattern, size} <- varints do
1121
    defp decode_dynamic(
1122
           <<0x12, unquote(pattern), tz::size(unquote(size))-bytes, rest::bytes>>,
1123
           dynamic,
1124
           types_rest,
1125
           row,
1126
           rows,
1127
           types
1128
         ) do
1129
      decode_dynamic_continue(rest, [{:datetime, tz} | dynamic], types_rest, row, rows, types)
1✔
1130
    end
1131
  end
1132

1133
  # DateTime64(P) 0x13 <uint8_precision>
1134
  defp decode_dynamic(
1135
         <<0x13, precision, rest::bytes>>,
1136
         dynamic,
1137
         types_rest,
1138
         row,
1139
         rows,
1140
         types
1141
       ) do
1142
    decode_dynamic_continue(
1✔
1143
      rest,
1144
      [decoding_type({:datetime64, precision}) | dynamic],
1145
      types_rest,
1146
      row,
1147
      rows,
1148
      types
1149
    )
1150
  end
1151

1152
  # DateTime64(P, time_zone) 0x14 <uint8_precision><var_uint_time_zone_name_size><time_zone_name_data>
1153
  for {pattern, size} <- varints do
1154
    defp decode_dynamic(
1155
           <<0x14, precision, unquote(pattern), tz::size(unquote(size))-bytes, rest::bytes>>,
1156
           dynamic,
1157
           types_rest,
1158
           row,
1159
           rows,
1160
           types
1161
         ) do
1162
      decode_dynamic_continue(
1✔
1163
        rest,
1164
        [decoding_type({:datetime64, precision, tz}) | dynamic],
1165
        types_rest,
1166
        row,
1167
        rows,
1168
        types
1169
      )
1170
    end
1171
  end
1172

1173
  # FixedString(N) 0x16 <var_uint_size>
1174
  for {pattern, size} <- varints do
1175
    defp decode_dynamic(
1176
           <<0x16, unquote(pattern), rest::bytes>>,
1177
           dynamic,
1178
           types_rest,
1179
           row,
1180
           rows,
1181
           types
1182
         ) do
1183
      decode_dynamic_continue(
2✔
1184
        rest,
1185
        [{:fixed_string, unquote(size)} | dynamic],
1186
        types_rest,
1187
        row,
1188
        rows,
1189
        types
1190
      )
1191
    end
1192
  end
1193

1194
  # Decimal32(P, S) 0x19 <uint8_precision><uint8_scale>
1195
  # Decimal64(P, S) 0x1A <uint8_precision><uint8_scale>
1196
  # Decimal128(P, S) 0x1B <uint8_precision><uint8_scale>
1197
  # Decimal256(P, S) 0x1C <uint8_precision><uint8_scale>
1198
  for {code, size} <- [{0x19, 32}, {0x1A, 64}, {0x1B, 128}, {0x1C, 256}] do
1199
    defp decode_dynamic(
1200
           <<unquote(code), _precision, scale, rest::bytes>>,
1201
           dynamic,
1202
           types_rest,
1203
           row,
1204
           rows,
1205
           types
1206
         ) do
1207
      decode_dynamic_continue(
4✔
1208
        rest,
1209
        [{:decimal, unquote(size), scale} | dynamic],
1210
        types_rest,
1211
        row,
1212
        rows,
1213
        types
1214
      )
1215
    end
1216
  end
1217

1218
  # Array(T) 0x1E <nested_type_encoding>
1219
  defp decode_dynamic(<<0x1E, rest::bytes>>, dynamic, types_rest, row, rows, types) do
1220
    decode_dynamic_continue(rest, [:array | dynamic], types_rest, row, rows, types)
29✔
1221
  end
1222

1223
  # Nullable(T)        0x23 <nested_type_encoding>
1224
  defp decode_dynamic(<<0x23, rest::bytes>>, dynamic, types_rest, row, rows, types) do
1225
    decode_dynamic_continue(rest, [:nullable | dynamic], types_rest, row, rows, types)
5✔
1226
  end
1227

1228
  # LowCardinality(T) 0x26 <nested_type_encoding>
1229
  defp decode_dynamic(<<0x26, rest::bytes>>, dynamic, types_rest, row, rows, types) do
1230
    decode_dynamic_continue(rest, [:low_cardinality | dynamic], types_rest, row, rows, types)
2✔
1231
  end
1232

1233
  # TODO
1234
  # Enum8        0x17 <var_uint_number_of_elements><var_uint_name_size_1><name_data_1><int8_value_1>...<var_uint_name_size_N><name_data_N><int8_value_N>
1235
  # Enum16        0x18 <var_uint_number_of_elements><var_uint_name_size_1><name_data_1><int16_little_endian_value_1>...><var_uint_name_size_N><name_data_N><int16_little_endian_value_N>
1236
  # Tuple(T1, ..., TN)        0x1F <var_uint_number_of_elements><nested_type_encoding_1>...<nested_type_encoding_N>
1237
  # Tuple(name1 T1, ..., nameN TN)        0x20 <var_uint_number_of_elements><var_uint_name_size_1><name_data_1><nested_type_encoding_1>...<var_uint_name_size_N><name_data_N><nested_type_encoding_N>
1238
  # Set        0x21
1239
  # Interval        0x22 <interval_kind> (see interval kind binary encoding)
1240
  # Function        0x24<var_uint_number_of_arguments><argument_type_encoding_1>...<argument_type_encoding_N><return_type_encoding>
1241
  # AggregateFunction(function_name(param_1, ..., param_N), arg_T1, ..., arg_TN)        0x25<var_uint_version><var_uint_function_name_size><function_name_data><var_uint_number_of_parameters><param_1>...<param_N><var_uint_number_of_arguments><argument_type_encoding_1>...<argument_type_encoding_N> (see aggregate function parameter binary encoding)
1242
  # Map(K, V)        0x27<key_type_encoding><value_type_encoding>
1243
  # Variant(T1, ..., TN)        0x2A<var_uint_number_of_variants><variant_type_encoding_1>...<variant_type_encoding_N>
1244
  # Dynamic(max_types=N)        0x2B<uint8_max_types>
1245
  # Custom type (Ring, Polygon, etc)        0x2C<var_uint_type_name_size><type_name_data>
1246
  # SimpleAggregateFunction(function_name(param_1, ..., param_N), arg_T1, ..., arg_TN)        0x2E<var_uint_function_name_size><function_name_data><var_uint_number_of_parameters><param_1>...<param_N><var_uint_number_of_arguments><argument_type_encoding_1>...<argument_type_encoding_N> (see aggregate function parameter binary encoding)
1247
  # Nested(name1 T1, ..., nameN TN)        0x2F<var_uint_number_of_elements><var_uint_name_size_1><name_data_1><nested_type_encoding_1>...<var_uint_name_size_N><name_data_N><nested_type_encoding_N>
1248
  # JSON(max_dynamic_paths=N, max_dynamic_types=M, path Type, SKIP skip_path, SKIP REGEXP skip_path_regexp)        0x30<uint8_serialization_version><var_int_max_dynamic_paths><uint8_max_dynamic_types><var_uint_number_of_typed_paths><var_uint_path_name_size_1><path_name_data_1><encoded_type_1>...<var_uint_number_of_skip_paths><var_uint_skip_path_size_1><skip_path_data_1>...<var_uint_number_of_skip_path_regexps><var_uint_skip_path_regexp_size_1><skip_path_data_regexp_1>...
1249

1250
  unsupported_dynamic_types = %{
1251
    "Enum8" => 0x17,
1252
    "Enum16" => 0x18,
1253
    "Tuple" => 0x1F,
1254
    "TupleWithNames" => 0x20,
1255
    "Set" => 0x21,
1256
    "Interval" => 0x22,
1257
    "Function" => 0x24,
1258
    "AggregateFunction" => 0x25,
1259
    "Map" => 0x27,
1260
    "Variant" => 0x2A,
1261
    "Dynamic" => 0x2B,
1262
    "CustomType" => 0x2C,
1263
    "SimpleAggregateFunction" => 0x2E,
1264
    "Nested" => 0x2F,
1265
    "JSON" => 0x30
1266
  }
1267

1268
  for {type, code} <- unsupported_dynamic_types do
1269
    defp decode_dynamic(<<unquote(code), _::bytes>>, _dynamic, _types_rest, _row, _rows, _types) do
1270
      raise ArgumentError, "unsupported dynamic type #{unquote(type)}"
9✔
1271
    end
1272
  end
1273

1274
  defp decode_dynamic(<<bin::bytes>>, dynamic, types_rest, row, rows, _types) do
1275
    to_be_continued(rows, bin, [{:dynamic, dynamic} | types_rest], row)
2✔
1276
  end
1277

1278
  @compile inline: [decode_dynamic_continue: 6]
1279

1280
  defp decode_dynamic_continue(<<rest::bytes>>, dynamic, types_rest, row, rows, types) do
1281
    continue? =
103✔
1282
      case dynamic do
1283
        [:array | _] -> true
29✔
1284
        [:nullable | _] -> true
5✔
1285
        [:low_cardinality | _] -> true
2✔
1286
        _ -> false
103✔
1287
      end
1288

1289
    if continue? do
103✔
1290
      decode_dynamic(rest, dynamic, types_rest, row, rows, types)
36✔
1291
    else
1292
      type = build_dynamic_type(:lists.reverse(dynamic))
103✔
1293
      decode_rows([type | types_rest], rest, row, rows, types)
103✔
1294
    end
1295
  end
1296

1297
  defp build_dynamic_type([type]), do: type
131✔
1298

1299
  defp build_dynamic_type(type) do
1300
    case type do
32✔
1301
      [:array | rest] -> {:array, build_dynamic_type(rest)}
25✔
1302
      [:nullable | rest] -> {:nullable, build_dynamic_type(rest)}
5✔
1303
      [:low_cardinality | rest] -> build_dynamic_type(rest)
2✔
1304
    end
1305
  end
1306

1307
  simple_types = %{
1308
    u8: %{pattern: quote(do: <<u>>), value: quote(do: u)},
1309
    u16: %{pattern: quote(do: <<u::16-little>>), value: quote(do: u)},
1310
    u32: %{pattern: quote(do: <<u::32-little>>), value: quote(do: u)},
1311
    u64: %{pattern: quote(do: <<u::64-little>>), value: quote(do: u)},
1312
    u128: %{pattern: quote(do: <<u::128-little>>), value: quote(do: u)},
1313
    u256: %{pattern: quote(do: <<u::256-little>>), value: quote(do: u)},
1314
    i8: %{pattern: quote(do: <<i::signed>>), value: quote(do: i)},
1315
    i16: %{pattern: quote(do: <<i::16-little-signed>>), value: quote(do: i)},
1316
    i32: %{pattern: quote(do: <<i::32-little-signed>>), value: quote(do: i)},
1317
    i64: %{pattern: quote(do: <<i::64-little-signed>>), value: quote(do: i)},
1318
    i128: %{pattern: quote(do: <<i::128-little-signed>>), value: quote(do: i)},
1319
    i256: %{pattern: quote(do: <<i::256-little-signed>>), value: quote(do: i)},
1320
    f32: [
1321
      %{pattern: quote(do: <<f::32-little-float>>), value: quote(do: f)},
1322
      %{pattern: quote(do: <<_nan_or_inf::32>>), value: quote(do: nil)}
1323
    ],
1324
    f64: [
1325
      %{pattern: quote(do: <<f::64-little-float>>), value: quote(do: f)},
1326
      %{pattern: quote(do: <<_nan_or_inf::64>>), value: quote(do: nil)}
1327
    ],
1328
    uuid: %{
1329
      pattern: quote(do: <<u1::64-little, u2::64-little>>),
1330
      value: quote(do: <<u1::64, u2::64>>)
1331
    },
1332
    date: %{
1333
      pattern: quote(do: <<d::16-little>>),
1334
      value: quote(do: Date.from_gregorian_days(d + @epoch_gregorian_days))
1335
    },
1336
    date32: %{
1337
      pattern: quote(do: <<d::32-little-signed>>),
1338
      value: quote(do: Date.from_gregorian_days(d + @epoch_gregorian_days))
1339
    },
1340
    time: %{
1341
      pattern: quote(do: <<s::32-little-signed>>),
1342
      value: quote(do: time_after_midnight(s, 1))
1343
    },
1344
    boolean: [
1345
      %{pattern: quote(do: <<0>>), value: quote(do: false)},
1346
      %{pattern: quote(do: <<1>>), value: quote(do: true)},
1347
      %{pattern: quote(do: <<b>>), value: quote(do: raise("invalid boolean value: #{b}"))}
1348
    ],
1349
    ipv4: %{
1350
      pattern: quote(do: <<b4, b3, b2, b1>>),
1351
      value: quote(do: {b1, b2, b3, b4})
1352
    },
1353
    ipv6: %{
1354
      pattern: quote(do: <<b1::16, b2::16, b3::16, b4::16, b5::16, b6::16, b7::16, b8::16>>),
1355
      value: quote(do: {b1, b2, b3, b4, b5, b6, b7, b8})
1356
    },
1357
    point: %{
1358
      pattern: quote(do: <<x::64-little-float, y::64-little-float>>),
1359
      value: quote(do: {x, y})
1360
    }
1361
  }
1362

1363
  for {type, clauses} <- simple_types do
1364
    fun = :"decode_#{type}_decode_rows"
1365
    @compile inline: [{fun, 5}]
1366

1367
    for %{pattern: pattern, value: value} <- List.wrap(clauses) do
1368
      defp unquote(fun)(<<unquote(pattern), rest::bytes>>, types_rest, row, rows, types) do
1369
        decode_rows(types_rest, rest, [unquote(value) | row], rows, types)
2,021,534✔
1370
      end
1371
    end
1372

1373
    defp unquote(fun)(<<bin::bytes>>, types_rest, row, rows, _types) do
1374
      to_be_continued(rows, bin, [unquote(type) | types_rest], row)
583✔
1375
    end
1376
  end
1377

1378
  defp decode_rows([type | types_rest], <<bin::bytes>>, row, rows, types) do
1379
    case type do
2,247,728✔
1380
      :u8 ->
1381
        decode_u8_decode_rows(bin, types_rest, row, rows, types)
5,874✔
1382

1383
      :u16 ->
1384
        decode_u16_decode_rows(bin, types_rest, row, rows, types)
10,302✔
1385

1386
      :u32 ->
1387
        decode_u32_decode_rows(bin, types_rest, row, rows, types)
124✔
1388

1389
      :u64 ->
1390
        decode_u64_decode_rows(bin, types_rest, row, rows, types)
2,000,238✔
1391

1392
      :u128 ->
1393
        decode_u128_decode_rows(bin, types_rest, row, rows, types)
63✔
1394

1395
      :u256 ->
1396
        decode_u256_decode_rows(bin, types_rest, row, rows, types)
95✔
1397

1398
      :i8 ->
1399
        decode_i8_decode_rows(bin, types_rest, row, rows, types)
165✔
1400

1401
      :i16 ->
1402
        decode_i16_decode_rows(bin, types_rest, row, rows, types)
175✔
1403

1404
      :i32 ->
1405
        decode_i32_decode_rows(bin, types_rest, row, rows, types)
71✔
1406

1407
      :i64 ->
1408
        decode_i64_decode_rows(bin, types_rest, row, rows, types)
170✔
1409

1410
      :i128 ->
1411
        decode_i128_decode_rows(bin, types_rest, row, rows, types)
89✔
1412

1413
      :i256 ->
1414
        decode_i256_decode_rows(bin, types_rest, row, rows, types)
89✔
1415

1416
      :f32 ->
1417
        decode_f32_decode_rows(bin, types_rest, row, rows, types)
836✔
1418

1419
      :f64 ->
1420
        decode_f64_decode_rows(bin, types_rest, row, rows, types)
939✔
1421

1422
      :string ->
1423
        decode_string_decode_rows(bin, types_rest, row, rows, types)
209,512✔
1424

1425
      :json ->
1426
        # assuming it arrives as text and not "native" binary JSON
1427
        # i.e. assumes `settings: [output_format_binary_write_json_as_string: 1]`
1428
        # TODO
1429
        decode_string_json_decode_rows(bin, types_rest, row, rows, types)
89✔
1430

1431
      :dynamic ->
1432
        decode_dynamic(bin, _dynamic = [], types_rest, row, rows, types)
140✔
1433

1434
      {:dynamic, dynamic} ->
1435
        decode_dynamic(bin, dynamic, types_rest, row, rows, types)
2✔
1436

1437
      {:fixed_string, size} ->
1438
        case bin do
4,408✔
1439
          <<s::size(^size)-bytes, rest::bytes>> ->
1440
            decode_rows(types_rest, rest, [s | row], rows, types)
4,394✔
1441

1442
          _ ->
1443
            to_be_continued(rows, bin, [type | types_rest], row)
14✔
1444
        end
1445

1446
      :boolean ->
1447
        decode_boolean_decode_rows(bin, types_rest, row, rows, types)
2,055✔
1448

1449
      :uuid ->
1450
        decode_uuid_decode_rows(bin, types_rest, row, rows, types)
182✔
1451

1452
      :date ->
1453
        decode_date_decode_rows(bin, types_rest, row, rows, types)
153✔
1454

1455
      :date32 ->
1456
        decode_date32_decode_rows(bin, types_rest, row, rows, types)
46✔
1457

1458
      :time ->
1459
        decode_time_decode_rows(bin, types_rest, row, rows, types)
230✔
1460

1461
      {:time64, time_unit} ->
1462
        case bin do
300✔
1463
          <<ticks::64-little-signed, bin::bytes>> ->
1464
            time = time_after_midnight(ticks, time_unit)
270✔
1465
            decode_rows(types_rest, bin, [time | row], rows, types)
265✔
1466

1467
          _ ->
1468
            to_be_continued(rows, bin, [type | types_rest], row)
30✔
1469
        end
1470

1471
      {:datetime, timezone} ->
1472
        case bin do
67✔
1473
          <<s::32-little, bin::bytes>> ->
1474
            dt = DateTime.from_unix!(s)
45✔
1475

1476
            dt =
45✔
1477
              case timezone do
1478
                nil -> DateTime.to_naive(dt)
16✔
1479
                "UTC" -> dt
23✔
1480
                _ -> DateTime.shift_zone!(dt, timezone)
6✔
1481
              end
1482

1483
            decode_rows(types_rest, bin, [dt | row], rows, types)
45✔
1484

1485
          _ ->
1486
            to_be_continued(rows, bin, [type | types_rest], row)
22✔
1487
        end
1488

1489
      {:decimal, size, scale} ->
1490
        case bin do
492✔
1491
          <<val::size(^size)-little-signed, bin::bytes>> ->
1492
            sign = if val < 0, do: -1, else: 1
374✔
1493
            d = Decimal.new(sign, abs(val), -scale)
374✔
1494
            decode_rows(types_rest, bin, [d | row], rows, types)
374✔
1495

1496
          _ ->
1497
            to_be_continued(rows, bin, [type | types_rest], row)
118✔
1498
        end
1499

1500
      {:nullable, inner_type} ->
1501
        case bin do
2,765✔
1502
          <<b, bin::bytes>> ->
1503
            case b do
2,762✔
1504
              0 -> decode_rows([inner_type | types_rest], bin, row, rows, types)
1,377✔
1505
              1 -> decode_rows(types_rest, bin, [nil | row], rows, types)
1,385✔
1506
            end
1507

1508
          _ ->
1509
            to_be_continued(rows, bin, [type | types_rest], row)
3✔
1510
        end
1511

1512
      :nothing ->
1513
        decode_rows(types_rest, bin, [nil | row], rows, types)
27✔
1514

1515
      {:array, inner_type} ->
1516
        decode_array_decode_rows(bin, inner_type, types_rest, row, rows, types)
3,100✔
1517

1518
      {:array_over, original_row} ->
1519
        decode_rows(types_rest, bin, [:lists.reverse(row) | original_row], rows, types)
2,686✔
1520

1521
      {:map, key_type, value_type} ->
1522
        decode_map_decode_rows(bin, key_type, value_type, types_rest, row, rows, types)
333✔
1523

1524
      {:map_over, original_row} ->
1525
        map = row |> Enum.chunk_every(2) |> Enum.map(fn [v, k] -> {k, v} end) |> Map.new()
291✔
1526
        decode_rows(types_rest, bin, [map | original_row], rows, types)
291✔
1527

1528
      {:tuple, tuple_types} ->
1529
        decode_rows(tuple_types ++ [{:tuple_over, row} | types_rest], bin, [], rows, types)
340✔
1530

1531
      {:tuple_over, original_row} ->
1532
        tuple = row |> :lists.reverse() |> List.to_tuple()
340✔
1533
        decode_rows(types_rest, bin, [tuple | original_row], rows, types)
340✔
1534

1535
      {:variant, variant_types} ->
1536
        case bin do
35✔
1537
          <<255, bin::bytes>> ->
1538
            # 255 is the variant type index for "nothing"
1539
            decode_rows(types_rest, bin, [nil | row], rows, types)
7✔
1540

1541
          # TODO varint?
1542
          <<variant_type_index::8, bin::bytes>>
1543
          when variant_type_index < tuple_size(variant_types) ->
1544
            variant_type = elem(variant_types, variant_type_index)
24✔
1545
            decode_rows([variant_type | types_rest], bin, row, rows, types)
24✔
1546

1547
          <<variant_type_index::8, _bin::bytes>> ->
1548
            raise ArgumentError, "invalid Variant type index: #{variant_type_index}"
1✔
1549

1550
          _ ->
1551
            to_be_continued(rows, bin, [type | types_rest], row)
3✔
1552
        end
1553

1554
      {:datetime64, time_unit, timezone} ->
1555
        case bin do
651✔
1556
          <<s::64-little-signed, bin::bytes>> ->
1557
            dt = DateTime.from_unix!(s, time_unit)
589✔
1558

1559
            dt =
589✔
1560
              case timezone do
1561
                nil -> DateTime.to_naive(dt)
7✔
1562
                "UTC" -> dt
576✔
1563
                _ -> DateTime.shift_zone!(dt, timezone)
6✔
1564
              end
1565

1566
            decode_rows(types_rest, bin, [dt | row], rows, types)
589✔
1567

1568
          _ ->
1569
            to_be_continued(rows, bin, [type | types_rest], row)
62✔
1570
        end
1571

1572
      {:enum8, mapping} ->
1573
        case bin do
27✔
1574
          <<v::signed, bin::bytes>> ->
1575
            decode_rows(types_rest, bin, [Map.fetch!(mapping, v) | row], rows, types)
26✔
1576

1577
          _ ->
1578
            to_be_continued(rows, bin, [type | types_rest], row)
1✔
1579
        end
1580

1581
      {:enum16, mapping} ->
1582
        case bin do
6✔
1583
          <<v::16-little-signed, bin::bytes>> ->
1584
            decode_rows(types_rest, bin, [Map.fetch!(mapping, v) | row], rows, types)
2✔
1585

1586
          _ ->
1587
            to_be_continued(rows, bin, [type | types_rest], row)
4✔
1588
        end
1589

1590
      :ipv4 ->
1591
        decode_ipv4_decode_rows(bin, types_rest, row, rows, types)
34✔
1592

1593
      :ipv6 ->
1594
        decode_ipv6_decode_rows(bin, types_rest, row, rows, types)
79✔
1595

1596
      :point ->
1597
        decode_point_decode_rows(bin, types_rest, row, rows, types)
108✔
1598
    end
1599
  end
1600

1601
  defp decode_rows([], <<>> = empty, row, rows, _types) do
1602
    rows = :lists.reverse([:lists.reverse(row) | rows])
3,477✔
1603
    {rows, empty, _no_state = nil}
3,477✔
1604
  end
1605

1606
  defp decode_rows([], <<bin::bytes>>, row, rows, types) do
1607
    row = :lists.reverse(row)
2,003,950✔
1608
    decode_rows(types, bin, [], [row | rows], types)
2,003,950✔
1609
  end
1610

1611
  @compile inline: [to_be_continued: 4]
1612
  defp to_be_continued(rows, bin, types_rest, row) do
1613
    {:lists.reverse(rows), bin, {:cont, types_rest, row}}
200,697✔
1614
  end
1615

1616
  @compile inline: [decimal_size: 1]
1617
  # https://clickhouse.com/docs/en/sql-reference/data-types/decimal/
1618
  defp decimal_size(precision) when is_integer(precision) do
1619
    cond do
364✔
1620
      precision >= 39 -> 256
207✔
1621
      precision >= 19 -> 128
156✔
1622
      precision >= 10 -> 64
150✔
1623
      true -> 32
18✔
1624
    end
1625
  end
1626

1627
  @compile inline: [time_unit: 1]
1628
  for precision <- 0..9 do
1629
    time_unit = Integer.pow(10, precision)
1630
    defp time_unit(unquote(precision)), do: unquote(time_unit)
692✔
1631
  end
1632

1633
  @compile inline: [time_after_midnight: 2]
1634
  defp time_after_midnight(ticks, time_unit) do
1635
    if ticks >= 0 and ticks < 86400 * time_unit do
486✔
1636
      ticks |> DateTime.from_unix!(time_unit) |> DateTime.to_time()
478✔
1637
    else
1638
      # since ClickHouse supports Time64 values of [-999:59:59.999999999, 999:59:59.999999999]
1639
      # and Elixir's Time supports values of [00:00:00.000000, 23:59:59.999999]
1640
      # we raise an error when ClickHouse's Time64 value is out of Elixir's Time range
1641
      raise ArgumentError,
8✔
1642
            "ClickHouse Time value #{:erlang.float_to_binary(ticks / time_unit, [:short])} (seconds) is out of Elixir's Time range (00:00:00.000000 - 23:59:59.999999)"
8✔
1643

1644
      # TODO: we could potentially decode ClickHouse's Time/Time64 values as Elixir's Duration when it's out of Elixir's Time range
1645
    end
1646
  end
1647
end
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc