Current section

51 Versions

Jump to

Compare versions

4 files changed
+44 additions
-25 deletions
  @@ -29,4 +29,4 @@
29 29 {<<"name">>,<<"hpack">>},
30 30 {<<"optional">>,false},
31 31 {<<"requirement">>,<<"~> 1.0.2">>}]]}.
32 - {<<"version">>,<<"0.0.3">>}.
32 + {<<"version">>,<<"0.0.4">>}.
  @@ -1,33 +1,35 @@
1 1 defmodule Ankh.Connection do
2 2 use GenServer
3 3
4 - alias Ankh.Frame
4 + alias Ankh.{Frame, Stream}
5 5 alias Ankh.Frame.{Encoder, Error, GoAway, Ping, Settings}
6 6
7 7 require Logger
8 8
9 - @ssl_opts binary: true, versions: [:"tlsv1.2"], secure_renegotiate: true,
10 - alpn_advertised_protocols: ["h2"], client_renegotiation: false,
9 + @default_ssl_opts binary: true, versions: [:"tlsv1.2"],
10 + secure_renegotiate: true, alpn_advertised_protocols: ["h2"],
11 + client_renegotiation: false,
11 12 ciphers: ["ECDHE-ECDSA-AES128-SHA256", "ECDHE-ECDSA-AES128-SHA"]
12 13 @preface "PRI * HTTP/2.0\r\n\r\nSM\r\n\r\n"
13 14 @frame_header_size 9
14 15 @max_stream_id 2_147_483_648
15 16
16 - def start_link(%URI{} = uri, stream \\ false, receiver \\ nil, options \\ [])
17 - do
17 + def start_link([uri: uri, receiver: receiver, stream: stream,
18 + ssl_options: ssl_options], options \\ []) do
18 19 target = if is_pid(receiver), do: receiver, else: self()
19 20 mode = if is_boolean(stream) && stream, do: :stream, else: :full
20 - GenServer.start_link(__MODULE__, [uri: uri, target: target, mode: mode],
21 - options)
21 + GenServer.start_link(__MODULE__, [uri: uri, target: target, mode: mode,
22 + ssl_options: ssl_options], options)
22 23 end
23 24
24 - def init([uri: uri, target: target, mode: mode]) do
25 + def init([uri: uri, target: target, mode: mode, ssl_options: ssl_opts]) do
25 26 with settings <- %Settings.Payload{},
26 27 {:ok, recv_ctx} = HPack.Table.start_link(settings.header_table_size),
27 28 {:ok, send_ctx} = HPack.Table.start_link(settings.header_table_size) do
28 - {:ok, %{uri: uri, target: target, mode: mode, socket: nil, streams: %{},
29 - last_stream_id: 0, buffer: <<>>, recv_ctx: recv_ctx, send_ctx: send_ctx,
30 - recv_settings: settings, send_settings: nil, window_size: 0}}
29 + {:ok, %{uri: uri, target: target, mode: mode, ssl_opts: ssl_opts,
30 + socket: nil, streams: %{}, last_stream_id: 0, buffer: <<>>,
31 + recv_ctx: recv_ctx, send_ctx: send_ctx, recv_settings: settings,
32 + send_settings: nil, window_size: 65_535}}
31 33 end
32 34 end
33 35
  @@ -38,9 +40,9 @@ defmodule Ankh.Connection do
38 40 def close(pid), do: GenServer.call(pid, {:close})
39 41
40 42 def handle_call({:send, frame}, from, %{socket: nil, uri: %URI{host: host,
41 - port: port}} = state) do
43 + port: port}, ssl_opts: ssl_opts} = state) do
42 44 hostname = String.to_charlist(host)
43 - ssl_options = @ssl_opts
45 + ssl_options = Keyword.merge(ssl_opts, @default_ssl_opts)
44 46 with {:ok, socket} <- :ssl.connect(hostname, port, ssl_options),
45 47 :ok <- :ssl.send(socket, @preface) do
46 48 state = %{state | socket: socket}
  @@ -97,7 +99,7 @@ defmodule Ankh.Connection do
97 99
98 100 defp send_frame(%{socket: socket, streams: streams, send_ctx: hpack}
99 101 = state, %Frame{stream_id: id} = frame) do
100 - stream = Map.get(streams, id, Ankh.Stream.new(id))
102 + stream = Map.get(streams, id, Stream.new(id))
101 103 Logger.debug "STREAM #{id} SEND #{inspect frame}"
102 104 {:ok, stream} = Ankh.Stream.send_frame(stream, frame)
103 105 Logger.debug "STREAM #{id} IS #{inspect stream}"
  @@ -251,11 +253,14 @@ defmodule Ankh.Connection do
251 253 end
252 254
253 255 defp decode_frame(%Frame{stream_id: id, type: :push_promise,
254 - flags: %{end_headers: false}, payload: %{header_block_fragment: hbf}} = frame,
256 + flags: %{end_headers: false}, payload: %{promised_stream_id: promised_id,
257 + header_block_fragment: hbf}} = frame,
255 258 %{streams: streams} = state) do
256 259 Logger.debug("STREAM #{id} RECEIVED PARTIAL HBF #{inspect hbf}")
257 260 stream = Map.get(streams, id)
258 - stream = %{stream | hbf: stream.hbf <> hbf}
261 +
262 + stream = %{stream | promised_id: promised_id, hbf_type: :push_promise,
263 + hbf: stream.hbf <> hbf}
259 264 {%{state | streams: Map.put(streams, id, stream)}, frame}
260 265 end
261 266
  @@ -279,13 +284,17 @@ defmodule Ankh.Connection do
279 284 end
280 285
281 286 defp decode_frame(%Frame{stream_id: id, type: :push_promise,
282 - flags: %{end_headers: true}, payload: %{header_block_fragment: hbf}} = frame,
287 + flags: %{end_headers: true}, payload: %{promised_stream_id: promised_id,
288 + header_block_fragment: hbf}} = frame,
283 289 %{streams: streams, recv_ctx: table, target: target} = state) do
284 290 stream = Map.get(streams, id)
285 291 headers = HPack.decode(stream.hbf <> hbf, table)
286 292 Logger.debug("STREAM #{id} RECEIVED HEADERS #{inspect headers}")
287 - Process.send(target, {:ankh, :headers, id, headers}, [])
288 - {%{state | streams: Map.put(streams, id, %{stream | hbf: <<>>})}, frame}
293 + Process.send(target, {:ankh, :push_promise, id, headers}, [])
294 + streams = streams
295 + |> Map.put(id, %{stream | hbf: <<>>})
296 + |> Map.put(promised_id, %{Stream.new(promised_id) | state: :reserved})
297 + {%{state | streams: streams}, frame}
289 298 end
290 299
291 300 defp decode_frame(%Frame{stream_id: id, type: :continuation,
  @@ -294,8 +303,17 @@ defmodule Ankh.Connection do
294 303 stream = Map.get(streams, id)
295 304 headers = HPack.decode(stream.hbf <> hbf, table)
296 305 Logger.debug("STREAM #{id} RECEIVED HEADERS #{inspect headers}")
297 - Process.send(target, {:ankh, :headers, id, headers}, [])
298 - {%{state | streams: Map.put(streams, id, %{stream | hbf: <<>>})}, frame}
306 + Process.send(target, {:ankh, stream.hbf_type, id, headers}, [])
307 + streams = case stream.hbf_type do
308 + :push_promise ->
309 + %{promised_id: promised_id} = stream
310 + streams
311 + |> Map.put(id, %{stream | hbf: <<>>})
312 + |> Map.put(promised_id, %{Stream.new(promised_id) | state: :reserved})
313 + _ ->
314 + Map.put(streams, id, %{stream | hbf: <<>>})
315 + end
316 + {%{state | streams: streams}, frame}
299 317 end
300 318
301 319 defp decode_frame(%Frame{stream_id: id, type: :data,
  @@ -317,7 +335,7 @@ defmodule Ankh.Connection do
317 335 %{streams: streams, target: target, mode: mode} = state) do
318 336 stream = Map.get(streams, id)
319 337 data = stream.data <> data
320 - Logger.debug("STREAM #{id} RECEIVED DATA #{data} SIZE #{byte_size data}")
338 + Logger.debug("STREAM #{id} RECEIVED #{byte_size data} BYTES DATA:\n#{data}")
321 339 case {mode, stream.data} do
322 340 {:full, _} ->
323 341 Process.send(target, {:ankh, :data, id, data}, [])
  @@ -1,7 +1,8 @@
1 1 defmodule Ankh.Stream do
2 2 alias Ankh.Frame
3 3
4 - defstruct [id: 0, state: :idle, hbf: <<>>, data: <<>>, window_size: 0]
4 + defstruct [id: 0, state: :idle, hbf_type: :headers, hbf: <<>>, data: <<>>,
5 + window_size: 65_535]
5 6
6 7 def new(id), do: %__MODULE__{id: id}
  @@ -4,7 +4,7 @@ defmodule Ankh.Mixfile do
4 4 def project do
5 5 [
6 6 app: :ankh,
7 - version: "0.0.3",
7 + version: "0.0.4",
8 8 elixir: "~> 1.3",
9 9 build_embedded: Mix.env == :prod,
10 10 start_permanent: Mix.env == :prod,