Skip to content

Repository files navigation

websocket

WebSocket library for fasthttp and net/http.

Checkout examples to inspire yourself.

Install

go get github.com/dgrr/websocket

Why another WebSocket package?

Other WebSocket packages DON'T allow concurrent Read/Write operations on servers and they do not provide low level access to WebSocket packet crafting. Those WebSocket packages try to emulate the Golang API by implementing io.Reader and io.Writer interfaces on their connections. io.Writer might be a good idea to use it, but no io.Reader, given that WebSocket is an async protocol by nature (all protocols are (?)).

Sometimes, WebSocket servers are just cumbersome when we want to handle a lot of clients in an async manner. For example, in other WebSocket packages to broadcast a message generated internally we'll need to do the following:

typeMyWebSocketServicestruct {
clients sync.Map
}
typeBlockingConnstruct {
lck sync.Mutexc websocketPackage.Conn
}
func (ws*MyWebSocketService) Broadcaster() {
formsg:=rangemessageProducerChannel {
ws.clients.Range(func(_, vinterface{}) bool {
c:=v.(*BlockingConn)
c.lck.Lock() // oh, we need to block, otherwise we can break the programerr:=c.Write(msg)
c.lck.Unlock()
iferr!=nil {
// we have an error, what can we do? Log it?// if the connection has been closed we'll receive that on// the Read call, so the connection will close automatically.
}
returntrue
})
}
}
func (ws*MyWebSocketService) Handle(request, response) {
c, err:=websocketPackage.Upgrade(request, response)
iferr!=nil {
// then it's clearly an error! Report back
}
bc:=&BlockingConn{
c: c,
} ws.clients.Store(bc, struct{}{})
// even though I just want to write, I need to block somehowfor {
content, err:=bc.Read()
iferr!=nil {
// handle the errorbreak
}
}
ws.clients.Delete(bc)
}

First, we need to store every client upon connection, and whenever we want to send data we need to iterate over a list, and send the message. If while, writing we get an error, then we need to handle that client's error What if the writing operation is happening at the same time in 2 different coroutines? Then we need a sync.Mutex and block until we finish writing.

To solve most of those problems websocket uses channels and separated coroutines, one for reading and another one for writing. By following the sharing principle.

Do not communicate by sharing memory; instead, share memory by communicating.

Following the fasthttp philosophy this library tries to take as much advantage of the Golang's multi-threaded model as possible, while keeping your code concurrently safe.

To see an example of what this package CAN do that others DONT checkout the broadcast example.

Server

How can I launch a server?

It's quite easy. You only need to create a Server, set your callbacks by calling the Handle* methods and then specify your fasthttp handler as Server.Upgrade.

package main
import (
"fmt""github.com/valyala/fasthttp""github.com/dgrr/websocket"
)
funcmain() {
ws:= websocket.Server{}
ws.HandleData(OnMessage)
fasthttp.ListenAndServe(":8080", ws.Upgrade)
}
funcOnMessage(c*websocket.Conn, isBinarybool, data []byte) {
fmt.Printf("Received data from %s: %s\n", c.RemoteAddr(), data)
}

How can I launch a server if I use net/http?

package main
import (
"fmt""net/http""github.com/dgrr/websocket"
)
funcmain() {
ws:= websocket.Server{}
ws.HandleData(OnMessage)
http.HandleFunc("/", ws.NetUpgrade)
http.ListenAndServe(":8080", nil)
}
funcOnMessage(c*websocket.Conn, isBinarybool, data []byte) {
fmt.Printf("Received data from %s: %s\n", c.RemoteAddr(), data)
}

How can I handle pings?

Pings are handle automatically by the library, but you can get the content of those pings setting the callback using HandlePing.

For example, let's try to get the round trip time to a client by using the PING frame. The website http2.gofiber.io uses this method to measure the round trip time displayed at the bottom of the webpage.

package main
import (
"sync""encoding/binary""log""time""github.com/valyala/fasthttp""github.com/dgrr/websocket"
)
// Struct to keep the clients connected//// it should be safe to access the clients concurrently from Open and Close.typeRTTMeasurestruct {
clients sync.Map
}
// just trigger the ping senderfunc (rtt*RTTMeasure) Start() {
time.AfterFunc(time.Second*2, rtt.sendPings)
}
func (rtt*RTTMeasure) sendPings() {
vardata [8]bytebinary.BigEndian.PutUint64(data[:], uint64(
time.Now().UnixNano()),
)
rtt.clients.Range(func(_, vinterface{}) bool {
c:=v.(*websocket.Conn)
c.Ping(data[:])
returntrue
})
rtt.Start()
}
// register a connection when it's openfunc (rtt*RTTMeasure) RegisterConn(c*websocket.Conn) {
rtt.clients.Store(c.ID(), c)
log.Printf("Client %s connected\n", c.RemoteAddr())
}
// remove the connection when receiving the closefunc (rtt*RTTMeasure) RemoveConn(c*websocket.Conn, errerror) {
rtt.clients.Delete(c.ID())
log.Printf("Client %s disconnected\n", c.RemoteAddr())
}
funcmain() {
rtt:=RTTMeasure{}
ws:= websocket.Server{}
ws.HandleOpen(rtt.RegisterConn)
ws.HandleClose(rtt.RemoveConn)
ws.HandlePong(OnPong)
// schedule the timerrtt.Start()
fasthttp.ListenAndServe(":8080", ws.Upgrade)
}
// handle the pong messagefuncOnPong(c*websocket.Conn, data []byte) {
iflen(data) ==8 {
n:=binary.BigEndian.Uint64(data)
ts:=time.Unix(0, int64(n))
log.Printf("RTT with %s is %s\n", c.RemoteAddr(), time.Now().Sub(ts))
}
}

websocket vs gorilla vs nhooyr vs gobwas

FeatureswebsocketGorillaNhooyrgowabs
Concurrent R/WYesNoNo. Only writesNo
Passes Autobahn Test SuiteMostlyYesYesMostly
Receive fragmented messageYesYesYesYes
Send close messageYesYesYesYes
Send pings and receive pongsYesYesYesYes
Get the type of a received data messageYesYesYesYes
Compression ExtensionsNoExperimentalYesNo (?)
Read message using io.ReaderNoYesNoNo (?)
Write message using io.WriteCloserYesYesNoNo (?)

Stress tests

The following stress test were performed without timeouts:

Executing tcpkali --ws -c 100 -m 'hello world!!13212312!' -r 10k localhost:8081 the tests shows the following:

Websocket:

Total data sent: 267.7 MiB (280678466 bytes)
Total data received: 229.5 MiB (240626600 bytes)
Bandwidth per channel: 4.167⇅ Mbps (520.9 kBps)
Aggregate bandwidth: 192.357↓, 224.375↑ Mbps
Packet rate estimate: 247050.1↓, 61842.9↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0075 s.

Websocket for net/http:

Total data sent: 267.3 MiB (280320124 bytes)
Total data received: 228.3 MiB (239396374 bytes)
Bandwidth per channel: 4.156⇅ Mbps (519.5 kBps)
Aggregate bandwidth: 191.442↓, 224.168↑ Mbps
Packet rate estimate: 188107.1↓, 52240.7↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0039 s.

Either for fasthttp and net/http should be quite close, the only difference is the way they both upgrade.

Gorilla:

Total data sent: 260.2 MiB (272886768 bytes)
Total data received: 109.3 MiB (114632982 bytes)
Bandwidth per channel: 3.097⇅ Mbps (387.1 kBps)
Aggregate bandwidth: 91.615↓, 218.092↑ Mbps
Packet rate estimate: 109755.3↓, 66807.4↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.01 s.

Nhooyr: (Don't know why is that low)

Total data sent: 224.3 MiB (235184096 bytes)
Total data received: 41.2 MiB (43209780 bytes)
Bandwidth per channel: 2.227⇅ Mbps (278.3 kBps)
Aggregate bandwidth: 34.559↓, 188.097↑ Mbps
Packet rate estimate: 88474.0↓, 55256.1↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0027 s.

Gobwas:

Total data sent: 265.8 MiB (278718160 bytes)
Total data received: 117.8 MiB (123548959 bytes)
Bandwidth per channel: 3.218⇅ Mbps (402.2 kBps)
Aggregate bandwidth: 98.825↓, 222.942↑ Mbps
Packet rate estimate: 148231.6↓, 72106.1↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0015 s.

The source files are in this folder.

About

WebSocket for fasthttp

Topics

Resources

Stars

67 stars

Watchers

2 watching

Forks

Releases

Packages

Used by

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Add copy buttons to all
 blocks
(function() {
function addCopyButtons() {
document.querySelectorAll('pre code').forEach(function(codeBlock) {
if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;
codeBlock.parentElement.setAttribute('data-copy-added', 'true');
var btn = document.createElement('button');
btn.textContent = 'Copy';
btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';
btn.onmouseover = function() { this.style.opacity = '1'; };
btn.onmouseout = function() { this.style.opacity = '0.7'; };
btn.onclick = function() {
navigator.clipboard.writeText(codeBlock.textContent).then(function() {
btn.textContent = 'Copied!';
setTimeout(function() { btn.textContent = 'Copy'; }, 1500);
});
};
codeBlock.parentElement.style.position = 'relative';
codeBlock.parentElement.appendChild(btn);
});
}
addCopyButtons();
// Re-run on dynamic content
var observer = new MutationObserver(addCopyButtons);
observer.observe(document.body, { childList: true, subtree: true });
})();
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
GitHub - dgrr/websocket: WebSocket for fasthttp · GitHub
Skip to content

Repository files navigation

websocket

WebSocket library for fasthttp and net/http.

Checkout examples to inspire yourself.

Install

go get github.com/dgrr/websocket

Why another WebSocket package?

Other WebSocket packages DON'T allow concurrent Read/Write operations on servers and they do not provide low level access to WebSocket packet crafting. Those WebSocket packages try to emulate the Golang API by implementing io.Reader and io.Writer interfaces on their connections. io.Writer might be a good idea to use it, but no io.Reader, given that WebSocket is an async protocol by nature (all protocols are (?)).

Sometimes, WebSocket servers are just cumbersome when we want to handle a lot of clients in an async manner. For example, in other WebSocket packages to broadcast a message generated internally we'll need to do the following:

typeMyWebSocketServicestruct {
clients sync.Map
}
typeBlockingConnstruct {
lck sync.Mutexc websocketPackage.Conn
}
func (ws*MyWebSocketService) Broadcaster() {
formsg:=rangemessageProducerChannel {
ws.clients.Range(func(_, vinterface{}) bool {
c:=v.(*BlockingConn)
c.lck.Lock() // oh, we need to block, otherwise we can break the programerr:=c.Write(msg)
c.lck.Unlock()
iferr!=nil {
// we have an error, what can we do? Log it?// if the connection has been closed we'll receive that on// the Read call, so the connection will close automatically.
}
returntrue
})
}
}
func (ws*MyWebSocketService) Handle(request, response) {
c, err:=websocketPackage.Upgrade(request, response)
iferr!=nil {
// then it's clearly an error! Report back
}
bc:=&BlockingConn{
c: c,
} ws.clients.Store(bc, struct{}{})
// even though I just want to write, I need to block somehowfor {
content, err:=bc.Read()
iferr!=nil {
// handle the errorbreak
}
}
ws.clients.Delete(bc)
}

First, we need to store every client upon connection, and whenever we want to send data we need to iterate over a list, and send the message. If while, writing we get an error, then we need to handle that client's error What if the writing operation is happening at the same time in 2 different coroutines? Then we need a sync.Mutex and block until we finish writing.

To solve most of those problems websocket uses channels and separated coroutines, one for reading and another one for writing. By following the sharing principle.

Do not communicate by sharing memory; instead, share memory by communicating.

Following the fasthttp philosophy this library tries to take as much advantage of the Golang's multi-threaded model as possible, while keeping your code concurrently safe.

To see an example of what this package CAN do that others DONT checkout the broadcast example.

Server

How can I launch a server?

It's quite easy. You only need to create a Server, set your callbacks by calling the Handle* methods and then specify your fasthttp handler as Server.Upgrade.

package main
import (
"fmt""github.com/valyala/fasthttp""github.com/dgrr/websocket"
)
funcmain() {
ws:= websocket.Server{}
ws.HandleData(OnMessage)
fasthttp.ListenAndServe(":8080", ws.Upgrade)
}
funcOnMessage(c*websocket.Conn, isBinarybool, data []byte) {
fmt.Printf("Received data from %s: %s\n", c.RemoteAddr(), data)
}

How can I launch a server if I use net/http?

package main
import (
"fmt""net/http""github.com/dgrr/websocket"
)
funcmain() {
ws:= websocket.Server{}
ws.HandleData(OnMessage)
http.HandleFunc("/", ws.NetUpgrade)
http.ListenAndServe(":8080", nil)
}
funcOnMessage(c*websocket.Conn, isBinarybool, data []byte) {
fmt.Printf("Received data from %s: %s\n", c.RemoteAddr(), data)
}

How can I handle pings?

Pings are handle automatically by the library, but you can get the content of those pings setting the callback using HandlePing.

For example, let's try to get the round trip time to a client by using the PING frame. The website http2.gofiber.io uses this method to measure the round trip time displayed at the bottom of the webpage.

package main
import (
"sync""encoding/binary""log""time""github.com/valyala/fasthttp""github.com/dgrr/websocket"
)
// Struct to keep the clients connected//// it should be safe to access the clients concurrently from Open and Close.typeRTTMeasurestruct {
clients sync.Map
}
// just trigger the ping senderfunc (rtt*RTTMeasure) Start() {
time.AfterFunc(time.Second*2, rtt.sendPings)
}
func (rtt*RTTMeasure) sendPings() {
vardata [8]bytebinary.BigEndian.PutUint64(data[:], uint64(
time.Now().UnixNano()),
)
rtt.clients.Range(func(_, vinterface{}) bool {
c:=v.(*websocket.Conn)
c.Ping(data[:])
returntrue
})
rtt.Start()
}
// register a connection when it's openfunc (rtt*RTTMeasure) RegisterConn(c*websocket.Conn) {
rtt.clients.Store(c.ID(), c)
log.Printf("Client %s connected\n", c.RemoteAddr())
}
// remove the connection when receiving the closefunc (rtt*RTTMeasure) RemoveConn(c*websocket.Conn, errerror) {
rtt.clients.Delete(c.ID())
log.Printf("Client %s disconnected\n", c.RemoteAddr())
}
funcmain() {
rtt:=RTTMeasure{}
ws:= websocket.Server{}
ws.HandleOpen(rtt.RegisterConn)
ws.HandleClose(rtt.RemoveConn)
ws.HandlePong(OnPong)
// schedule the timerrtt.Start()
fasthttp.ListenAndServe(":8080", ws.Upgrade)
}
// handle the pong messagefuncOnPong(c*websocket.Conn, data []byte) {
iflen(data) ==8 {
n:=binary.BigEndian.Uint64(data)
ts:=time.Unix(0, int64(n))
log.Printf("RTT with %s is %s\n", c.RemoteAddr(), time.Now().Sub(ts))
}
}

websocket vs gorilla vs nhooyr vs gobwas

FeatureswebsocketGorillaNhooyrgowabs
Concurrent R/WYesNoNo. Only writesNo
Passes Autobahn Test SuiteMostlyYesYesMostly
Receive fragmented messageYesYesYesYes
Send close messageYesYesYesYes
Send pings and receive pongsYesYesYesYes
Get the type of a received data messageYesYesYesYes
Compression ExtensionsNoExperimentalYesNo (?)
Read message using io.ReaderNoYesNoNo (?)
Write message using io.WriteCloserYesYesNoNo (?)

Stress tests

The following stress test were performed without timeouts:

Executing tcpkali --ws -c 100 -m 'hello world!!13212312!' -r 10k localhost:8081 the tests shows the following:

Websocket:

Total data sent: 267.7 MiB (280678466 bytes)
Total data received: 229.5 MiB (240626600 bytes)
Bandwidth per channel: 4.167⇅ Mbps (520.9 kBps)
Aggregate bandwidth: 192.357↓, 224.375↑ Mbps
Packet rate estimate: 247050.1↓, 61842.9↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0075 s.

Websocket for net/http:

Total data sent: 267.3 MiB (280320124 bytes)
Total data received: 228.3 MiB (239396374 bytes)
Bandwidth per channel: 4.156⇅ Mbps (519.5 kBps)
Aggregate bandwidth: 191.442↓, 224.168↑ Mbps
Packet rate estimate: 188107.1↓, 52240.7↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0039 s.

Either for fasthttp and net/http should be quite close, the only difference is the way they both upgrade.

Gorilla:

Total data sent: 260.2 MiB (272886768 bytes)
Total data received: 109.3 MiB (114632982 bytes)
Bandwidth per channel: 3.097⇅ Mbps (387.1 kBps)
Aggregate bandwidth: 91.615↓, 218.092↑ Mbps
Packet rate estimate: 109755.3↓, 66807.4↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.01 s.

Nhooyr: (Don't know why is that low)

Total data sent: 224.3 MiB (235184096 bytes)
Total data received: 41.2 MiB (43209780 bytes)
Bandwidth per channel: 2.227⇅ Mbps (278.3 kBps)
Aggregate bandwidth: 34.559↓, 188.097↑ Mbps
Packet rate estimate: 88474.0↓, 55256.1↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0027 s.

Gobwas:

Total data sent: 265.8 MiB (278718160 bytes)
Total data received: 117.8 MiB (123548959 bytes)
Bandwidth per channel: 3.218⇅ Mbps (402.2 kBps)
Aggregate bandwidth: 98.825↓, 222.942↑ Mbps
Packet rate estimate: 148231.6↓, 72106.1↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0015 s.

The source files are in this folder.

About

WebSocket for fasthttp

Topics

Resources

Stars

67 stars

Watchers

2 watching

Forks

Releases

Packages

Used by

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Force GitHub README to respect dark mode (function() { var style = document.createElement('style'); style.textContent = ' .markdown-body { color-scheme: dark light; } .markdown-body pre { background: #161b22 !important; } .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; } .markdown-body table th, .markdown-body table td { border-color: #30363d !important; } .markdown-body img { background: #0d1117; } .markdown-body blockquote { border-left-color: #8b949e; } .markdown-body hr { border-color: #30363d; } '; document.head.appendChild(style); })(); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' GitHub - dgrr/websocket: WebSocket for fasthttp · GitHub
Skip to content

Repository files navigation

websocket

WebSocket library for fasthttp and net/http.

Checkout examples to inspire yourself.

Install

go get github.com/dgrr/websocket

Why another WebSocket package?

Other WebSocket packages DON'T allow concurrent Read/Write operations on servers and they do not provide low level access to WebSocket packet crafting. Those WebSocket packages try to emulate the Golang API by implementing io.Reader and io.Writer interfaces on their connections. io.Writer might be a good idea to use it, but no io.Reader, given that WebSocket is an async protocol by nature (all protocols are (?)).

Sometimes, WebSocket servers are just cumbersome when we want to handle a lot of clients in an async manner. For example, in other WebSocket packages to broadcast a message generated internally we'll need to do the following:

typeMyWebSocketServicestruct {
clients sync.Map
}
typeBlockingConnstruct {
lck sync.Mutexc websocketPackage.Conn
}
func (ws*MyWebSocketService) Broadcaster() {
formsg:=rangemessageProducerChannel {
ws.clients.Range(func(_, vinterface{}) bool {
c:=v.(*BlockingConn)
c.lck.Lock() // oh, we need to block, otherwise we can break the programerr:=c.Write(msg)
c.lck.Unlock()
iferr!=nil {
// we have an error, what can we do? Log it?// if the connection has been closed we'll receive that on// the Read call, so the connection will close automatically.
}
returntrue
})
}
}
func (ws*MyWebSocketService) Handle(request, response) {
c, err:=websocketPackage.Upgrade(request, response)
iferr!=nil {
// then it's clearly an error! Report back
}
bc:=&BlockingConn{
c: c,
} ws.clients.Store(bc, struct{}{})
// even though I just want to write, I need to block somehowfor {
content, err:=bc.Read()
iferr!=nil {
// handle the errorbreak
}
}
ws.clients.Delete(bc)
}

First, we need to store every client upon connection, and whenever we want to send data we need to iterate over a list, and send the message. If while, writing we get an error, then we need to handle that client's error What if the writing operation is happening at the same time in 2 different coroutines? Then we need a sync.Mutex and block until we finish writing.

To solve most of those problems websocket uses channels and separated coroutines, one for reading and another one for writing. By following the sharing principle.

Do not communicate by sharing memory; instead, share memory by communicating.

Following the fasthttp philosophy this library tries to take as much advantage of the Golang's multi-threaded model as possible, while keeping your code concurrently safe.

To see an example of what this package CAN do that others DONT checkout the broadcast example.

Server

How can I launch a server?

It's quite easy. You only need to create a Server, set your callbacks by calling the Handle* methods and then specify your fasthttp handler as Server.Upgrade.

package main
import (
"fmt""github.com/valyala/fasthttp""github.com/dgrr/websocket"
)
funcmain() {
ws:= websocket.Server{}
ws.HandleData(OnMessage)
fasthttp.ListenAndServe(":8080", ws.Upgrade)
}
funcOnMessage(c*websocket.Conn, isBinarybool, data []byte) {
fmt.Printf("Received data from %s: %s\n", c.RemoteAddr(), data)
}

How can I launch a server if I use net/http?

package main
import (
"fmt""net/http""github.com/dgrr/websocket"
)
funcmain() {
ws:= websocket.Server{}
ws.HandleData(OnMessage)
http.HandleFunc("/", ws.NetUpgrade)
http.ListenAndServe(":8080", nil)
}
funcOnMessage(c*websocket.Conn, isBinarybool, data []byte) {
fmt.Printf("Received data from %s: %s\n", c.RemoteAddr(), data)
}

How can I handle pings?

Pings are handle automatically by the library, but you can get the content of those pings setting the callback using HandlePing.

For example, let's try to get the round trip time to a client by using the PING frame. The website http2.gofiber.io uses this method to measure the round trip time displayed at the bottom of the webpage.

package main
import (
"sync""encoding/binary""log""time""github.com/valyala/fasthttp""github.com/dgrr/websocket"
)
// Struct to keep the clients connected//// it should be safe to access the clients concurrently from Open and Close.typeRTTMeasurestruct {
clients sync.Map
}
// just trigger the ping senderfunc (rtt*RTTMeasure) Start() {
time.AfterFunc(time.Second*2, rtt.sendPings)
}
func (rtt*RTTMeasure) sendPings() {
vardata [8]bytebinary.BigEndian.PutUint64(data[:], uint64(
time.Now().UnixNano()),
)
rtt.clients.Range(func(_, vinterface{}) bool {
c:=v.(*websocket.Conn)
c.Ping(data[:])
returntrue
})
rtt.Start()
}
// register a connection when it's openfunc (rtt*RTTMeasure) RegisterConn(c*websocket.Conn) {
rtt.clients.Store(c.ID(), c)
log.Printf("Client %s connected\n", c.RemoteAddr())
}
// remove the connection when receiving the closefunc (rtt*RTTMeasure) RemoveConn(c*websocket.Conn, errerror) {
rtt.clients.Delete(c.ID())
log.Printf("Client %s disconnected\n", c.RemoteAddr())
}
funcmain() {
rtt:=RTTMeasure{}
ws:= websocket.Server{}
ws.HandleOpen(rtt.RegisterConn)
ws.HandleClose(rtt.RemoveConn)
ws.HandlePong(OnPong)
// schedule the timerrtt.Start()
fasthttp.ListenAndServe(":8080", ws.Upgrade)
}
// handle the pong messagefuncOnPong(c*websocket.Conn, data []byte) {
iflen(data) ==8 {
n:=binary.BigEndian.Uint64(data)
ts:=time.Unix(0, int64(n))
log.Printf("RTT with %s is %s\n", c.RemoteAddr(), time.Now().Sub(ts))
}
}

websocket vs gorilla vs nhooyr vs gobwas

FeatureswebsocketGorillaNhooyrgowabs
Concurrent R/WYesNoNo. Only writesNo
Passes Autobahn Test SuiteMostlyYesYesMostly
Receive fragmented messageYesYesYesYes
Send close messageYesYesYesYes
Send pings and receive pongsYesYesYesYes
Get the type of a received data messageYesYesYesYes
Compression ExtensionsNoExperimentalYesNo (?)
Read message using io.ReaderNoYesNoNo (?)
Write message using io.WriteCloserYesYesNoNo (?)

Stress tests

The following stress test were performed without timeouts:

Executing tcpkali --ws -c 100 -m 'hello world!!13212312!' -r 10k localhost:8081 the tests shows the following:

Websocket:

Total data sent: 267.7 MiB (280678466 bytes)
Total data received: 229.5 MiB (240626600 bytes)
Bandwidth per channel: 4.167⇅ Mbps (520.9 kBps)
Aggregate bandwidth: 192.357↓, 224.375↑ Mbps
Packet rate estimate: 247050.1↓, 61842.9↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0075 s.

Websocket for net/http:

Total data sent: 267.3 MiB (280320124 bytes)
Total data received: 228.3 MiB (239396374 bytes)
Bandwidth per channel: 4.156⇅ Mbps (519.5 kBps)
Aggregate bandwidth: 191.442↓, 224.168↑ Mbps
Packet rate estimate: 188107.1↓, 52240.7↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0039 s.

Either for fasthttp and net/http should be quite close, the only difference is the way they both upgrade.

Gorilla:

Total data sent: 260.2 MiB (272886768 bytes)
Total data received: 109.3 MiB (114632982 bytes)
Bandwidth per channel: 3.097⇅ Mbps (387.1 kBps)
Aggregate bandwidth: 91.615↓, 218.092↑ Mbps
Packet rate estimate: 109755.3↓, 66807.4↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.01 s.

Nhooyr: (Don't know why is that low)

Total data sent: 224.3 MiB (235184096 bytes)
Total data received: 41.2 MiB (43209780 bytes)
Bandwidth per channel: 2.227⇅ Mbps (278.3 kBps)
Aggregate bandwidth: 34.559↓, 188.097↑ Mbps
Packet rate estimate: 88474.0↓, 55256.1↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0027 s.

Gobwas:

Total data sent: 265.8 MiB (278718160 bytes)
Total data received: 117.8 MiB (123548959 bytes)
Bandwidth per channel: 3.218⇅ Mbps (402.2 kBps)
Aggregate bandwidth: 98.825↓, 222.942↑ Mbps
Packet rate estimate: 148231.6↓, 72106.1↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0015 s.

The source files are in this folder.

About

WebSocket for fasthttp

Topics

Resources

Stars

67 stars

Watchers

2 watching

Forks

Releases

Packages

Used by

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Highlight search terms from Google/DuckDuckGo/Bing referrer (function() { var ref = document.referrer; var terms = []; if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) { var url = new URL(ref); var q = url.searchParams.get('q') || url.searchParams.get('p'); if (q) { terms = q.split(/\s+/).filter(function(t) { return t.length > 2; }); } } if (terms.length === 0) return; var style = document.createElement('style'); style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }'; document.head.appendChild(style); function highlight(node) { if (node.nodeType === 3) { // text node var text = node.textContent; var found = false; terms.forEach(function(term) { var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\]\\]/g, '\\') + ')', 'gi'); if (regex.test(text)) { found = true; var frag = document.createDocumentFragment(); var parts = text.split(regex); parts.forEach(function(part, i) { if (i % 2 === 0) { frag.appendChild(document.createTextNode(part)); } else { var span = document.createElement('span'); span.className = 'userscript-highlight'; span.textContent = part; frag.appendChild(span); } }); node.parentNode.replaceChild(frag, node); } }); } else if (node.nodeType === 1 && node.childNodes) { // element var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT']; if (!skipTags.includes(node.tagName)) { Array.from(node.childNodes).forEach(highlight); } } } highlight(document.body); // Re-highlight on dynamic content var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1 || node.nodeType === 3) highlight(node); }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' GitHub - dgrr/websocket: WebSocket for fasthttp · GitHub
Skip to content

Repository files navigation

websocket

WebSocket library for fasthttp and net/http.

Checkout examples to inspire yourself.

Install

go get github.com/dgrr/websocket

Why another WebSocket package?

Other WebSocket packages DON'T allow concurrent Read/Write operations on servers and they do not provide low level access to WebSocket packet crafting. Those WebSocket packages try to emulate the Golang API by implementing io.Reader and io.Writer interfaces on their connections. io.Writer might be a good idea to use it, but no io.Reader, given that WebSocket is an async protocol by nature (all protocols are (?)).

Sometimes, WebSocket servers are just cumbersome when we want to handle a lot of clients in an async manner. For example, in other WebSocket packages to broadcast a message generated internally we'll need to do the following:

typeMyWebSocketServicestruct {
clients sync.Map
}
typeBlockingConnstruct {
lck sync.Mutexc websocketPackage.Conn
}
func (ws*MyWebSocketService) Broadcaster() {
formsg:=rangemessageProducerChannel {
ws.clients.Range(func(_, vinterface{}) bool {
c:=v.(*BlockingConn)
c.lck.Lock() // oh, we need to block, otherwise we can break the programerr:=c.Write(msg)
c.lck.Unlock()
iferr!=nil {
// we have an error, what can we do? Log it?// if the connection has been closed we'll receive that on// the Read call, so the connection will close automatically.
}
returntrue
})
}
}
func (ws*MyWebSocketService) Handle(request, response) {
c, err:=websocketPackage.Upgrade(request, response)
iferr!=nil {
// then it's clearly an error! Report back
}
bc:=&BlockingConn{
c: c,
} ws.clients.Store(bc, struct{}{})
// even though I just want to write, I need to block somehowfor {
content, err:=bc.Read()
iferr!=nil {
// handle the errorbreak
}
}
ws.clients.Delete(bc)
}

First, we need to store every client upon connection, and whenever we want to send data we need to iterate over a list, and send the message. If while, writing we get an error, then we need to handle that client's error What if the writing operation is happening at the same time in 2 different coroutines? Then we need a sync.Mutex and block until we finish writing.

To solve most of those problems websocket uses channels and separated coroutines, one for reading and another one for writing. By following the sharing principle.

Do not communicate by sharing memory; instead, share memory by communicating.

Following the fasthttp philosophy this library tries to take as much advantage of the Golang's multi-threaded model as possible, while keeping your code concurrently safe.

To see an example of what this package CAN do that others DONT checkout the broadcast example.

Server

How can I launch a server?

It's quite easy. You only need to create a Server, set your callbacks by calling the Handle* methods and then specify your fasthttp handler as Server.Upgrade.

package main
import (
"fmt""github.com/valyala/fasthttp""github.com/dgrr/websocket"
)
funcmain() {
ws:= websocket.Server{}
ws.HandleData(OnMessage)
fasthttp.ListenAndServe(":8080", ws.Upgrade)
}
funcOnMessage(c*websocket.Conn, isBinarybool, data []byte) {
fmt.Printf("Received data from %s: %s\n", c.RemoteAddr(), data)
}

How can I launch a server if I use net/http?

package main
import (
"fmt""net/http""github.com/dgrr/websocket"
)
funcmain() {
ws:= websocket.Server{}
ws.HandleData(OnMessage)
http.HandleFunc("/", ws.NetUpgrade)
http.ListenAndServe(":8080", nil)
}
funcOnMessage(c*websocket.Conn, isBinarybool, data []byte) {
fmt.Printf("Received data from %s: %s\n", c.RemoteAddr(), data)
}

How can I handle pings?

Pings are handle automatically by the library, but you can get the content of those pings setting the callback using HandlePing.

For example, let's try to get the round trip time to a client by using the PING frame. The website http2.gofiber.io uses this method to measure the round trip time displayed at the bottom of the webpage.

package main
import (
"sync""encoding/binary""log""time""github.com/valyala/fasthttp""github.com/dgrr/websocket"
)
// Struct to keep the clients connected//// it should be safe to access the clients concurrently from Open and Close.typeRTTMeasurestruct {
clients sync.Map
}
// just trigger the ping senderfunc (rtt*RTTMeasure) Start() {
time.AfterFunc(time.Second*2, rtt.sendPings)
}
func (rtt*RTTMeasure) sendPings() {
vardata [8]bytebinary.BigEndian.PutUint64(data[:], uint64(
time.Now().UnixNano()),
)
rtt.clients.Range(func(_, vinterface{}) bool {
c:=v.(*websocket.Conn)
c.Ping(data[:])
returntrue
})
rtt.Start()
}
// register a connection when it's openfunc (rtt*RTTMeasure) RegisterConn(c*websocket.Conn) {
rtt.clients.Store(c.ID(), c)
log.Printf("Client %s connected\n", c.RemoteAddr())
}
// remove the connection when receiving the closefunc (rtt*RTTMeasure) RemoveConn(c*websocket.Conn, errerror) {
rtt.clients.Delete(c.ID())
log.Printf("Client %s disconnected\n", c.RemoteAddr())
}
funcmain() {
rtt:=RTTMeasure{}
ws:= websocket.Server{}
ws.HandleOpen(rtt.RegisterConn)
ws.HandleClose(rtt.RemoveConn)
ws.HandlePong(OnPong)
// schedule the timerrtt.Start()
fasthttp.ListenAndServe(":8080", ws.Upgrade)
}
// handle the pong messagefuncOnPong(c*websocket.Conn, data []byte) {
iflen(data) ==8 {
n:=binary.BigEndian.Uint64(data)
ts:=time.Unix(0, int64(n))
log.Printf("RTT with %s is %s\n", c.RemoteAddr(), time.Now().Sub(ts))
}
}

websocket vs gorilla vs nhooyr vs gobwas

FeatureswebsocketGorillaNhooyrgowabs
Concurrent R/WYesNoNo. Only writesNo
Passes Autobahn Test SuiteMostlyYesYesMostly
Receive fragmented messageYesYesYesYes
Send close messageYesYesYesYes
Send pings and receive pongsYesYesYesYes
Get the type of a received data messageYesYesYesYes
Compression ExtensionsNoExperimentalYesNo (?)
Read message using io.ReaderNoYesNoNo (?)
Write message using io.WriteCloserYesYesNoNo (?)

Stress tests

The following stress test were performed without timeouts:

Executing tcpkali --ws -c 100 -m 'hello world!!13212312!' -r 10k localhost:8081 the tests shows the following:

Websocket:

Total data sent: 267.7 MiB (280678466 bytes)
Total data received: 229.5 MiB (240626600 bytes)
Bandwidth per channel: 4.167⇅ Mbps (520.9 kBps)
Aggregate bandwidth: 192.357↓, 224.375↑ Mbps
Packet rate estimate: 247050.1↓, 61842.9↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0075 s.

Websocket for net/http:

Total data sent: 267.3 MiB (280320124 bytes)
Total data received: 228.3 MiB (239396374 bytes)
Bandwidth per channel: 4.156⇅ Mbps (519.5 kBps)
Aggregate bandwidth: 191.442↓, 224.168↑ Mbps
Packet rate estimate: 188107.1↓, 52240.7↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0039 s.

Either for fasthttp and net/http should be quite close, the only difference is the way they both upgrade.

Gorilla:

Total data sent: 260.2 MiB (272886768 bytes)
Total data received: 109.3 MiB (114632982 bytes)
Bandwidth per channel: 3.097⇅ Mbps (387.1 kBps)
Aggregate bandwidth: 91.615↓, 218.092↑ Mbps
Packet rate estimate: 109755.3↓, 66807.4↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.01 s.

Nhooyr: (Don't know why is that low)

Total data sent: 224.3 MiB (235184096 bytes)
Total data received: 41.2 MiB (43209780 bytes)
Bandwidth per channel: 2.227⇅ Mbps (278.3 kBps)
Aggregate bandwidth: 34.559↓, 188.097↑ Mbps
Packet rate estimate: 88474.0↓, 55256.1↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0027 s.

Gobwas:

Total data sent: 265.8 MiB (278718160 bytes)
Total data received: 117.8 MiB (123548959 bytes)
Bandwidth per channel: 3.218⇅ Mbps (402.2 kBps)
Aggregate bandwidth: 98.825↓, 222.942↑ Mbps
Packet rate estimate: 148231.6↓, 72106.1↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0015 s.

The source files are in this folder.

About

WebSocket for fasthttp

Topics

Resources

Stars

67 stars

Watchers

2 watching

Forks

Releases

Packages

Used by

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Strip utm_, fbclid, gclid, etc. from all links on page (function() { var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content', 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid', 'ref', 'ref_src', 'source', 'medium', 'campaign']; function cleanUrl(url) { try { var u = new URL(url, window.location.origin); var changed = false; trackingParams.forEach(function(p) { if (u.searchParams.has(p)) { u.searchParams.delete(p); changed = true; } }); return changed ? u.toString() : url; } catch (e) { return url; } } function cleanLinks() { document.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } cleanLinks(); var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1) { if (node.tagName === 'A') cleanLinks(); node.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + ' GitHub - dgrr/websocket: WebSocket for fasthttp · GitHub
Skip to content

Repository files navigation

websocket

WebSocket library for fasthttp and net/http.

Checkout examples to inspire yourself.

Install

go get github.com/dgrr/websocket

Why another WebSocket package?

Other WebSocket packages DON'T allow concurrent Read/Write operations on servers and they do not provide low level access to WebSocket packet crafting. Those WebSocket packages try to emulate the Golang API by implementing io.Reader and io.Writer interfaces on their connections. io.Writer might be a good idea to use it, but no io.Reader, given that WebSocket is an async protocol by nature (all protocols are (?)).

Sometimes, WebSocket servers are just cumbersome when we want to handle a lot of clients in an async manner. For example, in other WebSocket packages to broadcast a message generated internally we'll need to do the following:

typeMyWebSocketServicestruct {
clients sync.Map
}
typeBlockingConnstruct {
lck sync.Mutexc websocketPackage.Conn
}
func (ws*MyWebSocketService) Broadcaster() {
formsg:=rangemessageProducerChannel {
ws.clients.Range(func(_, vinterface{}) bool {
c:=v.(*BlockingConn)
c.lck.Lock() // oh, we need to block, otherwise we can break the programerr:=c.Write(msg)
c.lck.Unlock()
iferr!=nil {
// we have an error, what can we do? Log it?// if the connection has been closed we'll receive that on// the Read call, so the connection will close automatically.
}
returntrue
})
}
}
func (ws*MyWebSocketService) Handle(request, response) {
c, err:=websocketPackage.Upgrade(request, response)
iferr!=nil {
// then it's clearly an error! Report back
}
bc:=&BlockingConn{
c: c,
} ws.clients.Store(bc, struct{}{})
// even though I just want to write, I need to block somehowfor {
content, err:=bc.Read()
iferr!=nil {
// handle the errorbreak
}
}
ws.clients.Delete(bc)
}

First, we need to store every client upon connection, and whenever we want to send data we need to iterate over a list, and send the message. If while, writing we get an error, then we need to handle that client's error What if the writing operation is happening at the same time in 2 different coroutines? Then we need a sync.Mutex and block until we finish writing.

To solve most of those problems websocket uses channels and separated coroutines, one for reading and another one for writing. By following the sharing principle.

Do not communicate by sharing memory; instead, share memory by communicating.

Following the fasthttp philosophy this library tries to take as much advantage of the Golang's multi-threaded model as possible, while keeping your code concurrently safe.

To see an example of what this package CAN do that others DONT checkout the broadcast example.

Server

How can I launch a server?

It's quite easy. You only need to create a Server, set your callbacks by calling the Handle* methods and then specify your fasthttp handler as Server.Upgrade.

package main
import (
"fmt""github.com/valyala/fasthttp""github.com/dgrr/websocket"
)
funcmain() {
ws:= websocket.Server{}
ws.HandleData(OnMessage)
fasthttp.ListenAndServe(":8080", ws.Upgrade)
}
funcOnMessage(c*websocket.Conn, isBinarybool, data []byte) {
fmt.Printf("Received data from %s: %s\n", c.RemoteAddr(), data)
}

How can I launch a server if I use net/http?

package main
import (
"fmt""net/http""github.com/dgrr/websocket"
)
funcmain() {
ws:= websocket.Server{}
ws.HandleData(OnMessage)
http.HandleFunc("/", ws.NetUpgrade)
http.ListenAndServe(":8080", nil)
}
funcOnMessage(c*websocket.Conn, isBinarybool, data []byte) {
fmt.Printf("Received data from %s: %s\n", c.RemoteAddr(), data)
}

How can I handle pings?

Pings are handle automatically by the library, but you can get the content of those pings setting the callback using HandlePing.

For example, let's try to get the round trip time to a client by using the PING frame. The website http2.gofiber.io uses this method to measure the round trip time displayed at the bottom of the webpage.

package main
import (
"sync""encoding/binary""log""time""github.com/valyala/fasthttp""github.com/dgrr/websocket"
)
// Struct to keep the clients connected//// it should be safe to access the clients concurrently from Open and Close.typeRTTMeasurestruct {
clients sync.Map
}
// just trigger the ping senderfunc (rtt*RTTMeasure) Start() {
time.AfterFunc(time.Second*2, rtt.sendPings)
}
func (rtt*RTTMeasure) sendPings() {
vardata [8]bytebinary.BigEndian.PutUint64(data[:], uint64(
time.Now().UnixNano()),
)
rtt.clients.Range(func(_, vinterface{}) bool {
c:=v.(*websocket.Conn)
c.Ping(data[:])
returntrue
})
rtt.Start()
}
// register a connection when it's openfunc (rtt*RTTMeasure) RegisterConn(c*websocket.Conn) {
rtt.clients.Store(c.ID(), c)
log.Printf("Client %s connected\n", c.RemoteAddr())
}
// remove the connection when receiving the closefunc (rtt*RTTMeasure) RemoveConn(c*websocket.Conn, errerror) {
rtt.clients.Delete(c.ID())
log.Printf("Client %s disconnected\n", c.RemoteAddr())
}
funcmain() {
rtt:=RTTMeasure{}
ws:= websocket.Server{}
ws.HandleOpen(rtt.RegisterConn)
ws.HandleClose(rtt.RemoveConn)
ws.HandlePong(OnPong)
// schedule the timerrtt.Start()
fasthttp.ListenAndServe(":8080", ws.Upgrade)
}
// handle the pong messagefuncOnPong(c*websocket.Conn, data []byte) {
iflen(data) ==8 {
n:=binary.BigEndian.Uint64(data)
ts:=time.Unix(0, int64(n))
log.Printf("RTT with %s is %s\n", c.RemoteAddr(), time.Now().Sub(ts))
}
}

websocket vs gorilla vs nhooyr vs gobwas

FeatureswebsocketGorillaNhooyrgowabs
Concurrent R/WYesNoNo. Only writesNo
Passes Autobahn Test SuiteMostlyYesYesMostly
Receive fragmented messageYesYesYesYes
Send close messageYesYesYesYes
Send pings and receive pongsYesYesYesYes
Get the type of a received data messageYesYesYesYes
Compression ExtensionsNoExperimentalYesNo (?)
Read message using io.ReaderNoYesNoNo (?)
Write message using io.WriteCloserYesYesNoNo (?)

Stress tests

The following stress test were performed without timeouts:

Executing tcpkali --ws -c 100 -m 'hello world!!13212312!' -r 10k localhost:8081 the tests shows the following:

Websocket:

Total data sent: 267.7 MiB (280678466 bytes)
Total data received: 229.5 MiB (240626600 bytes)
Bandwidth per channel: 4.167⇅ Mbps (520.9 kBps)
Aggregate bandwidth: 192.357↓, 224.375↑ Mbps
Packet rate estimate: 247050.1↓, 61842.9↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0075 s.

Websocket for net/http:

Total data sent: 267.3 MiB (280320124 bytes)
Total data received: 228.3 MiB (239396374 bytes)
Bandwidth per channel: 4.156⇅ Mbps (519.5 kBps)
Aggregate bandwidth: 191.442↓, 224.168↑ Mbps
Packet rate estimate: 188107.1↓, 52240.7↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0039 s.

Either for fasthttp and net/http should be quite close, the only difference is the way they both upgrade.

Gorilla:

Total data sent: 260.2 MiB (272886768 bytes)
Total data received: 109.3 MiB (114632982 bytes)
Bandwidth per channel: 3.097⇅ Mbps (387.1 kBps)
Aggregate bandwidth: 91.615↓, 218.092↑ Mbps
Packet rate estimate: 109755.3↓, 66807.4↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.01 s.

Nhooyr: (Don't know why is that low)

Total data sent: 224.3 MiB (235184096 bytes)
Total data received: 41.2 MiB (43209780 bytes)
Bandwidth per channel: 2.227⇅ Mbps (278.3 kBps)
Aggregate bandwidth: 34.559↓, 188.097↑ Mbps
Packet rate estimate: 88474.0↓, 55256.1↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0027 s.

Gobwas:

Total data sent: 265.8 MiB (278718160 bytes)
Total data received: 117.8 MiB (123548959 bytes)
Bandwidth per channel: 3.218⇅ Mbps (402.2 kBps)
Aggregate bandwidth: 98.825↓, 222.942↑ Mbps
Packet rate estimate: 148231.6↓, 72106.1↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0015 s.

The source files are in this folder.

About

WebSocket for fasthttp

Topics

Resources

Stars

67 stars

Watchers

2 watching

Forks

Releases

Packages

Used by

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Auto-enable theater mode on YouTube (function() { function tryTheater() { var btn = document.querySelector('button[aria-label="Theater mode"], ytd-player #player button[title="Theater mode"]'); if (btn && !btn.classList.contains('activated')) { btn.click(); } } // Try immediately tryTheater(); // Try after navigation (SPA) var lastUrl = location.href; setInterval(function() { if (location.href !== lastUrl) { lastUrl = location.href; setTimeout(tryTheater, 500); } }, 1000); // Also try on player load var observer = new MutationObserver(tryTheater); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' GitHub - dgrr/websocket: WebSocket for fasthttp · GitHub
Skip to content

Repository files navigation

websocket

WebSocket library for fasthttp and net/http.

Checkout examples to inspire yourself.

Install

go get github.com/dgrr/websocket

Why another WebSocket package?

Other WebSocket packages DON'T allow concurrent Read/Write operations on servers and they do not provide low level access to WebSocket packet crafting. Those WebSocket packages try to emulate the Golang API by implementing io.Reader and io.Writer interfaces on their connections. io.Writer might be a good idea to use it, but no io.Reader, given that WebSocket is an async protocol by nature (all protocols are (?)).

Sometimes, WebSocket servers are just cumbersome when we want to handle a lot of clients in an async manner. For example, in other WebSocket packages to broadcast a message generated internally we'll need to do the following:

typeMyWebSocketServicestruct {
clients sync.Map
}
typeBlockingConnstruct {
lck sync.Mutexc websocketPackage.Conn
}
func (ws*MyWebSocketService) Broadcaster() {
formsg:=rangemessageProducerChannel {
ws.clients.Range(func(_, vinterface{}) bool {
c:=v.(*BlockingConn)
c.lck.Lock() // oh, we need to block, otherwise we can break the programerr:=c.Write(msg)
c.lck.Unlock()
iferr!=nil {
// we have an error, what can we do? Log it?// if the connection has been closed we'll receive that on// the Read call, so the connection will close automatically.
}
returntrue
})
}
}
func (ws*MyWebSocketService) Handle(request, response) {
c, err:=websocketPackage.Upgrade(request, response)
iferr!=nil {
// then it's clearly an error! Report back
}
bc:=&BlockingConn{
c: c,
} ws.clients.Store(bc, struct{}{})
// even though I just want to write, I need to block somehowfor {
content, err:=bc.Read()
iferr!=nil {
// handle the errorbreak
}
}
ws.clients.Delete(bc)
}

First, we need to store every client upon connection, and whenever we want to send data we need to iterate over a list, and send the message. If while, writing we get an error, then we need to handle that client's error What if the writing operation is happening at the same time in 2 different coroutines? Then we need a sync.Mutex and block until we finish writing.

To solve most of those problems websocket uses channels and separated coroutines, one for reading and another one for writing. By following the sharing principle.

Do not communicate by sharing memory; instead, share memory by communicating.

Following the fasthttp philosophy this library tries to take as much advantage of the Golang's multi-threaded model as possible, while keeping your code concurrently safe.

To see an example of what this package CAN do that others DONT checkout the broadcast example.

Server

How can I launch a server?

It's quite easy. You only need to create a Server, set your callbacks by calling the Handle* methods and then specify your fasthttp handler as Server.Upgrade.

package main
import (
"fmt""github.com/valyala/fasthttp""github.com/dgrr/websocket"
)
funcmain() {
ws:= websocket.Server{}
ws.HandleData(OnMessage)
fasthttp.ListenAndServe(":8080", ws.Upgrade)
}
funcOnMessage(c*websocket.Conn, isBinarybool, data []byte) {
fmt.Printf("Received data from %s: %s\n", c.RemoteAddr(), data)
}

How can I launch a server if I use net/http?

package main
import (
"fmt""net/http""github.com/dgrr/websocket"
)
funcmain() {
ws:= websocket.Server{}
ws.HandleData(OnMessage)
http.HandleFunc("/", ws.NetUpgrade)
http.ListenAndServe(":8080", nil)
}
funcOnMessage(c*websocket.Conn, isBinarybool, data []byte) {
fmt.Printf("Received data from %s: %s\n", c.RemoteAddr(), data)
}

How can I handle pings?

Pings are handle automatically by the library, but you can get the content of those pings setting the callback using HandlePing.

For example, let's try to get the round trip time to a client by using the PING frame. The website http2.gofiber.io uses this method to measure the round trip time displayed at the bottom of the webpage.

package main
import (
"sync""encoding/binary""log""time""github.com/valyala/fasthttp""github.com/dgrr/websocket"
)
// Struct to keep the clients connected//// it should be safe to access the clients concurrently from Open and Close.typeRTTMeasurestruct {
clients sync.Map
}
// just trigger the ping senderfunc (rtt*RTTMeasure) Start() {
time.AfterFunc(time.Second*2, rtt.sendPings)
}
func (rtt*RTTMeasure) sendPings() {
vardata [8]bytebinary.BigEndian.PutUint64(data[:], uint64(
time.Now().UnixNano()),
)
rtt.clients.Range(func(_, vinterface{}) bool {
c:=v.(*websocket.Conn)
c.Ping(data[:])
returntrue
})
rtt.Start()
}
// register a connection when it's openfunc (rtt*RTTMeasure) RegisterConn(c*websocket.Conn) {
rtt.clients.Store(c.ID(), c)
log.Printf("Client %s connected\n", c.RemoteAddr())
}
// remove the connection when receiving the closefunc (rtt*RTTMeasure) RemoveConn(c*websocket.Conn, errerror) {
rtt.clients.Delete(c.ID())
log.Printf("Client %s disconnected\n", c.RemoteAddr())
}
funcmain() {
rtt:=RTTMeasure{}
ws:= websocket.Server{}
ws.HandleOpen(rtt.RegisterConn)
ws.HandleClose(rtt.RemoveConn)
ws.HandlePong(OnPong)
// schedule the timerrtt.Start()
fasthttp.ListenAndServe(":8080", ws.Upgrade)
}
// handle the pong messagefuncOnPong(c*websocket.Conn, data []byte) {
iflen(data) ==8 {
n:=binary.BigEndian.Uint64(data)
ts:=time.Unix(0, int64(n))
log.Printf("RTT with %s is %s\n", c.RemoteAddr(), time.Now().Sub(ts))
}
}

websocket vs gorilla vs nhooyr vs gobwas

FeatureswebsocketGorillaNhooyrgowabs
Concurrent R/WYesNoNo. Only writesNo
Passes Autobahn Test SuiteMostlyYesYesMostly
Receive fragmented messageYesYesYesYes
Send close messageYesYesYesYes
Send pings and receive pongsYesYesYesYes
Get the type of a received data messageYesYesYesYes
Compression ExtensionsNoExperimentalYesNo (?)
Read message using io.ReaderNoYesNoNo (?)
Write message using io.WriteCloserYesYesNoNo (?)

Stress tests

The following stress test were performed without timeouts:

Executing tcpkali --ws -c 100 -m 'hello world!!13212312!' -r 10k localhost:8081 the tests shows the following:

Websocket:

Total data sent: 267.7 MiB (280678466 bytes)
Total data received: 229.5 MiB (240626600 bytes)
Bandwidth per channel: 4.167⇅ Mbps (520.9 kBps)
Aggregate bandwidth: 192.357↓, 224.375↑ Mbps
Packet rate estimate: 247050.1↓, 61842.9↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0075 s.

Websocket for net/http:

Total data sent: 267.3 MiB (280320124 bytes)
Total data received: 228.3 MiB (239396374 bytes)
Bandwidth per channel: 4.156⇅ Mbps (519.5 kBps)
Aggregate bandwidth: 191.442↓, 224.168↑ Mbps
Packet rate estimate: 188107.1↓, 52240.7↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0039 s.

Either for fasthttp and net/http should be quite close, the only difference is the way they both upgrade.

Gorilla:

Total data sent: 260.2 MiB (272886768 bytes)
Total data received: 109.3 MiB (114632982 bytes)
Bandwidth per channel: 3.097⇅ Mbps (387.1 kBps)
Aggregate bandwidth: 91.615↓, 218.092↑ Mbps
Packet rate estimate: 109755.3↓, 66807.4↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.01 s.

Nhooyr: (Don't know why is that low)

Total data sent: 224.3 MiB (235184096 bytes)
Total data received: 41.2 MiB (43209780 bytes)
Bandwidth per channel: 2.227⇅ Mbps (278.3 kBps)
Aggregate bandwidth: 34.559↓, 188.097↑ Mbps
Packet rate estimate: 88474.0↓, 55256.1↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0027 s.

Gobwas:

Total data sent: 265.8 MiB (278718160 bytes)
Total data received: 117.8 MiB (123548959 bytes)
Bandwidth per channel: 3.218⇅ Mbps (402.2 kBps)
Aggregate bandwidth: 98.825↓, 222.942↑ Mbps
Packet rate estimate: 148231.6↓, 72106.1↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0015 s.

The source files are in this folder.

About

WebSocket for fasthttp

Topics

Resources

Stars

67 stars

Watchers

2 watching

Forks

Releases

Packages

Used by

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Remove or un-stick sticky/fixed headers that block content (function() { function unstick() { document.querySelectorAll('header, nav, [role="banner"], .header, .navbar, .sticky, .fixed-top, [style*="position: fixed"], [style*="position:sticky"]').forEach(function(el) { if (el.style.position === 'fixed' || el.style.position === 'sticky' || getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') { el.style.position = 'static'; el.style.top = 'auto'; el.style.zIndex = 'auto'; } }); } unstick(); var observer = new MutationObserver(unstick); observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] }); })(); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + ' GitHub - dgrr/websocket: WebSocket for fasthttp · GitHub
Skip to content

Repository files navigation

websocket

WebSocket library for fasthttp and net/http.

Checkout examples to inspire yourself.

Install

go get github.com/dgrr/websocket

Why another WebSocket package?

Other WebSocket packages DON'T allow concurrent Read/Write operations on servers and they do not provide low level access to WebSocket packet crafting. Those WebSocket packages try to emulate the Golang API by implementing io.Reader and io.Writer interfaces on their connections. io.Writer might be a good idea to use it, but no io.Reader, given that WebSocket is an async protocol by nature (all protocols are (?)).

Sometimes, WebSocket servers are just cumbersome when we want to handle a lot of clients in an async manner. For example, in other WebSocket packages to broadcast a message generated internally we'll need to do the following:

typeMyWebSocketServicestruct {
clients sync.Map
}
typeBlockingConnstruct {
lck sync.Mutexc websocketPackage.Conn
}
func (ws*MyWebSocketService) Broadcaster() {
formsg:=rangemessageProducerChannel {
ws.clients.Range(func(_, vinterface{}) bool {
c:=v.(*BlockingConn)
c.lck.Lock() // oh, we need to block, otherwise we can break the programerr:=c.Write(msg)
c.lck.Unlock()
iferr!=nil {
// we have an error, what can we do? Log it?// if the connection has been closed we'll receive that on// the Read call, so the connection will close automatically.
}
returntrue
})
}
}
func (ws*MyWebSocketService) Handle(request, response) {
c, err:=websocketPackage.Upgrade(request, response)
iferr!=nil {
// then it's clearly an error! Report back
}
bc:=&BlockingConn{
c: c,
} ws.clients.Store(bc, struct{}{})
// even though I just want to write, I need to block somehowfor {
content, err:=bc.Read()
iferr!=nil {
// handle the errorbreak
}
}
ws.clients.Delete(bc)
}

First, we need to store every client upon connection, and whenever we want to send data we need to iterate over a list, and send the message. If while, writing we get an error, then we need to handle that client's error What if the writing operation is happening at the same time in 2 different coroutines? Then we need a sync.Mutex and block until we finish writing.

To solve most of those problems websocket uses channels and separated coroutines, one for reading and another one for writing. By following the sharing principle.

Do not communicate by sharing memory; instead, share memory by communicating.

Following the fasthttp philosophy this library tries to take as much advantage of the Golang's multi-threaded model as possible, while keeping your code concurrently safe.

To see an example of what this package CAN do that others DONT checkout the broadcast example.

Server

How can I launch a server?

It's quite easy. You only need to create a Server, set your callbacks by calling the Handle* methods and then specify your fasthttp handler as Server.Upgrade.

package main
import (
"fmt""github.com/valyala/fasthttp""github.com/dgrr/websocket"
)
funcmain() {
ws:= websocket.Server{}
ws.HandleData(OnMessage)
fasthttp.ListenAndServe(":8080", ws.Upgrade)
}
funcOnMessage(c*websocket.Conn, isBinarybool, data []byte) {
fmt.Printf("Received data from %s: %s\n", c.RemoteAddr(), data)
}

How can I launch a server if I use net/http?

package main
import (
"fmt""net/http""github.com/dgrr/websocket"
)
funcmain() {
ws:= websocket.Server{}
ws.HandleData(OnMessage)
http.HandleFunc("/", ws.NetUpgrade)
http.ListenAndServe(":8080", nil)
}
funcOnMessage(c*websocket.Conn, isBinarybool, data []byte) {
fmt.Printf("Received data from %s: %s\n", c.RemoteAddr(), data)
}

How can I handle pings?

Pings are handle automatically by the library, but you can get the content of those pings setting the callback using HandlePing.

For example, let's try to get the round trip time to a client by using the PING frame. The website http2.gofiber.io uses this method to measure the round trip time displayed at the bottom of the webpage.

package main
import (
"sync""encoding/binary""log""time""github.com/valyala/fasthttp""github.com/dgrr/websocket"
)
// Struct to keep the clients connected//// it should be safe to access the clients concurrently from Open and Close.typeRTTMeasurestruct {
clients sync.Map
}
// just trigger the ping senderfunc (rtt*RTTMeasure) Start() {
time.AfterFunc(time.Second*2, rtt.sendPings)
}
func (rtt*RTTMeasure) sendPings() {
vardata [8]bytebinary.BigEndian.PutUint64(data[:], uint64(
time.Now().UnixNano()),
)
rtt.clients.Range(func(_, vinterface{}) bool {
c:=v.(*websocket.Conn)
c.Ping(data[:])
returntrue
})
rtt.Start()
}
// register a connection when it's openfunc (rtt*RTTMeasure) RegisterConn(c*websocket.Conn) {
rtt.clients.Store(c.ID(), c)
log.Printf("Client %s connected\n", c.RemoteAddr())
}
// remove the connection when receiving the closefunc (rtt*RTTMeasure) RemoveConn(c*websocket.Conn, errerror) {
rtt.clients.Delete(c.ID())
log.Printf("Client %s disconnected\n", c.RemoteAddr())
}
funcmain() {
rtt:=RTTMeasure{}
ws:= websocket.Server{}
ws.HandleOpen(rtt.RegisterConn)
ws.HandleClose(rtt.RemoveConn)
ws.HandlePong(OnPong)
// schedule the timerrtt.Start()
fasthttp.ListenAndServe(":8080", ws.Upgrade)
}
// handle the pong messagefuncOnPong(c*websocket.Conn, data []byte) {
iflen(data) ==8 {
n:=binary.BigEndian.Uint64(data)
ts:=time.Unix(0, int64(n))
log.Printf("RTT with %s is %s\n", c.RemoteAddr(), time.Now().Sub(ts))
}
}

websocket vs gorilla vs nhooyr vs gobwas

FeatureswebsocketGorillaNhooyrgowabs
Concurrent R/WYesNoNo. Only writesNo
Passes Autobahn Test SuiteMostlyYesYesMostly
Receive fragmented messageYesYesYesYes
Send close messageYesYesYesYes
Send pings and receive pongsYesYesYesYes
Get the type of a received data messageYesYesYesYes
Compression ExtensionsNoExperimentalYesNo (?)
Read message using io.ReaderNoYesNoNo (?)
Write message using io.WriteCloserYesYesNoNo (?)

Stress tests

The following stress test were performed without timeouts:

Executing tcpkali --ws -c 100 -m 'hello world!!13212312!' -r 10k localhost:8081 the tests shows the following:

Websocket:

Total data sent: 267.7 MiB (280678466 bytes)
Total data received: 229.5 MiB (240626600 bytes)
Bandwidth per channel: 4.167⇅ Mbps (520.9 kBps)
Aggregate bandwidth: 192.357↓, 224.375↑ Mbps
Packet rate estimate: 247050.1↓, 61842.9↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0075 s.

Websocket for net/http:

Total data sent: 267.3 MiB (280320124 bytes)
Total data received: 228.3 MiB (239396374 bytes)
Bandwidth per channel: 4.156⇅ Mbps (519.5 kBps)
Aggregate bandwidth: 191.442↓, 224.168↑ Mbps
Packet rate estimate: 188107.1↓, 52240.7↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0039 s.

Either for fasthttp and net/http should be quite close, the only difference is the way they both upgrade.

Gorilla:

Total data sent: 260.2 MiB (272886768 bytes)
Total data received: 109.3 MiB (114632982 bytes)
Bandwidth per channel: 3.097⇅ Mbps (387.1 kBps)
Aggregate bandwidth: 91.615↓, 218.092↑ Mbps
Packet rate estimate: 109755.3↓, 66807.4↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.01 s.

Nhooyr: (Don't know why is that low)

Total data sent: 224.3 MiB (235184096 bytes)
Total data received: 41.2 MiB (43209780 bytes)
Bandwidth per channel: 2.227⇅ Mbps (278.3 kBps)
Aggregate bandwidth: 34.559↓, 188.097↑ Mbps
Packet rate estimate: 88474.0↓, 55256.1↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0027 s.

Gobwas:

Total data sent: 265.8 MiB (278718160 bytes)
Total data received: 117.8 MiB (123548959 bytes)
Bandwidth per channel: 3.218⇅ Mbps (402.2 kBps)
Aggregate bandwidth: 98.825↓, 222.942↑ Mbps
Packet rate estimate: 148231.6↓, 72106.1↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0015 s.

The source files are in this folder.

About

WebSocket for fasthttp

Topics

Resources

Stars

67 stars

Watchers

2 watching

Forks

Releases

Packages

Used by

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { // Universal Dark Mode - works on any site (function() { var enabled = true; function applyDarkMode() { if (!enabled) return; // Create style element if it doesn't exist var style = document.getElementById('universal-dark-mode-style'); if (!style) { style = document.createElement('style'); style.id = 'universal-dark-mode-style'; document.head.appendChild(style); } // Dark mode CSS - inverts colors but preserves images/video style.textContent = ' /* Invert everything except media */ html { filter: invert(1) hue-rotate(180deg) !important; background: #1a1a2e !important; } /* Restore images, videos, iframes, canvas */ img, video, iframe, canvas, svg, picture, [style*="background-image"] { filter: invert(1) hue-rotate(180deg) !important; } /* Preserve specific elements that should not be inverted */ .no-dark-mode, .no-dark-mode *, [data-theme="light"], [data-theme="light"], .ace_editor, .ace_editor *, .CodeMirror, .CodeMirror *, .monaco-editor, .monaco-editor *, .markdown-body pre, .markdown-body pre *, .highlight, .highlight *, pre code, pre code * { filter: none !important; } /* Fix common UI elements */ .modal, .popup, .dropdown-menu, .tooltip, .popover { filter: invert(1) hue-rotate(180deg) !important; background: #2d2d44 !important; border-color: #444 !important; } /* Scrollbars */ ::-webkit-scrollbar { background: #1a1a2e !important; } ::-webkit-scrollbar-thumb { background: #444 !important; } ::-webkit-scrollbar-thumb:hover { background: #555 !important; } /* Selection */ ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; } ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; } '; } function removeDarkMode() { var style = document.getElementById('universal-dark-mode-style'); if (style) style.remove(); } // Toggle with Alt+Shift+D document.addEventListener('keydown', function(e) { if (e.altKey && e.shiftKey && e.key === 'D') { e.preventDefault(); enabled = !enabled; if (enabled) { applyDarkMode(); console.log('[Universal Dark Mode] Enabled'); } else { removeDarkMode(); console.log('[Universal Dark Mode] Disabled'); } } }); // Apply on load applyDarkMode(); // Re-apply on dynamic content var observer = new MutationObserver(function(mutations) { if (enabled && !document.getElementById('universal-dark-mode-style')) { applyDarkMode(); } }); observer.observe(document.head, { childList: true }); console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle'); })(); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })(); GitHub - dgrr/websocket: WebSocket for fasthttp · GitHub
Skip to content

Repository files navigation

websocket

WebSocket library for fasthttp and net/http.

Checkout examples to inspire yourself.

Install

go get github.com/dgrr/websocket

Why another WebSocket package?

Other WebSocket packages DON'T allow concurrent Read/Write operations on servers and they do not provide low level access to WebSocket packet crafting. Those WebSocket packages try to emulate the Golang API by implementing io.Reader and io.Writer interfaces on their connections. io.Writer might be a good idea to use it, but no io.Reader, given that WebSocket is an async protocol by nature (all protocols are (?)).

Sometimes, WebSocket servers are just cumbersome when we want to handle a lot of clients in an async manner. For example, in other WebSocket packages to broadcast a message generated internally we'll need to do the following:

typeMyWebSocketServicestruct {
clients sync.Map
}
typeBlockingConnstruct {
lck sync.Mutexc websocketPackage.Conn
}
func (ws*MyWebSocketService) Broadcaster() {
formsg:=rangemessageProducerChannel {
ws.clients.Range(func(_, vinterface{}) bool {
c:=v.(*BlockingConn)
c.lck.Lock() // oh, we need to block, otherwise we can break the programerr:=c.Write(msg)
c.lck.Unlock()
iferr!=nil {
// we have an error, what can we do? Log it?// if the connection has been closed we'll receive that on// the Read call, so the connection will close automatically.
}
returntrue
})
}
}
func (ws*MyWebSocketService) Handle(request, response) {
c, err:=websocketPackage.Upgrade(request, response)
iferr!=nil {
// then it's clearly an error! Report back
}
bc:=&BlockingConn{
c: c,
} ws.clients.Store(bc, struct{}{})
// even though I just want to write, I need to block somehowfor {
content, err:=bc.Read()
iferr!=nil {
// handle the errorbreak
}
}
ws.clients.Delete(bc)
}

First, we need to store every client upon connection, and whenever we want to send data we need to iterate over a list, and send the message. If while, writing we get an error, then we need to handle that client's error What if the writing operation is happening at the same time in 2 different coroutines? Then we need a sync.Mutex and block until we finish writing.

To solve most of those problems websocket uses channels and separated coroutines, one for reading and another one for writing. By following the sharing principle.

Do not communicate by sharing memory; instead, share memory by communicating.

Following the fasthttp philosophy this library tries to take as much advantage of the Golang's multi-threaded model as possible, while keeping your code concurrently safe.

To see an example of what this package CAN do that others DONT checkout the broadcast example.

Server

How can I launch a server?

It's quite easy. You only need to create a Server, set your callbacks by calling the Handle* methods and then specify your fasthttp handler as Server.Upgrade.

package main
import (
"fmt""github.com/valyala/fasthttp""github.com/dgrr/websocket"
)
funcmain() {
ws:= websocket.Server{}
ws.HandleData(OnMessage)
fasthttp.ListenAndServe(":8080", ws.Upgrade)
}
funcOnMessage(c*websocket.Conn, isBinarybool, data []byte) {
fmt.Printf("Received data from %s: %s\n", c.RemoteAddr(), data)
}

How can I launch a server if I use net/http?

package main
import (
"fmt""net/http""github.com/dgrr/websocket"
)
funcmain() {
ws:= websocket.Server{}
ws.HandleData(OnMessage)
http.HandleFunc("/", ws.NetUpgrade)
http.ListenAndServe(":8080", nil)
}
funcOnMessage(c*websocket.Conn, isBinarybool, data []byte) {
fmt.Printf("Received data from %s: %s\n", c.RemoteAddr(), data)
}

How can I handle pings?

Pings are handle automatically by the library, but you can get the content of those pings setting the callback using HandlePing.

For example, let's try to get the round trip time to a client by using the PING frame. The website http2.gofiber.io uses this method to measure the round trip time displayed at the bottom of the webpage.

package main
import (
"sync""encoding/binary""log""time""github.com/valyala/fasthttp""github.com/dgrr/websocket"
)
// Struct to keep the clients connected//// it should be safe to access the clients concurrently from Open and Close.typeRTTMeasurestruct {
clients sync.Map
}
// just trigger the ping senderfunc (rtt*RTTMeasure) Start() {
time.AfterFunc(time.Second*2, rtt.sendPings)
}
func (rtt*RTTMeasure) sendPings() {
vardata [8]bytebinary.BigEndian.PutUint64(data[:], uint64(
time.Now().UnixNano()),
)
rtt.clients.Range(func(_, vinterface{}) bool {
c:=v.(*websocket.Conn)
c.Ping(data[:])
returntrue
})
rtt.Start()
}
// register a connection when it's openfunc (rtt*RTTMeasure) RegisterConn(c*websocket.Conn) {
rtt.clients.Store(c.ID(), c)
log.Printf("Client %s connected\n", c.RemoteAddr())
}
// remove the connection when receiving the closefunc (rtt*RTTMeasure) RemoveConn(c*websocket.Conn, errerror) {
rtt.clients.Delete(c.ID())
log.Printf("Client %s disconnected\n", c.RemoteAddr())
}
funcmain() {
rtt:=RTTMeasure{}
ws:= websocket.Server{}
ws.HandleOpen(rtt.RegisterConn)
ws.HandleClose(rtt.RemoveConn)
ws.HandlePong(OnPong)
// schedule the timerrtt.Start()
fasthttp.ListenAndServe(":8080", ws.Upgrade)
}
// handle the pong messagefuncOnPong(c*websocket.Conn, data []byte) {
iflen(data) ==8 {
n:=binary.BigEndian.Uint64(data)
ts:=time.Unix(0, int64(n))
log.Printf("RTT with %s is %s\n", c.RemoteAddr(), time.Now().Sub(ts))
}
}

websocket vs gorilla vs nhooyr vs gobwas

FeatureswebsocketGorillaNhooyrgowabs
Concurrent R/WYesNoNo. Only writesNo
Passes Autobahn Test SuiteMostlyYesYesMostly
Receive fragmented messageYesYesYesYes
Send close messageYesYesYesYes
Send pings and receive pongsYesYesYesYes
Get the type of a received data messageYesYesYesYes
Compression ExtensionsNoExperimentalYesNo (?)
Read message using io.ReaderNoYesNoNo (?)
Write message using io.WriteCloserYesYesNoNo (?)

Stress tests

The following stress test were performed without timeouts:

Executing tcpkali --ws -c 100 -m 'hello world!!13212312!' -r 10k localhost:8081 the tests shows the following:

Websocket:

Total data sent: 267.7 MiB (280678466 bytes)
Total data received: 229.5 MiB (240626600 bytes)
Bandwidth per channel: 4.167⇅ Mbps (520.9 kBps)
Aggregate bandwidth: 192.357↓, 224.375↑ Mbps
Packet rate estimate: 247050.1↓, 61842.9↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0075 s.

Websocket for net/http:

Total data sent: 267.3 MiB (280320124 bytes)
Total data received: 228.3 MiB (239396374 bytes)
Bandwidth per channel: 4.156⇅ Mbps (519.5 kBps)
Aggregate bandwidth: 191.442↓, 224.168↑ Mbps
Packet rate estimate: 188107.1↓, 52240.7↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0039 s.

Either for fasthttp and net/http should be quite close, the only difference is the way they both upgrade.

Gorilla:

Total data sent: 260.2 MiB (272886768 bytes)
Total data received: 109.3 MiB (114632982 bytes)
Bandwidth per channel: 3.097⇅ Mbps (387.1 kBps)
Aggregate bandwidth: 91.615↓, 218.092↑ Mbps
Packet rate estimate: 109755.3↓, 66807.4↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.01 s.

Nhooyr: (Don't know why is that low)

Total data sent: 224.3 MiB (235184096 bytes)
Total data received: 41.2 MiB (43209780 bytes)
Bandwidth per channel: 2.227⇅ Mbps (278.3 kBps)
Aggregate bandwidth: 34.559↓, 188.097↑ Mbps
Packet rate estimate: 88474.0↓, 55256.1↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0027 s.

Gobwas:

Total data sent: 265.8 MiB (278718160 bytes)
Total data received: 117.8 MiB (123548959 bytes)
Bandwidth per channel: 3.218⇅ Mbps (402.2 kBps)
Aggregate bandwidth: 98.825↓, 222.942↑ Mbps
Packet rate estimate: 148231.6↓, 72106.1↑ (1↓, 1↑ TCP MSS/op)
Test duration: 10.0015 s.

The source files are in this folder.

About

WebSocket for fasthttp

Topics

Resources

Stars

67 stars

Watchers

2 watching

Forks

Releases

Packages

Used by

Contributors

Languages