diff --git a/CHANGELOG.md b/CHANGELOG.md index 89d0ce2..df975cf 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -31,6 +31,11 @@ and no init/loop callback for response streaming anymore. - Response streaming and server-sent events no longer spawn a process per response, they run in the connection process itself. +- `send_chunk` and `send_event` no longer return a `Result`. Writing to a client + that has gone now ends the handler where it stands rather than letting it + carry on producing a body with nowhere to go. A stream that ends this way is + not an error and is not reported as one while a handler that crashes for its + own reasons still does. `on_close` runs either way. - Rename `SSEConnection`, `SSEEvent` and `SSENext` to `SseConnection`, `SseEvent` and `SseNext`. - Add `comment`, for the server-sent events comment that keeps an idle stream diff --git a/src/ewe.gleam b/src/ewe.gleam index 74c59aa..37c9ff6 100644 --- a/src/ewe.gleam +++ b/src/ewe.gleam @@ -582,6 +582,11 @@ pub type ResponseWriter = /// Starts a streamed response. `handler` must end by calling `finish_chunk` or /// `finish_response` on it, since that's what closes the stream. +/// +/// A send to a client that has gone ends the handler there and then, rather +/// than letting it carry on producing a body with nowhere to go. Nothing after +/// that send runs, so hold anything that needs releasing in a way that survives +/// on the process ending rather than in code after the write. pub fn stream_response( response: response.Response(a), handler: fn(ResponseWriter) -> Nil, @@ -676,11 +681,9 @@ pub fn event_retry(event: SseEvent, retry: Int) -> SseEvent { sse.Event(..event, retry: Some(retry)) } -/// Sends event to the client. -pub fn send_event( - conn: SseConnection, - event: SseEvent, -) -> Result(Nil, socket.SocketReason) { +/// Sends event to the client. If the client has gone the stream ends here: +/// `on_close` runs and the handler is not called again. +pub fn send_event(conn: SseConnection, event: SseEvent) -> Nil { case conn { connection.Http1Sse(conn) -> http1_sse.send(conn, event) connection.Http2Sse -> todo as "HTTP/2 is not implemented yet!" diff --git a/src/ewe/internal/http1.gleam b/src/ewe/internal/http1.gleam index 6585e40..dcd393e 100644 --- a/src/ewe/internal/http1.gleam +++ b/src/ewe/internal/http1.gleam @@ -4,6 +4,7 @@ import ewe/internal/http1/body import ewe/internal/http1/connection as http1 import ewe/internal/http1/encoder import ewe/internal/http1/parser +import ewe/internal/stream import gleam/bytes_tree import gleam/erlang/process import gleam/http @@ -199,24 +200,32 @@ fn send_response( encoder.RemainderStream(handler: stream_handler, framing:) -> { use Nil <- result.try(transport.send(transport, socket, head)) - connection.Http1Writer(http1.ResponseWriter( - transport:, - socket:, - self:, - framing:, - keep_alive:, - )) - |> stream_handler - - let drained = drain_messages(self) - case drained.stream { - option.Some(http1.StreamFinished(keep_alive:)) -> - Ok(to_sent(keep_alive)) - // A handler that returns without finishing left the body unterminated, - // so close it out here and drop a connection we can no longer reuse. - option.None -> { - let _ = encoder.end_stream(transport, socket, framing) - Ok(SentClose) + let writer = + connection.Http1Writer(http1.ResponseWriter( + transport:, + socket:, + self:, + framing:, + keep_alive:, + )) + + case stream.rescue_dead(fn() { stream_handler(writer) }) { + // The client went away mid stream. There is nothing to terminate the + // body with and nothing to reuse. + Error(_reason) -> Ok(SentClose) + Ok(Nil) -> { + let drained = drain_messages(self) + case drained.stream { + option.Some(http1.StreamFinished(keep_alive:)) -> + Ok(to_sent(keep_alive)) + // A handler that returns without finishing left the body + // unterminated, so close it out here and drop a connection we can + // no longer reuse. + option.None -> { + let _ = encoder.end_stream(transport, socket, framing) + Ok(SentClose) + } + } } } } diff --git a/src/ewe/internal/http1/encoder.gleam b/src/ewe/internal/http1/encoder.gleam index 8b74f54..3e66293 100644 --- a/src/ewe/internal/http1/encoder.gleam +++ b/src/ewe/internal/http1/encoder.gleam @@ -3,6 +3,7 @@ import ewe/internal/connection import ewe/internal/file import ewe/internal/http1/connection as http1 import ewe/internal/http1/parser +import ewe/internal/stream import gleam/bit_array import gleam/bytes_tree import gleam/erlang/process @@ -260,8 +261,9 @@ pub fn frame( } pub fn send_chunk(writer: ResponseWriter, chunk: BitArray) -> ResponseWriter { - let bytes = frame(bytes_tree.from_bit_array(chunk), writer.framing) - let _ = transport.send(writer.transport, writer.socket, bytes) + frame(bytes_tree.from_bit_array(chunk), writer.framing) + |> write(writer, _) + writer } @@ -275,13 +277,22 @@ pub fn finish_chunk(writer: ResponseWriter, chunk: BitArray) -> Nil { ) http1.CloseDelimitedStream -> bytes_tree.from_bit_array(chunk) } - let _ = transport.send(writer.transport, writer.socket, bytes) + write(writer, bytes) finish(writer) } pub fn finish_response(writer: ResponseWriter) -> Nil { - let _ = end_stream(writer.transport, writer.socket, writer.framing) - finish(writer) + case end_stream(writer.transport, writer.socket, writer.framing) { + Ok(Nil) -> finish(writer) + Error(reason) -> stream.dead(reason) + } +} + +fn write(writer: ResponseWriter, bytes: bytes_tree.BytesTree) -> Nil { + case transport.send(writer.transport, writer.socket, bytes) { + Ok(Nil) -> Nil + Error(reason) -> stream.dead(reason) + } } pub fn end_stream( diff --git a/src/ewe/internal/http1/parser.gleam b/src/ewe/internal/http1/parser.gleam index b2a3767..1fa8404 100644 --- a/src/ewe/internal/http1/parser.gleam +++ b/src/ewe/internal/http1/parser.gleam @@ -1,7 +1,6 @@ import ewe/internal/http1/connection as http1 import gleam/bit_array import gleam/http -import gleam/int import gleam/list import gleam/option diff --git a/src/ewe/internal/http1/sse.gleam b/src/ewe/internal/http1/sse.gleam index 5b9d964..0222349 100644 --- a/src/ewe/internal/http1/sse.gleam +++ b/src/ewe/internal/http1/sse.gleam @@ -2,6 +2,7 @@ import ewe/internal/connection import ewe/internal/http1/connection as http1 import ewe/internal/http1/encoder import ewe/internal/sse +import ewe/internal/stream import gleam/dynamic import gleam/erlang/atom import gleam/erlang/process @@ -25,14 +26,28 @@ pub fn run( case activate(conn) { Ok(Nil) -> loop(conn, handle, selector(subject), state, Clean, step, on_close) - Error(reason) -> { - on_close(handle, state) - finished(conn, http1.CloseAfterResponse) - connection.StoppedAbnormal(socket.reason_to_string(reason)) - } + Error(reason) -> + socket.reason_to_string(reason) + |> connection.StoppedAbnormal + |> ended(conn, handle, state, on_close, http1.CloseAfterResponse, _) } } +/// Every way a stream ends runs the handler's `on_close`, reports whether the +/// connection survived it, and answers with the outcome. +fn ended( + conn: http1.SseConnection, + handle: connection.SseConnection, + state: user_state, + on_close: fn(connection.SseConnection, user_state) -> Nil, + keep_alive: http1.KeepAlive, + outcome: connection.Outcome, +) -> connection.Outcome { + on_close(handle, state) + finished(conn, keep_alive) + outcome +} + /// Whether anything happened during the stream that rules out handing the /// connection back for another request. type Reuse { @@ -54,25 +69,42 @@ fn loop( // Whatever the client sent has been taken off the socket and cannot be put // back, so the connection is no longer safe to reuse. ClientData -> loop(conn, handle, selector, state, Spoiled, step, on_close) - Disconnected -> { - on_close(handle, state) - finished(conn, http1.CloseAfterResponse) - connection.Stopped - } - Failed(reason) -> { - on_close(handle, state) - finished(conn, http1.CloseAfterResponse) + Disconnected -> + ended( + conn, + handle, + state, + on_close, + http1.CloseAfterResponse, + connection.Stopped, + ) + Failed(reason) -> connection.StoppedAbnormal(reason) - } + |> ended(conn, handle, state, on_close, http1.CloseAfterResponse, _) Message(message) -> - case step(handle, state, message) { - sse.Proceed(state) -> + // A send inside the handler can find the client gone before the socket + // has told us, so both routes out land on the same teardown. + case stream.rescue_dead(fn() { step(handle, state, message) }) { + Error(_reason) -> + ended( + conn, + handle, + state, + on_close, + http1.CloseAfterResponse, + connection.Stopped, + ) + Ok(sse.Proceed(state)) -> loop(conn, handle, selector, state, reuse, step, on_close) - sse.Halt(outcome) -> { - on_close(handle, state) - finished(conn, keep_alive(outcome, reuse)) - outcome - } + Ok(sse.Halt(outcome)) -> + ended( + conn, + handle, + state, + on_close, + keep_alive(outcome, reuse), + outcome, + ) } } } @@ -93,13 +125,13 @@ fn finished(conn: http1.SseConnection, keep_alive: http1.KeepAlive) -> Nil { |> process.send(conn.self, _) } -pub fn send( - conn: http1.SseConnection, - event: sse.Event, -) -> Result(Nil, socket.SocketReason) { - sse.encode(event) - |> encoder.frame(conn.framing) - |> transport.send(conn.transport, conn.socket, _) +pub fn send(conn: http1.SseConnection, event: sse.Event) -> Nil { + let bytes = sse.encode(event) |> encoder.frame(conn.framing) + + case transport.send(conn.transport, conn.socket, bytes) { + Ok(Nil) -> Nil + Error(reason) -> stream.dead(reason) + } } /// glisten re arms `{active, once}` only once its loop callback returns, and an diff --git a/src/ewe/internal/stream.gleam b/src/ewe/internal/stream.gleam new file mode 100644 index 0000000..fca9ca9 --- /dev/null +++ b/src/ewe/internal/stream.gleam @@ -0,0 +1,16 @@ +import glisten/socket + +/// Ends the handler writing this stream, because there is no longer anywhere +/// for it to write to. Never returns. +/// +/// A handler is straight line code that ewe has handed control to, so unwinding +/// is the only way to stop it doing work for a client that has gone. It also +/// keeps the one thing a handler could do about a failed write out of its way, +/// since stopping is the only sane answer. +@external(erlang, "stream_ffi", "dead") +pub fn dead(reason: socket.SocketReason) -> a + +/// Runs `handler`, catching only the end that `dead` raises. Anything else a +/// handler raises is a bug and is left to crash. +@external(erlang, "stream_ffi", "rescue_dead") +pub fn rescue_dead(handler: fn() -> a) -> Result(a, socket.SocketReason) diff --git a/src/ewe/internal/stream_ffi.erl b/src/ewe/internal/stream_ffi.erl new file mode 100644 index 0000000..0794ce8 --- /dev/null +++ b/src/ewe/internal/stream_ffi.erl @@ -0,0 +1,18 @@ +-module(stream_ffi). + +-export([dead/1, rescue_dead/1]). + +-define(STREAM_DEAD, ewe_stream_dead). + +dead(Reason) -> + erlang:error({?STREAM_DEAD, Reason}). + +%% Matches the sentinel exactly and reraises everything else with its own +%% stacktrace so a bug in a handler still crashes as a bug. +rescue_dead(Func) -> + try + {ok, Func()} + catch + error:{?STREAM_DEAD, Reason} -> {error, Reason}; + Class:Reason:Stacktrace -> erlang:raise(Class, Reason, Stacktrace) + end. diff --git a/test/ewe/internal/stream_test.gleam b/test/ewe/internal/stream_test.gleam new file mode 100644 index 0000000..25fdef3 --- /dev/null +++ b/test/ewe/internal/stream_test.gleam @@ -0,0 +1,41 @@ +import ewe/internal/stream +import gleam/erlang/process +import glisten/socket + +pub fn handler_running_to_completion_is_returned_test() { + assert stream.rescue_dead(fn() { 42 }) == Ok(42) +} + +pub fn dead_stream_unwinds_to_the_rescue_test() { + let handler = fn() { + stream.dead(socket.Closed) + panic as "a dead stream must not return to its handler" + } + + assert stream.rescue_dead(handler) == Error(socket.Closed) +} + +pub fn work_after_a_dead_stream_does_not_run_test() { + let subject = process.new_subject() + + let handler = fn() { + stream.dead(socket.Closed) + process.send(subject, "kept working") + } + + let _dead = stream.rescue_dead(handler) + + assert process.receive(subject, 0) == Error(Nil) + as "unwinding is what stops a handler producing a body nobody will read" +} + +pub fn a_handler_bug_is_not_mistaken_for_a_dead_stream_test() { + let crashed = + rescue(fn() { stream.rescue_dead(fn() { panic as "bug in the handler" }) }) + + assert crashed == Error(Nil) + as "swallowing a bug would report it as a client that went away" +} + +@external(erlang, "stream_test_ffi", "rescue") +fn rescue(handler: fn() -> a) -> Result(a, Nil) diff --git a/test/ewe/internal/stream_test_ffi.erl b/test/ewe/internal/stream_test_ffi.erl new file mode 100644 index 0000000..fe3e322 --- /dev/null +++ b/test/ewe/internal/stream_test_ffi.erl @@ -0,0 +1,11 @@ +-module(stream_test_ffi). + +-export([rescue/1]). + +%% Catches every class, so a test can assert that something crashed at all. +rescue(Func) -> + try + {ok, Func()} + catch + _Class:_Reason -> {error, nil} + end.