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

plausible / ch / 97d475eb0340305bbc2004d75aa8d9481ce8a97a-PR-410

03 Aug 2026 01:39PM UTC coverage: 97.538% (-0.5%) from 98.062%
97d475eb0340305bbc2004d75aa8d9481ce8a97a-PR-410

Pull #410

github

ruslandoga
Harden RowBinary boundary validation
Pull Request #410: Strengthen test ownership and RowBinary boundary safety

110 of 116 new or added lines in 2 files covered. (94.83%)

832 of 853 relevant lines covered (97.54%)

13943.39 hits per line

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

99.06
/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
  # @compile {:bin_opt_info, true}
5
  @dialyzer :no_improper_lists
6

7
  import Bitwise
8

9
  @epoch_gregorian_seconds 62_167_219_200
10
  @epoch_gregorian_days 719_528
11

12
  @doc false
13
  def encode_names_and_types(names, types) do
5✔
14
    [encode(:varint, length(names)), encode_many(names, :string), encode_types(types)]
15
  end
16

17
  defp encode_types([type | types]) do
12✔
18
    encoded =
12✔
19
      case type do
20
        _ when is_binary(type) -> type
11✔
21
        _ -> Ch.Types.encode(type)
1✔
22
      end
23

24
    [encode(:string, encoded) | encode_types(types)]
25
  end
26

27
  defp encode_types([] = done), do: done
5✔
28

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

32
  Examples:
33

34
      iex> encode_row([], [])
35
      []
36

37
      iex> encode_row([1], ["UInt8"])
38
      [1]
39

40
      iex> encode_row([3, "hello"], ["UInt8", "String"])
41
      [3, [5 | "hello"]]
42

43
  """
44
  def encode_row(row, types) do
45
    _encode_row(row, encoding_types(types))
23✔
46
  end
47

48
  defp _encode_row([el | els], [type | types]), do: [encode(type, el) | _encode_row(els, types)]
52✔
49
  defp _encode_row([] = done, []), do: done
23✔
50

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

54
  Examples:
55

56
      iex> encode_rows([], [])
57
      []
58

59
      iex> encode_rows([[1]], ["UInt8"])
60
      [1]
61

62
      iex> encode_rows([[3, "hello"], [4, "hi"]], ["UInt8", "String"])
63
      [3, [5 | "hello"], 4, [2 | "hi"]]
64

65
  """
66
  def encode_rows(rows, types) do
67
    _encode_rows(rows, encoding_types(types))
762✔
68
  end
69

70
  @doc false
71
  def _encode_rows([row | rows], types), do: _encode_rows(row, types, rows, types)
4,260✔
72
  def _encode_rows([] = done, _types), do: done
754✔
73

74
  defp _encode_rows([el | els], [t | ts], rows, types) do
10,779✔
75
    [encode(t, el) | _encode_rows(els, ts, rows, types)]
76
  end
77

78
  defp _encode_rows([], [], rows, types), do: _encode_rows(rows, types)
4,256✔
79

80
  @doc false
81
  def encoding_types([type | types]) do
1,930✔
82
    [encoding_type(type) | encoding_types(types)]
83
  end
84

85
  def encoding_types([] = done), do: done
789✔
86

87
  defp encoding_type(type) when is_binary(type) do
88
    encoding_type(Ch.Types.decode(type))
1,875✔
89
  end
90

91
  defp encoding_type(t)
92
       when t in [
93
              :string,
94
              :json,
95
              :dynamic,
96
              :boolean,
97
              :uuid,
98
              :date,
99
              :datetime,
100
              :date32,
101
              :time,
102
              :ipv4,
103
              :ipv6,
104
              :point,
105
              :nothing
106
            ],
107
       do: t
713✔
108

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

111
  defp encoding_type({:datetime, tz}) do
112
    raise ArgumentError, "can't encode DateTime with non-UTC timezone: #{inspect(tz)}"
1✔
113
  end
114

115
  defp encoding_type({:fixed_string, _len} = t), do: t
309✔
116

117
  for size <- [8, 16, 32, 64, 128, 256] do
118
    defp encoding_type(unquote(:"u#{size}") = u), do: u
543✔
119
    defp encoding_type(unquote(:"i#{size}") = i), do: i
22✔
120
  end
121

122
  for size <- [32, 64] do
123
    defp encoding_type(unquote(:"f#{size}") = f), do: f
220✔
124
  end
125

126
  defp encoding_type({:array = a, t}), do: {a, encoding_type(t)}
535✔
127

128
  defp encoding_type({:tuple = t, ts}) do
5✔
129
    {t, Enum.map(ts, &encoding_type/1)}
130
  end
131

132
  defp encoding_type({:variant = v, ts}) do
3✔
133
    {v, Enum.map(ts, &encoding_type/1)}
134
  end
135

136
  defp encoding_type({:map = m, kt, vt}) do
137
    {m, encoding_type(kt), encoding_type(vt)}
9✔
138
  end
139

140
  defp encoding_type({:nullable = n, t}), do: {n, encoding_type(t)}
113✔
141
  defp encoding_type({:low_cardinality, t}), do: encoding_type(t)
202✔
142

143
  defp encoding_type({:decimal, precision, scale} = type) do
144
    validate_decimal_type!(precision, scale)
10✔
145
    type
6✔
146
  end
147

148
  defp encoding_type({d, _scale} = t)
149
       when d in [:decimal32, :decimal64, :decimal128, :decimal256],
150
       do: t
9✔
151

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

154
  defp encoding_type({:datetime64 = t, p, "UTC"}), do: {t, time_unit(p)}
3✔
155

156
  defp encoding_type({:datetime64, _, tz}) do
157
    raise ArgumentError, "can't encode DateTime64 with non-UTC timezone: #{inspect(tz)}"
1✔
158
  end
159

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

162
  defp encoding_type({e, mappings}) when e in [:enum8, :enum16] do
5✔
163
    Ch.Types.encode({e, mappings})
5✔
164
    {e, Map.new(mappings)}
165
  end
166

167
  defp encoding_type({:simple_aggregate_function, _f, t}), do: encoding_type(t)
1✔
168

169
  defp encoding_type(:ring), do: {:array, :point}
1✔
170
  defp encoding_type(:polygon), do: {:array, {:array, :point}}
1✔
171
  defp encoding_type(:multipolygon), do: {:array, {:array, {:array, :point}}}
1✔
172

173
  defp encoding_type(type) do
174
    raise ArgumentError, "unsupported type for encoding: #{inspect(type)}"
1✔
175
  end
176

177
  @doc false
178
  def encode(type, value)
179

180
  def encode(:varint, i) when is_integer(i) and i >= 0 and i < 128, do: i
7,178✔
181
  def encode(:varint, i) when is_integer(i) and i >= 0, do: encode_varint_cont(i)
7✔
182

183
  def encode(:varint, i) when is_integer(i) do
184
    raise ArgumentError, "invalid varint: #{inspect(i)}"
1✔
185
  end
186

187
  def encode(:string, str) do
188
    case str do
5,191✔
189
      _ when is_binary(str) -> [encode(:varint, byte_size(str)) | str]
5,174✔
190
      _ when is_list(str) -> [encode(:varint, IO.iodata_length(str)) | str]
8✔
191
      nil -> 0
3✔
192
    end
193
  end
194

195
  def encode(:json, json) do
196
    # assuming it can be sent as text and not "native" binary JSON
197
    # i.e. assumes `settings: [input_format_binary_read_json_as_string: 1]`
198
    # TODO
199
    encode(:string, JSON.encode_to_iodata!(json))
5✔
200
  end
201

202
  def encode({:fixed_string, size}, str) when byte_size(str) == size do
203
    str
743✔
204
  end
205

206
  def encode({:fixed_string, size}, str) when byte_size(str) < size do
3,408✔
207
    to_pad = size - byte_size(str)
3,408✔
208
    [str | <<0::size(to_pad * 8)>>]
209
  end
210

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

213
  # UInt8 — [0 : 255]
214
  def encode(:u8, u) when is_integer(u) and u >= 0 and u <= 255, do: u
4,010✔
215
  def encode(:u8, nil), do: 0
4✔
216

217
  def encode(:u8, term) do
218
    raise ArgumentError, "invalid UInt8: #{inspect(term)}"
7✔
219
  end
220

221
  # Int8 — [-128 : 127]
222
  def encode(:i8, i) when is_integer(i) and i >= 0 and i <= 127, do: i
26✔
223
  def encode(:i8, i) when is_integer(i) and i < 0 and i >= -128, do: <<i::signed>>
16✔
224
  def encode(:i8, nil), do: 0
1✔
225

226
  def encode(:i8, term) do
227
    raise ArgumentError, "invalid Int8: #{inspect(term)}"
6✔
228
  end
229

230
  for size <- [16, 32, 64, 128, 256] do
231
    unsigned_max = (1 <<< size) - 1
232
    signed_min = -(1 <<< (size - 1))
233
    signed_max = (1 <<< (size - 1)) - 1
234
    uint = :"u#{size}"
235
    int = :"i#{size}"
236

237
    def encode(unquote(uint), u) when is_integer(u) and u >= 0 and u <= unquote(unsigned_max) do
238
      <<u::unquote(size)-little>>
128✔
239
    end
240

241
    def encode(unquote(int), i)
242
        when is_integer(i) and i >= unquote(signed_min) and i <= unquote(signed_max) do
243
      <<i::unquote(size)-little-signed>>
131✔
244
    end
245

246
    def encode(unquote(uint), nil), do: <<0::unquote(size)>>
3✔
247
    def encode(unquote(int), nil), do: <<0::unquote(size)>>
3✔
248

249
    def encode(unquote(uint), term) do
250
      raise ArgumentError, "invalid UInt#{unquote(size)}: #{inspect(term)}"
15✔
251
    end
252

253
    def encode(unquote(int), term) do
254
      raise ArgumentError, "invalid Int#{unquote(size)}: #{inspect(term)}"
15✔
255
    end
256
  end
257

258
  for {size, max} <- [{32, 3.4028234663852886e38}, {64, 1.7976931348623157e308}] do
259
    type = :"f#{size}"
260

261
    def encode(unquote(type), f) when is_number(f) and f >= -unquote(max) and f <= unquote(max) do
262
      <<f::unquote(size)-little-signed-float>>
1,276✔
263
    end
264

265
    def encode(unquote(type), nil), do: <<0::unquote(size)>>
4✔
266

267
    def encode(unquote(type), term) do
268
      raise ArgumentError, "invalid Float#{unquote(size)}: #{inspect(term)}"
5✔
269
    end
270
  end
271

272
  def encode({:decimal, precision, scale}, %Decimal{} = decimal) do
273
    validate_decimal_type!(precision, scale)
33✔
274
    size = decimal_size(precision)
33✔
275
    coefficient = decimal_coefficient!(decimal, scale, size)
33✔
276

277
    if abs(coefficient) >= Integer.pow(10, precision) do
33✔
278
      raise ArgumentError,
16✔
279
            "Decimal value #{Decimal.to_string(decimal)} exceeds precision #{precision}"
16✔
280
    end
281

282
    encode_decimal!(decimal, coefficient, size, scale)
17✔
283
  end
284

285
  def encode({:decimal, precision, scale}, nil) do
NEW
286
    validate_decimal_type!(precision, scale)
×
NEW
287
    <<0::size(decimal_size(precision))>>
×
288
  end
289

290
  for size <- [32, 64, 128, 256] do
291
    type = :"decimal#{size}"
292
    precision = %{32 => 9, 64 => 18, 128 => 38, 256 => 76}[size]
293

294
    def encode({unquote(type), scale}, %Decimal{} = decimal) do
295
      validate_decimal_scale!(unquote(type), scale, unquote(precision))
143✔
296
      coefficient = decimal_coefficient!(decimal, scale, unquote(size))
139✔
297
      encode_decimal!(decimal, coefficient, unquote(size), scale)
136✔
298
    end
299

300
    def encode({unquote(type), scale}, nil) do
301
      validate_decimal_scale!(unquote(type), scale, unquote(precision))
4✔
302
      <<0::unquote(size)>>
4✔
303
    end
304
  end
305

306
  def encode(:boolean, true), do: 1
751✔
307
  def encode(:boolean, false), do: 0
711✔
308
  def encode(:boolean, nil), do: 0
1✔
309

310
  def encode({:array, type}, [_ | _] = l) do
1,979✔
311
    [encode(:varint, length(l)) | encode_many(l, type)]
312
  end
313

314
  def encode({:array, _type}, []), do: 0
286✔
315
  def encode({:array, _type}, nil), do: 0
4✔
316

317
  def encode({:map, k, v}, [_ | _] = m) do
1✔
318
    [encode(:varint, length(m)) | encode_many_kv(m, k, v)]
319
  end
320

321
  def encode({:map, k, v}, m) when is_map(m) do
17✔
322
    [
323
      encode(:varint, map_size(m))
324
      | :maps.fold(fn key, value, acc -> [encode(k, key), encode(v, value) | acc] end, [], m)
18✔
325
    ]
326
  end
327

328
  def encode({:map, _k, _v}, []), do: 0
1✔
329
  def encode({:map, _k, _v}, nil), do: 0
1✔
330

331
  def encode({:tuple, _types} = t, v) when is_tuple(v) do
332
    encode(t, Tuple.to_list(v))
12✔
333
  end
334

335
  def encode({:tuple, types}, values) when is_list(types) and is_list(values) do
336
    encode_row(values, types)
12✔
337
  end
338

339
  def encode({:tuple, types}, nil) when is_list(types) do
340
    Enum.map(types, fn type -> encode(type, nil) end)
1✔
341
  end
342

343
  def encode({:variant, _types}, nil), do: 255
3✔
344

345
  def encode({:variant, types}, value) do
346
    try_encode_variant(types, 0, value)
8✔
347
  end
348

349
  def encode(:datetime, %NaiveDateTime{} = datetime) do
350
    {seconds, _micros} = NaiveDateTime.to_gregorian_seconds(datetime)
14✔
351
    encode_fixed_integer!(seconds - @epoch_gregorian_seconds, 32, :unsigned, "DateTime")
14✔
352
  end
353

354
  def encode(:datetime, %DateTime{} = datetime) do
355
    datetime
356
    |> DateTime.to_unix(:second)
357
    |> encode_fixed_integer!(32, :unsigned, "DateTime")
9✔
358
  end
359

360
  def encode(:datetime, nil), do: <<0::32>>
1✔
361

362
  def encode({:datetime64, time_unit}, %NaiveDateTime{} = datetime) do
363
    {seconds, micros} = NaiveDateTime.to_gregorian_seconds(datetime)
5✔
364
    ticks = (seconds - @epoch_gregorian_seconds) * time_unit + div(micros * time_unit, 1_000_000)
5✔
365
    encode_fixed_integer!(ticks, 64, :signed, "DateTime64")
5✔
366
  end
367

368
  def encode({:datetime64, time_unit}, %DateTime{} = datetime) do
369
    datetime
370
    |> DateTime.to_unix(time_unit)
371
    |> encode_fixed_integer!(64, :signed, "DateTime64")
9✔
372
  end
373

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

376
  def encode(:date, %Date{} = date) do
377
    date
378
    |> Date.to_gregorian_days()
379
    |> Kernel.-(@epoch_gregorian_days)
380
    |> encode_fixed_integer!(16, :unsigned, "Date")
16✔
381
  end
382

383
  def encode(:date, nil), do: <<0::16>>
1✔
384

385
  def encode(:date32, %Date{} = date) do
386
    <<Date.to_gregorian_days(date) - @epoch_gregorian_days::32-little-signed>>
6✔
387
  end
388

389
  def encode(:date32, nil), do: <<0::32>>
1✔
390

391
  def encode(:time, %Time{} = time) do
392
    {s, _micros} = Time.to_seconds_after_midnight(time)
107✔
393
    <<s::32-little-signed>>
107✔
394
  end
395

396
  def encode(:time, nil), do: <<0::32>>
1✔
397

398
  def encode({:time64, time_unit}, %Time{} = time) do
399
    {s, micros} = Time.to_seconds_after_midnight(time)
117✔
400

401
    micros_as_ticks =
117✔
402
      cond do
403
        time_unit < 1_000_000 -> div(micros, div(1_000_000, time_unit))
74✔
404
        time_unit == 1_000_000 -> micros
43✔
405
        true -> micros * div(time_unit, 1_000_000)
31✔
406
      end
407

408
    ticks = s * time_unit + micros_as_ticks
117✔
409
    <<ticks::64-little-signed>>
117✔
410
  end
411

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

414
  def encode(:uuid, <<u1::64, u2::64>>), do: <<u1::64-little, u2::64-little>>
15✔
415

416
  def encode(
417
        :uuid,
418
        <<a1, a2, a3, a4, a5, a6, a7, a8, ?-, b1, b2, b3, b4, ?-, c1, c2, c3, c4, ?-, d1, d2, d3,
419
          d4, ?-, e1, e2, e3, e4, e5, e6, e7, e8, e9, e10, e11, e12>>
420
      ) do
421
    raw =
2✔
422
      <<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,
423
        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,
424
        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,
425
        d(e8)::4, d(e9)::4, d(e10)::4, d(e11)::4, d(e12)::4>>
426

427
    encode(:uuid, raw)
2✔
428
  end
429

430
  def encode(:uuid, nil), do: <<0::128>>
1✔
431

432
  def encode(:ipv4, {a, b, c, d})
106✔
433
      when is_integer(a) and a in 0..255 and is_integer(b) and b in 0..255 and is_integer(c) and
434
             c in 0..255 and is_integer(d) and d in 0..255,
435
      do: [d, c, b, a]
436

437
  def encode(:ipv4, {_, _, _, _} = address) do
438
    raise ArgumentError, "invalid IPv4 address: #{inspect(address)}"
4✔
439
  end
440

441
  def encode(:ipv4, nil), do: <<0::32>>
1✔
442

443
  def encode(:ipv6, {b1, b2, b3, b4, b5, b6, b7, b8})
444
      when is_integer(b1) and b1 in 0..65_535 and is_integer(b2) and b2 in 0..65_535 and
445
             is_integer(b3) and b3 in 0..65_535 and is_integer(b4) and b4 in 0..65_535 and
446
             is_integer(b5) and b5 in 0..65_535 and is_integer(b6) and b6 in 0..65_535 and
447
             is_integer(b7) and b7 in 0..65_535 and is_integer(b8) and b8 in 0..65_535 do
448
    <<b1::16, b2::16, b3::16, b4::16, b5::16, b6::16, b7::16, b8::16>>
106✔
449
  end
450

451
  def encode(:ipv6, {_, _, _, _, _, _, _, _} = address) do
452
    raise ArgumentError, "invalid IPv6 address: #{inspect(address)}"
4✔
453
  end
454

455
  def encode(:ipv6, <<_::128>> = encoded), do: encoded
1✔
456
  def encode(:ipv6, nil), do: <<0::128>>
1✔
457

458
  def encode(:point, {x, y}), do: [encode(:f64, x) | encode(:f64, y)]
21✔
459
  def encode(:point, nil), do: <<0::128>>
1✔
460
  def encode(:ring, points), do: encode({:array, :point}, points)
1✔
461
  def encode(:polygon, rings), do: encode({:array, :ring}, rings)
1✔
462
  def encode(:multipolygon, polygons), do: encode({:array, :polygon}, polygons)
1✔
463

464
  # TODO
465
  def encode(:dynamic, value) do
466
    case value do
12✔
467
      _ when is_binary(value) -> [0x15 | encode(:string, value)]
2✔
468
      _ when is_integer(value) and value >= 0 -> [0x04 | encode(:u64, value)]
3✔
469
      _ when is_integer(value) -> [0x0A | encode(:i64, value)]
1✔
470
      _ when is_float(value) -> [0x0E | encode(:f64, value)]
2✔
471
      %Date{} -> [0x0F | encode(:date, value)]
2✔
472
      %NaiveDateTime{} -> [0x11 | encode(:datetime, value)]
1✔
473
      [] -> [0x1E, 0x00]
1✔
474
    end
475
  end
476

477
  # TODO enum8 and enum16 nil
478
  for size <- [8, 16] do
479
    enum_t = :"enum#{size}"
480
    int_t = :"i#{size}"
481

482
    def encode({unquote(enum_t), mapping}, e) do
483
      i =
17✔
484
        case e do
485
          _ when is_integer(e) ->
486
            e
2✔
487

488
          _ when is_binary(e) ->
489
            case Map.fetch(mapping, e) do
15✔
490
              {:ok, res} ->
491
                res
14✔
492

493
              :error ->
494
                raise ArgumentError,
1✔
495
                      "enum value #{inspect(e)} not found in mapping: #{inspect(mapping)}"
496
            end
497
        end
498

499
      encode(unquote(int_t), i)
16✔
500
    end
501
  end
502

503
  def encode({:nullable, _type}, nil), do: 1
889✔
504

505
  def encode({:nullable, type}, value) do
506
    case encode(type, value) do
888✔
507
      e when is_list(e) or is_binary(e) -> [0 | e]
887✔
508
      e -> [0, e]
1✔
509
    end
510
  end
511

512
  defp encode_varint_cont(i) when i < 128, do: <<i>>
7✔
513

514
  defp encode_varint_cont(i) do
12✔
515
    [(i &&& 0b0111_1111) ||| 0b1000_0000 | encode_varint_cont(i >>> 7)]
516
  end
517

518
  defp encode_many([el | rest], type), do: [encode(type, el) | encode_many(rest, type)]
8,696✔
519
  defp encode_many([] = done, _type), do: done
1,984✔
520

521
  defp encode_many_kv([{key, value} | rest], key_type, value_type) do
1✔
522
    [
523
      encode(key_type, key),
524
      encode(value_type, value)
525
      | encode_many_kv(rest, key_type, value_type)
526
    ]
527
  end
528

529
  defp encode_many_kv([] = done, _key_type, _value_type), do: done
1✔
530

531
  # TODO find a better way than try/rescue
532
  defp try_encode_variant([type | types], idx, value) do
533
    try do
13✔
534
      encode(type, value)
13✔
535
    else
536
      encoded -> [idx | encoded]
7✔
537
    rescue
538
      _e -> try_encode_variant(types, idx + 1, value)
6✔
539
    end
540
  end
541

542
  defp try_encode_variant([], _idx, value) do
543
    raise ArgumentError, "no matching type found for encoding #{inspect(value)} as Variant"
1✔
544
  end
545

546
  @compile {:inline, d: 1}
547

548
  defp d(?0), do: 0
1✔
549
  defp d(?1), do: 1
1✔
550
  defp d(?2), do: 2
3✔
551
  defp d(?3), do: 3
1✔
552
  defp d(?4), do: 4
1✔
553
  defp d(?5), do: 5
3✔
554
  defp d(?6), do: 6
1✔
555
  defp d(?7), do: 7
1✔
556
  defp d(?8), do: 8
1✔
557
  defp d(?9), do: 9
1✔
558
  defp d(?A), do: 10
1✔
559
  defp d(?B), do: 11
1✔
560
  defp d(?C), do: 12
1✔
561
  defp d(?D), do: 13
1✔
562
  defp d(?E), do: 14
1✔
563
  defp d(?F), do: 15
1✔
564
  defp d(?a), do: 10
1✔
565
  defp d(?b), do: 11
1✔
566
  defp d(?c), do: 12
1✔
567
  defp d(?d), do: 13
1✔
568
  defp d(?e), do: 14
1✔
569
  defp d(?f), do: 15
1✔
570

571
  varints = [
572
    {_pattern = quote(do: <<0::1, v1::7>>), _value = quote(do: v1)},
573
    {quote(do: <<1::1, v1::7, 0::1, v2::7>>), quote(do: (v2 <<< 7) + v1)},
574
    {quote(do: <<1::1, v1::7, 1::1, v2::7, 0::1, v3::7>>),
575
     quote(do: (v3 <<< 14) + (v2 <<< 7) + v1)},
576
    {quote(do: <<1::1, v1::7, 1::1, v2::7, 1::1, v3::7, 0::1, v4::7>>),
577
     quote(do: (v4 <<< 21) + (v3 <<< 14) + (v2 <<< 7) + v1)},
578
    {quote(do: <<1::1, v1::7, 1::1, v2::7, 1::1, v3::7, 1::1, v4::7, 0::1, v5::7>>),
579
     quote(do: (v5 <<< 28) + (v4 <<< 21) + (v3 <<< 14) + (v2 <<< 7) + v1)},
580
    {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>>),
581
     quote(do: (v6 <<< 35) + (v5 <<< 28) + (v4 <<< 21) + (v3 <<< 14) + (v2 <<< 7) + v1)},
582
    {quote do
583
       <<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,
584
         v7::7>>
585
     end,
586
     quote do
587
       (v7 <<< 42) + (v6 <<< 35) + (v5 <<< 28) + (v4 <<< 21) + (v3 <<< 14) + (v2 <<< 7) + v1
588
     end},
589
    {quote do
590
       <<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,
591
         v7::7, 0::1, v8::7>>
592
     end,
593
     quote do
594
       (v8 <<< 49) + (v7 <<< 42) + (v6 <<< 35) + (v5 <<< 28) + (v4 <<< 21) + (v3 <<< 14) +
595
         (v2 <<< 7) + v1
596
     end}
597
  ]
598

599
  @doc false
600
  @spec decode_header(binary()) ::
601
          {:ok, names :: [String.t()], types :: [term], rest :: binary} | :more
602
  def decode_header(row_binary_with_names_and_types)
603

604
  for {pattern, value} <- varints do
605
    def decode_header(<<unquote(pattern), rest::bytes>>) do
606
      decode_header_names(rest, unquote(value), unquote(value), _acc = [])
36✔
607
    end
608
  end
609

610
  def decode_header(<<_bin::bytes>>) do
1✔
611
    :more
612
  end
613

614
  defp decode_header_names(<<rest::bytes>>, 0, count, names) do
615
    decode_header_types(rest, count, _acc = [], :lists.reverse(names))
21✔
616
  end
617

618
  for {pattern, value} <- varints do
619
    defp decode_header_names(
620
           <<unquote(pattern), name::size(unquote(value))-bytes, rest::bytes>>,
621
           left,
622
           count,
623
           acc
624
         ) do
625
      decode_header_names(rest, left - 1, count, [name | acc])
78✔
626
    end
627
  end
628

629
  defp decode_header_names(<<_bin::bytes>>, _left, _count, _acc) do
15✔
630
    :more
631
  end
632

633
  defp decode_header_types(<<rest::bytes>>, 0, types, names) do
634
    {:ok, names, decoding_types_reverse(types), rest}
1✔
635
  end
636

637
  for {pattern, value} <- varints do
638
    defp decode_header_types(
639
           <<unquote(pattern), type::size(unquote(value))-bytes, rest::bytes>>,
640
           count,
641
           acc,
642
           names
643
         ) do
644
      decode_header_types(rest, count - 1, [type | acc], names)
24✔
645
    end
646
  end
647

648
  defp decode_header_types(<<_bin::bytes>>, _count, _acc, _names) do
20✔
649
    :more
650
  end
651

652
  @doc """
653
  Decodes [RowBinaryWithNamesAndTypes](https://clickhouse.com/docs/en/interfaces/formats/RowBinaryWithNamesAndTypes) into rows.
654

655
  Example:
656

657
      iex> decode_rows(<<1, 3, "1+1"::bytes, 5, "UInt8"::bytes, 2>>)
658
      [[2]]
659

660
  """
661
  def decode_rows(row_binary_with_names_and_types)
662
  def decode_rows(<<>>), do: []
1✔
663

664
  for {pattern, value} <- varints do
665
    def decode_rows(<<unquote(pattern), rest::bytes>>) do
666
      skip_names(rest, unquote(value), unquote(value))
5✔
667
    end
668
  end
669

670
  @doc """
671
  Same as `decode_rows/1` but the first element is a list of column names.
672

673
  Example:
674

675
      iex> decode_names_and_rows(<<1, 3, "1+1"::bytes, 5, "UInt8"::bytes, 2>>)
676
      [["1+1"], [2]]
677

678
  """
679
  def decode_names_and_rows(row_binary_with_names_and_types)
680

681
  for {pattern, value} <- varints do
682
    def decode_names_and_rows(<<unquote(pattern), rest::bytes>>) do
683
      decode_names(rest, unquote(value), unquote(value), _acc = [])
2,316✔
684
    end
685
  end
686

687
  @doc """
688
  Decodes [RowBinary](https://clickhouse.com/docs/en/interfaces/formats/RowBinary) into rows.
689

690
  Example:
691

692
      iex> decode_rows(<<1>>, ["UInt8"])
693
      [[1]]
694

695
  """
696
  def decode_rows(row_binary, types)
697
  def decode_rows(<<>>, _types), do: []
1✔
698

699
  def decode_rows(<<data::bytes>>, types) do
700
    decode_rows!(data, decoding_types(types))
748✔
701
  end
702

703
  defp decode_rows!(data, types) do
704
    {rows, remaining_data, state} = decode_rows(types, data, [], [], types)
3,047✔
705

706
    case state do
3,029✔
707
      nil ->
708
        rows
3,027✔
709

710
      {:cont, types_rest, row} ->
711
        raise ArgumentError, """
2✔
712
        incomplete RowBinary data: ran out of bytes while decoding
713

714
        Expected to decode: #{inspect(types_rest)}
715
        Remaining bytes: #{byte_size(remaining_data)} bytes
2✔
716
        Partial row: #{inspect(row)}
717
        Completed rows: #{length(rows)}
2✔
718
        """
719
    end
720
  end
721

722
  @doc false
723
  def decode_rows_continue(<<data::bytes>>, types, state) do
724
    case state do
201,104✔
725
      {:cont, types_rest, row} -> decode_rows(types_rest, data, row, [], types)
201,046✔
726
      nil -> decode_rows(types, data, [], [], types)
58✔
727
    end
728
  end
729

730
  @doc false
731
  def decoding_types([type | types]) do
865✔
732
    [decoding_type(type) | decoding_types(types)]
733
  end
734

735
  def decoding_types([] = done), do: done
820✔
736

737
  defp decoding_types_reverse(types), do: decoding_types_reverse(types, [])
2,301✔
738

739
  defp decoding_types_reverse([type | types], acc) do
740
    decoding_types_reverse(types, [decoding_type(type) | acc])
15,429✔
741
  end
742

743
  defp decoding_types_reverse([], acc), do: acc
2,301✔
744

745
  defp decoding_type(t) when is_binary(t) do
746
    decoding_type(Ch.Types.decode(t))
15,746✔
747
  end
748

749
  defp decoding_type(t)
750
       when t in [
751
              :string,
752
              :json,
753
              :dynamic,
754
              :boolean,
755
              :uuid,
756
              :date,
757
              :date32,
758
              :time,
759
              :time64,
760
              :ipv4,
761
              :ipv6,
762
              :point,
763
              :nothing
764
            ],
765
       do: t
2,909✔
766

767
  defp decoding_type({:datetime, _tz} = t), do: t
13✔
768
  defp decoding_type({:fixed_string, _len} = t), do: t
419✔
769

770
  for size <- [8, 16, 32, 64, 128, 256] do
771
    defp decoding_type(unquote(:"u#{size}") = u), do: u
11,894✔
772
    defp decoding_type(unquote(:"i#{size}") = i), do: i
218✔
773
  end
774

775
  for size <- [32, 64] do
776
    defp decoding_type(unquote(:"f#{size}") = f), do: f
453✔
777
  end
778

779
  defp decoding_type(:datetime = t), do: {t, _tz = nil}
12✔
780

781
  defp decoding_type({:array = a, t}), do: {a, decoding_type(t)}
1,123✔
782

783
  defp decoding_type({:tuple = t, ts}) do
221✔
784
    {t, Enum.map(ts, &decoding_type/1)}
785
  end
786

787
  defp decoding_type({:variant = v, ts}) do
17✔
788
    {v, ts |> Enum.map(&decoding_type/1) |> List.to_tuple()}
789
  end
790

791
  defp decoding_type({:map = m, kt, vt}) do
792
    {m, decoding_type(kt), decoding_type(vt)}
226✔
793
  end
794

795
  defp decoding_type({:nullable = n, t}), do: {n, decoding_type(t)}
546✔
796
  defp decoding_type({:low_cardinality, t}), do: decoding_type(t)
277✔
797

798
  defp decoding_type({:decimal = t, p, s}), do: {t, decimal_size(p), s}
250✔
799
  defp decoding_type({:decimal32, s}), do: {:decimal, 32, s}
26✔
800
  defp decoding_type({:decimal64, s}), do: {:decimal, 64, s}
22✔
801
  defp decoding_type({:decimal128, s}), do: {:decimal, 128, s}
26✔
802
  defp decoding_type({:decimal256, s}), do: {:decimal, 256, s}
30✔
803

804
  defp decoding_type({:datetime64 = t, p}), do: {t, time_unit(p), _tz = nil}
6✔
805
  defp decoding_type({:datetime64 = t, p, tz}), do: {t, time_unit(p), tz}
222✔
806

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

809
  defp decoding_type({e, mappings}) when e in [:enum8, :enum16] do
13✔
810
    {e, Map.new(mappings, fn {k, v} -> {v, k} end)}
29✔
811
  end
812

813
  defp decoding_type({:simple_aggregate_function, _f, t}), do: decoding_type(t)
6✔
814

815
  defp decoding_type(:ring), do: {:array, :point}
1✔
816
  defp decoding_type(:polygon), do: {:array, {:array, :point}}
1✔
817
  defp decoding_type(:multipolygon), do: {:array, {:array, {:array, :point}}}
1✔
818

819
  defp decoding_type(type) do
820
    raise ArgumentError, "unsupported type for decoding: #{inspect(type)}"
1✔
821
  end
822

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

825
  for {pattern, value} <- varints do
826
    defp skip_names(<<unquote(pattern), _::size(unquote(value))-bytes, rest::bytes>>, left, count) do
827
      skip_names(rest, left - 1, count)
19✔
828
    end
829
  end
830

831
  defp decode_names(<<rest::bytes>>, 0, count, names) do
2,316✔
832
    [:lists.reverse(names) | decode_types(rest, count, _acc = [])]
833
  end
834

835
  for {pattern, value} <- varints do
836
    defp decode_names(
837
           <<unquote(pattern), name::size(unquote(value))-bytes, rest::bytes>>,
838
           left,
839
           count,
840
           acc
841
         ) do
842
      decode_names(rest, left - 1, count, [name | acc])
15,460✔
843
    end
844
  end
845

846
  defp decode_types(<<>>, 0, _types), do: []
21✔
847

848
  defp decode_types(<<rest::bytes>>, 0, types) do
849
    decode_rows!(rest, decoding_types_reverse(types))
2,300✔
850
  end
851

852
  for {pattern, value} <- varints do
853
    defp decode_types(
854
           <<unquote(pattern), type::size(unquote(value))-bytes, rest::bytes>>,
855
           count,
856
           acc
857
         ) do
858
      decode_types(rest, count - 1, [type | acc])
15,479✔
859
    end
860
  end
861

862
  @compile inline: [decode_string_decode_rows: 5]
863

864
  for {pattern, size} <- varints do
865
    defp decode_string_decode_rows(
866
           <<unquote(pattern), s::size(unquote(size))-bytes, bin::bytes>>,
867
           types_rest,
868
           row,
869
           rows,
870
           types
871
         ) do
872
      decode_rows(types_rest, bin, [s | row], rows, types)
8,068✔
873
    end
874
  end
875

876
  defp decode_string_decode_rows(<<bin::bytes>>, types_rest, row, rows, _types) do
877
    to_be_continued(rows, bin, [:string | types_rest], row)
200,144✔
878
  end
879

880
  @compile inline: [decode_string_json_decode_rows: 5]
881

882
  for {pattern, size} <- varints do
883
    defp decode_string_json_decode_rows(
884
           <<unquote(pattern), s::size(unquote(size))-bytes, bin::bytes>>,
885
           types_rest,
886
           row,
887
           rows,
888
           types
889
         ) do
890
      decode_rows(types_rest, bin, [JSON.decode!(s) | row], rows, types)
43✔
891
    end
892
  end
893

894
  defp decode_string_json_decode_rows(<<bin::bytes>>, types_rest, row, rows, _types) do
895
    to_be_continued(rows, bin, [:json | types_rest], row)
46✔
896
  end
897

898
  @compile inline: [decode_array_decode_rows: 6]
899
  defp decode_array_decode_rows(<<0, bin::bytes>>, _type, types_rest, row, rows, types) do
900
    decode_rows(types_rest, bin, [[] | row], rows, types)
385✔
901
  end
902

903
  for {pattern, size} <- varints do
904
    defp decode_array_decode_rows(
905
           <<unquote(pattern), bin::bytes>>,
906
           type,
907
           types_rest,
908
           row,
909
           rows,
910
           types
911
         ) do
912
      array_types = List.duplicate(type, unquote(size))
2,535✔
913
      types_rest = array_types ++ [{:array_over, row} | types_rest]
2,535✔
914
      decode_rows(types_rest, bin, [], rows, types)
2,535✔
915
    end
916
  end
917

918
  defp decode_array_decode_rows(<<bin::bytes>>, type, types_rest, row, rows, _types) do
919
    to_be_continued(rows, bin, [{:array, type} | types_rest], row)
12✔
920
  end
921

922
  @compile inline: [decode_map_decode_rows: 7]
923
  defp decode_map_decode_rows(
924
         <<0, bin::bytes>>,
925
         _key_type,
926
         _value_type,
927
         types_rest,
928
         row,
929
         rows,
930
         types
931
       ) do
932
    decode_rows(types_rest, bin, [%{} | row], rows, types)
26✔
933
  end
934

935
  for {pattern, size} <- varints do
936
    defp decode_map_decode_rows(
937
           <<unquote(pattern), bin::bytes>>,
938
           key_type,
939
           value_type,
940
           types_rest,
941
           row,
942
           rows,
943
           types
944
         ) do
945
      types_rest =
212✔
946
        map_types(unquote(size), key_type, value_type) ++ [{:map_over, row} | types_rest]
947

948
      decode_rows(types_rest, bin, [], rows, types)
212✔
949
    end
950
  end
951

952
  defp decode_map_decode_rows(<<bin::bytes>>, key_type, value_type, types_rest, row, rows, _types) do
953
    to_be_continued(rows, bin, [{:map, key_type, value_type} | types_rest], row)
6✔
954
  end
955

956
  defp map_types(count, key_type, value_type) when count > 0 do
851✔
957
    [key_type, value_type | map_types(count - 1, key_type, value_type)]
958
  end
959

960
  defp map_types(0, _key_type, _value_types), do: []
212✔
961

962
  # https://clickhouse.com/docs/sql-reference/data-types/data-types-binary-encoding
963
  dynamic_types = [
964
    nothing: 0x00,
965
    u8: 0x01,
966
    u16: 0x02,
967
    u32: 0x03,
968
    u64: 0x04,
969
    u128: 0x05,
970
    u256: 0x06,
971
    i8: 0x07,
972
    i16: 0x08,
973
    i32: 0x09,
974
    i64: 0x0A,
975
    i128: 0x0B,
976
    i256: 0x0C,
977
    f32: 0x0D,
978
    f64: 0x0E,
979
    date: 0x0F,
980
    date32: 0x10,
981
    string: 0x15,
982
    uuid: 0x1D,
983
    ipv4: 0x28,
984
    ipv6: 0x29,
985
    boolean: 0x2D
986
  ]
987

988
  # TODO compile inline?
989

990
  for {type, code} <- dynamic_types do
991
    defp decode_dynamic(
992
           <<unquote(code), rest::bytes>>,
993
           dynamic,
994
           types_rest,
995
           row,
996
           rows,
997
           types
998
         ) do
999
      decode_dynamic_continue(rest, [unquote(type) | dynamic], types_rest, row, rows, types)
120✔
1000
    end
1001
  end
1002

1003
  # DateTime 0x11
1004
  defp decode_dynamic(<<0x11, rest::bytes>>, dynamic, types_rest, row, rows, types) do
1005
    decode_dynamic_continue(rest, [{:datetime, nil} | dynamic], types_rest, row, rows, types)
2✔
1006
  end
1007

1008
  # DateTime(time_zone) 0x12 <var_uint_time_zone_name_size><time_zone_name_data>
1009
  for {pattern, size} <- varints do
1010
    defp decode_dynamic(
1011
           <<0x12, unquote(pattern), tz::size(unquote(size))-bytes, rest::bytes>>,
1012
           dynamic,
1013
           types_rest,
1014
           row,
1015
           rows,
1016
           types
1017
         ) do
1018
      decode_dynamic_continue(rest, [{:datetime, tz} | dynamic], types_rest, row, rows, types)
1✔
1019
    end
1020
  end
1021

1022
  # DateTime64(P) 0x13 <uint8_precision>
1023
  defp decode_dynamic(
1024
         <<0x13, precision, rest::bytes>>,
1025
         dynamic,
1026
         types_rest,
1027
         row,
1028
         rows,
1029
         types
1030
       ) do
1031
    decode_dynamic_continue(
1✔
1032
      rest,
1033
      [decoding_type({:datetime64, precision}) | dynamic],
1034
      types_rest,
1035
      row,
1036
      rows,
1037
      types
1038
    )
1039
  end
1040

1041
  # DateTime64(P, time_zone) 0x14 <uint8_precision><var_uint_time_zone_name_size><time_zone_name_data>
1042
  for {pattern, size} <- varints do
1043
    defp decode_dynamic(
1044
           <<0x14, precision, unquote(pattern), tz::size(unquote(size))-bytes, rest::bytes>>,
1045
           dynamic,
1046
           types_rest,
1047
           row,
1048
           rows,
1049
           types
1050
         ) do
1051
      decode_dynamic_continue(
1✔
1052
        rest,
1053
        [decoding_type({:datetime64, precision, tz}) | dynamic],
1054
        types_rest,
1055
        row,
1056
        rows,
1057
        types
1058
      )
1059
    end
1060
  end
1061

1062
  # FixedString(N) 0x16 <var_uint_size>
1063
  for {pattern, size} <- varints do
1064
    defp decode_dynamic(
1065
           <<0x16, unquote(pattern), rest::bytes>>,
1066
           dynamic,
1067
           types_rest,
1068
           row,
1069
           rows,
1070
           types
1071
         ) do
1072
      decode_dynamic_continue(
2✔
1073
        rest,
1074
        [{:fixed_string, unquote(size)} | dynamic],
1075
        types_rest,
1076
        row,
1077
        rows,
1078
        types
1079
      )
1080
    end
1081
  end
1082

1083
  # Decimal32(P, S) 0x19 <uint8_precision><uint8_scale>
1084
  # Decimal64(P, S) 0x1A <uint8_precision><uint8_scale>
1085
  # Decimal128(P, S) 0x1B <uint8_precision><uint8_scale>
1086
  # Decimal256(P, S) 0x1C <uint8_precision><uint8_scale>
1087
  for {code, size} <- [{0x19, 32}, {0x1A, 64}, {0x1B, 128}, {0x1C, 256}] do
1088
    defp decode_dynamic(
1089
           <<unquote(code), _precision, scale, rest::bytes>>,
1090
           dynamic,
1091
           types_rest,
1092
           row,
1093
           rows,
1094
           types
1095
         ) do
1096
      decode_dynamic_continue(
4✔
1097
        rest,
1098
        [{:decimal, unquote(size), scale} | dynamic],
1099
        types_rest,
1100
        row,
1101
        rows,
1102
        types
1103
      )
1104
    end
1105
  end
1106

1107
  # Array(T) 0x1E <nested_type_encoding>
1108
  defp decode_dynamic(<<0x1E, rest::bytes>>, dynamic, types_rest, row, rows, types) do
1109
    decode_dynamic_continue(rest, [:array | dynamic], types_rest, row, rows, types)
29✔
1110
  end
1111

1112
  # Nullable(T)        0x23 <nested_type_encoding>
1113
  defp decode_dynamic(<<0x23, rest::bytes>>, dynamic, types_rest, row, rows, types) do
1114
    decode_dynamic_continue(rest, [:nullable | dynamic], types_rest, row, rows, types)
5✔
1115
  end
1116

1117
  # LowCardinality(T) 0x26 <nested_type_encoding>
1118
  defp decode_dynamic(<<0x26, rest::bytes>>, dynamic, types_rest, row, rows, types) do
1119
    decode_dynamic_continue(rest, [:low_cardinality | dynamic], types_rest, row, rows, types)
2✔
1120
  end
1121

1122
  # TODO
1123
  # 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>
1124
  # 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>
1125
  # Tuple(T1, ..., TN)        0x1F <var_uint_number_of_elements><nested_type_encoding_1>...<nested_type_encoding_N>
1126
  # 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>
1127
  # Set        0x21
1128
  # Interval        0x22 <interval_kind> (see interval kind binary encoding)
1129
  # Function        0x24<var_uint_number_of_arguments><argument_type_encoding_1>...<argument_type_encoding_N><return_type_encoding>
1130
  # 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)
1131
  # Map(K, V)        0x27<key_type_encoding><value_type_encoding>
1132
  # Variant(T1, ..., TN)        0x2A<var_uint_number_of_variants><variant_type_encoding_1>...<variant_type_encoding_N>
1133
  # Dynamic(max_types=N)        0x2B<uint8_max_types>
1134
  # Custom type (Ring, Polygon, etc)        0x2C<var_uint_type_name_size><type_name_data>
1135
  # 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)
1136
  # 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>
1137
  # 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>...
1138

1139
  unsupported_dynamic_types = %{
1140
    "Enum8" => 0x17,
1141
    "Enum16" => 0x18,
1142
    "Tuple" => 0x1F,
1143
    "TupleWithNames" => 0x20,
1144
    "Set" => 0x21,
1145
    "Interval" => 0x22,
1146
    "Function" => 0x24,
1147
    "AggregateFunction" => 0x25,
1148
    "Map" => 0x27,
1149
    "Variant" => 0x2A,
1150
    "Dynamic" => 0x2B,
1151
    "CustomType" => 0x2C,
1152
    "SimpleAggregateFunction" => 0x2E,
1153
    "Nested" => 0x2F,
1154
    "JSON" => 0x30
1155
  }
1156

1157
  for {type, code} <- unsupported_dynamic_types do
1158
    defp decode_dynamic(<<unquote(code), _::bytes>>, _dynamic, _types_rest, _row, _rows, _types) do
1159
      raise ArgumentError, "unsupported dynamic type #{unquote(type)}"
9✔
1160
    end
1161
  end
1162

1163
  defp decode_dynamic(<<bin::bytes>>, dynamic, types_rest, row, rows, _types) do
1164
    to_be_continued(rows, bin, [{:dynamic, dynamic} | types_rest], row)
2✔
1165
  end
1166

1167
  @compile inline: [decode_dynamic_continue: 6]
1168

1169
  defp decode_dynamic_continue(<<rest::bytes>>, dynamic, types_rest, row, rows, types) do
1170
    continue? =
103✔
1171
      case dynamic do
1172
        [:array | _] -> true
29✔
1173
        [:nullable | _] -> true
5✔
1174
        [:low_cardinality | _] -> true
2✔
1175
        _ -> false
103✔
1176
      end
1177

1178
    if continue? do
103✔
1179
      decode_dynamic(rest, dynamic, types_rest, row, rows, types)
36✔
1180
    else
1181
      type = build_dynamic_type(:lists.reverse(dynamic))
103✔
1182
      decode_rows([type | types_rest], rest, row, rows, types)
103✔
1183
    end
1184
  end
1185

1186
  defp build_dynamic_type([type]), do: type
131✔
1187

1188
  defp build_dynamic_type(type) do
1189
    case type do
32✔
1190
      [:array | rest] -> {:array, build_dynamic_type(rest)}
25✔
1191
      [:nullable | rest] -> {:nullable, build_dynamic_type(rest)}
5✔
1192
      [:low_cardinality | rest] -> build_dynamic_type(rest)
2✔
1193
    end
1194
  end
1195

1196
  simple_types = %{
1197
    u8: %{pattern: quote(do: <<u>>), value: quote(do: u)},
1198
    u16: %{pattern: quote(do: <<u::16-little>>), value: quote(do: u)},
1199
    u32: %{pattern: quote(do: <<u::32-little>>), value: quote(do: u)},
1200
    u64: %{pattern: quote(do: <<u::64-little>>), value: quote(do: u)},
1201
    u128: %{pattern: quote(do: <<u::128-little>>), value: quote(do: u)},
1202
    u256: %{pattern: quote(do: <<u::256-little>>), value: quote(do: u)},
1203
    i8: %{pattern: quote(do: <<i::signed>>), value: quote(do: i)},
1204
    i16: %{pattern: quote(do: <<i::16-little-signed>>), value: quote(do: i)},
1205
    i32: %{pattern: quote(do: <<i::32-little-signed>>), value: quote(do: i)},
1206
    i64: %{pattern: quote(do: <<i::64-little-signed>>), value: quote(do: i)},
1207
    i128: %{pattern: quote(do: <<i::128-little-signed>>), value: quote(do: i)},
1208
    i256: %{pattern: quote(do: <<i::256-little-signed>>), value: quote(do: i)},
1209
    f32: [
1210
      %{pattern: quote(do: <<f::32-little-float>>), value: quote(do: f)},
1211
      %{pattern: quote(do: <<_nan_or_inf::32>>), value: quote(do: nil)}
1212
    ],
1213
    f64: [
1214
      %{pattern: quote(do: <<f::64-little-float>>), value: quote(do: f)},
1215
      %{pattern: quote(do: <<_nan_or_inf::64>>), value: quote(do: nil)}
1216
    ],
1217
    uuid: %{
1218
      pattern: quote(do: <<u1::64-little, u2::64-little>>),
1219
      value: quote(do: <<u1::64, u2::64>>)
1220
    },
1221
    date: %{
1222
      pattern: quote(do: <<d::16-little>>),
1223
      value: quote(do: Date.from_gregorian_days(d + @epoch_gregorian_days))
1224
    },
1225
    date32: %{
1226
      pattern: quote(do: <<d::32-little-signed>>),
1227
      value: quote(do: Date.from_gregorian_days(d + @epoch_gregorian_days))
1228
    },
1229
    time: %{
1230
      pattern: quote(do: <<s::32-little-signed>>),
1231
      value: quote(do: time_after_midnight(s, 1))
1232
    },
1233
    boolean: [
1234
      %{pattern: quote(do: <<0>>), value: quote(do: false)},
1235
      %{pattern: quote(do: <<1>>), value: quote(do: true)},
1236
      %{pattern: quote(do: <<b>>), value: quote(do: raise("invalid boolean value: #{b}"))}
1237
    ],
1238
    ipv4: %{
1239
      pattern: quote(do: <<b4, b3, b2, b1>>),
1240
      value: quote(do: {b1, b2, b3, b4})
1241
    },
1242
    ipv6: %{
1243
      pattern: quote(do: <<b1::16, b2::16, b3::16, b4::16, b5::16, b6::16, b7::16, b8::16>>),
1244
      value: quote(do: {b1, b2, b3, b4, b5, b6, b7, b8})
1245
    },
1246
    point: %{
1247
      pattern: quote(do: <<x::64-little-float, y::64-little-float>>),
1248
      value: quote(do: {x, y})
1249
    }
1250
  }
1251

1252
  for {type, clauses} <- simple_types do
1253
    fun = :"decode_#{type}_decode_rows"
1254
    @compile inline: [{fun, 5}]
1255

1256
    for %{pattern: pattern, value: value} <- List.wrap(clauses) do
1257
      defp unquote(fun)(<<unquote(pattern), rest::bytes>>, types_rest, row, rows, types) do
1258
        decode_rows(types_rest, rest, [unquote(value) | row], rows, types)
2,019,837✔
1259
      end
1260
    end
1261

1262
    defp unquote(fun)(<<bin::bytes>>, types_rest, row, rows, _types) do
1263
      to_be_continued(rows, bin, [unquote(type) | types_rest], row)
583✔
1264
    end
1265
  end
1266

1267
  defp decode_rows([type | types_rest], <<bin::bytes>>, row, rows, types) do
1268
    case type do
2,243,931✔
1269
      :u8 ->
1270
        decode_u8_decode_rows(bin, types_rest, row, rows, types)
4,958✔
1271

1272
      :u16 ->
1273
        decode_u16_decode_rows(bin, types_rest, row, rows, types)
10,275✔
1274

1275
      :u32 ->
1276
        decode_u32_decode_rows(bin, types_rest, row, rows, types)
85✔
1277

1278
      :u64 ->
1279
        decode_u64_decode_rows(bin, types_rest, row, rows, types)
2,000,217✔
1280

1281
      :u128 ->
1282
        decode_u128_decode_rows(bin, types_rest, row, rows, types)
58✔
1283

1284
      :u256 ->
1285
        decode_u256_decode_rows(bin, types_rest, row, rows, types)
90✔
1286

1287
      :i8 ->
1288
        decode_i8_decode_rows(bin, types_rest, row, rows, types)
57✔
1289

1290
      :i16 ->
1291
        decode_i16_decode_rows(bin, types_rest, row, rows, types)
65✔
1292

1293
      :i32 ->
1294
        decode_i32_decode_rows(bin, types_rest, row, rows, types)
51✔
1295

1296
      :i64 ->
1297
        decode_i64_decode_rows(bin, types_rest, row, rows, types)
127✔
1298

1299
      :i128 ->
1300
        decode_i128_decode_rows(bin, types_rest, row, rows, types)
59✔
1301

1302
      :i256 ->
1303
        decode_i256_decode_rows(bin, types_rest, row, rows, types)
91✔
1304

1305
      :f32 ->
1306
        decode_f32_decode_rows(bin, types_rest, row, rows, types)
892✔
1307

1308
      :f64 ->
1309
        decode_f64_decode_rows(bin, types_rest, row, rows, types)
965✔
1310

1311
      :string ->
1312
        decode_string_decode_rows(bin, types_rest, row, rows, types)
208,212✔
1313

1314
      :json ->
1315
        # assuming it arrives as text and not "native" binary JSON
1316
        # i.e. assumes `settings: [output_format_binary_write_json_as_string: 1]`
1317
        # TODO
1318
        decode_string_json_decode_rows(bin, types_rest, row, rows, types)
89✔
1319

1320
      :dynamic ->
1321
        decode_dynamic(bin, _dynamic = [], types_rest, row, rows, types)
140✔
1322

1323
      {:dynamic, dynamic} ->
1324
        decode_dynamic(bin, dynamic, types_rest, row, rows, types)
2✔
1325

1326
      {:fixed_string, size} ->
1327
        case bin do
4,278✔
1328
          <<s::size(^size)-bytes, rest::bytes>> ->
1329
            decode_rows(types_rest, rest, [s | row], rows, types)
4,264✔
1330

1331
          _ ->
1332
            to_be_continued(rows, bin, [type | types_rest], row)
14✔
1333
        end
1334

1335
      :boolean ->
1336
        decode_boolean_decode_rows(bin, types_rest, row, rows, types)
1,507✔
1337

1338
      :uuid ->
1339
        decode_uuid_decode_rows(bin, types_rest, row, rows, types)
200✔
1340

1341
      :date ->
1342
        decode_date_decode_rows(bin, types_rest, row, rows, types)
46✔
1343

1344
      :date32 ->
1345
        decode_date32_decode_rows(bin, types_rest, row, rows, types)
32✔
1346

1347
      :time ->
1348
        decode_time_decode_rows(bin, types_rest, row, rows, types)
230✔
1349

1350
      {:time64, time_unit} ->
1351
        case bin do
300✔
1352
          <<ticks::64-little-signed, bin::bytes>> ->
1353
            time = time_after_midnight(ticks, time_unit)
270✔
1354
            decode_rows(types_rest, bin, [time | row], rows, types)
265✔
1355

1356
          _ ->
1357
            to_be_continued(rows, bin, [type | types_rest], row)
30✔
1358
        end
1359

1360
      {:datetime, timezone} ->
1361
        case bin do
49✔
1362
          <<s::32-little, bin::bytes>> ->
1363
            dt = DateTime.from_unix!(s)
27✔
1364

1365
            dt =
27✔
1366
              case timezone do
1367
                nil -> DateTime.to_naive(dt)
13✔
1368
                "UTC" -> dt
8✔
1369
                _ -> DateTime.shift_zone!(dt, timezone)
6✔
1370
              end
1371

1372
            decode_rows(types_rest, bin, [dt | row], rows, types)
27✔
1373

1374
          _ ->
1375
            to_be_continued(rows, bin, [type | types_rest], row)
22✔
1376
        end
1377

1378
      {:decimal, size, scale} ->
1379
        case bin do
484✔
1380
          <<val::size(^size)-little-signed, bin::bytes>> ->
1381
            sign = if val < 0, do: -1, else: 1
366✔
1382
            d = Decimal.new(sign, abs(val), -scale)
366✔
1383
            decode_rows(types_rest, bin, [d | row], rows, types)
366✔
1384

1385
          _ ->
1386
            to_be_continued(rows, bin, [type | types_rest], row)
118✔
1387
        end
1388

1389
      {:nullable, inner_type} ->
1390
        case bin do
2,824✔
1391
          <<b, bin::bytes>> ->
1392
            case b do
2,821✔
1393
              0 -> decode_rows([inner_type | types_rest], bin, row, rows, types)
1,406✔
1394
              1 -> decode_rows(types_rest, bin, [nil | row], rows, types)
1,415✔
1395
            end
1396

1397
          _ ->
1398
            to_be_continued(rows, bin, [type | types_rest], row)
3✔
1399
        end
1400

1401
      :nothing ->
1402
        decode_rows(types_rest, bin, [nil | row], rows, types)
27✔
1403

1404
      {:array, inner_type} ->
1405
        decode_array_decode_rows(bin, inner_type, types_rest, row, rows, types)
2,932✔
1406

1407
      {:array_over, original_row} ->
1408
        decode_rows(types_rest, bin, [:lists.reverse(row) | original_row], rows, types)
2,534✔
1409

1410
      {:map, key_type, value_type} ->
1411
        decode_map_decode_rows(bin, key_type, value_type, types_rest, row, rows, types)
244✔
1412

1413
      {:map_over, original_row} ->
1414
        map = row |> Enum.chunk_every(2) |> Enum.map(fn [v, k] -> {k, v} end) |> Map.new()
212✔
1415
        decode_rows(types_rest, bin, [map | original_row], rows, types)
212✔
1416

1417
      {:tuple, tuple_types} ->
1418
        decode_rows(tuple_types ++ [{:tuple_over, row} | types_rest], bin, [], rows, types)
234✔
1419

1420
      {:tuple_over, original_row} ->
1421
        tuple = row |> :lists.reverse() |> List.to_tuple()
234✔
1422
        decode_rows(types_rest, bin, [tuple | original_row], rows, types)
234✔
1423

1424
      {:variant, variant_types} ->
1425
        case bin do
35✔
1426
          <<255, bin::bytes>> ->
1427
            # 255 is the variant type index for "nothing"
1428
            decode_rows(types_rest, bin, [nil | row], rows, types)
7✔
1429

1430
          # TODO varint?
1431
          <<variant_type_index::8, bin::bytes>>
1432
          when variant_type_index < tuple_size(variant_types) ->
1433
            variant_type = elem(variant_types, variant_type_index)
24✔
1434
            decode_rows([variant_type | types_rest], bin, row, rows, types)
24✔
1435

1436
          <<variant_type_index::8, _bin::bytes>> ->
1437
            raise ArgumentError, "invalid Variant type index: #{variant_type_index}"
1✔
1438

1439
          _ ->
1440
            to_be_continued(rows, bin, [type | types_rest], row)
3✔
1441
        end
1442

1443
      {:datetime64, time_unit, timezone} ->
1444
        case bin do
642✔
1445
          <<s::64-little-signed, bin::bytes>> ->
1446
            dt = DateTime.from_unix!(s, time_unit)
580✔
1447

1448
            dt =
580✔
1449
              case timezone do
1450
                nil -> DateTime.to_naive(dt)
7✔
1451
                "UTC" -> dt
567✔
1452
                _ -> DateTime.shift_zone!(dt, timezone)
6✔
1453
              end
1454

1455
            decode_rows(types_rest, bin, [dt | row], rows, types)
580✔
1456

1457
          _ ->
1458
            to_be_continued(rows, bin, [type | types_rest], row)
62✔
1459
        end
1460

1461
      {:enum8, mapping} ->
1462
        case bin do
33✔
1463
          <<v::signed, bin::bytes>> ->
1464
            decode_rows(types_rest, bin, [Map.fetch!(mapping, v) | row], rows, types)
32✔
1465

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

1470
      {:enum16, mapping} ->
1471
        case bin do
6✔
1472
          <<v::16-little-signed, bin::bytes>> ->
1473
            decode_rows(types_rest, bin, [Map.fetch!(mapping, v) | row], rows, types)
2✔
1474

1475
          _ ->
1476
            to_be_continued(rows, bin, [type | types_rest], row)
4✔
1477
        end
1478

1479
      :ipv4 ->
1480
        decode_ipv4_decode_rows(bin, types_rest, row, rows, types)
133✔
1481

1482
      :ipv6 ->
1483
        decode_ipv6_decode_rows(bin, types_rest, row, rows, types)
175✔
1484

1485
      :point ->
1486
        decode_point_decode_rows(bin, types_rest, row, rows, types)
107✔
1487
    end
1488
  end
1489

1490
  defp decode_rows([], <<>> = empty, row, rows, _types) do
1491
    rows = :lists.reverse([:lists.reverse(row) | rows])
3,083✔
1492
    {rows, empty, _no_state = nil}
3,083✔
1493
  end
1494

1495
  defp decode_rows([], <<bin::bytes>>, row, rows, types) do
1496
    row = :lists.reverse(row)
2,003,565✔
1497
    decode_rows(types, bin, [], [row | rows], types)
2,003,565✔
1498
  end
1499

1500
  @compile inline: [to_be_continued: 4]
1501
  defp to_be_continued(rows, bin, types_rest, row) do
1502
    {:lists.reverse(rows), bin, {:cont, types_rest, row}}
200,697✔
1503
  end
1504

1505
  defp decimal_coefficient!(decimal, scale, size) do
1506
    unless is_integer(decimal.coef) do
172✔
1507
      raise ArgumentError, "ClickHouse Decimal values must be finite"
3✔
1508
    end
1509

1510
    %Decimal{sign: sign, coef: coefficient, exp: exponent} = decimal
169✔
1511
    shift = exponent + scale
169✔
1512
    max_digits = size |> then(&((1 <<< (&1 - 1)) - 1)) |> Integer.digits() |> length()
169✔
1513

1514
    cond do
169✔
1515
      coefficient == 0 ->
1✔
1516
        0
1517

1518
      shift >= 0 and length(Integer.digits(coefficient)) + shift > max_digits ->
168✔
NEW
1519
        raise ArgumentError,
×
NEW
1520
              "Decimal#{size}(#{scale}) value #{Decimal.to_string(decimal)} is out of range"
×
1521

1522
      shift >= 0 ->
168✔
1523
        sign * coefficient * Integer.pow(10, shift)
164✔
1524

1525
      -shift > length(Integer.digits(coefficient)) ->
4✔
1526
        0
1527

1528
      true ->
4✔
1529
        divisor = Integer.pow(10, -shift)
4✔
1530
        quotient = div(coefficient, divisor)
4✔
1531
        remainder = rem(coefficient, divisor)
4✔
1532
        rounded = if remainder * 2 >= divisor, do: quotient + 1, else: quotient
4✔
1533
        sign * rounded
4✔
1534
    end
1535
  end
1536

1537
  defp encode_decimal!(decimal, coefficient, size, scale) do
1538
    encode_fixed_integer!(
153✔
1539
      coefficient,
1540
      size,
1541
      :signed,
1542
      "Decimal#{size}(#{scale}) value #{Decimal.to_string(decimal)}"
153✔
1543
    )
1544
  end
1545

1546
  defp validate_decimal_type!(precision, scale)
39✔
1547
       when is_integer(precision) and precision in 1..76 and is_integer(scale) and scale >= 0 and
1548
              scale <= precision,
1549
       do: :ok
1550

1551
  defp validate_decimal_type!(precision, scale) do
1552
    raise ArgumentError,
4✔
1553
          "invalid Decimal precision and scale: precision=#{inspect(precision)}, scale=#{inspect(scale)}"
1554
  end
1555

1556
  defp validate_decimal_scale!(_type, scale, precision)
143✔
1557
       when is_integer(scale) and scale >= 0 and scale <= precision,
1558
       do: :ok
1559

1560
  defp validate_decimal_scale!(type, scale, precision) do
1561
    raise ArgumentError,
4✔
1562
          "invalid #{decimal_type_name(type)} scale #{inspect(scale)}; expected 0..#{precision}"
4✔
1563
  end
1564

1565
  defp decimal_type_name(type) do
1566
    type |> Atom.to_string() |> String.replace_prefix("decimal", "Decimal")
4✔
1567
  end
1568

1569
  defp encode_fixed_integer!(integer, size, :unsigned, _type)
1570
       when is_integer(integer) and integer >= 0 and integer < 1 <<< size do
1571
    <<integer::little-unsigned-size(size)>>
33✔
1572
  end
1573

1574
  defp encode_fixed_integer!(integer, size, :signed, _type)
1575
       when is_integer(integer) and integer >= -(1 <<< (size - 1)) and
1576
              integer < 1 <<< (size - 1) do
1577
    <<integer::little-signed-size(size)>>
152✔
1578
  end
1579

1580
  defp encode_fixed_integer!(integer, size, signedness, type) do
1581
    raise ArgumentError,
21✔
1582
          "#{type} is out of range for a #{size}-bit #{signedness} integer: #{inspect(integer)}"
21✔
1583
  end
1584

1585
  @compile inline: [decimal_size: 1]
1586
  # https://clickhouse.com/docs/en/sql-reference/data-types/decimal/
1587
  defp decimal_size(precision) when is_integer(precision) do
1588
    cond do
283✔
1589
      precision >= 39 -> 256
215✔
1590
      precision >= 19 -> 128
68✔
1591
      precision >= 10 -> 64
55✔
1592
      true -> 32
26✔
1593
    end
1594
  end
1595

1596
  @compile inline: [time_unit: 1]
1597
  for precision <- 0..9 do
1598
    time_unit = Integer.pow(10, precision)
1599
    defp time_unit(unquote(precision)), do: unquote(time_unit)
599✔
1600
  end
1601

1602
  @compile inline: [time_after_midnight: 2]
1603
  defp time_after_midnight(ticks, time_unit) do
1604
    if ticks >= 0 and ticks < 86400 * time_unit do
486✔
1605
      ticks |> DateTime.from_unix!(time_unit) |> DateTime.to_time()
478✔
1606
    else
1607
      # since ClickHouse supports Time64 values of [-999:59:59.999999999, 999:59:59.999999999]
1608
      # and Elixir's Time supports values of [00:00:00.000000, 23:59:59.999999]
1609
      # we raise an error when ClickHouse's Time64 value is out of Elixir's Time range
1610
      raise ArgumentError,
8✔
1611
            "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✔
1612

1613
      # TODO: we could potentially decode ClickHouse's Time/Time64 values as Elixir's Duration when it's out of Elixir's Time range
1614
    end
1615
  end
1616
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