diff --git a/benchmark/ewe@5/src/app.gleam b/benchmark/ewe@5/src/app.gleam index ba05aff..fece45b 100644 --- a/benchmark/ewe@5/src/app.gleam +++ b/benchmark/ewe@5/src/app.gleam @@ -48,6 +48,7 @@ fn handle_request( // head -c 100K /dev/urandom > file_100kb.bin let assert Ok(file) = ewe.file( + request.body, "../priv/file_100kb.bin", offset: option.None, limit: option.None, @@ -63,6 +64,7 @@ fn handle_request( // head -c 1G /dev/urandom > file_1gb.bin let assert Ok(file) = ewe.file( + request.body, "../priv/file_1gb.bin", offset: option.None, limit: option.None, diff --git a/dev/serve.gleam b/dev/serve.gleam index f294f47..97b8a7c 100644 --- a/dev/serve.gleam +++ b/dev/serve.gleam @@ -45,6 +45,7 @@ fn handle_request( // head -c 100K /dev/urandom > file_100kb.bin let assert Ok(file) = ewe.file( + request.body, "./dev/priv/file_100kb.bin", offset: option.None, limit: option.None, @@ -60,6 +61,7 @@ fn handle_request( // head -c 1G /dev/urandom > file_1gb.bin let assert Ok(file) = ewe.file( + request.body, "./dev/priv/file_1gb.bin", offset: option.None, limit: option.None, diff --git a/src/ewe.gleam b/src/ewe.gleam index bba8dd8..01f04ba 100644 --- a/src/ewe.gleam +++ b/src/ewe.gleam @@ -392,15 +392,20 @@ fn from_internal_file_error(error: file.FileError) -> FileError { } } -/// Prepares a file to be streamed as a response body. `offset` and `limit` in -/// bytes let you serve a byte range from the file. leave either as `None` to +/// Prepares a file to be streamed as a response body. `offset` and `limit` in +/// bytes let you serve a byte range from the file. leave either as `None` to /// serve from the start or through the end. +/// +/// On HTTP/1 this opens the file, and the returned body holds it open until the +/// response is written. Put it on a response you go on to return. A body that +/// is built and then discarded keeps its file open until it is collected. pub fn file( + connection: Connection, path: String, offset offset: Option(Int), limit limit: Option(Int), ) -> Result(Body, FileError) { - case file.resolve(path, offset, limit) { + case file.resolve(connection, path, offset, limit) { Ok(file) -> Ok(File(file)) Error(error) -> Error(from_internal_file_error(error)) } diff --git a/src/ewe/internal/connection.gleam b/src/ewe/internal/connection.gleam index 16fd5a8..92c3396 100644 --- a/src/ewe/internal/connection.gleam +++ b/src/ewe/internal/connection.gleam @@ -27,10 +27,18 @@ pub type Streaming { StreamingMetadata(handler: fn(ResponseWriter) -> Nil) } +/// A raw descriptor belongs to the process that opened it, so whether the +/// handler can carry one depends on the protocol: an HTTP/1 handler runs in the +/// process that writes the socket, an HTTP/2 stream handler does not. pub type File { - FileMetadata(path: String, offset: Int, length: Int) + /// Already open, and closed by whoever writes or drops the response. + OpenFile(handle: FileDescriptor, offset: Int, length: Int) + /// Sized but not yet open; the connection process opens it at send time. + PendingFile(path: String, offset: Int, length: Int) } +pub type FileDescriptor + pub type ResponseWriter { Http1Writer(http1.ResponseWriter) Http2Writer diff --git a/src/ewe/internal/file.gleam b/src/ewe/internal/file.gleam index 5669ba2..34a5430 100644 --- a/src/ewe/internal/file.gleam +++ b/src/ewe/internal/file.gleam @@ -16,11 +16,40 @@ pub type FileError { } pub fn resolve( + conn: connection.Connection, path: String, offset: option.Option(Int), limit: option.Option(Int), ) -> Result(connection.File, FileError) { - use size <- result.try(stat(path)) + case conn { + // The handler already runs in the process that writes the socket, so it can + // hold the descriptor itself. + connection.Http1(_conn) -> { + use handle <- result.try(open(path)) + + case size(handle) |> result.try(range(_, offset, limit)) { + Ok(#(offset, length)) -> + Ok(connection.OpenFile(handle:, offset:, length:)) + Error(error) -> { + close(handle) + Error(error) + } + } + } + connection.Http2 -> { + use size <- result.try(stat(path)) + use #(offset, length) <- result.map(range(size, offset, limit)) + + connection.PendingFile(path:, offset:, length:) + } + } +} + +fn range( + size: Int, + offset: option.Option(Int), + limit: option.Option(Int), +) -> Result(#(Int, Int), FileError) { let offset = option.unwrap(offset, 0) case offset >= 0 && offset <= size { @@ -29,37 +58,64 @@ pub fn resolve( let available = size - offset case limit { - option.None -> - Ok(connection.FileMetadata(path:, offset:, length: available)) + option.None -> Ok(#(offset, available)) option.Some(limit) if limit < 0 -> Error(InvalidLimit) - option.Some(limit) -> - Ok(connection.FileMetadata( - path:, - offset:, - length: int.min(limit, available), - )) + option.Some(limit) -> Ok(#(offset, int.min(limit, available))) } } } } +/// Hands back a descriptor. +pub fn release(file: connection.File) -> Nil { + case file { + connection.OpenFile(handle:, ..) -> close(handle) + connection.PendingFile(..) -> Nil + } +} + +pub fn release_body(body: connection.Body) -> Nil { + case body { + connection.File(file) -> release(file) + connection.Bytes(..) + | connection.Text(..) + | connection.Empty + | connection.Streaming(..) + | connection.Sse(..) -> Nil + } +} + pub fn send( transport: transport.Transport, socket: socket.Socket, file: connection.File, ) -> Result(Nil, socket.SocketReason) { - case file.length { + case file { + connection.OpenFile(handle:, offset:, length:) -> + send_handle(transport, socket, handle, offset, length) + connection.PendingFile(..) -> todo as "HTTP/2 is not implemented yet!" + } +} + +/// Owns the descriptor from here on, so it is closed however the write ends. +fn send_handle( + transport: transport.Transport, + socket: socket.Socket, + handle: connection.FileDescriptor, + offset: Int, + length: Int, +) -> Result(Nil, socket.SocketReason) { + let sent = case length { 0 -> Ok(Nil) - length -> { - use fd <- result.try(open(file.path)) - let result = case transport { - transport.Tcp -> do_sendfile(fd, socket, file.offset, length) - transport.Ssl -> send_chunks(transport, socket, fd, file.offset, length) + _length -> + case transport { + transport.Tcp -> do_sendfile(handle, socket, offset, length) + transport.Ssl -> send_chunks(transport, socket, handle, offset, length) } - close(fd) - result - } } + + close(handle) + sent } const chunk_size = 65_536 @@ -67,7 +123,7 @@ const chunk_size = 65_536 fn send_chunks( transport: transport.Transport, socket: socket.Socket, - fd: FileDescriptor, + handle: connection.FileDescriptor, offset: Int, remaining: Int, ) -> Result(Nil, socket.SocketReason) { @@ -75,40 +131,47 @@ fn send_chunks( 0 -> Ok(Nil) _remaining -> { let amount = int.min(remaining, chunk_size) - use data <- result.try(pread(fd, offset, amount)) + use data <- result.try(pread(handle, offset, amount)) use Nil <- result.try(transport.send( transport, socket, bytes_tree.from_bit_array(data), )) - send_chunks(transport, socket, fd, offset + amount, remaining - amount) + send_chunks( + transport, + socket, + handle, + offset + amount, + remaining - amount, + ) } } } -pub type FileDescriptor - @external(erlang, "file_ffi", "stat") fn stat(path: String) -> Result(Int, FileError) @external(erlang, "file_ffi", "sendfile") fn do_sendfile( - fd: FileDescriptor, + handle: connection.FileDescriptor, socket: socket.Socket, offset: Int, bytes: Int, ) -> Result(Nil, socket.SocketReason) @external(erlang, "file_ffi", "open") -fn open(path: String) -> Result(FileDescriptor, socket.SocketReason) +fn open(path: String) -> Result(connection.FileDescriptor, FileError) + +@external(erlang, "file_ffi", "size") +fn size(handle: connection.FileDescriptor) -> Result(Int, FileError) @external(erlang, "file_ffi", "pread") fn pread( - fd: FileDescriptor, + handle: connection.FileDescriptor, offset: Int, length: Int, ) -> Result(BitArray, socket.SocketReason) @external(erlang, "file_ffi", "close") -fn close(fd: FileDescriptor) -> Nil +fn close(handle: connection.FileDescriptor) -> Nil diff --git a/src/ewe/internal/file_ffi.erl b/src/ewe/internal/file_ffi.erl index 0e510f6..6ab9546 100644 --- a/src/ewe/internal/file_ffi.erl +++ b/src/ewe/internal/file_ffi.erl @@ -2,7 +2,7 @@ -include_lib("kernel/include/file.hrl"). --export([stat/1, sendfile/4, open/1, pread/3, close/1]). +-export([stat/1, sendfile/4, open/1, size/1, pread/3, close/1]). stat(Path) -> case file:read_file_info(Path, [raw, {time, posix}]) of @@ -20,7 +20,21 @@ sendfile(Fd, Socket, Offset, Bytes) -> end. open(Path) -> - file:open(Path, [raw, binary, read]). + case file:open(Path, [raw, binary, read]) of + {ok, Fd} -> {ok, Fd}; + {error, enoent} -> {error, not_found}; + {error, eisdir} -> {error, is_directory}; + {error, eacces} -> {error, access_denied}; + {error, _Reason} -> {error, unknown_error} + end. + +%% Sizes the open handle, so the bytes framed are the ones about to be sent +%% rather than whatever the path pointed at a moment ago. +size(Fd) -> + case file:position(Fd, eof) of + {ok, Size} -> {ok, Size}; + {error, _Reason} -> {error, unknown_error} + end. pread(Fd, Offset, Length) -> case file:pread(Fd, Offset, Length) of diff --git a/src/ewe/internal/http1.gleam b/src/ewe/internal/http1.gleam index 9c3138e..5cc3471 100644 --- a/src/ewe/internal/http1.gleam +++ b/src/ewe/internal/http1.gleam @@ -77,6 +77,7 @@ pub fn handle_message( logging.Error, "Handler produced an unsafe response header: " <> name, ) + file.release_body(response.body) transport.send( connection.transport, @@ -170,11 +171,19 @@ fn send_response( )) Ok(to_sent(keep_alive)) } - encoder.RemainderFile(data) -> { - use Nil <- result.try(transport.send(transport, socket, head)) - use Nil <- result.try(file.send(transport, socket, data)) - Ok(to_sent(keep_alive)) - } + // `file.send` owns the descriptor once it is reached, so only a head that + // never made it to the socket leaves one to hand back. + encoder.RemainderFile(data) -> + case transport.send(transport, socket, head) { + Error(reason) -> { + file.release(data) + Error(reason) + } + Ok(Nil) -> { + use Nil <- result.try(file.send(transport, socket, data)) + Ok(to_sent(keep_alive)) + } + } encoder.RemainderStream(handler: stream_handler, framing:) -> { use Nil <- result.try(transport.send(transport, socket, head)) diff --git a/src/ewe/internal/http1/encoder.gleam b/src/ewe/internal/http1/encoder.gleam index 4d5d243..94f382c 100644 --- a/src/ewe/internal/http1/encoder.gleam +++ b/src/ewe/internal/http1/encoder.gleam @@ -1,5 +1,6 @@ import ewe/internal/clock import ewe/internal/connection +import ewe/internal/file import ewe/internal/http1/connection as http1 import ewe/internal/http1/parser import gleam/bit_array @@ -86,11 +87,24 @@ pub fn encode_response( // A HEAD response keeps the framing headers it would have had, minus the body. Ok(case method { - http.Head -> Encoded(..encoded, remainder: NoRemainder) + http.Head -> drop_body(encoded) _method -> encoded }) } +/// Anything the body was holding is let go here. +fn drop_body(encoded: Encoded) -> Encoded { + case encoded.remainder { + RemainderFile(data) -> file.release(data) + NoRemainder + | RemainderInline(..) + | RemainderStream(..) + | RemainderSse(..) -> Nil + } + + Encoded(..encoded, remainder: NoRemainder) +} + fn sized( state: EncodeState, status: Int,