package kafka-eio

  1. Overview
  2. Docs
Eio-native Kafka client for OCaml 5, built on librdkafka

Install

dune-project
 Dependency

Authors

Maintainers

Sources

v0.1.1.tar.gz
md5=88c2b9c6d263a1eb114cc5d1b570288a
sha512=44f5395fcc58b8d4303bc0b36113a17193da3a1925125cb0742fabeb96f32eb6b243c98d2fde283b6c95872ff36ddbf96de6e98b7612799e2aefeb2f85a1db95

doc/src/kafka-eio.core/kafka_security.ml.html

Source file kafka_security.ml

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
type protocol =
  [ `Plaintext
  | `Ssl
  | `Sasl_plaintext
  | `Sasl_ssl
  ]

type sasl_mechanism =
  | Plain
  | Scram_sha256
  | Scram_sha512

type sasl = {
  mechanism : sasl_mechanism;
  username  : string;
  password  : string;
}

let mechanism_to_string = function
  | Plain         -> "PLAIN"
  | Scram_sha256  -> "SCRAM-SHA-256"
  | Scram_sha512  -> "SCRAM-SHA-512"

let mechanism_of_string value =
  match String.uppercase_ascii value with
  | "PLAIN"          -> Ok Plain
  | "SCRAM-SHA-256"  -> Ok Scram_sha256
  | "SCRAM-SHA-512"  -> Ok Scram_sha512
  | other ->
    Error
      (Printf.sprintf
         "kafka security: unknown KAFKA_SASL_MECHANISM %S \
          (expected PLAIN, SCRAM-SHA-256, or SCRAM-SHA-512)"
         other)

type t =
  | Plaintext
  | Ssl of { ssl_ca_location : string option }
  | Sasl_plaintext of sasl
  | Sasl_ssl of { ssl_ca_location : string option; sasl : sasl }

let default = Plaintext

let protocol_to_string = function
  | `Plaintext      -> "plaintext"
  | `Ssl            -> "ssl"
  | `Sasl_plaintext -> "sasl_plaintext"
  | `Sasl_ssl       -> "sasl_ssl"

let protocol_of_string value =
  match String.lowercase_ascii value with
  | "plaintext"      -> Ok `Plaintext
  | "ssl"            -> Ok `Ssl
  | "sasl_plaintext" -> Ok `Sasl_plaintext
  | "sasl_ssl"       -> Ok `Sasl_ssl
  | other ->
    Error
      (Printf.sprintf
         "kafka security: unknown KAFKA_SECURITY_PROTOCOL %S \
          (expected plaintext, ssl, sasl_plaintext, or sasl_ssl)"
         other)

let env_opt name =
  match Sys.getenv_opt name with Some v when v <> "" -> Some v | _ -> None

let required_env name =
  match env_opt name with
  | Some value -> Ok value
  | None       -> Error ("kafka security: " ^ name ^ " required for SASL protocols")

let sasl_of_env () =
  let ( let* ) = Result.bind in
  let* mechanism_raw = required_env "KAFKA_SASL_MECHANISM" in
  let* mechanism     = mechanism_of_string mechanism_raw in
  let* username      = required_env "KAFKA_SASL_USERNAME" in
  let* password      = required_env "KAFKA_SASL_PASSWORD" in
  Ok { mechanism; username; password }

let of_env () =
  let ( let* ) = Result.bind in
  let* protocol =
    match env_opt "KAFKA_SECURITY_PROTOCOL" with
    | None       -> Ok `Plaintext
    | Some value -> protocol_of_string value
  in
  let ssl_ca_location = env_opt "KAFKA_SSL_CA_LOCATION" in
  match protocol with
  | `Plaintext      -> Ok Plaintext
  | `Ssl            -> Ok (Ssl { ssl_ca_location })
  | `Sasl_plaintext ->
    let* sasl = sasl_of_env () in
    Ok (Sasl_plaintext sasl)
  | `Sasl_ssl ->
    let* sasl = sasl_of_env () in
    Ok (Sasl_ssl { ssl_ca_location; sasl })

let apply conf t =
  let errs = ref [] in
  let set k v = match Kafka_raw.conf_set conf k v with
    | Ok ()   -> ()
    | Error s -> errs := ("kafka security conf " ^ k ^ ": " ^ s) :: !errs
  in
  let set_sasl sasl =
    set "sasl.mechanism" (mechanism_to_string sasl.mechanism);
    set "sasl.username" sasl.username;
    set "sasl.password" sasl.password
  in
  (match t with
   | Plaintext ->
     set "security.protocol" (protocol_to_string `Plaintext)
   | Ssl { ssl_ca_location } ->
     set "security.protocol" (protocol_to_string `Ssl);
     Option.iter (set "ssl.ca.location") ssl_ca_location
   | Sasl_plaintext sasl ->
     set "security.protocol" (protocol_to_string `Sasl_plaintext);
     set_sasl sasl
   | Sasl_ssl { ssl_ca_location; sasl } ->
     set "security.protocol" (protocol_to_string `Sasl_ssl);
     Option.iter (set "ssl.ca.location") ssl_ca_location;
     set_sasl sasl);
  match !errs with
  | []   -> Ok ()
  | errs -> Error (String.concat "; " (List.rev errs))