ledisdb/server/client_resp.go

275 lines
5.0 KiB
Go
Raw Normal View History

package server
import (
"bufio"
"errors"
2014-10-29 19:01:19 +03:00
"github.com/siddontang/go/arena"
2014-09-24 08:34:21 +04:00
"github.com/siddontang/go/hack"
2014-09-24 05:46:36 +04:00
"github.com/siddontang/go/log"
2014-09-24 09:29:27 +04:00
"github.com/siddontang/go/num"
"github.com/siddontang/ledisdb/ledis"
"io"
"net"
"runtime"
"strconv"
2014-10-30 06:11:45 +03:00
"time"
)
var errReadRequest = errors.New("invalid request protocol")
type respClient struct {
2014-08-25 10:18:23 +04:00
*client
conn net.Conn
rb *bufio.Reader
2014-10-29 19:01:19 +03:00
ar *arena.Arena
}
type respWriter struct {
buff *bufio.Writer
}
func newClientRESP(conn net.Conn, app *App) {
c := new(respClient)
2014-08-25 10:18:23 +04:00
c.client = newClient(app)
c.conn = conn
2014-10-25 10:48:52 +04:00
if tcpConn, ok := conn.(*net.TCPConn); ok {
tcpConn.SetReadBuffer(app.cfg.ConnReadBufferSize)
tcpConn.SetWriteBuffer(app.cfg.ConnWriteBufferSize)
}
2014-10-28 12:58:37 +03:00
2014-10-25 10:48:52 +04:00
c.rb = bufio.NewReaderSize(conn, app.cfg.ConnReadBufferSize)
2014-10-25 10:48:52 +04:00
c.resp = newWriterRESP(conn, app.cfg.ConnWriteBufferSize)
2014-08-25 10:18:23 +04:00
c.remoteAddr = conn.RemoteAddr().String()
2014-10-29 19:01:19 +03:00
//maybe another config?
c.ar = arena.NewArena(app.cfg.ConnReadBufferSize)
go c.run()
}
func (c *respClient) run() {
2014-10-30 05:51:02 +03:00
c.app.info.addClients(1)
defer func() {
c.client.close()
2014-10-30 05:51:02 +03:00
c.app.info.addClients(-1)
if e := recover(); e != nil {
buf := make([]byte, 4096)
n := runtime.Stack(buf, false)
buf = buf[0:n]
log.Fatal("client run panic %s:%v", buf, e)
}
2014-10-21 18:51:17 +04:00
handleQuit := true
2014-08-25 10:18:23 +04:00
if c.conn != nil {
2014-10-21 18:51:17 +04:00
//if handle quit command before, conn is nil
handleQuit = false
2014-08-25 10:18:23 +04:00
c.conn.Close()
}
2014-09-27 16:11:36 +04:00
if c.tx != nil {
c.tx.Rollback()
c.tx = nil
}
2014-10-21 18:51:17 +04:00
c.app.removeSlave(c.client, handleQuit)
}()
2014-10-30 06:11:45 +03:00
kc := time.Duration(c.app.cfg.ConnKeepaliveInterval) * time.Second
for {
2014-10-30 06:11:45 +03:00
if kc > 0 {
c.conn.SetReadDeadline(time.Now().Add(kc))
}
reqData, err := c.readRequest()
if err == nil {
c.handleRequest(reqData)
}
2014-10-30 04:03:58 +03:00
if err != nil {
return
}
2014-10-30 04:03:58 +03:00
if c.conn == nil {
return
}
}
}
func (c *respClient) readRequest() ([][]byte, error) {
2014-10-29 19:01:19 +03:00
return ReadRequest(c.rb, c.ar)
}
func (c *respClient) handleRequest(reqData [][]byte) {
if len(reqData) == 0 {
2014-08-25 10:18:23 +04:00
c.cmd = ""
c.args = reqData[0:0]
} else {
2014-10-30 07:48:52 +03:00
c.cmd = hack.String(lowerSlice(reqData[0]))
2014-08-25 10:18:23 +04:00
c.args = reqData[1:]
}
2014-08-25 10:18:23 +04:00
if c.cmd == "quit" {
c.resp.writeStatus(OK)
c.resp.flush()
2014-08-04 07:06:28 +04:00
c.conn.Close()
2014-10-21 18:51:17 +04:00
c.conn = nil
2014-08-04 07:37:21 +04:00
return
2014-08-04 07:06:28 +04:00
}
2014-08-25 10:18:23 +04:00
c.perform()
2014-08-01 05:38:08 +04:00
2014-10-29 19:01:19 +03:00
c.cmd = ""
c.args = nil
c.ar.Reset()
2014-08-01 05:38:08 +04:00
return
}
// response writer
2014-10-25 10:48:52 +04:00
func newWriterRESP(conn net.Conn, size int) *respWriter {
w := new(respWriter)
2014-10-25 10:48:52 +04:00
w.buff = bufio.NewWriterSize(conn, size)
return w
}
func (w *respWriter) writeError(err error) {
2014-09-24 08:34:21 +04:00
w.buff.Write(hack.Slice("-ERR"))
if err != nil {
w.buff.WriteByte(' ')
2014-09-24 08:34:21 +04:00
w.buff.Write(hack.Slice(err.Error()))
}
w.buff.Write(Delims)
}
func (w *respWriter) writeStatus(status string) {
w.buff.WriteByte('+')
2014-09-24 08:34:21 +04:00
w.buff.Write(hack.Slice(status))
w.buff.Write(Delims)
}
func (w *respWriter) writeInteger(n int64) {
w.buff.WriteByte(':')
2014-09-24 09:29:27 +04:00
w.buff.Write(num.FormatInt64ToSlice(n))
w.buff.Write(Delims)
}
func (w *respWriter) writeBulk(b []byte) {
w.buff.WriteByte('$')
if b == nil {
w.buff.Write(NullBulk)
} else {
2014-09-24 08:34:21 +04:00
w.buff.Write(hack.Slice(strconv.Itoa(len(b))))
w.buff.Write(Delims)
w.buff.Write(b)
}
w.buff.Write(Delims)
}
func (w *respWriter) writeArray(lst []interface{}) {
w.buff.WriteByte('*')
if lst == nil {
w.buff.Write(NullArray)
w.buff.Write(Delims)
} else {
2014-09-24 08:34:21 +04:00
w.buff.Write(hack.Slice(strconv.Itoa(len(lst))))
w.buff.Write(Delims)
for i := 0; i < len(lst); i++ {
switch v := lst[i].(type) {
case []interface{}:
w.writeArray(v)
2014-08-26 19:21:45 +04:00
case [][]byte:
w.writeSliceArray(v)
case []byte:
w.writeBulk(v)
case nil:
w.writeBulk(nil)
case int64:
w.writeInteger(v)
default:
panic("invalid array type")
}
}
}
}
func (w *respWriter) writeSliceArray(lst [][]byte) {
w.buff.WriteByte('*')
if lst == nil {
w.buff.Write(NullArray)
w.buff.Write(Delims)
} else {
2014-09-24 08:34:21 +04:00
w.buff.Write(hack.Slice(strconv.Itoa(len(lst))))
w.buff.Write(Delims)
for i := 0; i < len(lst); i++ {
w.writeBulk(lst[i])
}
}
}
func (w *respWriter) writeFVPairArray(lst []ledis.FVPair) {
w.buff.WriteByte('*')
if lst == nil {
w.buff.Write(NullArray)
w.buff.Write(Delims)
} else {
2014-09-24 08:34:21 +04:00
w.buff.Write(hack.Slice(strconv.Itoa(len(lst) * 2)))
w.buff.Write(Delims)
for i := 0; i < len(lst); i++ {
w.writeBulk(lst[i].Field)
w.writeBulk(lst[i].Value)
}
}
}
func (w *respWriter) writeScorePairArray(lst []ledis.ScorePair, withScores bool) {
w.buff.WriteByte('*')
if lst == nil {
w.buff.Write(NullArray)
w.buff.Write(Delims)
} else {
if withScores {
2014-09-24 08:34:21 +04:00
w.buff.Write(hack.Slice(strconv.Itoa(len(lst) * 2)))
w.buff.Write(Delims)
} else {
2014-09-24 08:34:21 +04:00
w.buff.Write(hack.Slice(strconv.Itoa(len(lst))))
w.buff.Write(Delims)
}
for i := 0; i < len(lst); i++ {
w.writeBulk(lst[i].Member)
if withScores {
2014-09-24 09:29:27 +04:00
w.writeBulk(num.FormatInt64ToSlice(lst[i].Score))
}
}
}
}
func (w *respWriter) writeBulkFrom(n int64, rb io.Reader) {
w.buff.WriteByte('$')
2014-09-24 08:34:21 +04:00
w.buff.Write(hack.Slice(strconv.FormatInt(n, 10)))
w.buff.Write(Delims)
io.Copy(w.buff, rb)
w.buff.Write(Delims)
}
func (w *respWriter) flush() {
w.buff.Flush()
}