~/bend-docscommunity

transport.bend source

transport.bend on the hub · documented module

import Baseimport ./json.bend as Jtype Connection is Type:  Connection{socket: Socket, buffered: String}type Read is Type:  Frame{connection: Connection, value: J.Json}  Broken{message: String}type Cut is Data:  Complete{line: String, rest: String}  Partial{}def cut(s: String, acc: String) -> Cut:  match s:    case SNil{}: Partial{}    case SCon{'\n', rest}: Complete{String.reverse(acc), rest}    case SCon{c, rest}: cut(rest, SCon{c, acc})def broken(socket: Socket, message: String) -> IO(Read):  do IO<Read>:    Socket.close(socket)    return Broken{message}def decoded(result: Result<&1, &1, String, J.Json>, socket: Socket, rest: String) -> IO(Read):  match result:    case Fail{error}: broken(socket, "Invalid bridge JSON: " ++ error)    case Done{value}: IO.pure(Read, Frame{Connection{socket, rest}, value})def chunk_nonempty(empty: Bool, socket: Socket, chunk: String, buffered: String,                   next: Connection -> IO(Read)) -> IO(Read):  match empty:    case True{}: broken(socket, "Bridge closed before a complete frame")    case False{}: next(Connection{socket, buffered ++ chunk})def chunk(result: Socket & Result<&1, &1, U32 & String, String>, buffered: String,          next: Connection -> IO(Read)) -> IO(Read):  (socket, received) = result  match received:    case Fail{(code, message)}: broken(socket, "Bridge receive failed: " ++ message)    case Done{+data}: chunk_nonempty(String.is_empty(data), socket, data, buffered, next)def read_cut(part: Cut, socket: Socket, buffered: String,             next: Connection -> IO(Read)) -> IO(Read):  match part:    case Complete{line, rest}: decoded(J.read(line), socket, rest)    case Partial{}:      do IO<Read>:        received : Socket & Result<&1, &1, U32 & String, String> <- TCP.recv(socket, 4096)        chunk(received, buffered, next)def read_limit(large: Bool, socket: Socket, +buffered: String,               next: Connection -> IO(Read)) -> IO(Read):  match large:    case True{}: broken(socket, "Bridge frame exceeds 4 MiB limit")    case False{}: read_cut(cut(buffered, ""), socket, buffered, next)# Socket progress depends on the external peer, not structural recursion.@unsafedef read(connection: Connection) -> IO(Read):  match connection:    case Connection{socket, +buffered}:      read_limit((U32.from_nat(String.length(buffered)) > 4198400 : U32), socket, buffered, c => read(c))def sent(result: Socket & Result<&1, &1, U32 & String, Unit>) -> IO(Result<&1, &1, String, Connection>):  (socket, result) = result  match result:    case Fail{(code, message)}:      do IO<Result<&1, &1, String, Connection>>:        Socket.close(socket)        return Fail{message}    case Done{unit}: IO.pure(Result<&1, &1, String, Connection>, Done{Connection{socket, ""}})def connected(result: Result<&1, &1, U32 & String, Socket>, body: String) -> IO(Result<&1, &1, String, Connection>):  match result:    case Fail{(code, message)}: IO.pure(Result<&1, &1, String, Connection>, Fail{message})    case Done{socket}:      do IO<Result<&1, &1, String, Connection>>:        result : Socket & Result<&1, &1, U32 & String, Unit> <- TCP.send(socket, body ++ "\n")        sent(result)def connect(port: U32, body: String) -> IO(Result<&1, &1, String, Connection>):  do IO<Result<&1, &1, String, Connection>>:    result : Result<&1, &1, U32 & String, Socket> <- TCP.connect("127.0.0.1", port)    connected(result, body)def close(connection: Connection) -> IO(Unit):  match connection:    case Connection{socket, buffered}: Socket.close(socket)