src/server.bend source
src/server.bend on the hub · documented module
import Baseimport ./json.bend as Jsontype ServerHeader is Data: ServerHeader{name: String, value: String}# One listener lifetime per process. Zero port asks the OS for a free port.type Config is Data: Config{address: String, port: U32, max_body: U32, max_pending: U32, timeout_ms: U32, grace_ms: U32}def config(address: String, port: U32) -> Config: Config{address, port, 1048576, 128, 10000, 5000}# Accepted socket budget and absolute time to receive the complete request.type TransportLimits is Data: TransportLimits{max_connections: U32, read_timeout_ms: U32}type Started is Data: Listening{port: U32} ListenError{code: String}type Incoming is Data: Incoming{id: U32, method: String, path: String, target: String, headers: List<&2, ServerHeader>, body: String}type Next is Data: Received{request: Incoming} Stopped{}# Opt-in upload streaming exposes validated request metadata before the body.type IncomingHead is Data: IncomingHead{id: U32, method: String, path: String, target: String, headers: List<&2, ServerHeader>}type StreamNext is Data: StreamReceived{request: IncomingHead} StreamStopped{}type BodyNext is Data: BodyChunk{value: String} BodyEnd{} BodyFailure{code: String}type Reply is Data: Reply{status: U32, headers: List<&2, ServerHeader>, body: String}type Replied is Data: Sent{} ReplyError{code: String}# Streaming is opt-in. Ordinary Reply keeps its connection-closing behavior.type ConnectionMode is Data: CloseAfter{} KeepAlive{}type StreamHead is Data: StreamHead{status: U32, headers: List<&2, ServerHeader>, connection: ConnectionMode}def json(status: U32, body: String) -> Reply: Reply{status, Con{ServerHeader{"Content-Type", "application/json"}, Nil{}}, body}def encoded_reply(status: U32, result: Json.EncodeResult) -> IO(Reply): match result: case Json.JsonEncoded{text}: IO.pure(Reply, json(status, text)) case Json.JsonEncodeFailure{code, message}: IO.pure(Reply, json(500, "{\"error\":\"JSON encoding failed\"}"))def json_value(status: U32, value: Json.Json) -> IO(Reply): IO.bind(Json.EncodeResult, Reply, Json.Json.stringify(value), encoded_reply(status))def with_header(name: String, value: String, reply: Reply) -> Reply: match reply: case Reply{status, headers, body}: Reply{status, Con{ServerHeader{name, value}, headers}, body}# Incoming field names are normalized to lowercase by the native boundary.def header_step(equal: Bool, value: String, rest: Unit -> Maybe<String>) -> Maybe<String>: match equal: case True{}: Some{value} case False{}: rest(Unit{})def header(+name: String, headers: List<&2, ServerHeader>) -> Maybe<String>: match headers: case Nil{}: None{} case Con{ServerHeader{key, value}, tail}: header_step(String.eq(name, key), value, u => header(name, tail))def Server.listen(config: Config) -> IO(Started): import "./effects/server.c"def Server.listen_with_limits(config: Config, limits: TransportLimits) -> IO(Started): import "./effects/server.c"# Request streaming is a listener-wide opt-in. Buffered listeners retain their# existing Incoming/next API and allocation behavior.def Server.listen_streaming(config: Config) -> IO(Started): import "./effects/server.c"def Server.listen_streaming_with_limits(config: Config, limits: TransportLimits) -> IO(Started): import "./effects/server.c"def Server.next() -> IO(Next): import "./effects/server.c"def Server.next_stream() -> IO(StreamNext): import "./effects/server.c"# Bounded pending text is retained per request while the consumer works.def Server.body_next(id: U32) -> IO(BodyNext): import "./effects/server.c"def Server.reply(id: U32, reply: Reply) -> IO(Replied): import "./effects/server.c"# Native boundary uses a numeric flag instead of relying on the compiler's# private representation of a nullary data constructor.def Server.stream_start_raw(id: U32, status: U32, headers: List<&2, ServerHeader>, keep_alive: U32) -> IO(Replied): import "./effects/server.c"def stream_start_mode(id: U32, status: U32, headers: List<&2, ServerHeader>, connection: ConnectionMode) -> IO(Replied): match connection: case CloseAfter{}: Server.stream_start_raw(id, status, headers, 0) case KeepAlive{}: Server.stream_start_raw(id, status, headers, 1)# Begin one bounded chunked response. Statuses that forbid bodies are rejected.def Server.stream_start(id: U32, head: StreamHead) -> IO(Replied): match head: case StreamHead{status, headers, connection}: stream_start_mode(id, status, headers, connection)# The aggregate body across writes is limited to 16 MiB.def Server.stream_write(id: U32, chunk: String) -> IO(Replied): import "./effects/server.c"def Server.stream_end(id: U32) -> IO(Replied): import "./effects/server.c"def sse_head(connection: ConnectionMode) -> StreamHead: StreamHead{200, [ ServerHeader{"Content-Type", "text/event-stream"}, ServerHeader{"Cache-Control", "no-cache"}], connection}# Cooperative checkpoint; does not preempt work or undo effects.def Server.active(id: U32) -> IO(Bool): import "./effects/server.c"# Manual dispatch must release its handler budget after all work finishes.def Server.finish(id: U32) -> IO(Unit): import "./effects/server.c"def Server.stop() -> IO(Unit): import "./effects/server.c"def finish_reply(id: U32, result: Replied) -> IO(Unit): match result: case Sent{}: IO.pure(Unit, Unit{}) case ReplyError{code}: do IO<Unit>: ignored : Replied <- Server.reply(id, json(500, "{\"error\":\"response rejected\"}")) IO.pure(Unit, Unit{})def dispatch_live(handler: Incoming -> IO(Reply), request: Incoming) -> IO(Unit): match request: case Incoming{+id, method, path, target, headers, body}: do IO<Unit>: reply : Reply <- handler(Incoming{id, method, path, target, headers, body}) result : Replied <- Server.reply(id, reply) finish_reply(id, result) Server.finish(id)def discard_request(request: Incoming) -> IO(Unit): match request: case Incoming{id, method, path, target, headers, body}: Server.finish(id)def dispatch_active(active: Bool, handler: Incoming -> IO(Reply), request: Incoming) -> IO(Unit): match active: case True{}: dispatch_live(handler, request) case False{}: discard_request(request)def dispatch(handler: Incoming -> IO(Reply), request: Incoming) -> IO(Unit): match request: case Incoming{+id, method, path, target, headers, body}: IO.bind(Bool, Unit, Server.active(id), active => dispatch_active(active, handler, Incoming{id, method, path, target, headers, body}))def serve_next(handler: Incoming -> IO(Reply), rest: Unit -> IO(Unit), next: Next) -> IO(Unit): match next: case Received{request}: do IO<Unit>: IO.spawn(Unit, dispatch(handler, request)) rest(Unit{}) case Stopped{}: IO.pure(Unit, Unit{})# The accept loop ends on external shutdown, not a pure termination proof.@unsafedef serve(~handler: Incoming -> IO(Reply)) -> IO(Unit): IO.bind(Next, Unit, Server.next(), serve_next(handler, u => serve(~handler)))# Process/listener lifetime counters and current bounded gauges. No reset API.type Metrics is Data: Metrics{admitted: U32, completed: U32, failed: U32, rejected: U32, expired: U32, read_expired: U32, disconnected: U32, connections: U32, pending: U32, handlers: U32}def Server.metrics() -> IO(Metrics): import "./effects/server.c"def metrics_json(metrics: Metrics) -> Json.Json: match metrics: case Metrics{admitted, completed, failed, rejected, expired, read_expired, disconnected, connections, pending, handlers}: Json.object([ ("admitted", Json.u32(admitted)), ("completed", Json.u32(completed)), ("failed", Json.u32(failed)), ("rejected", Json.u32(rejected)), ("expired", Json.u32(expired)), ("read_expired", Json.u32(read_expired)), ("disconnected", Json.u32(disconnected)), ("connections", Json.u32(connections)), ("pending", Json.u32(pending)), ("handlers", Json.u32(handlers))])