基于 epoll/kqueue 的网络库,callback 写法。起点是 WebSocket,已经有 HTTP/1.1、HTTP/2、HTTP/3、gRPC 和 TLS。
- 支持 epoll/kqueue
- 低内存占用
- 高tps(详见 docs/bench-vs-fnet.md)
- WebSocket 完整实现 rfc6455 / rfc7692,Autobahn 517 用例零失败
- 多协议共用一套引擎:HTTP/1.1、HTTP/2、HTTP/3、gRPC、TLS
| 协议 | 包 | 状态 |
|---|---|---|
| WebSocket (rfc6455 / rfc7692) | websocket/ |
完整实现,Autobahn 517 用例零失败 |
| HTTP/1.1 | http/ |
端到端可用(net/http 的 handler 直接能用) |
| HTTP/2 (RFC 9113 + HPACK) | http2/ |
帧层 + 流 + HPACK,和官方 Framer 双向对测通过 |
| gRPC | grpc/ |
消息分帧 + trailer 状态,端到端可用 |
| TLS 1.3 | tls/ |
状态机实现,不启动 goroutine、不阻塞事件循环 |
| HTTP/3 | http3/ |
端到端可用(QUIC 用 quic-go,帧层自己实现) |
引擎在 engine/(epoll/kqueue + 非阻塞 io + Handler 接口),
每个协议包实现它。分层和边界见 docs/architecture.md。
tls/ 是自己实现的状态机,不是包 crypto/tls。原因:crypto/tls
的握手不支持分步——它内部的 handshakeErr 一旦置上,后面每次
Handshake() 都直接返回那个错误,不会再尝试读。套在非阻塞 fd 上,
第一次"数据不够"那条连接就废了。
状态机版是"喂一段密文、推进一点、吐一段密文",每一步都能停,所以能和
普通连接一样跑在同一个 epoll 上(海量连接场景需要的)。已经和标准库的
crypto/tls 服务端握手成功。
- windows
- io-uring
- QUIC 的完整实现(
http3/用 quic-go 做传输,自己实现的是帧层)
早期阶段,暂时不建议生产使用
go get github.com/antlabs/fio/websocketpackage main
import (
"fmt"
"github.com/antlabs/fio/websocket"
)
type echoHandler struct{}
func (e *echoHandler) OnOpen(c *websocket.Conn) {
// fmt.Printf("OnOpen: %p\n", c)
}
func (e *echoHandler) OnMessage(c *websocket.Conn, op websocket.Opcode, msg []byte) {
if err := c.WriteTimeout(op, msg, 3*time.Second); err != nil {
fmt.Println("write fail:", err)
}
// if err := c.WriteMessage(op, msg); err != nil {
// slog.Error("write fail:", err)
// }
}
func (e *echoHandler) OnClose(c *websocket.Conn, err error) {
errMsg := ""
if err != nil {
errMsg = err.Error()
}
slog.Error("OnClose:", errMsg)
}
type handler struct {
m *websocket.MultiEventLoop
}
func (h *handler) echo(w http.ResponseWriter, r *http.Request) {
c, err := websocket.Upgrade(w, r,
websocket.WithServerReplyPing(),
// websocket.WithServerDecompression(),
websocket.WithServerIgnorePong(),
websocket.WithServerCallback(&echoHandler{}),
// websocket.WithServerEnableUTF8Check(),
websocket.WithServerReadTimeout(5*time.Second),
websocket.WithServerMultiEventLoop(h.m),
)
if err != nil {
slog.Error("Upgrade fail:", "err", err.Error())
}
_ = c
}
func main() {
var h handler
h.m = websocket.NewMultiEventLoopMust(websocket.WithEventLoops(0), websocket.WithMaxEventNum(256), websocket.WithLogLevel(slog.LevelError)) // epoll, kqueue
h.m.Start()
fmt.Printf("apiname:%s\n", h.m.GetApiName())
mux := &http.ServeMux{}
mux.HandleFunc("/autobahn", h.echo)
rawTCP, err := net.Listen("tcp", ":9001")
if err != nil {
fmt.Println("Listen fail:", err)
return
}
log.Println("non-tls server exit:", http.Serve(rawTCP, mux))
}package main
import (
"fmt"
"github.com/antlabs/fio/websocket"
"github.com/gin-gonic/gin"
)
type handler struct{
m *websocket.MultiEventLoop
}
func (h *handler) OnOpen(c *websocket.Conn) {
fmt.Printf("服务端收到一个新的连接")
}
func (h *handler) OnMessage(c *websocket.Conn, op websocket.Opcode, msg []byte) {
// 如果msg的生命周期不是在OnMessage中结束,需要拷贝一份
// newMsg := make([]byte, len(msg))
// copy(newMsg, msg)
fmt.Printf("收到客户端消息:%s\n", msg)
c.WriteMessage(op, msg)
// os.Stdout.Write(msg)
}
func (h *handler) OnClose(c *websocket.Conn, err error) {
fmt.Printf("服务端连接关闭:%v\n", err)
}
func main() {
r := gin.Default()
var h handler
h.m = websocket.NewMultiEventLoopMust(websocket.WithEventLoops(0), websocket.WithMaxEventNum(256), websocket.WithLogLevel(slog.LevelError)) // epoll, kqueue
h.m.Start()
r.GET("/", func(c *gin.Context) {
con, err := websocket.Upgrade(c.Writer, c.Request, websocket.WithServerCallback(h.m), websocket.WithServerMultiEventLoop(h.m))
if err != nil {
return
}
con.StartReadLoop()
})
r.Run()
}package main
import (
"fmt"
"time"
"github.com/antlabs/fio/websocket"
)
var m *websocket.MultiEventLoop
type handler struct{}
func (h *handler) OnOpen(c *websocket.Conn) {
fmt.Printf("客户端连接成功\n")
}
func (h *handler) OnMessage(c *websocket.Conn, op websocket.Opcode, msg []byte) {
// 如果msg的生命周期不是在OnMessage中结束,需要拷贝一份
// newMsg := make([]byte, len(msg))
// copy(newMsg, msg)
fmt.Printf("收到服务端消息:%s\n", msg)
c.WriteMessage(op, msg)
time.Sleep(time.Second)
}
func (h *handler) OnClose(c *websocket.Conn, err error) {
fmt.Printf("客户端端连接关闭:%v\n", err)
}
func main() {
m = websocket.NewMultiEventLoopMust(websocket.WithEventLoops(0), websocket.WithMaxEventNum(256), websocket.WithLogLevel(slog.LevelError)) // epoll, kqueue
m.Start()
c, err := websocket.Dial("ws://127.0.0.1:8080/", websocket.WithClientCallback(&handler{}), websocket.WithServerMultiEventLoop(h.m))
if err != nil {
fmt.Printf("连接失败:%v\n", err)
return
}
c.WriteMessage(opcode.Text, []byte("hello"))
time.Sleep(time.Hour) //demo里面等待下OnMessage 看下执行效果,因为websocket.Dial和WriteMessage都是非阻塞的函数调用,不会卡住主go程
}func main() {
websocket.Dial("ws://127.0.0.1:12345/test", websocket.WithClientHTTPHeader(http.Header{
"h1": "v1",
"h2":"v2",
}))
}func main() {
websocket.Dial("ws://127.0.0.1:12345/test", websocket.WithClientDialTimeout(2 * time.Second))
}func main() {
websocket.Dial("ws://127.0.0.1:12345/test", websocket.WithClientReplyPing())
} // 限制客户端最大服务返回返回的最大包是1024,如果超过这个大小报错
websocket.Dial("ws://127.0.0.1:12345/test", websocket.WithClientReadMaxMessage(1024))func main() {
websocket.Dial("ws://127.0.0.1:12345/test", websocket.WithClientDecompressAndCompress())
}func main() {
websocket.Dial("ws://127.0.0.1:12345/test", websocket.WithClientContextTakeover())
}func main() {
c, err := websocket.Upgrade(w, r, websocket.WithServerReplyPing())
if err != nil {
fmt.Println("Upgrade fail:", err)
return
}
}func main() {
// 配置服务端读取客户端最大的包是1024大小, 超过该值报错
c, err := websocket.Upgrade(w, r, websocket.WithServerReadMaxMessage(1024))
if err != nil {
fmt.Println("Upgrade fail:", err)
return
}
}func main() {
// 配置服务端读取客户端最大的包是1024大小, 超过该值报错
c, err := websocket.Upgrade(w, r, websocket.WithServerDecompression())
if err != nil {
fmt.Println("Upgrade fail:", err)
return
}
}func main() {
c, err := websocket.Upgrade(w, r, websocket.WithServerDecompressAndCompress())
if err != nil {
fmt.Println("Upgrade fail:", err)
return
}
}func main() {
// 配置服务端读取客户端最大的包是1024大小, 超过该值报错
c, err := websocket.Upgrade(w, r, websocket.WithServerContextTakeover)
if err != nil {
fmt.Println("Upgrade fail:", err)
return
}
}- cpu=e5 2686(单路)
- memory=32GB
BenchType : BenchEcho
Framework : fio
TPS : 106014
EER : 218.54
Min : 49.26us
Avg : 94.08ms
Max : 954.33ms
TP50 : 45.76ms
TP75 : 52.27ms
TP90 : 336.85ms
TP95 : 427.07ms
TP99 : 498.66ms
Used : 18.87s
Total : 2000000
Success : 2000000
Failed : 0
Conns : 1000000
Concurrency: 10000
Payload : 1024
CPU Min : 184.90%
CPU Avg : 485.10%
CPU Max : 588.31%
MEM Min : 563.40M
MEM Avg : 572.40M
MEM Max : 594.48M
- cpu=5800h
- memory=64GB
BenchType : BenchEcho
Framework : fio
TPS : 103544
EER : 397.07
Min : 26.51us
Avg : 95.79ms
Max : 1.34s
TP50 : 58.26ms
TP75 : 60.94ms
TP90 : 62.50ms
TP95 : 63.04ms
TP99 : 63.47ms
Used : 40.76s
Total : 5000000
Success : 4220634
Failed : 779366
Conns : 1000000
Concurrency: 10000
Payload : 1024
CPU Min : 30.54%
CPU Avg : 260.77%
CPU Max : 335.88%
MEM Min : 432.25M
MEM Avg : 439.71M
MEM Max : 449.62M
