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

plausible / ch / fe719c22b4b5e378e075209d6c75dcbb68a25720-PR-403

03 Aug 2026 12:43PM UTC coverage: 96.944% (-1.1%) from 98.06%
fe719c22b4b5e378e075209d6c75dcbb68a25720-PR-403

Pull #403

github

ruslandoga
Restrict late exceptions to RowBinary responses
Pull Request #403: Handle late ClickHouse HTTP exceptions

39 of 46 new or added lines in 2 files covered. (84.78%)

14 existing lines in 2 files now uncovered.

793 of 818 relevant lines covered (96.94%)

14649.68 hits per line

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

98.58
/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
  @exception_marker "__exception__"
12
  @exception_tag_length 16
13
  @max_exception_size 16 * 1024
14
  @max_exception_length_digits 8
15

16
  @doc false
17
  def encode_names_and_types(names, types) do
5✔
18
    [encode(:varint, length(names)), encode_many(names, :string), encode_types(types)]
19
  end
20

21
  defp encode_types([type | types]) do
12✔
22
    encoded =
12✔
23
      case type do
24
        _ when is_binary(type) -> type
11✔
25
        _ -> Ch.Types.encode(type)
1✔
26
      end
27

28
    [encode(:string, encoded) | encode_types(types)]
29
  end
30

31
  defp encode_types([] = done), do: done
5✔
32

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

36
  Examples:
37

38
      iex> encode_row([], [])
39
      []
40

41
      iex> encode_row([1], ["UInt8"])
42
      [1]
43

44
      iex> encode_row([3, "hello"], ["UInt8", "String"])
45
      [3, [5 | "hello"]]
46

47
  """
48
  def encode_row(row, types) do
49
    _encode_row(row, encoding_types(types))
22✔
50
  end
51

52
  defp _encode_row([el | els], [type | types]), do: [encode(type, el) | _encode_row(els, types)]
106✔
53
  defp _encode_row([] = done, []), do: done
22✔
54

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

58
  Examples:
59

60
      iex> encode_rows([], [])
61
      []
62

63
      iex> encode_rows([[1]], ["UInt8"])
64
      [1]
65

66
      iex> encode_rows([[3, "hello"], [4, "hi"]], ["UInt8", "String"])
67
      [3, [5 | "hello"], 4, [2 | "hi"]]
68

69
  """
70
  def encode_rows(rows, types) do
71
    _encode_rows(rows, encoding_types(types))
854✔
72
  end
73

74
  @doc false
75
  def _encode_rows([row | rows], types), do: _encode_rows(row, types, rows, types)
4,860✔
76
  def _encode_rows([] = done, _types), do: done
854✔
77

78
  defp _encode_rows([el | els], [t | ts], rows, types) do
12,524✔
79
    [encode(t, el) | _encode_rows(els, ts, rows, types)]
80
  end
81

82
  defp _encode_rows([], [], rows, types), do: _encode_rows(rows, types)
4,860✔
83

84
  @doc false
85
  def encoding_types([type | types]) do
2,262✔
86
    [encoding_type(type) | encoding_types(types)]
87
  end
88

89
  def encoding_types([] = done), do: done
884✔
90

91
  defp encoding_type(type) when is_binary(type) do
92
    encoding_type(Ch.Types.decode(type))
2,162✔
93
  end
94

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

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

115
  defp encoding_type({:datetime, tz}) do
116
    raise ArgumentError, "can't encode DateTime with non-UTC timezone: #{inspect(tz)}"
1✔
117
  end
118

119
  defp encoding_type({:fixed_string, _len} = t), do: t
311✔
120

121
  for size <- [8, 16, 32, 64, 128, 256] do
122
    defp encoding_type(unquote(:"u#{size}") = u), do: u
655✔
123
    defp encoding_type(unquote(:"i#{size}") = i), do: i
40✔
124
  end
125

126
  for size <- [32, 64] do
127
    defp encoding_type(unquote(:"f#{size}") = f), do: f
222✔
128
  end
129

130
  defp encoding_type({:array = a, t}), do: {a, encoding_type(t)}
548✔
131

132
  defp encoding_type({:tuple = t, ts}) do
5✔
133
    {t, Enum.map(ts, &encoding_type/1)}
134
  end
135

136
  defp encoding_type({:variant = v, ts}) do
3✔
137
    {v, Enum.map(ts, &encoding_type/1)}
138
  end
139

140
  defp encoding_type({:map = m, kt, vt}) do
141
    {m, encoding_type(kt), encoding_type(vt)}
8✔
142
  end
143

144
  defp encoding_type({:nullable = n, t}), do: {n, encoding_type(t)}
113✔
145
  defp encoding_type({:low_cardinality, t}), do: encoding_type(t)
202✔
146

147
  defp encoding_type({:decimal, p, s}) do
148
    case decimal_size(p) do
5✔
149
      32 -> {:decimal32, s}
1✔
150
      64 -> {:decimal64, s}
2✔
151
      128 -> {:decimal128, s}
1✔
152
      256 -> {:decimal256, s}
1✔
153
    end
154
  end
155

156
  defp encoding_type({d, _scale} = t)
157
       when d in [:decimal32, :decimal64, :decimal128, :decimal256],
158
       do: t
5✔
159

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

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

164
  defp encoding_type({:datetime64, _, tz}) do
165
    raise ArgumentError, "can't encode DateTime64 with non-UTC timezone: #{inspect(tz)}"
1✔
166
  end
167

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

170
  defp encoding_type({e, mappings}) when e in [:enum8, :enum16] do
4✔
171
    {e, Map.new(mappings)}
172
  end
173

174
  defp encoding_type({:simple_aggregate_function, _f, t}), do: encoding_type(t)
1✔
175

176
  defp encoding_type(:ring), do: {:array, :point}
1✔
177
  defp encoding_type(:polygon), do: {:array, {:array, :point}}
1✔
178
  defp encoding_type(:multipolygon), do: {:array, {:array, {:array, :point}}}
1✔
179

180
  defp encoding_type(type) do
181
    raise ArgumentError, "unsupported type for encoding: #{inspect(type)}"
1✔
182
  end
183

184
  @doc false
185
  def encode(type, value)
186

187
  def encode(:varint, i) when is_integer(i) and i < 128, do: i
8,313✔
188
  def encode(:varint, i) when is_integer(i), do: encode_varint_cont(i)
14✔
189

190
  def encode(:string, str) do
191
    case str do
1✔
192
      _ when is_binary(str) -> [encode(:varint, byte_size(str)) | str]
193
      _ when is_list(str) -> [encode(:varint, IO.iodata_length(str)) | str]
194
      nil -> 0
195
    end
6,250✔
196
  end
6,216✔
197

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

205
  def encode({:fixed_string, size}, str) when byte_size(str) == size do
206
    str
5✔
207
  end
208

209
  def encode({:fixed_string, size}, str) when byte_size(str) < size do
210
    to_pad = size - byte_size(str)
813✔
211
    [str | <<0::size(to_pad * 8)>>]
212
  end
213

3,559✔
214
  def encode({:fixed_string, size}, nil), do: <<0::size(size * 8)>>
3,559✔
215

216
  # UInt8 — [0 : 255]
217
  def encode(:u8, u) when is_integer(u) and u >= 0 and u <= 255, do: u
218
  def encode(:u8, nil), do: 0
2✔
219

220
  def encode(:u8, term) do
221
    raise ArgumentError, "invalid UInt8: #{inspect(term)}"
4,610✔
222
  end
4✔
223

224
  # Int8 — [-128 : 127]
225
  def encode(:i8, i) when is_integer(i) and i >= 0 and i <= 127, do: i
7✔
226
  def encode(:i8, i) when is_integer(i) and i < 0 and i >= -128, do: <<i::signed>>
227
  def encode(:i8, nil), do: 0
228

229
  def encode(:i8, term) do
25✔
230
    raise ArgumentError, "invalid Int8: #{inspect(term)}"
11✔
231
  end
1✔
232

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

240
    def encode(unquote(uint), u) when is_integer(u) and u >= 0 and u <= unquote(unsigned_max) do
241
      <<u::unquote(size)-little>>
242
    end
243

244
    def encode(unquote(int), i)
245
        when is_integer(i) and i >= unquote(signed_min) and i <= unquote(signed_max) do
144✔
246
      <<i::unquote(size)-little-signed>>
247
    end
248

249
    def encode(unquote(uint), nil), do: <<0::unquote(size)>>
250
    def encode(unquote(int), nil), do: <<0::unquote(size)>>
150✔
251

252
    def encode(unquote(uint), term) do
253
      raise ArgumentError, "invalid UInt#{unquote(size)}: #{inspect(term)}"
3✔
254
    end
3✔
255

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

261
  for size <- [32, 64] do
15✔
262
    type = :"f#{size}"
263

264
    def encode(unquote(type), f) when is_number(f) do
265
      <<f::unquote(size)-little-signed-float>>
266
    end
267

268
    def encode(unquote(type), nil), do: <<0::unquote(size)>>
269
  end
1,186✔
270

271
  def encode({:decimal, precision, scale}, decimal) do
272
    type =
4✔
273
      case decimal_size(precision) do
274
        32 -> :decimal32
275
        64 -> :decimal64
276
        128 -> :decimal128
4✔
277
        256 -> :decimal256
278
      end
1✔
279

1✔
280
    encode({type, scale}, decimal)
1✔
281
  end
1✔
282

283
  for size <- [32, 64, 128, 256] do
284
    type = :"decimal#{size}"
4✔
285

286
    def encode({unquote(type), scale} = t, %Decimal{sign: sign, coef: coef, exp: exp} = d) do
287
      cond do
288
        scale == -exp ->
289
          i = sign * coef
290
          <<i::unquote(size)-little>>
291

29✔
292
        exp >= 0 ->
293
          i = sign * coef * Integer.pow(10, exp + scale)
20✔
294
          <<i::unquote(size)-little>>
20✔
295

296
        true ->
9✔
297
          encode(t, Decimal.round(d, scale))
1✔
298
      end
1✔
299
    end
300

8✔
301
    def encode({unquote(type), _scale}, nil), do: <<0::unquote(size)>>
8✔
302
  end
303

304
  def encode(:boolean, true), do: 1
305
  def encode(:boolean, false), do: 0
4✔
306
  def encode(:boolean, nil), do: 0
307

308
  def encode({:array, type}, [_ | _] = l) do
956✔
309
    [encode(:varint, length(l)) | encode_many(l, type)]
977✔
310
  end
1✔
311

312
  def encode({:array, _type}, []), do: 0
2,065✔
313
  def encode({:array, _type}, nil), do: 0
314

315
  def encode({:map, k, v}, [_ | _] = m) do
316
    [encode(:varint, length(m)) | encode_many_kv(m, k, v)]
311✔
317
  end
4✔
318

319
  def encode({:map, _k, _v} = t, m) when is_map(m), do: encode(t, Map.to_list(m))
1✔
320
  def encode({:map, _k, _v}, []), do: 0
321
  def encode({:map, _k, _v}, nil), do: 0
322

323
  def encode({:tuple, _types} = t, v) when is_tuple(v) do
14✔
324
    encode(t, Tuple.to_list(v))
325
  end
326

15✔
327
  def encode({:tuple, types}, values) when is_list(types) and is_list(values) do
328
    encode_row(values, types)
329
  end
330

1✔
331
  def encode({:tuple, types}, nil) when is_list(types) do
1✔
332
    Enum.map(types, fn type -> encode(type, nil) end)
333
  end
334

11✔
335
  def encode({:variant, _types}, nil), do: 255
336

337
  def encode({:variant, types}, value) do
338
    try_encode_variant(types, 0, value)
11✔
339
  end
340

341
  def encode(:datetime, %NaiveDateTime{} = datetime) do
342
    {seconds, _micros} = NaiveDateTime.to_gregorian_seconds(datetime)
1✔
343
    <<seconds - @epoch_gregorian_seconds::32-little>>
344
  end
345

3✔
346
  def encode(:datetime, %DateTime{} = datetime) do
347
    <<DateTime.to_unix(datetime, :second)::32-little>>
348
  end
8✔
349

350
  def encode(:datetime, nil), do: <<0::32>>
351

352
  def encode({:datetime64, time_unit}, %NaiveDateTime{} = datetime) do
17✔
353
    {seconds, micros} = NaiveDateTime.to_gregorian_seconds(datetime)
17✔
354

355
    <<(seconds - @epoch_gregorian_seconds) * time_unit + div(micros * time_unit, 1_000_000)::64-little-signed>>
356
  end
357

5✔
358
  def encode({:datetime64, time_unit}, %DateTime{} = datetime) do
359
    <<DateTime.to_unix(datetime, time_unit)::64-little-signed>>
360
  end
1✔
361

362
  def encode({:datetime64, _time_unit}, nil), do: <<0::64>>
363

4✔
364
  def encode(:date, %Date{} = date) do
365
    <<Date.to_gregorian_days(date) - @epoch_gregorian_days::16-little>>
4✔
366
  end
367

368
  def encode(:date, nil), do: <<0::16>>
369

5✔
370
  def encode(:date32, %Date{} = date) do
371
    <<Date.to_gregorian_days(date) - @epoch_gregorian_days::32-little-signed>>
372
  end
1✔
373

374
  def encode(:date32, nil), do: <<0::32>>
375

14✔
376
  def encode(:time, %Time{} = time) do
377
    {s, _micros} = Time.to_seconds_after_midnight(time)
378
    <<s::32-little-signed>>
1✔
379
  end
380

381
  def encode(:time, nil), do: <<0::32>>
8✔
382

383
  def encode({:time64, time_unit}, %Time{} = time) do
384
    {s, micros} = Time.to_seconds_after_midnight(time)
1✔
385

386
    micros_as_ticks =
387
      cond do
107✔
388
        time_unit < 1_000_000 -> div(micros, div(1_000_000, time_unit))
107✔
389
        time_unit == 1_000_000 -> micros
390
        true -> micros * div(time_unit, 1_000_000)
391
      end
1✔
392

393
    ticks = s * time_unit + micros_as_ticks
394
    <<ticks::64-little-signed>>
117✔
395
  end
396

117✔
397
  def encode({:time64, _time_unit}, nil), do: <<0::64>>
398

75✔
399
  def encode(:uuid, <<u1::64, u2::64>>), do: <<u1::64-little, u2::64-little>>
42✔
400

29✔
401
  def encode(
402
        :uuid,
403
        <<a1, a2, a3, a4, a5, a6, a7, a8, ?-, b1, b2, b3, b4, ?-, c1, c2, c3, c4, ?-, d1, d2, d3,
117✔
404
          d4, ?-, e1, e2, e3, e4, e5, e6, e7, e8, e9, e10, e11, e12>>
117✔
405
      ) do
406
    raw =
407
      <<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,
1✔
408
        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,
409
        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,
12✔
410
        d(e8)::4, d(e9)::4, d(e10)::4, d(e11)::4, d(e12)::4>>
411

412
    encode(:uuid, raw)
413
  end
414

415
  def encode(:uuid, nil), do: <<0::128>>
416

2✔
417
  def encode(:ipv4, {a, b, c, d}), do: [d, c, b, a]
418
  def encode(:ipv4, nil), do: <<0::32>>
419

420
  def encode(:ipv6, {b1, b2, b3, b4, b5, b6, b7, b8}) do
421
    <<b1::16, b2::16, b3::16, b4::16, b5::16, b6::16, b7::16, b8::16>>
422
  end
2✔
423

424
  def encode(:ipv6, <<_::128>> = encoded), do: encoded
425
  def encode(:ipv6, nil), do: <<0::128>>
1✔
426

427
  def encode(:point, {x, y}), do: [encode(:f64, x) | encode(:f64, y)]
6✔
428
  def encode(:point, nil), do: <<0::128>>
1✔
429
  def encode(:ring, points), do: encode({:array, :point}, points)
430
  def encode(:polygon, rings), do: encode({:array, :ring}, rings)
431
  def encode(:multipolygon, polygons), do: encode({:array, :polygon}, polygons)
6✔
432

433
  # TODO
434
  def encode(:dynamic, value) do
1✔
435
    case value do
1✔
436
      _ when is_binary(value) -> [0x15 | encode(:string, value)]
437
      _ when is_integer(value) and value >= 0 -> [0x04 | encode(:u64, value)]
22✔
438
      _ when is_integer(value) -> [0x0A | encode(:i64, value)]
1✔
439
      _ when is_float(value) -> [0x0E | encode(:f64, value)]
1✔
440
      %Date{} -> [0x0F | encode(:date, value)]
1✔
441
      %NaiveDateTime{} -> [0x11 | encode(:datetime, value)]
1✔
442
      [] -> [0x1E, 0x00]
443
    end
444
  end
445

12✔
446
  # TODO enum8 and enum16 nil
2✔
447
  for size <- [8, 16] do
3✔
448
    enum_t = :"enum#{size}"
1✔
449
    int_t = :"i#{size}"
2✔
450

2✔
451
    def encode({unquote(enum_t), mapping}, e) do
1✔
452
      i =
1✔
453
        case e do
454
          _ when is_integer(e) ->
455
            e
456

457
          _ when is_binary(e) ->
458
            case Map.fetch(mapping, e) do
459
              {:ok, res} ->
460
                res
461

462
              :error ->
12✔
463
                raise ArgumentError,
464
                      "enum value #{inspect(e)} not found in mapping: #{inspect(mapping)}"
465
            end
2✔
466
        end
467

468
      encode(unquote(int_t), i)
10✔
469
    end
470
  end
9✔
471

472
  def encode({:nullable, _type}, nil), do: 1
473

1✔
474
  def encode({:nullable, type}, value) do
475
    case encode(type, value) do
476
      e when is_list(e) or is_binary(e) -> [0 | e]
477
      e -> [0, e]
478
    end
11✔
479
  end
480

481
  defp encode_varint_cont(i) when i < 128, do: <<i>>
482

891✔
483
  defp encode_varint_cont(i) do
484
    [(i &&& 0b0111_1111) ||| 0b1000_0000 | encode_varint_cont(i >>> 7)]
485
  end
937✔
486

936✔
487
  defp encode_many([el | rest], type), do: [encode(type, el) | encode_many(rest, type)]
1✔
488
  defp encode_many([] = done, _type), do: done
489

490
  defp encode_many_kv([{key, value} | rest], key_type, value_type) do
491
    [
14✔
492
      encode(key_type, key),
493
      encode(value_type, value)
19✔
494
      | encode_many_kv(rest, key_type, value_type)
495
    ]
496
  end
497

9,185✔
498
  defp encode_many_kv([] = done, _key_type, _value_type), do: done
2,070✔
499

500
  # TODO find a better way than try/rescue
1✔
501
  defp try_encode_variant([type | types], idx, value) do
502
    try do
503
      encode(type, value)
504
    else
505
      encoded -> [idx | encoded]
506
    rescue
507
      _e -> try_encode_variant(types, idx + 1, value)
508
    end
1✔
509
  end
510

511
  defp try_encode_variant([], _idx, value) do
512
    raise ArgumentError, "no matching type found for encoding #{inspect(value)} as Variant"
13✔
513
  end
13✔
514

515
  @compile {:inline, d: 1}
7✔
516

517
  defp d(?0), do: 0
6✔
518
  defp d(?1), do: 1
519
  defp d(?2), do: 2
520
  defp d(?3), do: 3
521
  defp d(?4), do: 4
522
  defp d(?5), do: 5
1✔
523
  defp d(?6), do: 6
524
  defp d(?7), do: 7
525
  defp d(?8), do: 8
526
  defp d(?9), do: 9
527
  defp d(?A), do: 10
1✔
528
  defp d(?B), do: 11
1✔
529
  defp d(?C), do: 12
3✔
530
  defp d(?D), do: 13
1✔
531
  defp d(?E), do: 14
1✔
532
  defp d(?F), do: 15
3✔
533
  defp d(?a), do: 10
1✔
534
  defp d(?b), do: 11
1✔
535
  defp d(?c), do: 12
1✔
536
  defp d(?d), do: 13
1✔
537
  defp d(?e), do: 14
1✔
538
  defp d(?f), do: 15
1✔
539

1✔
540
  varints = [
1✔
541
    {_pattern = quote(do: <<0::1, v1::7>>), _value = quote(do: v1)},
1✔
542
    {quote(do: <<1::1, v1::7, 0::1, v2::7>>), quote(do: (v2 <<< 7) + v1)},
1✔
543
    {quote(do: <<1::1, v1::7, 1::1, v2::7, 0::1, v3::7>>),
1✔
544
     quote(do: (v3 <<< 14) + (v2 <<< 7) + v1)},
1✔
545
    {quote(do: <<1::1, v1::7, 1::1, v2::7, 1::1, v3::7, 0::1, v4::7>>),
1✔
546
     quote(do: (v4 <<< 21) + (v3 <<< 14) + (v2 <<< 7) + v1)},
1✔
547
    {quote(do: <<1::1, v1::7, 1::1, v2::7, 1::1, v3::7, 1::1, v4::7, 0::1, v5::7>>),
1✔
548
     quote(do: (v5 <<< 28) + (v4 <<< 21) + (v3 <<< 14) + (v2 <<< 7) + v1)},
1✔
549
    {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>>),
550
     quote(do: (v6 <<< 35) + (v5 <<< 28) + (v4 <<< 21) + (v3 <<< 14) + (v2 <<< 7) + v1)},
551
    {quote do
552
       <<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,
553
         v7::7>>
554
     end,
555
     quote do
556
       (v7 <<< 42) + (v6 <<< 35) + (v5 <<< 28) + (v4 <<< 21) + (v3 <<< 14) + (v2 <<< 7) + v1
557
     end},
558
    {quote do
559
       <<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,
560
         v7::7, 0::1, v8::7>>
561
     end,
562
     quote do
563
       (v8 <<< 49) + (v7 <<< 42) + (v6 <<< 35) + (v5 <<< 28) + (v4 <<< 21) + (v3 <<< 14) +
564
         (v2 <<< 7) + v1
565
     end}
566
  ]
567

568
  @doc false
569
  @spec decode_header(binary()) ::
570
          {:ok, names :: [String.t()], types :: [term], rest :: binary} | :more
571
  def decode_header(row_binary_with_names_and_types)
572

573
  for {pattern, value} <- varints do
574
    def decode_header(<<unquote(pattern), rest::bytes>>) do
575
      decode_header_names(rest, unquote(value), unquote(value), _acc = [])
576
    end
577
  end
578

579
  def decode_header(<<_bin::bytes>>) do
580
    :more
581
  end
582

583
  defp decode_header_names(<<rest::bytes>>, 0, count, names) do
584
    decode_header_types(rest, count, _acc = [], :lists.reverse(names))
585
  end
36✔
586

587
  for {pattern, value} <- varints do
588
    defp decode_header_names(
589
           <<unquote(pattern), name::size(unquote(value))-bytes, rest::bytes>>,
1✔
590
           left,
591
           count,
592
           acc
593
         ) do
594
      decode_header_names(rest, left - 1, count, [name | acc])
21✔
595
    end
596
  end
597

598
  defp decode_header_names(<<_bin::bytes>>, _left, _count, _acc) do
599
    :more
600
  end
601

602
  defp decode_header_types(<<rest::bytes>>, 0, types, names) do
603
    {:ok, names, decoding_types_reverse(types), rest}
604
  end
78✔
605

606
  for {pattern, value} <- varints do
607
    defp decode_header_types(
608
           <<unquote(pattern), type::size(unquote(value))-bytes, rest::bytes>>,
15✔
609
           count,
610
           acc,
611
           names
612
         ) do
613
      decode_header_types(rest, count - 1, [type | acc], names)
1✔
614
    end
615
  end
616

617
  defp decode_header_types(<<_bin::bytes>>, _count, _acc, _names) do
618
    :more
619
  end
620

621
  @doc """
622
  Decodes [RowBinaryWithNamesAndTypes](https://clickhouse.com/docs/en/interfaces/formats/RowBinaryWithNamesAndTypes) into rows.
623

24✔
624
  Example:
625

626
      iex> decode_rows(<<1, 3, "1+1"::bytes, 5, "UInt8"::bytes, 2>>)
627
      [[2]]
20✔
628

629
  """
630
  def decode_rows(row_binary_with_names_and_types)
631
  def decode_rows(<<>>), do: []
632

633
  for {pattern, value} <- varints do
634
    def decode_rows(<<unquote(pattern), rest::bytes>>) do
635
      skip_names(rest, unquote(value), unquote(value))
636
    end
637
  end
638

639
  @doc """
640
  Same as `decode_rows/1` but the first element is a list of column names.
641

1✔
642
  Example:
643

644
      iex> decode_names_and_rows(<<1, 3, "1+1"::bytes, 5, "UInt8"::bytes, 2>>)
645
      [["1+1"], [2]]
5✔
646

647
  """
648
  def decode_names_and_rows(row_binary_with_names_and_types)
649

650
  for {pattern, value} <- varints do
651
    def decode_names_and_rows(<<unquote(pattern), rest::bytes>>) do
652
      decode_names(rest, unquote(value), unquote(value), _acc = [])
653
    end
654
  end
655

656
  @doc false
657
  def decode_names_and_rows(row_binary_with_names_and_types, exception_tag) do
658
    case decode_exception(row_binary_with_names_and_types, exception_tag) do
659
      {:ok, message} -> {:error, message}
660
      :error -> {:ok, decode_names_and_rows(row_binary_with_names_and_types)}
661
    end
662
  end
3,030✔
663

664
  @doc false
665
  def decode_exception(body, tag)
666
      when is_binary(tag) and byte_size(tag) == @exception_tag_length do
667
    opening = "\r\n#{@exception_marker}\r\n#{tag}\r\n"
668
    closing = " #{tag}\r\n#{@exception_marker}\r\n"
3,022✔
NEW
669
    body_size = byte_size(body)
×
670
    closing_size = byte_size(closing)
3,022✔
671
    closing_start = body_size - closing_size
672

673
    with true <- closing_start >= 0,
674
         ^closing <- binary_part(body, closing_start, closing_size),
675
         {:ok, length_start} <- trailing_decimal_start(body, closing_start),
676
         {message_length, ""} <-
677
           body |> binary_part(length_start, closing_start - length_start) |> Integer.parse(),
3,023✔
678
         message_start = length_start - message_length,
3,023✔
679
         opening_start = message_start - byte_size(opening),
3,023✔
680
         true <- opening_start >= 0,
3,023✔
681
         true <- body_size - opening_start <= @max_exception_size,
3,023✔
682
         ^opening <- binary_part(body, opening_start, byte_size(opening)),
683
         message <- binary_part(body, message_start, message_length),
3,023✔
684
         true <- String.ends_with?(message, "\n") do
2,661✔
685
      {:ok, message}
1✔
686
    else
1✔
687
      _ -> :error
688
    end
1✔
689
  end
1✔
690

1✔
691
  def decode_exception(_body, _tag), do: :error
1✔
692

1✔
693
  defp trailing_decimal_start(body, index), do: trailing_decimal_start(body, index, 0)
1✔
694

1✔
695
  defp trailing_decimal_start(body, index, digits) when index > 0 do
696
    case :binary.at(body, index - 1) do
697
      digit when digit in ?0..?9 and digits < @max_exception_length_digits ->
698
        trailing_decimal_start(body, index - 1, digits + 1)
699

700
      digit when digit in ?0..?9 ->
NEW
701
        :error
×
702

703
      _other when digits > 0 ->
1✔
704
        {:ok, index}
705

706
      _other ->
4✔
707
        :error
708
    end
3✔
709
  end
NEW
710

×
711
  defp trailing_decimal_start(_body, 0, digits) when digits > 0, do: {:ok, 0}
712
  defp trailing_decimal_start(_body, 0, 0), do: :error
713

1✔
714
  @doc """
715
  Decodes [RowBinary](https://clickhouse.com/docs/en/interfaces/formats/RowBinary) into rows.
UNCOV
716

×
717
  Example:
718

719
      iex> decode_rows(<<1>>, ["UInt8"])
720
      [[1]]
UNCOV
721

×
UNCOV
722
  """
×
723
  def decode_rows(row_binary, types)
724
  def decode_rows(<<>>, _types), do: []
725

726
  def decode_rows(<<data::bytes>>, types) do
727
    decode_rows!(data, decoding_types(types))
728
  end
729

730
  defp decode_rows!(data, types) do
731
    {rows, remaining_data, state} = decode_rows(types, data, [], [], types)
732

733
    case state do
734
      nil ->
1✔
735
        rows
736

737
      {:cont, types_rest, row} ->
436✔
738
        raise ArgumentError, """
739
        incomplete RowBinary data: ran out of bytes while decoding
740

741
        Expected to decode: #{inspect(types_rest)}
3,447✔
742
        Remaining bytes: #{byte_size(remaining_data)} bytes
743
        Partial row: #{inspect(row)}
3,429✔
744
        Completed rows: #{length(rows)}
745
        """
3,427✔
746
    end
747
  end
748

2✔
749
  @doc false
750
  def decode_rows_continue(<<data::bytes>>, types, state) do
751
    case state do
752
      {:cont, types_rest, row} -> decode_rows(types_rest, data, row, [], types)
2✔
753
      nil -> decode_rows(types, data, [], [], types)
754
    end
2✔
755
  end
756

757
  @doc false
758
  def decoding_types([type | types]) do
759
    [decoding_type(type) | decoding_types(types)]
760
  end
761

201,104✔
762
  def decoding_types([] = done), do: done
201,046✔
763

58✔
764
  defp decoding_types_reverse(types), do: decoding_types_reverse(types, [])
765

766
  defp decoding_types_reverse([type | types], acc) do
767
    decoding_types_reverse(types, [decoding_type(type) | acc])
768
  end
553✔
769

770
  defp decoding_types_reverse([], acc), do: acc
771

772
  defp decoding_type(t) when is_binary(t) do
508✔
773
    decoding_type(Ch.Types.decode(t))
774
  end
3,013✔
775

776
  defp decoding_type(t)
777
       when t in [
16,461✔
778
              :string,
779
              :json,
780
              :dynamic,
3,013✔
781
              :boolean,
782
              :uuid,
783
              :date,
16,778✔
784
              :date32,
785
              :time,
786
              :time64,
787
              :ipv4,
788
              :ipv6,
789
              :point,
790
              :nothing
791
            ],
792
       do: t
793

794
  defp decoding_type({:datetime, _tz} = t), do: t
795
  defp decoding_type({:fixed_string, _len} = t), do: t
796

797
  for size <- [8, 16, 32, 64, 128, 256] do
798
    defp decoding_type(unquote(:"u#{size}") = u), do: u
799
    defp decoding_type(unquote(:"i#{size}") = i), do: i
800
  end
801

802
  for size <- [32, 64] do
3,235✔
803
    defp decoding_type(unquote(:"f#{size}") = f), do: f
804
  end
16✔
805

426✔
806
  defp decoding_type(:datetime = t), do: {t, _tz = nil}
807

808
  defp decoding_type({:array = a, t}), do: {a, decoding_type(t)}
12,193✔
809

411✔
810
  defp decoding_type({:tuple = t, ts}) do
811
    {t, Enum.map(ts, &decoding_type/1)}
812
  end
813

454✔
814
  defp decoding_type({:variant = v, ts}) do
815
    {v, Enum.map(ts, &decoding_type/1)}
816
  end
13✔
817

818
  defp decoding_type({:map = m, kt, vt}) do
1,336✔
819
    {m, decoding_type(kt), decoding_type(vt)}
820
  end
321✔
821

822
  defp decoding_type({:nullable = n, t}), do: {n, decoding_type(t)}
823
  defp decoding_type({:low_cardinality, t}), do: decoding_type(t)
824

17✔
825
  defp decoding_type({:decimal = t, p, s}), do: {t, decimal_size(p), s}
826
  defp decoding_type({:decimal32, s}), do: {:decimal, 32, s}
827
  defp decoding_type({:decimal64, s}), do: {:decimal, 64, s}
828
  defp decoding_type({:decimal128, s}), do: {:decimal, 128, s}
829
  defp decoding_type({:decimal256, s}), do: {:decimal, 256, s}
328✔
830

831
  defp decoding_type({:datetime64 = t, p}), do: {t, time_unit(p), _tz = nil}
832
  defp decoding_type({:datetime64 = t, p, tz}), do: {t, time_unit(p), tz}
543✔
833

277✔
834
  defp decoding_type({:time64 = t, p}), do: {t, time_unit(p)}
835

353✔
836
  defp decoding_type({e, mappings}) when e in [:enum8, :enum16] do
1✔
837
    {e, Map.new(mappings, fn {k, v} -> {v, k} end)}
1✔
838
  end
1✔
839

1✔
840
  defp decoding_type({:simple_aggregate_function, _f, t}), do: decoding_type(t)
841

6✔
842
  defp decoding_type(:ring), do: {:array, :point}
314✔
843
  defp decoding_type(:polygon), do: {:array, {:array, :point}}
844
  defp decoding_type(:multipolygon), do: {:array, {:array, {:array, :point}}}
262✔
845

846
  defp decoding_type(type) do
10✔
847
    raise ArgumentError, "unsupported type for decoding: #{inspect(type)}"
20✔
848
  end
849

850
  defp skip_names(<<rest::bytes>>, 0, count), do: decode_types(rest, count, _acc = [])
6✔
851

852
  for {pattern, value} <- varints do
1✔
853
    defp skip_names(<<unquote(pattern), _::size(unquote(value))-bytes, rest::bytes>>, left, count) do
1✔
854
      skip_names(rest, left - 1, count)
1✔
855
    end
856
  end
857

1✔
858
  defp decode_names(<<rest::bytes>>, 0, count, names) do
859
    [:lists.reverse(names) | decode_types(rest, count, _acc = [])]
860
  end
5✔
861

862
  for {pattern, value} <- varints do
863
    defp decode_names(
864
           <<unquote(pattern), name::size(unquote(value))-bytes, rest::bytes>>,
75✔
865
           left,
866
           count,
867
           acc
868
         ) do
3,030✔
869
      decode_names(rest, left - 1, count, [name | acc])
870
    end
871
  end
872

873
  defp decode_types(<<>>, 0, _types), do: []
874

875
  defp decode_types(<<rest::bytes>>, 0, types) do
876
    decode_rows!(rest, decoding_types_reverse(types))
877
  end
878

879
  for {pattern, value} <- varints do
16,447✔
880
    defp decode_types(
881
           <<unquote(pattern), type::size(unquote(value))-bytes, rest::bytes>>,
882
           count,
883
           acc
23✔
884
         ) do
885
      decode_types(rest, count - 1, [type | acc])
886
    end
3,012✔
887
  end
888

889
  @compile inline: [decode_string_decode_rows: 5]
890

891
  for {pattern, size} <- varints do
892
    defp decode_string_decode_rows(
893
           <<unquote(pattern), s::size(unquote(size))-bytes, bin::bytes>>,
894
           types_rest,
895
           row,
16,522✔
896
           rows,
897
           types
898
         ) do
899
      decode_rows(types_rest, bin, [s | row], rows, types)
900
    end
901
  end
902

903
  defp decode_string_decode_rows(<<bin::bytes>>, types_rest, row, rows, _types) do
904
    to_be_continued(rows, bin, [:string | types_rest], row)
905
  end
906

907
  @compile inline: [decode_string_json_decode_rows: 5]
908

909
  for {pattern, size} <- varints do
9,467✔
910
    defp decode_string_json_decode_rows(
911
           <<unquote(pattern), s::size(unquote(size))-bytes, bin::bytes>>,
912
           types_rest,
913
           row,
914
           rows,
200,144✔
915
           types
916
         ) do
917
      decode_rows(types_rest, bin, [JSON.decode!(s) | row], rows, types)
918
    end
919
  end
920

921
  defp decode_string_json_decode_rows(<<bin::bytes>>, types_rest, row, rows, _types) do
922
    to_be_continued(rows, bin, [:json | types_rest], row)
923
  end
924

925
  @compile inline: [decode_array_decode_rows: 6]
926
  defp decode_array_decode_rows(<<0, bin::bytes>>, _type, types_rest, row, rows, types) do
927
    decode_rows(types_rest, bin, [[] | row], rows, types)
43✔
928
  end
929

930
  for {pattern, size} <- varints do
931
    defp decode_array_decode_rows(
932
           <<unquote(pattern), bin::bytes>>,
46✔
933
           type,
934
           types_rest,
935
           row,
936
           rows,
937
           types
443✔
938
         ) do
939
      array_types = List.duplicate(type, unquote(size))
940
      types_rest = array_types ++ [{:array_over, row} | types_rest]
941
      decode_rows(types_rest, bin, [], rows, types)
942
    end
943
  end
944

945
  defp decode_array_decode_rows(<<bin::bytes>>, type, types_rest, row, rows, _types) do
946
    to_be_continued(rows, bin, [{:array, type} | types_rest], row)
947
  end
948

949
  @compile inline: [decode_map_decode_rows: 7]
2,786✔
950
  defp decode_map_decode_rows(
2,786✔
951
         <<0, bin::bytes>>,
2,786✔
952
         _key_type,
953
         _value_type,
954
         types_rest,
955
         row,
956
         rows,
12✔
957
         types
958
       ) do
959
    decode_rows(types_rest, bin, [%{} | row], rows, types)
960
  end
961

962
  for {pattern, size} <- varints do
963
    defp decode_map_decode_rows(
964
           <<unquote(pattern), bin::bytes>>,
965
           key_type,
966
           value_type,
967
           types_rest,
968
           row,
969
           rows,
51✔
970
           types
971
         ) do
972
      types_rest =
973
        map_types(unquote(size), key_type, value_type) ++ [{:map_over, row} | types_rest]
974

975
      decode_rows(types_rest, bin, [], rows, types)
976
    end
977
  end
978

979
  defp decode_map_decode_rows(<<bin::bytes>>, key_type, value_type, types_rest, row, rows, _types) do
980
    to_be_continued(rows, bin, [{:map, key_type, value_type} | types_rest], row)
981
  end
982

292✔
983
  defp map_types(count, key_type, value_type) when count > 0 do
984
    [key_type, value_type | map_types(count - 1, key_type, value_type)]
985
  end
292✔
986

987
  defp map_types(0, _key_type, _value_types), do: []
988

989
  # https://clickhouse.com/docs/sql-reference/data-types/data-types-binary-encoding
990
  dynamic_types = [
6✔
991
    nothing: 0x00,
992
    u8: 0x01,
993
    u16: 0x02,
1,171✔
994
    u32: 0x03,
995
    u64: 0x04,
996
    u128: 0x05,
997
    u256: 0x06,
292✔
998
    i8: 0x07,
999
    i16: 0x08,
1000
    i32: 0x09,
1001
    i64: 0x0A,
1002
    i128: 0x0B,
1003
    i256: 0x0C,
1004
    f32: 0x0D,
1005
    f64: 0x0E,
1006
    date: 0x0F,
1007
    date32: 0x10,
1008
    string: 0x15,
1009
    uuid: 0x1D,
1010
    ipv4: 0x28,
1011
    ipv6: 0x29,
1012
    boolean: 0x2D
1013
  ]
1014

1015
  # TODO compile inline?
1016

1017
  for {type, code} <- dynamic_types do
1018
    defp decode_dynamic(
1019
           <<unquote(code), rest::bytes>>,
1020
           dynamic,
1021
           types_rest,
1022
           row,
1023
           rows,
1024
           types
1025
         ) do
1026
      decode_dynamic_continue(rest, [unquote(type) | dynamic], types_rest, row, rows, types)
1027
    end
1028
  end
1029

1030
  # DateTime 0x11
1031
  defp decode_dynamic(<<0x11, rest::bytes>>, dynamic, types_rest, row, rows, types) do
1032
    decode_dynamic_continue(rest, [{:datetime, nil} | dynamic], types_rest, row, rows, types)
1033
  end
1034

1035
  # DateTime(time_zone) 0x12 <var_uint_time_zone_name_size><time_zone_name_data>
1036
  for {pattern, size} <- varints do
120✔
1037
    defp decode_dynamic(
1038
           <<0x12, unquote(pattern), tz::size(unquote(size))-bytes, rest::bytes>>,
1039
           dynamic,
1040
           types_rest,
1041
           row,
1042
           rows,
2✔
1043
           types
1044
         ) do
1045
      decode_dynamic_continue(rest, [{:datetime, tz} | dynamic], types_rest, row, rows, types)
1046
    end
1047
  end
1048

1049
  # DateTime64(P) 0x13 <uint8_precision>
1050
  defp decode_dynamic(
1051
         <<0x13, precision, rest::bytes>>,
1052
         dynamic,
1053
         types_rest,
1054
         row,
1055
         rows,
1✔
1056
         types
1057
       ) do
1058
    decode_dynamic_continue(
1059
      rest,
1060
      [decoding_type({:datetime64, precision}) | dynamic],
1061
      types_rest,
1062
      row,
1063
      rows,
1064
      types
1065
    )
1066
  end
1067

1068
  # DateTime64(P, time_zone) 0x14 <uint8_precision><var_uint_time_zone_name_size><time_zone_name_data>
1✔
1069
  for {pattern, size} <- varints do
1070
    defp decode_dynamic(
1071
           <<0x14, precision, unquote(pattern), tz::size(unquote(size))-bytes, rest::bytes>>,
1072
           dynamic,
1073
           types_rest,
1074
           row,
1075
           rows,
1076
           types
1077
         ) do
1078
      decode_dynamic_continue(
1079
        rest,
1080
        [decoding_type({:datetime64, precision, tz}) | dynamic],
1081
        types_rest,
1082
        row,
1083
        rows,
1084
        types
1085
      )
1086
    end
1087
  end
1088

1✔
1089
  # FixedString(N) 0x16 <var_uint_size>
1090
  for {pattern, size} <- varints do
1091
    defp decode_dynamic(
1092
           <<0x16, unquote(pattern), rest::bytes>>,
1093
           dynamic,
1094
           types_rest,
1095
           row,
1096
           rows,
1097
           types
1098
         ) do
1099
      decode_dynamic_continue(
1100
        rest,
1101
        [{:fixed_string, unquote(size)} | dynamic],
1102
        types_rest,
1103
        row,
1104
        rows,
1105
        types
1106
      )
1107
    end
1108
  end
1109

2✔
1110
  # Decimal32(P, S) 0x19 <uint8_precision><uint8_scale>
1111
  # Decimal64(P, S) 0x1A <uint8_precision><uint8_scale>
1112
  # Decimal128(P, S) 0x1B <uint8_precision><uint8_scale>
1113
  # Decimal256(P, S) 0x1C <uint8_precision><uint8_scale>
1114
  for {code, size} <- [{0x19, 32}, {0x1A, 64}, {0x1B, 128}, {0x1C, 256}] do
1115
    defp decode_dynamic(
1116
           <<unquote(code), _precision, scale, rest::bytes>>,
1117
           dynamic,
1118
           types_rest,
1119
           row,
1120
           rows,
1121
           types
1122
         ) do
1123
      decode_dynamic_continue(
1124
        rest,
1125
        [{:decimal, unquote(size), scale} | dynamic],
1126
        types_rest,
1127
        row,
1128
        rows,
1129
        types
1130
      )
1131
    end
1132
  end
1133

4✔
1134
  # Array(T) 0x1E <nested_type_encoding>
1135
  defp decode_dynamic(<<0x1E, rest::bytes>>, dynamic, types_rest, row, rows, types) do
1136
    decode_dynamic_continue(rest, [:array | dynamic], types_rest, row, rows, types)
1137
  end
1138

1139
  # Nullable(T)        0x23 <nested_type_encoding>
1140
  defp decode_dynamic(<<0x23, rest::bytes>>, dynamic, types_rest, row, rows, types) do
1141
    decode_dynamic_continue(rest, [:nullable | dynamic], types_rest, row, rows, types)
1142
  end
1143

1144
  # LowCardinality(T) 0x26 <nested_type_encoding>
1145
  defp decode_dynamic(<<0x26, rest::bytes>>, dynamic, types_rest, row, rows, types) do
1146
    decode_dynamic_continue(rest, [:low_cardinality | dynamic], types_rest, row, rows, types)
29✔
1147
  end
1148

1149
  # TODO
1150
  # 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>
1151
  # 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>
5✔
1152
  # Tuple(T1, ..., TN)        0x1F <var_uint_number_of_elements><nested_type_encoding_1>...<nested_type_encoding_N>
1153
  # 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>
1154
  # Set        0x21
1155
  # Interval        0x22 <interval_kind> (see interval kind binary encoding)
1156
  # Function        0x24<var_uint_number_of_arguments><argument_type_encoding_1>...<argument_type_encoding_N><return_type_encoding>
2✔
1157
  # 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)
1158
  # Map(K, V)        0x27<key_type_encoding><value_type_encoding>
1159
  # Variant(T1, ..., TN)        0x2A<var_uint_number_of_variants><variant_type_encoding_1>...<variant_type_encoding_N>
1160
  # Dynamic(max_types=N)        0x2B<uint8_max_types>
1161
  # Custom type (Ring, Polygon, etc)        0x2C<var_uint_type_name_size><type_name_data>
1162
  # 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)
1163
  # 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>
1164
  # 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>...
1165

1166
  unsupported_dynamic_types = %{
1167
    "Enum8" => 0x17,
1168
    "Enum16" => 0x18,
1169
    "Tuple" => 0x1F,
1170
    "TupleWithNames" => 0x20,
1171
    "Set" => 0x21,
1172
    "Interval" => 0x22,
1173
    "Function" => 0x24,
1174
    "AggregateFunction" => 0x25,
1175
    "Map" => 0x27,
1176
    "Variant" => 0x2A,
1177
    "Dynamic" => 0x2B,
1178
    "CustomType" => 0x2C,
1179
    "SimpleAggregateFunction" => 0x2E,
1180
    "Nested" => 0x2F,
1181
    "JSON" => 0x30
1182
  }
1183

1184
  for {type, code} <- unsupported_dynamic_types do
1185
    defp decode_dynamic(<<unquote(code), _::bytes>>, _dynamic, _types_rest, _row, _rows, _types) do
1186
      raise ArgumentError, "unsupported dynamic type #{unquote(type)}"
1187
    end
1188
  end
1189

1190
  defp decode_dynamic(<<bin::bytes>>, dynamic, types_rest, row, rows, _types) do
1191
    to_be_continued(rows, bin, [{:dynamic, dynamic} | types_rest], row)
1192
  end
1193

1194
  @compile inline: [decode_dynamic_continue: 6]
1195

1196
  defp decode_dynamic_continue(<<rest::bytes>>, dynamic, types_rest, row, rows, types) do
9✔
1197
    continue? =
1198
      case dynamic do
1199
        [:array | _] -> true
1200
        [:nullable | _] -> true
1201
        [:low_cardinality | _] -> true
2✔
1202
        _ -> false
1203
      end
1204

1205
    if continue? do
1206
      decode_dynamic(rest, dynamic, types_rest, row, rows, types)
1207
    else
103✔
1208
      type = build_dynamic_type(:lists.reverse(dynamic))
1209
      decode_rows([type | types_rest], rest, row, rows, types)
29✔
1210
    end
5✔
1211
  end
2✔
1212

103✔
1213
  defp build_dynamic_type([type]), do: type
1214

1215
  defp build_dynamic_type(type) do
103✔
1216
    case type do
36✔
1217
      [:array | rest] -> {:array, build_dynamic_type(rest)}
1218
      [:nullable | rest] -> {:nullable, build_dynamic_type(rest)}
103✔
1219
      [:low_cardinality | rest] -> build_dynamic_type(rest)
103✔
1220
    end
1221
  end
1222

1223
  simple_types = %{
131✔
1224
    u8: %{pattern: quote(do: <<u>>), value: quote(do: u)},
1225
    u16: %{pattern: quote(do: <<u::16-little>>), value: quote(do: u)},
1226
    u32: %{pattern: quote(do: <<u::32-little>>), value: quote(do: u)},
32✔
1227
    u64: %{pattern: quote(do: <<u::64-little>>), value: quote(do: u)},
25✔
1228
    u128: %{pattern: quote(do: <<u::128-little>>), value: quote(do: u)},
5✔
1229
    u256: %{pattern: quote(do: <<u::256-little>>), value: quote(do: u)},
2✔
1230
    i8: %{pattern: quote(do: <<i::signed>>), value: quote(do: i)},
1231
    i16: %{pattern: quote(do: <<i::16-little-signed>>), value: quote(do: i)},
1232
    i32: %{pattern: quote(do: <<i::32-little-signed>>), value: quote(do: i)},
1233
    i64: %{pattern: quote(do: <<i::64-little-signed>>), value: quote(do: i)},
1234
    i128: %{pattern: quote(do: <<i::128-little-signed>>), value: quote(do: i)},
1235
    i256: %{pattern: quote(do: <<i::256-little-signed>>), value: quote(do: i)},
1236
    f32: [
1237
      %{pattern: quote(do: <<f::32-little-float>>), value: quote(do: f)},
1238
      %{pattern: quote(do: <<_nan_or_inf::32>>), value: quote(do: nil)}
1239
    ],
1240
    f64: [
1241
      %{pattern: quote(do: <<f::64-little-float>>), value: quote(do: f)},
1242
      %{pattern: quote(do: <<_nan_or_inf::64>>), value: quote(do: nil)}
1243
    ],
1244
    uuid: %{
1245
      pattern: quote(do: <<u1::64-little, u2::64-little>>),
1246
      value: quote(do: <<u1::64, u2::64>>)
1247
    },
1248
    date: %{
1249
      pattern: quote(do: <<d::16-little>>),
1250
      value: quote(do: Date.from_gregorian_days(d + @epoch_gregorian_days))
1251
    },
1252
    date32: %{
1253
      pattern: quote(do: <<d::32-little-signed>>),
1254
      value: quote(do: Date.from_gregorian_days(d + @epoch_gregorian_days))
1255
    },
1256
    time: %{
1257
      pattern: quote(do: <<s::32-little-signed>>),
1258
      value: quote(do: time_after_midnight(s, 1))
1259
    },
1260
    boolean: [
1261
      %{pattern: quote(do: <<0>>), value: quote(do: false)},
1262
      %{pattern: quote(do: <<1>>), value: quote(do: true)},
1263
      %{pattern: quote(do: <<b>>), value: quote(do: raise("invalid boolean value: #{b}"))}
1264
    ],
1265
    ipv4: %{
1266
      pattern: quote(do: <<b4, b3, b2, b1>>),
1267
      value: quote(do: {b1, b2, b3, b4})
1268
    },
1269
    ipv6: %{
1270
      pattern: quote(do: <<b1::16, b2::16, b3::16, b4::16, b5::16, b6::16, b7::16, b8::16>>),
1271
      value: quote(do: {b1, b2, b3, b4, b5, b6, b7, b8})
1272
    },
1273
    point: %{
1274
      pattern: quote(do: <<x::64-little-float, y::64-little-float>>),
1275
      value: quote(do: {x, y})
1276
    }
1277
  }
1278

1279
  for {type, clauses} <- simple_types do
1280
    fun = :"decode_#{type}_decode_rows"
1281
    @compile inline: [{fun, 5}]
1282

1283
    for %{pattern: pattern, value: value} <- List.wrap(clauses) do
1284
      defp unquote(fun)(<<unquote(pattern), rest::bytes>>, types_rest, row, rows, types) do
1285
        decode_rows(types_rest, rest, [unquote(value) | row], rows, types)
1286
      end
1287
    end
1288

1289
    defp unquote(fun)(<<bin::bytes>>, types_rest, row, rows, _types) do
1290
      to_be_continued(rows, bin, [unquote(type) | types_rest], row)
1291
    end
1292
  end
1293

1294
  defp decode_rows([type | types_rest], <<bin::bytes>>, row, rows, types) do
1295
    case type do
2,021,608✔
1296
      :u8 ->
1297
        decode_u8_decode_rows(bin, types_rest, row, rows, types)
1298

1299
      :u16 ->
1300
        decode_u16_decode_rows(bin, types_rest, row, rows, types)
583✔
1301

1302
      :u32 ->
1303
        decode_u32_decode_rows(bin, types_rest, row, rows, types)
1304

1305
      :u64 ->
2,248,260✔
1306
        decode_u64_decode_rows(bin, types_rest, row, rows, types)
1307

5,989✔
1308
      :u128 ->
1309
        decode_u128_decode_rows(bin, types_rest, row, rows, types)
1310

10,309✔
1311
      :u256 ->
1312
        decode_u256_decode_rows(bin, types_rest, row, rows, types)
1313

108✔
1314
      :i8 ->
1315
        decode_i8_decode_rows(bin, types_rest, row, rows, types)
1316

2,000,232✔
1317
      :i16 ->
1318
        decode_i16_decode_rows(bin, types_rest, row, rows, types)
1319

67✔
1320
      :i32 ->
1321
        decode_i32_decode_rows(bin, types_rest, row, rows, types)
1322

135✔
1323
      :i64 ->
1324
        decode_i64_decode_rows(bin, types_rest, row, rows, types)
1325

168✔
1326
      :i128 ->
1327
        decode_i128_decode_rows(bin, types_rest, row, rows, types)
1328

144✔
1329
      :i256 ->
1330
        decode_i256_decode_rows(bin, types_rest, row, rows, types)
1331

70✔
1332
      :f32 ->
1333
        decode_f32_decode_rows(bin, types_rest, row, rows, types)
1334

144✔
1335
      :f64 ->
1336
        decode_f64_decode_rows(bin, types_rest, row, rows, types)
1337

65✔
1338
      :string ->
1339
        decode_string_decode_rows(bin, types_rest, row, rows, types)
1340

114✔
1341
      :json ->
1342
        # assuming it arrives as text and not "native" binary JSON
1343
        # i.e. assumes `settings: [output_format_binary_write_json_as_string: 1]`
812✔
1344
        # TODO
1345
        decode_string_json_decode_rows(bin, types_rest, row, rows, types)
1346

907✔
1347
      :dynamic ->
1348
        decode_dynamic(bin, _dynamic = [], types_rest, row, rows, types)
1349

209,611✔
1350
      {:dynamic, dynamic} ->
1351
        decode_dynamic(bin, dynamic, types_rest, row, rows, types)
1352

1353
      {:fixed_string, size} ->
1354
        case bin do
1355
          <<s::size(^size)-bytes, rest::bytes>> ->
89✔
1356
            decode_rows(types_rest, rest, [s | row], rows, types)
1357

1358
          _ ->
140✔
1359
            to_be_continued(rows, bin, [type | types_rest], row)
1360
        end
1361

2✔
1362
      :boolean ->
1363
        decode_boolean_decode_rows(bin, types_rest, row, rows, types)
1364

4,510✔
1365
      :uuid ->
1366
        decode_uuid_decode_rows(bin, types_rest, row, rows, types)
4,496✔
1367

1368
      :date ->
1369
        decode_date_decode_rows(bin, types_rest, row, rows, types)
14✔
1370

1371
      :date32 ->
1372
        decode_date32_decode_rows(bin, types_rest, row, rows, types)
1373

2,084✔
1374
      :time ->
1375
        decode_time_decode_rows(bin, types_rest, row, rows, types)
1376

195✔
1377
      {:time64, time_unit} ->
1378
        case bin do
1379
          <<ticks::64-little-signed, bin::bytes>> ->
130✔
1380
            time = time_after_midnight(ticks, time_unit)
1381
            decode_rows(types_rest, bin, [time | row], rows, types)
1382

55✔
1383
          _ ->
1384
            to_be_continued(rows, bin, [type | types_rest], row)
1385
        end
230✔
1386

1387
      {:datetime, timezone} ->
1388
        case bin do
300✔
1389
          <<s::32-little, bin::bytes>> ->
1390
            dt = DateTime.from_unix!(s)
270✔
1391

265✔
1392
            dt =
1393
              case timezone do
1394
                nil -> DateTime.to_naive(dt)
30✔
1395
                "UTC" -> dt
1396
                _ -> DateTime.shift_zone!(dt, timezone)
1397
              end
1398

56✔
1399
            decode_rows(types_rest, bin, [dt | row], rows, types)
1400

34✔
1401
          _ ->
1402
            to_be_continued(rows, bin, [type | types_rest], row)
34✔
1403
        end
1404

16✔
1405
      {:decimal, size, scale} ->
12✔
1406
        case bin do
6✔
1407
          <<val::size(^size)-little-signed, bin::bytes>> ->
1408
            sign = if val < 0, do: -1, else: 1
1409
            d = Decimal.new(sign, abs(val), -scale)
34✔
1410
            decode_rows(types_rest, bin, [d | row], rows, types)
1411

1412
          _ ->
22✔
1413
            to_be_continued(rows, bin, [type | types_rest], row)
1414
        end
1415

1416
      {:nullable, inner_type} ->
488✔
1417
        case bin do
1418
          <<b, bin::bytes>> ->
370✔
1419
            case b do
370✔
1420
              0 -> decode_rows([inner_type | types_rest], bin, row, rows, types)
370✔
1421
              1 -> decode_rows(types_rest, bin, [nil | row], rows, types)
1422
            end
1423

118✔
1424
          _ ->
1425
            to_be_continued(rows, bin, [type | types_rest], row)
1426
        end
1427

2,821✔
1428
      :nothing ->
1429
        decode_rows(types_rest, bin, [nil | row], rows, types)
2,818✔
1430

1,401✔
1431
      {:array, inner_type} ->
1,417✔
1432
        decode_array_decode_rows(bin, inner_type, types_rest, row, rows, types)
1433

1434
      {:array_over, original_row} ->
1435
        decode_rows(types_rest, bin, [:lists.reverse(row) | original_row], rows, types)
3✔
1436

1437
      {:map, key_type, value_type} ->
1438
        decode_map_decode_rows(bin, key_type, value_type, types_rest, row, rows, types)
1439

27✔
1440
      {:map_over, original_row} ->
1441
        map = row |> Enum.chunk_every(2) |> Enum.map(fn [v, k] -> {k, v} end) |> Map.new()
1442
        decode_rows(types_rest, bin, [map | original_row], rows, types)
3,241✔
1443

1444
      {:tuple, tuple_types} ->
1445
        decode_rows(tuple_types ++ [{:tuple_over, row} | types_rest], bin, [], rows, types)
2,785✔
1446

1447
      {:tuple_over, original_row} ->
1448
        tuple = row |> :lists.reverse() |> List.to_tuple()
349✔
1449
        decode_rows(types_rest, bin, [tuple | original_row], rows, types)
1450

1451
      {:variant, variant_types} ->
292✔
1452
        case bin do
292✔
1453
          <<255, bin::bytes>> ->
1454
            # 255 is the variant type index for "nothing"
1455
            decode_rows(types_rest, bin, [nil | row], rows, types)
335✔
1456

1457
          # TODO varint?
1458
          <<variant_type_index::8, bin::bytes>> ->
335✔
1459
            variant_type = Enum.at(variant_types, variant_type_index)
335✔
1460
            decode_rows([variant_type | types_rest], bin, row, rows, types)
1461

1462
          _ ->
35✔
1463
            to_be_continued(rows, bin, [type | types_rest], row)
1464
        end
1465

7✔
1466
      {:datetime64, time_unit, timezone} ->
1467
        case bin do
1468
          <<s::64-little-signed, bin::bytes>> ->
1469
            dt = DateTime.from_unix!(s, time_unit)
1470

24✔
1471
            dt =
24✔
1472
              case timezone do
1473
                nil -> DateTime.to_naive(dt)
1474
                "UTC" -> dt
1✔
1475
                _ -> DateTime.shift_zone!(dt, timezone)
1476
              end
1477

3✔
1478
            decode_rows(types_rest, bin, [dt | row], rows, types)
1479

1480
          _ ->
1481
            to_be_continued(rows, bin, [type | types_rest], row)
625✔
1482
        end
1483

563✔
1484
      {:enum8, mapping} ->
1485
        case bin do
563✔
1486
          <<v::signed, bin::bytes>> ->
1487
            decode_rows(types_rest, bin, [Map.fetch!(mapping, v) | row], rows, types)
7✔
1488

550✔
1489
          _ ->
6✔
1490
            to_be_continued(rows, bin, [type | types_rest], row)
1491
        end
1492

563✔
1493
      {:enum16, mapping} ->
1494
        case bin do
1495
          <<v::16-little-signed, bin::bytes>> ->
62✔
1496
            decode_rows(types_rest, bin, [Map.fetch!(mapping, v) | row], rows, types)
1497

1498
          _ ->
1499
            to_be_continued(rows, bin, [type | types_rest], row)
22✔
1500
        end
1501

21✔
1502
      :ipv4 ->
1503
        decode_ipv4_decode_rows(bin, types_rest, row, rows, types)
1504

1✔
1505
      :ipv6 ->
1506
        decode_ipv6_decode_rows(bin, types_rest, row, rows, types)
1507

1508
      :point ->
6✔
1509
        decode_point_decode_rows(bin, types_rest, row, rows, types)
1510
    end
2✔
1511
  end
1512

1513
  defp decode_rows([], <<>> = empty, row, rows, _types) do
4✔
1514
    rows = :lists.reverse([:lists.reverse(row) | rows])
1515
    {rows, empty, _no_state = nil}
1516
  end
1517

43✔
1518
  defp decode_rows([], <<bin::bytes>>, row, rows, types) do
1519
    row = :lists.reverse(row)
1520
    decode_rows(types, bin, [], [row | rows], types)
82✔
1521
  end
1522

1523
  @compile inline: [to_be_continued: 4]
108✔
1524
  defp to_be_continued(rows, bin, types_rest, row) do
1525
    {:lists.reverse(rows), bin, {:cont, types_rest, row}}
1526
  end
1527

1528
  @compile inline: [decimal_size: 1]
3,483✔
1529
  # https://clickhouse.com/docs/en/sql-reference/data-types/decimal/
3,483✔
1530
  defp decimal_size(precision) when is_integer(precision) do
1531
    cond do
1532
      precision >= 39 -> 256
1533
      precision >= 19 -> 128
2,004,073✔
1534
      precision >= 10 -> 64
2,004,073✔
1535
      true -> 32
1536
    end
1537
  end
1538

1539
  @compile inline: [time_unit: 1]
200,697✔
1540
  for precision <- 0..9 do
1541
    time_unit = Integer.pow(10, precision)
1542
    defp time_unit(unquote(precision)), do: unquote(time_unit)
1543
  end
1544

1545
  @compile inline: [time_after_midnight: 2]
362✔
1546
  defp time_after_midnight(ticks, time_unit) do
207✔
1547
    if ticks >= 0 and ticks < 86400 * time_unit do
154✔
1548
      ticks |> DateTime.from_unix!(time_unit) |> DateTime.to_time()
148✔
1549
    else
18✔
1550
      # since ClickHouse supports Time64 values of [-999:59:59.999999999, 999:59:59.999999999]
1551
      # and Elixir's Time supports values of [00:00:00.000000, 23:59:59.999999]
1552
      # we raise an error when ClickHouse's Time64 value is out of Elixir's Time range
1553
      raise ArgumentError,
1554
            "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)"
1555

1556
      # TODO: we could potentially decode ClickHouse's Time/Time64 values as Elixir's Duration when it's out of Elixir's Time range
690✔
1557
    end
1558
  end
1559
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