๐Ÿ‘ a fluffy Gleam web server
23

Configure Feed

Select the types of activity you want to include in your feed.

Replace `gramps` with `websocks` (#8)

authored by

Vladislav Shakitskiy and committed by
GitHub
(Nov 6, 2025, 4:03 PM +0300) 064ad0b8 29726fa3

+312 -1195
+7
CHANGELOG.md
··· 1 1 # Changelog 2 2 3 + # v2.0.3 4 + 5 + - Replace `gramps` with `websocks` package 6 + - Remove alias names for internal stream modules 7 + - Improve script that change documentation 8 + - Move `gleam_crypto` package as dev dependency 9 + 3 10 # v2.0.2 4 11 5 12 - Refactor internal handler code
+2 -2
gleam.toml
··· 17 17 gleam_http = ">= 4.3.0 and < 5.0.0" 18 18 logging = ">= 1.3.0 and < 2.0.0" 19 19 gleam_erlang = ">= 1.3.0 and < 2.0.0" 20 - gleam_crypto = ">= 1.5.1 and < 2.0.0" 21 20 compresso = "0.1.0" 22 - # gramps = ">= 6.0.0 and < 7.0.0" 21 + websocks = ">= 2.0.0 and < 3.0.0" 23 22 24 23 [dev-dependencies] 25 24 gleeunit = ">= 1.0.0 and < 2.0.0" 26 25 gleam_httpc = ">= 5.0.0 and < 6.0.0" 27 26 gleam_json = ">= 3.0.2 and < 4.0.0" 27 + gleam_crypto = ">= 1.5.1 and < 2.0.0" 28 28 29 29 [erlang] 30 30 application_start_module = "ewe@internal@clock"
+3 -1
manifest.toml
··· 16 16 { name = "glisten", version = "8.0.1", build_tools = ["gleam"], requirements = ["gleam_erlang", "gleam_otp", "gleam_stdlib", "logging", "telemetry"], otp_app = "glisten", source = "hex", outer_checksum = "534BB27C71FB9E506345A767C0D76B17A9E9199934340C975DC003C710E3692D" }, 17 17 { name = "logging", version = "1.3.0", build_tools = ["gleam"], requirements = ["gleam_stdlib"], otp_app = "logging", source = "hex", outer_checksum = "1098FBF10B54B44C2C7FDF0B01C1253CAFACDACABEFB4B0D027803246753E06D" }, 18 18 { name = "telemetry", version = "1.3.0", build_tools = ["rebar3"], requirements = [], otp_app = "telemetry", source = "hex", outer_checksum = "7015FC8919DBE63764F4B4B87A95B7C0996BD539E0D499BE6EC9D7F3875B79E6" }, 19 + { name = "websocks", version = "2.0.0", build_tools = ["gleam"], requirements = ["gleam_crypto", "gleam_erlang", "gleam_stdlib"], otp_app = "websocks", source = "hex", outer_checksum = "A13BF89A8AC63C478C0E9FE502C81A6CBB0513F427FA5C1DFD96383BF46D501D" }, 19 20 ] 20 21 21 22 [requirements] 22 23 compresso = { version = "0.1.0" } 23 - gleam_crypto = { version = ">= 1.5.1 and < 2.0.0" } 24 24 gleam_erlang = { version = ">= 1.3.0 and < 2.0.0" } 25 25 gleam_http = { version = ">= 4.3.0 and < 5.0.0" } 26 26 gleam_httpc = { version = ">= 5.0.0 and < 6.0.0" } ··· 30 30 gleeunit = { version = ">= 1.0.0 and < 2.0.0" } 31 31 glisten = { version = ">= 8.0.1 and < 9.0.0" } 32 32 logging = { version = ">= 1.3.0 and < 2.0.0" } 33 + websocks = { version = ">= 2.0.0 and < 3.0.0" } 34 + gleam_crypto = { version = ">= 1.5.1 and < 2.0.0" }
+164 -129
src/ewe.gleam
··· 1 - //// <style> 2 - //// .content > h4, 3 - //// .content > ul { 4 - //// display: none; 1 + //// <script> 2 + //// const docs = [ 3 + //// { 4 + //// header: "IP Address", 5 + //// functions: ["ip_address_to_string"] 6 + //// }, 7 + //// { 8 + //// header: "Information", 9 + //// functions: [ 10 + //// "get_client_info", 11 + //// "get_server_info" 12 + //// ] 13 + //// }, 14 + //// { 15 + //// header: "Builder", 16 + //// functions: [ 17 + //// "new", 18 + //// "bind", 19 + //// "bind_all", 20 + //// "listening", 21 + //// "listening_random", 22 + //// "enable_ipv6", 23 + //// "enable_tls", 24 + //// "with_name", 25 + //// "quiet", 26 + //// "idle_timeout", 27 + //// "on_start", 28 + //// "on_crash" 29 + //// ] 30 + //// }, 31 + //// { 32 + //// header: "Server", 33 + //// functions: [ 34 + //// "start", 35 + //// "supervised" 36 + //// ] 37 + //// }, 38 + //// { 39 + //// header: "Request", 40 + //// functions: [ 41 + //// "read_body", 42 + //// "stream_body" 43 + //// ] 44 + //// }, 45 + //// { 46 + //// header: "Response", 47 + //// functions: ["file"] 48 + //// }, 49 + //// { 50 + //// header: "Chunked Response", 51 + //// functions: [ 52 + //// "chunked_body", 53 + //// "send_chunk", 54 + //// "chunked_continue", 55 + //// "chunked_stop", 56 + //// "chunked_stop_abnormal" 57 + //// ] 58 + //// }, 59 + //// { 60 + //// header: "Websocket", 61 + //// functions: [ 62 + //// "upgrade_websocket", 63 + //// "send_binary_frame", 64 + //// "send_text_frame", 65 + //// "websocket_continue", 66 + //// "websocket_continue_with_selector", 67 + //// "websocket_stop", 68 + //// "websocket_stop_abnormal" 69 + //// ] 70 + //// }, 71 + //// { 72 + //// header: "Server-Sent Events", 73 + //// functions: [ 74 + //// "sse", 75 + //// "event", 76 + //// "event_name", 77 + //// "event_id", 78 + //// "event_retry", 79 + //// "send_event", 80 + //// "sse_continue", 81 + //// "sse_stop", 82 + //// "sse_stop_abnormal" 83 + //// ] 5 84 //// } 6 - //// </style> 7 - //// <script> 8 - //// // https://gitlab.com/arkandos/smol/-/blob/main/src/smol.gleam?ref_type=heads 9 - //// (callback => document.readyState !== 'loading' ? callback() : document.addEventListener('DOMContentLoaded', callback, { once: true }))(() => { 10 - //// const list = document.querySelector('.sidebar > ul:last-of-type') 85 + //// ] 86 + //// 87 + //// const callback = () => { 88 + //// const list = document.querySelector(".sidebar > ul:last-of-type") 11 89 //// const sortedLists = document.createDocumentFragment() 12 90 //// const sortedMembers = document.createDocumentFragment() 13 - //// 14 - //// for (const header of document.querySelectorAll('main > h4')) { 91 + //// 92 + //// for (const section of docs) { 15 93 //// sortedLists.append((() => { 16 - //// const node = document.createElement('h3') 17 - //// node.append(header.textContent) 94 + //// const node = document.createElement("h3") 95 + //// node.append(section.header) 18 96 //// return node 19 97 //// })()) 20 98 //// sortedMembers.append((() => { 21 - //// const node = document.createElement('h2') 22 - //// node.append(header.textContent) 99 + //// const node = document.createElement("h2") 100 + //// node.append(section.header) 23 101 //// return node 24 102 //// })()) 25 - //// 26 - //// const sortedList = document.createElement('ul') 103 + //// 104 + //// const sortedList = document.createElement("ul") 27 105 //// sortedLists.append(sortedList) 28 - //// 29 - //// for (const anchor of header.nextElementSibling.querySelectorAll('a')) { 30 - //// const href = anchor.getAttribute('href') 31 - //// const member = document.querySelector(`.member:has(h2 > a[href="${href}"])`) 106 + //// 107 + //// const sortedFunctions = [...section.functions].sort() 108 + //// 109 + //// for (const funcName of sortedFunctions) { 110 + //// const href = `#${funcName}` 111 + //// const member = document.querySelector( 112 + //// `.member:has(h2 > a[href="${href}"])` 113 + //// ) 32 114 //// const sidebar = list.querySelector(`li:has(a[href="${href}"])`) 33 115 //// sortedList.append(sidebar) 34 116 //// sortedMembers.append(member) 35 117 //// } 36 118 //// } 37 - //// 38 - //// document.querySelector('.sidebar').insertBefore(sortedLists, list) 39 - //// document.querySelector('.module-members:has(#module-values)').insertBefore(sortedMembers, document.querySelector('#module-values').nextSibling) 40 - //// }) 119 + //// 120 + //// document.querySelector(".sidebar").insertBefore(sortedLists, list) 121 + //// document 122 + //// .querySelector(".module-members:has(#module-values)") 123 + //// .insertBefore( 124 + //// sortedMembers, 125 + //// document.querySelector("#module-values").nextSibling 126 + //// ) 127 + //// } 128 + //// 129 + //// document.readyState !== "loading" 130 + //// ? callback() 131 + //// : document.addEventListener( 132 + //// "DOMContentLoaded", 133 + //// callback, 134 + //// { once: true } 135 + //// ) 41 136 //// </script> 42 - //// #### IP Address 43 - //// - [ip_address_to_string](#ip_address_to_string) 44 - //// #### Information 45 - //// - [get_client_info](#get_client_info) 46 - //// - [get_server_info](#get_server_info) 47 - //// #### Builder 48 - //// - [new](#new) 49 - //// - [bind](#bind) 50 - //// - [bind_all](#bind_all) 51 - //// - [listening](#listening) 52 - //// - [listening_random](#listening_random) 53 - //// - [enable_ipv6](#enable_ipv6) 54 - //// - [enable_tls](#enable_tls) 55 - //// - [with_name](#with_name) 56 - //// - [quiet](#quiet) 57 - //// - [idle_timeout](#idle_timeout) 58 - //// - [on_start](#on_start) 59 - //// - [on_crash](#on_crash) 60 - //// #### Server 61 - //// - [start](#start) 62 - //// - [supervised](#supervised) 63 - //// #### Request 64 - //// - [read_body](#read_body) 65 - //// - [stream_body](#stream_body) 66 - //// #### Response 67 - //// - [file](#file) 68 - //// #### Chunked Response 69 - //// - [chunked_body](#chunked_body) 70 - //// - [send_chunk](#send_chunk) 71 - //// - [chunked_continue](#chunked_continue) 72 - //// - [chunked_stop](#chunked_stop) 73 - //// - [chunked_stop_abnormal](#chunked_stop_abnormal) 74 - //// #### Websocket 75 - //// - [upgrade_websocket](#upgrade_websocket) 76 - //// - [send_binary_frame](#send_binary_frame) 77 - //// - [send_text_frame](#send_text_frame) 78 - //// - [websocket_continue](#websocket_continue) 79 - //// - [websocket_continue_with_selector](#websocket_continue_with_selector) 80 - //// - [websocket_stop](#websocket_stop) 81 - //// - [websocket_stop_abnormal](#websocket_stop_abnormal) 82 - //// #### Server-Sent Events 83 - //// - [sse](#sse) 84 - //// - [event](#event) 85 - //// - [event_name](#event_name) 86 - //// - [event_id](#event_id) 87 - //// - [event_retry](#event_retry) 88 - //// - [send_event](#send_event) 89 - //// - [sse_continue](#sse_continue) 90 - //// - [sse_stop](#sse_stop) 91 - //// - [sse_stop_abnormal](#sse_stop_abnormal) 92 137 93 138 // ----------------------------------------------------------------------------- 94 139 // IMPORTS ··· 111 156 import gleam/string_tree.{type StringTree} 112 157 import logging 113 158 159 + import websocks 160 + 114 161 import glisten 115 162 import glisten/internal/listener 116 163 import glisten/socket/options as glisten_options 117 164 import glisten/transport 118 165 119 - // TODO: replace this once gramps changes are published 120 - import ewe/internal/gramps/websocket as ws 121 - 122 166 import ewe/internal/file 123 167 import ewe/internal/handler 124 168 import ewe/internal/http1 as ewe_http 125 - import ewe/internal/stream/chunked as ewe_chunked 126 - import ewe/internal/stream/sse as ewe_sse 127 - import ewe/internal/stream/websocket as ewe_ws 169 + import ewe/internal/stream/chunked 170 + import ewe/internal/stream/sse 171 + import ewe/internal/stream/websocket 128 172 129 173 // ----------------------------------------------------------------------------- 130 174 // CONNECTION ··· 665 709 666 710 /// Represents a chunked response body. This type is used to send a chunked response to the client. 667 711 pub type ChunkedBody = 668 - ewe_chunked.ChunkedBody 712 + chunked.ChunkedBody 669 713 670 714 /// Represents an instruction on how chunked response should be processed. 671 715 /// ··· 699 743 700 744 fn to_internal_chunked_next( 701 745 next: ChunkedNext(user_state), 702 - ) -> ewe_chunked.ChunkedNext(user_state) { 746 + ) -> chunked.ChunkedNext(user_state) { 703 747 case next { 704 - ChunkedContinue(user_state) -> ewe_chunked.Continue(user_state) 705 - ChunkedStop -> ewe_chunked.NormalStop 706 - ChunkedAbnormalStop(reason) -> ewe_chunked.AbnormalStop(reason) 748 + ChunkedContinue(user_state) -> chunked.Continue(user_state) 749 + ChunkedStop -> chunked.NormalStop 750 + ChunkedAbnormalStop(reason) -> chunked.AbnormalStop(reason) 707 751 } 708 752 } 709 753 ··· 735 779 let socket = req.body.socket 736 780 let factory_name = req.body.factory_name 737 781 738 - case ewe_chunked.send_response(resp, transport, socket) { 782 + case chunked.send_response(resp, transport, socket) { 739 783 Ok(Nil) -> { 740 784 let supervisor = factory.get_by_name(factory_name) 741 785 742 786 let start_result = 743 787 factory.start_child(supervisor, fn() { 744 - ewe_chunked.start(transport, socket, on_init, handler, on_close) 788 + chunked.start(transport, socket, on_init, handler, on_close) 745 789 }) 746 790 747 791 case start_result { ··· 762 806 body: ChunkedBody, 763 807 chunk: BitArray, 764 808 ) -> Result(Nil, glisten.SocketReason) { 765 - ewe_chunked.send_chunk(body.transport, body.socket, chunk) 809 + chunked.send_chunk(body.transport, body.socket, chunk) 766 810 } 767 811 768 812 // ----------------------------------------------------------------------------- ··· 772 816 /// Represents a WebSocket connection between a client and a server. 773 817 /// 774 818 pub type WebsocketConnection = 775 - ewe_ws.WebsocketConnection 819 + websocket.WebsocketConnection 776 820 777 821 /// Represents an instruction on how WebSocket connection should proceed. 778 822 /// ··· 822 866 823 867 fn to_internal_websocket_next( 824 868 next: WebsocketNext(user_state, user_message), 825 - ) -> ewe_ws.WebsocketNext(user_state, user_message) { 869 + ) -> websocket.WebsocketNext(user_state, user_message) { 826 870 case next { 827 871 WebsocketContinue(user_state, selector) -> 828 - ewe_ws.Continue(user_state, selector) 829 - WebsocketNormalStop -> ewe_ws.NormalStop 830 - WebsocketAbnormalStop(reason) -> ewe_ws.AbnormalStop(reason) 872 + websocket.Continue(user_state, selector) 873 + WebsocketNormalStop -> websocket.NormalStop 874 + WebsocketAbnormalStop(reason) -> websocket.AbnormalStop(reason) 831 875 } 832 876 } 833 877 ··· 846 890 } 847 891 848 892 fn transform_websocket_message( 849 - message: ewe_ws.WebsocketMessage(user_message), 893 + message: websocket.WebsocketMessage(user_message), 850 894 ) -> Result(WebsocketMessage(user_message), Nil) { 851 - // NOTE: see "https://github.com/rawhat/gramps/pull/7" 852 895 case message { 853 - ewe_ws.WebsocketFrame(ws.Data(frame)) -> { 854 - ws.match_data_frame( 855 - frame, 856 - on_text: fn(payload, _) { 857 - bit_array.to_string(payload) |> result.map(Text) 858 - }, 859 - on_binary: fn(payload, _) { Ok(Binary(payload)) }, 860 - ) 861 - } 862 - ewe_ws.UserMessage(user_message) -> Ok(User(user_message)) 896 + websocket.Frame(websocks.Text(payload)) -> 897 + bit_array.to_string(payload) |> result.map(Text) 898 + websocket.Frame(websocks.Binary(payload)) -> Ok(Binary(payload)) 899 + websocket.UserMessage(user_message) -> Ok(User(user_message)) 863 900 _ -> Error(Nil) 864 901 } 865 902 } ··· 905 942 let supervisor = factory.get_by_name(factory_name) 906 943 let start_result = 907 944 factory.start_child(supervisor, fn() { 908 - ewe_ws.start( 945 + websocket.start( 909 946 transport, 910 947 socket, 911 948 on_init, ··· 935 972 conn: WebsocketConnection, 936 973 bits: BitArray, 937 974 ) -> Result(Nil, glisten.SocketReason) { 938 - ewe_ws.send_frame( 939 - ws.encode_binary_frame, 975 + websocket.send_frame( 976 + websocks.encode_binary_frame, 940 977 conn.transport, 941 978 conn.socket, 942 - conn.deflate, 979 + conn.context, 943 980 bits, 944 981 ) 945 982 } ··· 950 987 conn: WebsocketConnection, 951 988 text: String, 952 989 ) -> Result(Nil, glisten.SocketReason) { 953 - ewe_ws.send_frame( 954 - ws.encode_text_frame, 990 + websocket.send_frame( 991 + websocks.encode_text_frame, 955 992 conn.transport, 956 993 conn.socket, 957 - conn.deflate, 958 - text, 994 + conn.context, 995 + bit_array.from_string(text), 959 996 ) 960 997 } 961 998 ··· 966 1003 /// Represents a Server-Sent Events connection between a client and a server. 967 1004 /// 968 1005 pub type SSEConnection = 969 - ewe_sse.SSEConnection 1006 + sse.SSEConnection 970 1007 971 1008 /// Represents an instruction on how Server-Sent Events connection should 972 1009 /// proceed. ··· 999 1036 SSEAbnormalStop(reason) 1000 1037 } 1001 1038 1002 - fn to_internal_sse_next( 1003 - next: SSENext(user_state), 1004 - ) -> ewe_sse.SSENext(user_state) { 1039 + fn to_internal_sse_next(next: SSENext(user_state)) -> sse.SSENext(user_state) { 1005 1040 case next { 1006 - SSEContinue(user_state) -> ewe_sse.Continue(user_state) 1007 - SSENormalStop -> ewe_sse.NormalStop 1008 - SSEAbnormalStop(reason) -> ewe_sse.AbnormalStop(reason) 1041 + SSEContinue(user_state) -> sse.Continue(user_state) 1042 + SSENormalStop -> sse.NormalStop 1043 + SSEAbnormalStop(reason) -> sse.AbnormalStop(reason) 1009 1044 } 1010 1045 } 1011 1046 ··· 1020 1055 /// `ewe.event_id`, and `ewe.event_retry`. 1021 1056 /// 1022 1057 pub type SSEEvent = 1023 - ewe_sse.SSEEvent 1058 + sse.SSEEvent 1024 1059 1025 1060 /// Creates a new SSE event with the given data. Use `ewe.event_name`, 1026 1061 /// `ewe.event_id`, and `ewe.event_retry` to modify other fields of the event. 1027 1062 /// 1028 1063 pub fn event(data: String) -> SSEEvent { 1029 - ewe_sse.SSEEvent(event: None, data:, id: None, retry: None) 1064 + sse.SSEEvent(event: None, data:, id: None, retry: None) 1030 1065 } 1031 1066 1032 1067 /// Sets the name of the event. 1033 1068 /// 1034 1069 pub fn event_name(event: SSEEvent, name: String) -> SSEEvent { 1035 - ewe_sse.SSEEvent(..event, event: Some(name)) 1070 + sse.SSEEvent(..event, event: Some(name)) 1036 1071 } 1037 1072 1038 1073 /// Sets the ID of the event. 1039 1074 /// 1040 1075 pub fn event_id(event: SSEEvent, id: String) -> SSEEvent { 1041 - ewe_sse.SSEEvent(..event, id: Some(id)) 1076 + sse.SSEEvent(..event, id: Some(id)) 1042 1077 } 1043 1078 1044 1079 /// Sets the retry time of the event. 1045 1080 /// 1046 1081 pub fn event_retry(event: SSEEvent, retry: Int) -> SSEEvent { 1047 - ewe_sse.SSEEvent(..event, retry: Some(retry)) 1082 + sse.SSEEvent(..event, retry: Some(retry)) 1048 1083 } 1049 1084 1050 1085 /// Sets up the connection for Server-Sent Events. ··· 1074 1109 let socket = req.body.socket 1075 1110 let factory_name = req.body.factory_name 1076 1111 1077 - case ewe_sse.send_response(transport, socket) { 1112 + case sse.send_response(transport, socket) { 1078 1113 Ok(Nil) -> { 1079 1114 let supervisor = factory.get_by_name(factory_name) 1080 1115 let start_result = 1081 1116 factory.start_child(supervisor, fn() { 1082 - ewe_sse.start(transport, socket, on_init, handler, on_close) 1117 + sse.start(transport, socket, on_init, handler, on_close) 1083 1118 }) 1084 1119 1085 1120 case start_result { ··· 1100 1135 conn: SSEConnection, 1101 1136 event: SSEEvent, 1102 1137 ) -> Result(Nil, glisten.SocketReason) { 1103 - ewe_sse.send_event(conn.transport, conn.socket, event) 1138 + sse.send_event(conn.transport, conn.socket, event) 1104 1139 }
-678
src/ewe/internal/gramps/websocket.gleam
··· 1 - // TODO: remove this once gramps changes are published 2 - // See https://github.com/rawhat/gramps 3 - 4 - import ewe/internal/gramps/websocket/compression.{ 5 - type Context, type ContextTakeover, ContextTakeover, 6 - } 7 - import gleam/bit_array 8 - import gleam/bool 9 - import gleam/bytes_tree.{type BytesTree} 10 - import gleam/crypto 11 - import gleam/list 12 - import gleam/option.{type Option, None, Some} 13 - import gleam/result 14 - import gleam/string 15 - 16 - pub opaque type DataFrame { 17 - TextFrame(payload: BitArray) 18 - BinaryFrame(payload: BitArray) 19 - 20 - CompressedTextFrame(payload: BitArray) 21 - CompressedBinaryFrame(payload: BitArray) 22 - } 23 - 24 - pub fn text_frame(payload: BitArray) -> DataFrame { 25 - TextFrame(payload) 26 - } 27 - 28 - pub fn binary_frame(payload: BitArray) -> DataFrame { 29 - BinaryFrame(payload) 30 - } 31 - 32 - pub fn match_data_frame( 33 - data_frame: DataFrame, 34 - on_text on_text: fn(BitArray, Bool) -> a, 35 - on_binary on_binary: fn(BitArray, Bool) -> a, 36 - ) -> a { 37 - case data_frame { 38 - TextFrame(payload) -> on_text(payload, False) 39 - CompressedTextFrame(payload) -> on_text(payload, True) 40 - BinaryFrame(payload) -> on_binary(payload, False) 41 - CompressedBinaryFrame(payload) -> on_binary(payload, True) 42 - } 43 - } 44 - 45 - pub type CloseReason { 46 - NotProvided 47 - Normal(body: BitArray) 48 - GoingAway(body: BitArray) 49 - ProtocolError(body: BitArray) 50 - UnexpectedDataType(body: BitArray) 51 - InconsistentDataType(body: BitArray) 52 - PolicyViolation(body: BitArray) 53 - MessageTooBig(body: BitArray) 54 - MissingExtensions(body: BitArray) 55 - UnexpectedCondition(body: BitArray) 56 - /// Usually used for `4000` codes. 57 - CustomCloseReason( 58 - /// If `code >= 5000`, it will be the same as a `Normal` close reason. 59 - code: Int, 60 - body: BitArray, 61 - ) 62 - } 63 - 64 - pub type ControlFrame { 65 - CloseFrame(reason: CloseReason) 66 - PingFrame(payload: BitArray) 67 - PongFrame(payload: BitArray) 68 - } 69 - 70 - pub type Frame { 71 - Data(DataFrame) 72 - Control(ControlFrame) 73 - Continuation(length: Int, payload: BitArray) 74 - } 75 - 76 - @external(erlang, "crypto", "exor") 77 - fn crypto_exor(a a: BitArray, b b: BitArray) -> BitArray 78 - 79 - fn mask_data(data: BitArray, masks: List(BitArray)) -> BitArray { 80 - let assert [m1, m2, m3, m4] = masks 81 - let mask_key = <<m1:bits, m2:bits, m3:bits, m4:bits>> 82 - 83 - let payload_size = bit_array.byte_size(data) 84 - let full_mask = create_repeating_mask(mask_key, payload_size) 85 - crypto_exor(data, full_mask) 86 - } 87 - 88 - fn create_repeating_mask(mask_key: BitArray, size: Int) -> BitArray { 89 - case size { 90 - 1 | 2 | 3 | 4 -> bit_array.slice(mask_key, 0, size) |> result.unwrap(<<>>) 91 - 92 - _ -> { 93 - let repetitions = size / 4 94 - let remainder = size % 4 95 - let base = list.repeat(mask_key, repetitions) |> bit_array.concat 96 - 97 - case remainder { 98 - 0 -> base 99 - n -> { 100 - let partial = bit_array.slice(mask_key, 0, n) |> result.unwrap(<<>>) 101 - <<base:bits, partial:bits>> 102 - } 103 - } 104 - } 105 - } 106 - } 107 - 108 - pub type FrameParseError { 109 - NeedMoreData(BitArray) 110 - InvalidFrame 111 - } 112 - 113 - pub type ParsedFrame { 114 - Complete(Frame) 115 - Incomplete(Frame) 116 - } 117 - 118 - pub fn decode_frame( 119 - message: BitArray, 120 - context: Option(Context), 121 - ) -> Result(#(ParsedFrame, BitArray), FrameParseError) { 122 - case message { 123 - << 124 - complete:1, 125 - compressed:1, 126 - rsv2:1, 127 - rsv3:1, 128 - opcode:int-size(4), 129 - masked:1, 130 - payload_length:int-size(7), 131 - rest:bits, 132 - >> -> { 133 - let compressed = compressed == 1 134 - let masked = masked == 1 135 - 136 - use <- bool.guard( 137 - when: compressed && option.is_none(context), 138 - return: Error(InvalidFrame), 139 - ) 140 - 141 - use <- bool.guard(rsv2 == 1 || rsv3 == 1, return: Error(InvalidFrame)) 142 - 143 - use <- bool.guard( 144 - when: { 145 - let is_control_frame = opcode >= 8 && opcode <= 10 146 - let is_fragmented = complete == 0 147 - is_control_frame && is_fragmented 148 - }, 149 - return: Error(InvalidFrame), 150 - ) 151 - 152 - let payload_size = case payload_length { 153 - 126 -> 16 154 - 127 -> 64 155 - _ -> 0 156 - } 157 - 158 - let maybe_pair = case masked, rest { 159 - True, 160 - << 161 - length:int-size(payload_size), 162 - mask1:bytes-size(1), 163 - mask2:bytes-size(1), 164 - mask3:bytes-size(1), 165 - mask4:bytes-size(1), 166 - rest:bits, 167 - >> 168 - -> { 169 - let payload_byte_size = case length { 170 - 0 -> payload_length 171 - n -> n 172 - } 173 - 174 - case bit_array.byte_size(rest) >= payload_byte_size, rest { 175 - True, <<payload:bytes-size(payload_byte_size), remaining:bits>> -> { 176 - let data = mask_data(payload, [mask1, mask2, mask3, mask4]) 177 - Ok(#(data, remaining)) 178 - } 179 - _, _ -> Error(NeedMoreData(message)) 180 - } 181 - } 182 - True, _rest -> Error(NeedMoreData(message)) 183 - False, <<length:int-size(payload_size), rest:bits>> -> { 184 - let payload_byte_size = case length { 185 - 0 -> payload_length 186 - n -> n 187 - } 188 - case rest { 189 - <<payload:bytes-size(payload_byte_size), rest:bits>> -> { 190 - Ok(#(payload, rest)) 191 - } 192 - _ -> { 193 - Error(NeedMoreData(message)) 194 - } 195 - } 196 - } 197 - _, _ -> Error(InvalidFrame) 198 - } 199 - 200 - use #(data, rest) <- result.try(maybe_pair) 201 - case opcode { 202 - 0 -> Ok(Continuation(payload_size, data)) 203 - 1 -> { 204 - case compressed { 205 - True -> Ok(Data(CompressedTextFrame(data))) 206 - False -> Ok(Data(TextFrame(data))) 207 - } 208 - } 209 - 2 -> { 210 - case compressed { 211 - True -> Ok(Data(CompressedBinaryFrame(data))) 212 - False -> Ok(Data(BinaryFrame(data))) 213 - } 214 - } 215 - 8 -> { 216 - case data { 217 - <<>> -> Ok(Control(CloseFrame(NotProvided))) 218 - <<code:16, rest:bits>> -> { 219 - use <- bool.guard( 220 - when: !bit_array.is_utf8(rest), 221 - return: Error(InvalidFrame), 222 - ) 223 - 224 - case code { 225 - 1000 -> Ok(Control(CloseFrame(Normal(rest)))) 226 - 1001 -> Ok(Control(CloseFrame(GoingAway(rest)))) 227 - 1002 -> Ok(Control(CloseFrame(ProtocolError(rest)))) 228 - 1003 -> Ok(Control(CloseFrame(UnexpectedDataType(rest)))) 229 - 1007 -> Ok(Control(CloseFrame(InconsistentDataType(rest)))) 230 - 1008 -> Ok(Control(CloseFrame(PolicyViolation(rest)))) 231 - 1009 -> Ok(Control(CloseFrame(MessageTooBig(rest)))) 232 - 1010 -> Ok(Control(CloseFrame(MissingExtensions(rest)))) 233 - 1011 -> Ok(Control(CloseFrame(UnexpectedCondition(rest)))) 234 - code if code >= 3000 && code <= 4999 -> 235 - Ok(Control(CloseFrame(CustomCloseReason(code, rest)))) 236 - _ -> Error(InvalidFrame) 237 - } 238 - } 239 - _ -> Error(InvalidFrame) 240 - } 241 - } 242 - 9 -> Ok(Control(PingFrame(data))) 243 - 10 -> Ok(Control(PongFrame(data))) 244 - _ -> Error(InvalidFrame) 245 - } 246 - |> result.try(fn(frame) { 247 - case complete { 248 - 1 -> Ok(#(Complete(frame), rest)) 249 - 0 -> Ok(#(Incomplete(frame), rest)) 250 - _ -> Error(InvalidFrame) 251 - } 252 - }) 253 - } 254 - _ -> Error(NeedMoreData(message)) 255 - } 256 - } 257 - 258 - pub fn encode_text_frame( 259 - data: String, 260 - context: Option(Context), 261 - mask: Option(BitArray), 262 - ) -> BytesTree { 263 - to_frame(bit_array.from_string(data), context, mask, TextFrame, Data) 264 - } 265 - 266 - pub fn encode_binary_frame( 267 - data: BitArray, 268 - context: Option(Context), 269 - mask: Option(BitArray), 270 - ) -> BytesTree { 271 - to_frame(data, context, mask, BinaryFrame, Data) 272 - } 273 - 274 - pub fn encode_close_frame( 275 - reason: CloseReason, 276 - mask: Option(BitArray), 277 - ) -> BytesTree { 278 - encode_frame(Control(CloseFrame(reason)), Uncompressed, mask) 279 - } 280 - 281 - pub fn encode_ping_frame(data: BitArray, mask: Option(BitArray)) -> BytesTree { 282 - to_frame(data, None, mask, PingFrame, Control) 283 - } 284 - 285 - pub fn encode_pong_frame(data: BitArray, mask: Option(BitArray)) -> BytesTree { 286 - to_frame(data, None, mask, PongFrame, Control) 287 - } 288 - 289 - pub fn encode_continuation_frame( 290 - data: BitArray, 291 - total_size: Int, 292 - mask: Option(BitArray), 293 - ) -> BytesTree { 294 - let payload = apply_mask(data, mask) 295 - encode_frame(Continuation(total_size, payload), Uncompressed, mask) 296 - } 297 - 298 - fn encode_frame( 299 - frame: Frame, 300 - compressed: Compression, 301 - mask: Option(BitArray), 302 - ) -> BytesTree { 303 - case frame { 304 - Data(TextFrame(payload)) | Data(CompressedTextFrame(payload)) -> { 305 - let payload_length = bit_array.byte_size(payload) 306 - make_frame(1, payload_length, payload, compressed, mask) 307 - } 308 - 309 - Data(BinaryFrame(payload)) | Data(CompressedBinaryFrame(payload)) -> { 310 - let payload_length = bit_array.byte_size(payload) 311 - make_frame(2, payload_length, payload, compressed, mask) 312 - } 313 - 314 - Control(CloseFrame(reason)) -> { 315 - let #(payload_length, payload) = case reason { 316 - NotProvided -> #(0, <<>>) 317 - GoingAway(body:) -> { 318 - let payload_size = bit_array.byte_size(body) + 2 319 - #(payload_size, <<1001:16, body:bits>>) 320 - } 321 - InconsistentDataType(body:) -> { 322 - let payload_size = bit_array.byte_size(body) + 2 323 - #(payload_size, <<1007:16, body:bits>>) 324 - } 325 - MessageTooBig(body:) -> { 326 - let payload_size = bit_array.byte_size(body) + 2 327 - #(payload_size, <<1009:16, body:bits>>) 328 - } 329 - MissingExtensions(body:) -> { 330 - let payload_size = bit_array.byte_size(body) + 2 331 - #(payload_size, <<1010:16, body:bits>>) 332 - } 333 - Normal(body:) -> { 334 - let payload_size = bit_array.byte_size(body) + 2 335 - #(payload_size, <<1000:16, body:bits>>) 336 - } 337 - PolicyViolation(body:) -> { 338 - let payload_size = bit_array.byte_size(body) + 2 339 - #(payload_size, <<1008:16, body:bits>>) 340 - } 341 - ProtocolError(body:) -> { 342 - let payload_size = bit_array.byte_size(body) + 2 343 - #(payload_size, <<1002:16, body:bits>>) 344 - } 345 - UnexpectedCondition(body:) -> { 346 - let payload_size = bit_array.byte_size(body) + 2 347 - #(payload_size, <<1011:16, body:bits>>) 348 - } 349 - UnexpectedDataType(body:) -> { 350 - let payload_size = bit_array.byte_size(body) + 2 351 - #(payload_size, <<1003:16, body:bits>>) 352 - } 353 - CustomCloseReason(code:, body:) -> { 354 - let payload_size = bit_array.byte_size(body) + 2 355 - // Prevents integer overflow and changes the status code to `Normal` for invalid codes. 356 - let code = case code < 5000 { 357 - True -> code 358 - False -> 1000 359 - } 360 - #(payload_size, <<code:16, body:bits>>) 361 - } 362 - } 363 - make_frame(8, payload_length, apply_mask(payload, mask), compressed, mask) 364 - } 365 - Control(PongFrame(payload)) -> { 366 - let payload_length = bit_array.byte_size(payload) 367 - make_frame(10, payload_length, payload, compressed, mask) 368 - } 369 - Control(PingFrame(payload)) -> { 370 - let payload_length = bit_array.byte_size(payload) 371 - make_frame(9, payload_length, payload, compressed, mask) 372 - } 373 - Continuation(length, payload) -> 374 - make_frame(0, length, payload, compressed, mask) 375 - } 376 - } 377 - 378 - fn make_length(length: Int) -> BitArray { 379 - case length { 380 - length if length > 65_535 -> <<127:7, length:int-size(64)>> 381 - length if length >= 126 -> <<126:7, length:int-size(16)>> 382 - _length -> <<length:7>> 383 - } 384 - } 385 - 386 - type Compression { 387 - Compressed 388 - Uncompressed 389 - } 390 - 391 - fn make_frame( 392 - opcode: Int, 393 - length: Int, 394 - payload: BitArray, 395 - compressed: Compression, 396 - mask: Option(BitArray), 397 - ) -> BytesTree { 398 - let length_section = make_length(length) 399 - 400 - let masked = case option.is_some(mask) { 401 - True -> 1 402 - False -> 0 403 - } 404 - 405 - let mask_key = option.unwrap(mask, <<>>) 406 - 407 - let compressed = case compressed { 408 - Compressed -> 1 409 - Uncompressed -> 0 410 - } 411 - 412 - << 413 - 1:1, 414 - compressed:1, 415 - 0:2, 416 - opcode:4, 417 - masked:1, 418 - length_section:bits, 419 - mask_key:bits, 420 - payload:bits, 421 - >> 422 - |> bytes_tree.from_bit_array 423 - } 424 - 425 - pub fn apply_mask(data: BitArray, mask: Option(BitArray)) -> BitArray { 426 - case mask { 427 - Some(mask) -> { 428 - let assert << 429 - mask1:bytes-size(1), 430 - mask2:bytes-size(1), 431 - mask3:bytes-size(1), 432 - mask4:bytes-size(1), 433 - >> = mask 434 - mask_data(data, [mask1, mask2, mask3, mask4]) 435 - } 436 - None -> data 437 - } 438 - } 439 - 440 - pub fn apply_deflate(data: BitArray, context: Option(Context)) -> BitArray { 441 - case context { 442 - Some(context) -> compression.deflate(context, data) 443 - _ -> data 444 - } 445 - } 446 - 447 - pub fn apply_inflate(data: BitArray, context: Option(Context)) -> BitArray { 448 - case context { 449 - Some(context) -> compression.inflate(context, data) 450 - _ -> data 451 - } 452 - } 453 - 454 - fn to_frame( 455 - data: BitArray, 456 - context: Option(Context), 457 - mask: Option(BitArray), 458 - create_inner_frame: fn(BitArray) -> a, 459 - create_frame: fn(a) -> Frame, 460 - ) -> BytesTree { 461 - let frame = 462 - data 463 - |> apply_deflate(context) 464 - |> apply_mask(mask) 465 - |> create_inner_frame 466 - |> create_frame 467 - let compress = case context { 468 - Some(_context) -> Compressed 469 - _ -> Uncompressed 470 - } 471 - encode_frame(frame, compress, mask) 472 - } 473 - 474 - pub fn decode_many_frames( 475 - data: BitArray, 476 - context: Option(Context), 477 - frames: List(ParsedFrame), 478 - ) -> #(List(ParsedFrame), BitArray) { 479 - case decode_frame(data, context) { 480 - Ok(#(frame, <<>>)) -> #(list.reverse([frame, ..frames]), <<>>) 481 - Ok(#(frame, rest)) -> decode_many_frames(rest, context, [frame, ..frames]) 482 - Error(NeedMoreData(rest)) -> #(list.reverse(frames), rest) 483 - Error(InvalidFrame) -> #(list.reverse(frames), data) 484 - } 485 - } 486 - 487 - pub type ManyFramesParseError { 488 - NeedMoreDataAccumulated(parsed: List(ParsedFrame), rest: BitArray) 489 - ContainsInvalidFrame 490 - } 491 - 492 - pub fn decode_many_frames_result( 493 - data: BitArray, 494 - context: Option(Context), 495 - frames: List(ParsedFrame), 496 - ) -> Result(#(List(ParsedFrame), BitArray), ManyFramesParseError) { 497 - case decode_frame(data, context) { 498 - Ok(#(frame, <<>>)) -> Ok(#(list.reverse([frame, ..frames]), <<>>)) 499 - Ok(#(frame, rest)) -> 500 - decode_many_frames_result(rest, context, [frame, ..frames]) 501 - Error(NeedMoreData(rest)) -> 502 - Error(NeedMoreDataAccumulated(list.reverse(frames), rest)) 503 - Error(InvalidFrame) -> Error(ContainsInvalidFrame) 504 - } 505 - } 506 - 507 - pub fn aggregate_frames( 508 - frames: List(ParsedFrame), 509 - accumulated: Option(#(Frame, List(BitArray))), 510 - joined: List(Frame), 511 - context: Option(Context), 512 - ) -> Result(List(Frame), Nil) { 513 - case frames, accumulated { 514 - // No more frames - we are done 515 - [], _ -> Ok(list.reverse(joined)) 516 - 517 - // Complete standalone frame 518 - [Complete(Data(CompressedTextFrame(data))), ..rest], None -> { 519 - case context { 520 - Some(ctx) -> { 521 - let decompressed = compression.inflate(ctx, data) 522 - case bit_array.is_utf8(decompressed) { 523 - True -> { 524 - let final_frame = Data(TextFrame(decompressed)) 525 - aggregate_frames(rest, None, [final_frame, ..joined], context) 526 - } 527 - False -> Error(Nil) 528 - } 529 - } 530 - None -> Error(Nil) 531 - } 532 - } 533 - [Complete(Data(CompressedBinaryFrame(data))), ..rest], None -> { 534 - case context { 535 - Some(ctx) -> { 536 - let decompressed = compression.inflate(ctx, data) 537 - let final_frame = Data(BinaryFrame(decompressed)) 538 - aggregate_frames(rest, None, [final_frame, ..joined], context) 539 - } 540 - None -> Error(Nil) 541 - } 542 - } 543 - [Complete(Data(TextFrame(data))), ..rest], None -> { 544 - case bit_array.is_utf8(data) { 545 - True -> 546 - aggregate_frames( 547 - rest, 548 - None, 549 - [Data(TextFrame(data)), ..joined], 550 - context, 551 - ) 552 - False -> Error(Nil) 553 - } 554 - } 555 - [Complete(Data(BinaryFrame(data))), ..rest], None -> 556 - aggregate_frames(rest, None, [Data(BinaryFrame(data)), ..joined], context) 557 - [Complete(Continuation(..)), ..], None -> Error(Nil) 558 - [Complete(frame), ..rest], None -> 559 - aggregate_frames(rest, None, [frame, ..joined], context) 560 - 561 - // Incomplete frame starting fragmentation 562 - [Incomplete(frame), ..rest], None -> { 563 - let initial_payload = case frame { 564 - Data(TextFrame(payload)) -> payload 565 - Data(BinaryFrame(payload)) -> payload 566 - Data(CompressedTextFrame(payload)) -> payload 567 - Data(CompressedBinaryFrame(payload)) -> payload 568 - Continuation(_, payload) -> payload 569 - Control(_) -> <<>> 570 - } 571 - aggregate_frames(rest, Some(#(frame, [initial_payload])), joined, context) 572 - } 573 - 574 - // Complete continuation; finish fragmented message 575 - [Complete(Continuation(payload: data, ..)), ..rest], 576 - Some(#(initial_frame, payloads)) 577 - -> { 578 - let all_payloads = [data, ..payloads] |> list.reverse() 579 - let final_payload = bit_array.concat(all_payloads) 580 - 581 - case initial_frame { 582 - Data(CompressedTextFrame(_)) -> { 583 - case context { 584 - Some(ctx) -> { 585 - let decompressed = compression.inflate(ctx, final_payload) 586 - case bit_array.is_utf8(decompressed) { 587 - True -> { 588 - let final_frame = Data(TextFrame(decompressed)) 589 - aggregate_frames(rest, None, [final_frame, ..joined], context) 590 - } 591 - False -> Error(Nil) 592 - } 593 - } 594 - None -> Error(Nil) 595 - } 596 - } 597 - Data(CompressedBinaryFrame(_)) -> { 598 - case context { 599 - Some(ctx) -> { 600 - let decompressed = compression.inflate(ctx, final_payload) 601 - let final_frame = Data(BinaryFrame(decompressed)) 602 - aggregate_frames(rest, None, [final_frame, ..joined], context) 603 - } 604 - None -> Error(Nil) 605 - } 606 - } 607 - Data(TextFrame(_)) -> { 608 - case bit_array.is_utf8(final_payload) { 609 - True -> { 610 - let final_frame = Data(TextFrame(final_payload)) 611 - aggregate_frames(rest, None, [final_frame, ..joined], context) 612 - } 613 - False -> Error(Nil) 614 - } 615 - } 616 - Data(BinaryFrame(_)) -> { 617 - let final_frame = Data(BinaryFrame(final_payload)) 618 - aggregate_frames(rest, None, [final_frame, ..joined], context) 619 - } 620 - Control(_) -> Error(Nil) 621 - Continuation(..) -> Error(Nil) 622 - } 623 - } 624 - 625 - // Incomplete continuation; keep building the message 626 - [Incomplete(Continuation(payload: data, ..)), ..rest], 627 - Some(#(initial_frame, payloads)) 628 - -> { 629 - aggregate_frames( 630 - rest, 631 - Some(#(initial_frame, [data, ..payloads])), 632 - joined, 633 - context, 634 - ) 635 - } 636 - 637 - _, _ -> Error(Nil) 638 - } 639 - } 640 - 641 - const websocket_key = "258EAFA5-E914-47DA-95CA-C5AB0DC85B11" 642 - 643 - pub fn make_client_key() -> String { 644 - let bytes = crypto.strong_random_bytes(16) 645 - bit_array.base64_encode(bytes, True) 646 - } 647 - 648 - type ShaHash { 649 - Sha 650 - } 651 - 652 - pub fn parse_websocket_key(key: String) -> String { 653 - key 654 - |> string.append(websocket_key) 655 - |> crypto_hash(Sha, _) 656 - |> base64_encode 657 - } 658 - 659 - @external(erlang, "crypto", "hash") 660 - fn crypto_hash(hash hash: ShaHash, data data: String) -> String 661 - 662 - @external(erlang, "base64", "encode") 663 - fn base64_encode(data data: String) -> String 664 - 665 - pub fn has_deflate(extensions: List(String)) -> Bool { 666 - list.any(extensions, fn(str) { str == "permessage-deflate" }) 667 - } 668 - 669 - pub fn get_context_takeovers(extensions: List(String)) -> ContextTakeover { 670 - let no_client_context_takeover = 671 - list.any(extensions, fn(str) { str == "client_no_context_takeover" }) 672 - let no_server_context_takeover = 673 - list.any(extensions, fn(str) { str == "server_no_context_takeover" }) 674 - ContextTakeover( 675 - no_client: no_client_context_takeover, 676 - no_server: no_server_context_takeover, 677 - ) 678 - }
-124
src/ewe/internal/gramps/websocket/compression.gleam
··· 1 - // TODO: remove this once gramps changes are published 2 - // See https://github.com/rawhat/gramps 3 - 4 - import gleam/bit_array 5 - import gleam/bytes_tree.{type BytesTree} 6 - import gleam/erlang/atom.{type Atom} 7 - import gleam/erlang/process.{type Pid} 8 - 9 - pub type CompressionContext 10 - 11 - pub type Context { 12 - Context(context: CompressionContext, no_takeover: Bool) 13 - } 14 - 15 - type Flush { 16 - Sync 17 - } 18 - 19 - type Deflated { 20 - Deflated 21 - } 22 - 23 - type Default { 24 - Default 25 - } 26 - 27 - pub type ContextTakeover { 28 - ContextTakeover(no_client: Bool, no_server: Bool) 29 - } 30 - 31 - pub type Compression { 32 - Compression(inflate: Context, deflate: Context) 33 - } 34 - 35 - pub fn init(takeover: ContextTakeover) -> Compression { 36 - let inflate = open() 37 - let inflate_context = 38 - Context(context: inflate, no_takeover: takeover.no_client) 39 - 40 - inflate_init(inflate, -15) 41 - let deflate = open() 42 - let deflate_context = 43 - Context(context: deflate, no_takeover: takeover.no_server) 44 - deflate_init(deflate, Default, Deflated, -15, 8, Default) 45 - 46 - Compression(inflate: inflate_context, deflate: deflate_context) 47 - } 48 - 49 - @external(erlang, "zlib", "inflateInit") 50 - fn inflate_init(context: CompressionContext, bits: Int) -> Atom 51 - 52 - @external(erlang, "zlib", "deflateInit") 53 - fn deflate_init( 54 - context: CompressionContext, 55 - level: Default, 56 - deflated: Deflated, 57 - bits: Int, 58 - mem_level: Int, 59 - strategy: Default, 60 - ) -> Atom 61 - 62 - @external(erlang, "zlib", "open") 63 - fn open() -> CompressionContext 64 - 65 - @external(erlang, "zlib", "inflate") 66 - fn do_inflate(context: CompressionContext, data: BitArray) -> BytesTree 67 - 68 - pub fn inflate(context: Context, data: BitArray) -> BitArray { 69 - let output = 70 - context.context 71 - |> do_inflate(<<data:bits, 0x00, 0x00, 0xFF, 0xFF>>) 72 - |> bytes_tree.to_bit_array 73 - 74 - let _ = case context.no_takeover { 75 - True -> inflate_reset(context.context) 76 - False -> Nil 77 - } 78 - 79 - output 80 - } 81 - 82 - @external(erlang, "zlib", "deflate") 83 - fn do_deflate( 84 - context: CompressionContext, 85 - data: BitArray, 86 - flush: Flush, 87 - ) -> BytesTree 88 - 89 - pub fn deflate(context: Context, data: BitArray) -> BitArray { 90 - let data = 91 - context.context 92 - |> do_deflate(data, Sync) 93 - |> bytes_tree.to_bit_array 94 - 95 - let size = bit_array.byte_size(data) - 4 96 - 97 - let return = case data { 98 - <<value:bytes-size(size), 0x00, 0x00, 0xFF, 0xFF>> -> value 99 - _ -> data 100 - } 101 - 102 - let _ = case context.no_takeover { 103 - True -> deflate_reset(context.context) 104 - False -> Nil 105 - } 106 - 107 - return 108 - } 109 - 110 - @external(erlang, "zlib", "set_controlling_process") 111 - pub fn set_controlling_process(context: Context, pid: Pid) -> Atom 112 - 113 - pub fn close(context: Context) -> Nil { 114 - do_close(context.context) 115 - } 116 - 117 - @external(erlang, "zlib", "close") 118 - fn do_close(context: CompressionContext) -> Nil 119 - 120 - @external(erlang, "zlib", "inflateReset") 121 - fn inflate_reset(context: CompressionContext) -> Nil 122 - 123 - @external(erlang, "zlib", "deflateReset") 124 - fn deflate_reset(context: CompressionContext) -> Nil
+4 -5
src/ewe/internal/http1.gleam
··· 21 21 import gleam/string_tree.{type StringTree} 22 22 import gleam/uri 23 23 24 + import websocks 25 + 24 26 import glisten 25 27 import glisten/socket.{type Socket} 26 28 import glisten/transport.{type Transport} 27 - 28 - // TODO: replace this once gramps changes are published 29 - import ewe/internal/gramps/websocket as ws 30 29 31 30 import ewe/internal/buffer.{type Buffer} 32 31 import ewe/internal/clock ··· 404 403 |> result.replace_error(MissingWebsocketKey), 405 404 ) 406 405 407 - let accept_key = ws.parse_websocket_key(key) 406 + let accept_key = websocks.compute_accept(key) 408 407 409 408 let extensions = 410 409 request.get_header(req, "sec-websocket-extensions") 411 410 |> result.map(string.split(_, ";")) 412 411 |> result.unwrap([]) 413 412 414 - let permessage_deflate = ws.has_deflate(extensions) 413 + let permessage_deflate = websocks.has_deflate(extensions) 415 414 416 415 let resp = 417 416 response.new(101)
+132 -256
src/ewe/internal/stream/websocket.gleam
··· 2 2 // IMPORTS 3 3 // ----------------------------------------------------------------------------- 4 4 import gleam/bit_array 5 - import gleam/bytes_tree.{type BytesTree} 5 + import gleam/bytes_tree 6 6 import gleam/dynamic/decode 7 7 import gleam/erlang/atom 8 8 import gleam/erlang/process.{type Selector, type Subject} 9 - import gleam/list 10 9 import gleam/option.{type Option, None, Some} 11 10 import gleam/otp/actor 12 11 import gleam/result 13 12 import gleam/string 14 13 import logging 15 14 15 + import websocks 16 + 16 17 import glisten/socket.{type Socket, type SocketReason} 17 18 import glisten/socket/options.{ActiveMode, Count} 18 19 import glisten/transport.{type Transport} 19 20 20 - // TODO: replace this once gramps changes are published 21 - import ewe/internal/gramps/websocket.{type Frame, CloseFrame, PingFrame} 22 - import ewe/internal/gramps/websocket/compression 23 - 24 21 import ewe/internal/exception 25 22 26 23 // ----------------------------------------------------------------------------- ··· 32 29 WebsocketConnection( 33 30 transport: Transport, 34 31 socket: Socket, 35 - deflate: Option(compression.Context), 32 + context: websocks.Context, 36 33 ) 37 34 } 38 35 39 36 // Messages that can be sent to or received from the WebSocket 40 37 pub type WebsocketMessage(user_message) { 41 - WebsocketFrame(Frame) 38 + Frame(websocks.Frame) 42 39 UserMessage(user_message) 43 40 } 44 41 45 42 // Control flow for WebSocket message handling 46 43 pub type WebsocketNext(user_state, user_message) { 47 - Continue(user_state, Option(Selector(user_message))) 44 + Continue(user_state: user_state, selector: Option(Selector(user_message))) 48 45 NormalStop 49 46 AbnormalStop(reason: String) 50 47 } ··· 55 52 56 53 // Internal state maintained by the WebSocket actor 57 54 type WebsocketState(user_state) { 58 - WebsocketState( 59 - user_state: user_state, 60 - per_message_deflate: Option(compression.Compression), 61 - buffer: BitArray, 62 - awaiting_frames: List(websocket.ParsedFrame), 63 - ) 55 + WebsocketState(user_state: user_state, context: websocks.Context) 64 56 } 65 57 66 58 // Type alias for actor next steps ··· 110 102 // COMPRESSION UTILITIES 111 103 // ----------------------------------------------------------------------------- 112 104 113 - /// Gets the deflate context from the compression option 114 - fn get_deflate( 115 - compression: Option(compression.Compression), 116 - ) -> Option(compression.Context) { 117 - option.map(compression, fn(compression) { compression.deflate }) 118 - } 105 + // /// Gets the deflate context from the compression option 106 + // fn get_deflate( 107 + // compression: Option(compression.Compression), 108 + // ) -> Option(compression.Context) { 109 + // option.map(compression, fn(compression) { compression.deflate }) 110 + // } 119 111 120 - /// Gets the inflate context from the compression option 121 - fn get_inflate( 122 - compression: Option(compression.Compression), 123 - ) -> Option(compression.Context) { 124 - option.map(compression, fn(compression) { compression.inflate }) 125 - } 112 + // /// Gets the inflate context from the compression option 113 + // fn get_inflate( 114 + // compression: Option(compression.Compression), 115 + // ) -> Option(compression.Context) { 116 + // option.map(compression, fn(compression) { compression.inflate }) 117 + // } 126 118 127 119 // ----------------------------------------------------------------------------- 128 120 // SELECTOR UTILITIES ··· 167 159 // SOCKET UTILITIES 168 160 // ----------------------------------------------------------------------------- 169 161 170 - /// Sets socket to active mode for one message delivery 171 - // fn set_socket_active_once(transport: Transport, socket: Socket) -> Nil { 172 - // // echo transport.get_socket_opts(transport, socket, [atom.create("active")]) 173 - 174 - // let _ = transport.set_opts(transport, socket, [ActiveMode(Once)]) 175 - 176 - // // echo transport.get_socket_opts(transport, socket, [atom.create("active")]) 177 - 178 - // Nil 179 - // } 180 - 181 162 const socket_active_count = 100 182 163 183 - // fn set_socket_active_smart( 184 - // transport: Transport, 185 - // socket: Socket, 186 - // count: Int, 187 - // ) -> Int { 188 - // echo #( 189 - // count, 190 - // transport.get_socket_opts(transport, socket, [atom.create("active")]), 191 - // ) 192 - 193 - // case count { 194 - // 0 -> { 195 - // let _ = 196 - // transport.set_opts(transport, socket, [ 197 - // ActiveMode(Count(socket_active_count)), 198 - // ]) 199 - // socket_active_count 200 - // } 201 - // _ -> count - 1 202 - // } 203 - // } 204 - 205 164 // ----------------------------------------------------------------------------- 206 165 // PUBLIC API 207 166 // ----------------------------------------------------------------------------- ··· 217 176 permessage_deflate: Bool, 218 177 ) -> Result(actor.Started(Nil), actor.StartError) { 219 178 actor.new_with_initialiser(1000, fn(subject) { 220 - let takeovers = websocket.get_context_takeovers(extensions) 221 - let deflate = case permessage_deflate { 222 - True -> Some(compression.init(takeovers)) 179 + let context_takeovers = websocks.get_context_takeovers(extensions) 180 + let compression = case permessage_deflate { 181 + True -> Some(context_takeovers) 223 182 False -> None 224 183 } 225 - 226 - let conn = WebsocketConnection(transport, socket, get_deflate(deflate)) 184 + let context = websocks.create_context(compression) 227 185 228 - let #(user_state, user_selector) = on_init(conn, process.new_selector()) 186 + let #(user_state, user_selector) = 187 + WebsocketConnection(transport, socket, context) 188 + |> on_init(process.new_selector()) 229 189 230 190 let selector = 231 191 process.map_selector(user_selector, User) 232 192 |> process.merge_selector(create_socket_selector()) 233 193 234 - let ws_state = 235 - WebsocketState( 236 - user_state:, 237 - per_message_deflate: deflate, 238 - buffer: <<>>, 239 - awaiting_frames: [], 240 - ) 241 - 242 - actor.initialised(ws_state) 194 + WebsocketState(user_state:, context:) 195 + |> actor.initialised() 243 196 |> actor.selecting(selector) 244 197 |> actor.returning(subject) 245 198 |> Ok 246 199 }) 247 200 |> actor.on_message(fn(state, msg) { 248 - let conn = 249 - WebsocketConnection( 250 - transport, 251 - socket, 252 - get_deflate(state.per_message_deflate), 253 - ) 254 - 255 201 case msg { 256 - Packet(data) -> handle_valid_packet(state, conn, data, handler, on_close) 257 - Close -> handle_close(on_close, state, conn, None) 202 + Packet(data) -> 203 + handle_valid_packet(transport, socket, state, data, handler, on_close) 258 204 User(user_message) -> 259 - handle_user_message(state, conn, user_message, handler, on_close) 260 - Invalid -> handle_close(on_close, state, conn, Some(malformed)) 205 + handle_user_message( 206 + transport, 207 + socket, 208 + state, 209 + user_message, 210 + handler, 211 + on_close, 212 + ) 213 + Close -> { 214 + let conn = WebsocketConnection(transport, socket, state.context) 215 + handle_close(on_close, state, conn, None) 216 + } 217 + Invalid -> { 218 + let conn = WebsocketConnection(transport, socket, state.context) 219 + handle_close(on_close, state, conn, Some(malformed)) 220 + } 261 221 TcpPassive -> { 262 222 let _ = 263 223 transport.set_opts(transport, socket, [ ··· 273 233 274 234 /// Sends a frame to the WebSocket 275 235 pub fn send_frame( 276 - encoder: fn(data, Option(compression.Context), Option(BitArray)) -> BytesTree, 236 + encoder: fn(BitArray, websocks.Context, Option(BitArray)) -> BitArray, 277 237 transport: Transport, 278 238 socket: Socket, 279 - deflate: Option(compression.Context), 280 - data: data, 239 + context: websocks.Context, 240 + payload: BitArray, 281 241 ) -> Result(Nil, SocketReason) { 282 242 let frame = 283 243 exception.rescue(fn() { 284 - encoder(data, deflate, option.None) 244 + encoder(payload, context, option.None) 245 + |> bytes_tree.from_bit_array() 285 246 |> transport.send(transport, socket, _) 286 247 }) 287 248 ··· 304 265 305 266 /// Handles incoming packet data, decoding frames and processing them 306 267 fn handle_valid_packet( 268 + transport: Transport, 269 + socket: Socket, 307 270 state: WebsocketState(user_state), 308 - conn: WebsocketConnection, 309 271 data: BitArray, 310 272 handler: Handler(user_state, user_message), 311 273 on_close: OnClose(user_state), 312 274 ) -> ActorNext(user_state, user_message) { 313 - let buffer = <<state.buffer:bits, data:bits>> 314 - 315 - let decoded = 316 - websocket.decode_many_frames_result( 317 - buffer, 318 - get_inflate(state.per_message_deflate), 319 - [], 275 + let conn = WebsocketConnection(transport, socket, state.context) 276 + let result = 277 + websocks.process_incoming_frames( 278 + data, 279 + state.context, 280 + ResolveState( 281 + socket:, 282 + transport:, 283 + handler:, 284 + next: Continue(state.user_state, None), 285 + ), 286 + handle_frame, 320 287 ) 321 288 322 - case decoded { 323 - Ok(#(frames, rest)) -> 324 - handle_frames_processing(state, conn, frames, rest, handler, on_close) 325 - Error(websocket.NeedMoreDataAccumulated(parsed, rest)) -> { 326 - actor.continue( 327 - WebsocketState( 328 - ..state, 329 - buffer: rest, 330 - // NOTE: idk if its correct 331 - awaiting_frames: list.append(state.awaiting_frames, parsed), 332 - ), 333 - ) 334 - } 335 - Error(websocket.ContainsInvalidFrame) -> { 336 - handle_close(on_close, state, conn, Some(malformed)) 337 - } 338 - } 339 - } 340 - 341 - /// Handles frames processing 342 - fn handle_frames_processing( 343 - state: WebsocketState(user_state), 344 - conn: WebsocketConnection, 345 - frames: List(websocket.ParsedFrame), 346 - rest: BitArray, 347 - handler: Handler(user_state, user_message), 348 - on_close: OnClose(user_state), 349 - ) { 350 - let frames = list.append(state.awaiting_frames, frames) 351 - 352 - let #(data_frames, control_frames) = separate_frames(frames, [], []) 289 + case result { 290 + Ok(#(resolved_state, context)) -> { 291 + case resolved_state.next { 292 + Continue(user_state, selector) -> { 293 + let next = actor.continue(WebsocketState(user_state:, context:)) 353 294 354 - let control_result = case control_frames { 355 - [] -> Continue(state.user_state, None) 356 - _ -> 357 - loop_by_frames( 358 - control_frames, 359 - conn, 360 - handler, 361 - Continue(state.user_state, None), 362 - ) 363 - } 364 - 365 - case control_result { 366 - NormalStop -> handle_close(on_close, state, conn, None) 367 - AbnormalStop(reason) -> handle_close(on_close, state, conn, Some(reason)) 368 - Continue(_, _) -> { 369 - let aggregated = 370 - websocket.aggregate_frames( 371 - data_frames, 372 - None, 373 - [], 374 - get_inflate(state.per_message_deflate), 375 - ) 376 - 377 - case aggregated { 378 - Ok([]) -> { 379 - actor.continue( 380 - WebsocketState(..state, buffer: rest, awaiting_frames: data_frames), 381 - ) 382 - } 383 - Ok(data_frames) -> { 384 - let next = 385 - loop_by_frames( 386 - data_frames, 387 - conn, 388 - handler, 389 - Continue(state.user_state, None), 390 - ) 391 - 392 - case next { 393 - Continue(user_state, selector) -> { 394 - let next = 395 - actor.continue( 396 - WebsocketState( 397 - ..state, 398 - user_state:, 399 - buffer: rest, 400 - awaiting_frames: [], 401 - ), 402 - ) 403 - 404 - case selector { 405 - Some(selector) -> actor.with_selector(next, selector) 406 - None -> next 407 - } 408 - } 409 - NormalStop -> handle_close(on_close, state, conn, None) 410 - AbnormalStop(reason) -> 411 - handle_close(on_close, state, conn, Some(reason)) 295 + case selector { 296 + Some(selector) -> actor.with_selector(next, selector) 297 + None -> next 412 298 } 413 299 } 414 - Error(Nil) -> handle_close(on_close, state, conn, Some(malformed)) 300 + NormalStop -> handle_close(on_close, state, conn, None) 301 + AbnormalStop(reason) -> 302 + handle_close(on_close, state, conn, Some(reason)) 415 303 } 416 304 } 305 + Error(violation) -> { 306 + echo violation as "violation during frame resolving" 307 + handle_close(on_close, state, conn, Some(malformed)) 308 + } 417 309 } 418 310 } 419 311 420 - /// Separates frames into data and control frames 421 - fn separate_frames( 422 - frames: List(websocket.ParsedFrame), 423 - data_frames: List(websocket.ParsedFrame), 424 - control_frames: List(websocket.Frame), 425 - ) -> #(List(websocket.ParsedFrame), List(websocket.Frame)) { 426 - case frames { 427 - [] -> #(list.reverse(data_frames), list.reverse(control_frames)) 428 - [websocket.Complete(websocket.Control(control_frame)), ..rest] -> 429 - separate_frames(rest, data_frames, [ 430 - websocket.Control(control_frame), 431 - ..control_frames 432 - ]) 433 - [data_frame, ..rest] -> 434 - separate_frames(rest, [data_frame, ..data_frames], control_frames) 435 - } 312 + type ResolveState(user_state, user_message) { 313 + ResolveState( 314 + socket: Socket, 315 + transport: Transport, 316 + handler: Handler(user_state, user_message), 317 + next: WebsocketNext(user_state, InternalMessage(user_message)), 318 + ) 436 319 } 437 320 438 321 /// Processes a list of frames sequentially 439 - fn loop_by_frames( 440 - frames: List(Frame), 441 - conn: WebsocketConnection, 442 - handler: Handler(user_state, user_message), 443 - next: WebsocketNext(user_state, InternalMessage(user_message)), 444 - ) -> WebsocketNext(user_state, InternalMessage(user_message)) { 445 - case frames, next { 446 - // Early termination cases 447 - _, NormalStop -> NormalStop 448 - _, AbnormalStop(reason) -> AbnormalStop(reason) 449 - 450 - // No more frames - finish 451 - [], next -> next 452 - 453 - // Control frames 454 - [websocket.Control(PingFrame(payload)), ..rest], Continue(user_state, _) -> { 322 + fn handle_frame( 323 + state: ResolveState(user_state, user_message), 324 + context: websocks.Context, 325 + frame: websocks.Frame, 326 + ) -> websocks.ResolveNext(ResolveState(user_state, user_message)) { 327 + case frame { 328 + websocks.Control(websocks.Ping(payload)) -> { 455 329 case bit_array.byte_size(payload) { 456 330 size if size > 125 -> 457 - AbnormalStop( 458 - "control frames are only allowed to have payload up to and including 125 octets", 331 + websocks.Stop( 332 + ResolveState( 333 + ..state, 334 + next: AbnormalStop( 335 + "control frames are only allowed to have payload up to and including 125 octets", 336 + ), 337 + ), 459 338 ) 460 339 _ -> { 461 340 let sent = 462 341 transport.send( 463 - conn.transport, 464 - conn.socket, 465 - websocket.encode_pong_frame(payload, None), 342 + state.transport, 343 + state.socket, 344 + websocks.encode_pong_frame(payload, None) 345 + |> bytes_tree.from_bit_array(), 466 346 ) 467 347 468 348 case sent { 469 - Ok(Nil) -> 470 - loop_by_frames(rest, conn, handler, Continue(user_state, None)) 471 - Error(_) -> AbnormalStop(failed_pong) 349 + Ok(Nil) -> websocks.Continue(state) 350 + Error(_) -> 351 + websocks.Stop( 352 + ResolveState(..state, next: AbnormalStop(failed_pong)), 353 + ) 472 354 } 473 355 } 474 356 } 475 357 } 476 - [websocket.Control(CloseFrame(reason)), ..], Continue(_, _) -> { 358 + websocks.Control(websocks.Close(reason)) -> { 477 359 let _ = 478 360 transport.send( 479 - conn.transport, 480 - conn.socket, 481 - websocket.encode_close_frame(reason, None), 361 + state.transport, 362 + state.socket, 363 + websocks.encode_close_frame(reason, None) 364 + |> bytes_tree.from_bit_array(), 482 365 ) 483 366 484 - NormalStop 367 + websocks.Stop(ResolveState(..state, next: NormalStop)) 485 368 } 486 369 487 - // NOTE: unsure if its should be here 488 - [websocket.Continuation(_, _), ..], Continue(_, _) -> { 489 - AbnormalStop("Unexpected continuation frame") 490 - } 370 + frame -> { 371 + let assert Continue(user_state, selector) = state.next 491 372 492 - // Data frames 493 - [frame, ..rest], Continue(user_state, selector) -> { 373 + let conn = WebsocketConnection(state.transport, state.socket, context) 374 + 494 375 let call = 495 - exception.rescue(fn() { 496 - handler(conn, user_state, WebsocketFrame(frame)) 497 - }) 376 + exception.rescue(fn() { state.handler(conn, user_state, Frame(frame)) }) 498 377 499 378 case call { 500 379 Ok(Continue(user_state, new_selector)) -> { ··· 503 382 |> option.or(selector) 504 383 |> option.map(process.merge_selector(create_socket_selector(), _)) 505 384 506 - loop_by_frames( 507 - rest, 508 - conn, 509 - handler, 510 - Continue(user_state, next_selector), 385 + websocks.Continue( 386 + ResolveState(..state, next: Continue(user_state, next_selector)), 511 387 ) 512 388 } 513 - Ok(NormalStop) -> NormalStop 514 - Ok(AbnormalStop(reason)) -> AbnormalStop(reason) 515 - Error(_) -> AbnormalStop(crashed) 389 + Ok(NormalStop) -> websocks.Stop(ResolveState(..state, next: NormalStop)) 390 + Ok(AbnormalStop(reason)) -> 391 + websocks.Stop(ResolveState(..state, next: AbnormalStop(reason))) 392 + Error(_) -> 393 + websocks.Stop(ResolveState(..state, next: AbnormalStop(crashed))) 516 394 } 517 395 } 518 396 } ··· 520 398 521 399 /// Handles user messages sent to the WebSocket 522 400 fn handle_user_message( 401 + transport: Transport, 402 + socket: Socket, 523 403 state: WebsocketState(user_state), 524 - conn: WebsocketConnection, 525 404 user_message: user_message, 526 405 handler: Handler(user_state, user_message), 527 406 on_close: OnClose(user_state), 528 407 ) -> ActorNext(user_state, user_message) { 408 + let conn = WebsocketConnection(transport, socket, state.context) 529 409 let call = 530 410 exception.rescue(fn() { 531 411 handler(conn, state.user_state, UserMessage(user_message)) ··· 563 443 conn: WebsocketConnection, 564 444 abnormal_reason: Option(String), 565 445 ) -> actor.Next(WebsocketState(user_state), InternalMessage(user_message)) { 566 - option.map(state.per_message_deflate, fn(compression) { 567 - compression.close(compression.deflate) 568 - compression.close(compression.inflate) 569 - }) 570 - 446 + websocks.close_context(state.context) 571 447 on_close(conn, state.user_state) 572 448 573 449 case abnormal_reason {