Do not use websocket for device

Signed-off-by: Jianhui Zhao <zhaojh329@gmail.com>
This commit is contained in:
Jianhui Zhao
2020-01-30 22:14:34 +08:00
parent ebd9b0d095
commit 0846fceb99
19 changed files with 900 additions and 809 deletions
+9 -5
View File
@@ -1,12 +1,12 @@
# rttys([中文](/README_ZH.md))
[1]: https://img.shields.io/badge/license-LGPL2-brightgreen.svg?style=plastic
[1]: https://img.shields.io/badge/license-MIT-brightgreen.svg?style=plastic
[2]: /LICENSE
[3]: https://img.shields.io/badge/PRs-welcome-brightgreen.svg?style=plastic
[4]: https://github.com/zhaojh329/rttys/pulls
[5]: https://img.shields.io/badge/Issues-welcome-brightgreen.svg?style=plastic
[6]: https://github.com/zhaojh329/rttys/issues/new
[7]: https://img.shields.io/badge/release-2.10.2-blue.svg?style=plastic
[7]: https://img.shields.io/badge/release-3.0.0-blue.svg?style=plastic
[8]: https://github.com/zhaojh329/rttys/releases
[9]: https://travis-ci.org/zhaojh329/rttys.svg?branch=master
[10]: https://travis-ci.org/zhaojh329/rttys
@@ -27,9 +27,13 @@ This is the server program of [rtty](https://github.com/zhaojh329/rtty)
## See Supported Command Line Parameters
./rttys -h
Usage of ./rttys:
-addr string
address to listen (default ":5912")
Usage of rttys:
-addr-dev string
address to listen device (default ":5912")
-addr-user string
address to listen user (default ":5913")
-base-url string
base url to serve on (default "/")
-conf string
config file to load (default "./rttys.conf")
-gen-token
+9 -5
View File
@@ -1,12 +1,12 @@
# rttys
[1]: https://img.shields.io/badge/license-LGPL2-brightgreen.svg?style=plastic
[1]: https://img.shields.io/badge/license-MIT-brightgreen.svg?style=plastic
[2]: /LICENSE
[3]: https://img.shields.io/badge/PRs-welcome-brightgreen.svg?style=plastic
[4]: https://github.com/zhaojh329/rttys/pulls
[5]: https://img.shields.io/badge/Issues-welcome-brightgreen.svg?style=plastic
[6]: https://github.com/zhaojh329/rttys/issues/new
[7]: https://img.shields.io/badge/release-2.10.2-blue.svg?style=plastic
[7]: https://img.shields.io/badge/release-3.0.0-blue.svg?style=plastic
[8]: https://github.com/zhaojh329/rttys/releases
[9]: https://travis-ci.org/zhaojh329/rttys.svg?branch=master
[10]: https://travis-ci.org/zhaojh329/rttys
@@ -27,9 +27,13 @@
## 查看支持哪些命令行参数
./rttys -h
Usage of ./rttys:
-addr string
address to listen (default ":5912")
Usage of rttys:
-addr-dev string
address to listen device (default ":5912")
-addr-user string
address to listen user (default ":5913")
-base-url string
base url to serve on (default "/")
-conf string
config file to load (default "./rttys.conf")
-gen-token
+121 -115
View File
@@ -1,159 +1,165 @@
package main
import (
"fmt"
"github.com/gorilla/websocket"
"github.com/json-iterator/go"
jsoniter "github.com/json-iterator/go"
log "github.com/sirupsen/logrus"
"time"
)
const (
RTTY_MESSAGE_VERSION = 2
RTTY_MAX_SESSION_ID = 1000000
)
type Broker struct {
devices map[string]*Device
sessions map[string]*Session
connecting chan *Device /* Connecting requests from Device. */
disconnecting chan *Device /* Disconnecting requests from Device. */
logining chan *User /* Login requests from the User. */
logouting chan *User /* Logout requests from the User. */
inDevMessage chan *DevMessage /* Buffered channel of inbound messages from device. */
inUsrMessage chan *UsrMessage /* Buffered channel of inbound messages from user. */
}
type Session struct {
dev *Device
user *User
devsid uint8
devsid byte
}
func newBroker() *Broker {
type Broker struct {
token string
login chan *User
logout chan string
register chan *Device
unregister chan *Device
devices map[string]*Device
sessions map[string]*Session
commands map[string]*CommandStatus
newSession chan *Session
cmdReq chan *CommandReq
devMessage chan *DevMessage
userMessage chan *UsrMessage
cmdMessage chan []byte
clearCmd chan string
}
func newBroker(token string) *Broker {
return &Broker{
connecting: make(chan *Device, 100),
disconnecting: make(chan *Device, 100),
logining: make(chan *User, 100),
logouting: make(chan *User, 100),
devices: make(map[string]*Device),
sessions: make(map[string]*Session),
inDevMessage: make(chan *DevMessage, 1000),
inUsrMessage: make(chan *UsrMessage, 1000),
}
}
func (br *Broker) newSession(user *User) bool {
devid := user.devid
sid := genUniqueID("tty")
if dev, ok := br.devices[devid]; ok {
devsid := dev.getFreeSid()
if devsid < 1 {
log.Warn("Not found available devsid")
return false
}
br.sessions[sid] = &Session{dev, user, devsid}
dev.sessions[devsid] = sid
user.sid = sid
msg := fmt.Sprintf(`{"type":"login","sid":%d}`, devsid)
// Notify the device to create a pty and associate it with a session id
dev.wsWrite(websocket.TextMessage, []byte(msg))
log.Info("New session:", sid)
return true
} else {
// Notify the user that the device is offline
msg := `{"type":"login","err":1,"msg":"offline"}`
user.wsWrite(websocket.TextMessage, []byte(msg))
log.Info("Device", devid, "offline")
return false
token: token,
login: make(chan *User, 10),
logout: make(chan string, 10),
register: make(chan *Device, 1000),
unregister: make(chan *Device, 1000),
devices: make(map[string]*Device),
sessions: make(map[string]*Session),
newSession: make(chan *Session, 10),
commands: make(map[string]*CommandStatus),
cmdReq: make(chan *CommandReq, 1000),
devMessage: make(chan *DevMessage, 1000),
userMessage: make(chan *UsrMessage, 1000),
cmdMessage: make(chan []byte, 1000),
clearCmd: make(chan string, 1000),
}
}
func (br *Broker) run() {
for {
select {
case dev := <-br.connecting:
if _, ok := br.devices[dev.devid]; ok {
log.Warn("ID conflicting:", dev.devid)
dev.Close()
case dev := <-br.register:
err := byte(0)
msg := "OK"
if _, ok := br.devices[dev.id]; ok {
log.Error("Device ID conflicting: ", dev.id)
msg = "ID conflicting"
err = 1
} else if dev.token != br.token {
log.Error("Invalid token from terminal device")
msg = "Invalid token"
err = 1
} else {
br.devices[dev.devid] = dev
log.Info("New device:", dev.devid)
br.devices[dev.id] = dev
log.Info("New device: ", dev.id)
}
case dev := <-br.disconnecting:
if dev, ok := br.devices[dev.devid]; ok {
delete(br.devices, dev.devid)
dev.writeMsg(MsgTypeRegister, append([]byte{err}, msg...))
if err == 1 {
dev.close()
}
log.Info("Died device:", dev.devid)
case dev := <-br.unregister:
if _, ok := br.devices[dev.id]; ok {
delete(br.devices, dev.id)
}
for sid, session := range br.sessions {
if session.dev.devid == dev.devid {
session.user.Close()
delete(br.sessions, sid)
log.Info("Delete session: ", sid)
}
for sid, session := range br.sessions {
if session.dev == dev {
session.user.close()
delete(br.sessions, sid)
log.Info("Delete session: ", sid)
}
}
case user := <-br.logining:
if !br.newSession(user) {
time.AfterFunc(500*time.Millisecond, user.Close)
case user := <-br.login:
if dev, ok := br.devices[user.devid]; ok {
if !dev.login(user) {
user.loginAck(LoginErrorBusy)
log.Errorf("Device '%s' is busy", dev.id)
}
} else {
user.loginAck(LoginErrorOffline)
log.Errorf("Not found the device '%s'", user.devid)
}
case user := <-br.logouting:
if session, ok := br.sessions[user.sid]; ok {
devsid := session.devsid
dev := session.dev
sid := user.sid
msg := fmt.Sprintf(`{"type":"logout","sid":%d}`, devsid)
dev.wsWrite(websocket.TextMessage, []byte(msg))
case sid := <-br.logout:
if session, ok := br.sessions[sid]; ok {
delete(br.sessions, sid)
delete(session.dev.sessions, devsid)
session.user.close()
session.dev.logout(sid[len(sid)-1] - '0')
log.Info("Delete session: ", sid)
}
case msg := <-br.inDevMessage:
msgType := msg.msgType
data := msg.data
devsid := uint8(0)
if msgType == websocket.BinaryMessage {
devsid = data[0]
data = data[1:]
} else {
typ := jsoniter.Get(data, "type").ToString()
if typ == "cmd" {
handleCmdResp(data)
continue
}
val := jsoniter.Get(data, "sid").ToInt()
devsid = uint8(val)
}
sid := msg.dev.sessions[devsid]
case session := <-br.newSession:
sid := session.dev.id + string(session.devsid+'0')
session.user.sid = sid
session.user.loginAck(LoginErrorNone)
br.sessions[sid] = session
log.Infof("New session: %s", sid)
case msg := <-br.devMessage:
sid := msg.devid + string(msg.sid+'0')
if session, ok := br.sessions[sid]; ok {
session.user.wsWrite(msgType, data)
data := []byte{0}
if msg.isFileMsg {
data[0] = 1
}
session.user.writeMessage(websocket.BinaryMessage, append(data, msg.data...))
}
case msg := <-br.inUsrMessage:
case msg := <-br.userMessage:
msgType := msg.msgType
data := msg.data
if session, ok := br.sessions[msg.user.sid]; ok {
if session, ok := br.sessions[msg.sid]; ok {
devsid := msg.sid[len(msg.sid)-1] - '0'
if msgType == websocket.BinaryMessage {
data = append([]byte{session.devsid}, data...)
isFileMsg := data[0] == 1
data = data[1:]
if isFileMsg {
session.dev.writeMsg(MsgTypeFile, data)
} else {
session.dev.writeMsg(MsgTypeTermData, append([]byte{devsid}, data...))
}
} else {
typ := jsoniter.Get(msg.data, "type").ToString()
switch typ {
case "winsize":
cols := jsoniter.Get(msg.data, "cols").ToInt()
rows := jsoniter.Get(msg.data, "rows").ToInt()
data = append([]byte{devsid}, intToBytes(cols, 2)...)
data = append(data, intToBytes(rows, 2)...)
session.dev.writeMsg(MsgTypeWinsize, data)
}
}
session.dev.wsWrite(msgType, data)
} else {
log.Error("Not found sid: ", msg.sid)
}
case cmdReq := <-br.cmdReq:
handleCmdReq(br, cmdReq)
case data := <-br.cmdMessage:
handleCmdResp(br, data)
case token := <-br.clearCmd:
if cmd, ok := br.commands[token]; ok {
delete(br.commands, token)
cmd.tmr.Stop()
}
}
}
-128
View File
@@ -1,128 +0,0 @@
package main
import (
"net/http"
"strconv"
"sync"
"time"
"github.com/gorilla/websocket"
log "github.com/sirupsen/logrus"
)
var upgrader = websocket.Upgrader{
CheckOrigin: func(r *http.Request) bool {
return true
},
}
type Client struct {
br *Broker
ws *websocket.Conn
devid string
mutex sync.Mutex /* Avoid repeated closes */
closed bool
closeChan chan byte
outMessage chan *wsOutMessage /* Buffered channel of outbound messages */
}
type wsOutMessage struct {
msgType int
data []byte
}
func (c *Client) Close() {
defer c.mutex.Unlock()
c.mutex.Lock()
if !c.closed {
c.ws.Close()
c.closed = true
close(c.closeChan)
}
}
func (c *Client) wsWrite(msgType int, data []byte) {
c.outMessage <- &wsOutMessage{msgType, data}
}
func (c *Client) writePump() {
defer c.Close()
for {
select {
case msg := <-c.outMessage:
if err := c.ws.WriteMessage(msg.msgType, msg.data); err != nil {
return
}
case <-c.closeChan:
return
}
}
}
/* serveWs handles websocket requests from the device or user. */
func serveWs(br *Broker, w http.ResponseWriter, r *http.Request, cfg *RttysConfig) {
isDev := r.URL.Query().Get("device") != ""
if isDev {
token := r.Header.Get("Authorization")
if token != cfg.token {
log.Error("Invalid token from terminal device")
http.Error(w, "Forbidden", http.StatusForbidden)
return
}
} else if _, ok := httpSessions.Get(r.URL.Query().Get("sid")); !ok {
log.Error("Invalid sid from client")
http.Error(w, "Forbidden", http.StatusForbidden)
return
}
keepalive, _ := strconv.Atoi(r.URL.Query().Get("keepalive"))
devid := r.URL.Query().Get("devid")
conn, err := upgrader.Upgrade(w, r, nil)
if err != nil {
log.Println(err)
return
}
if devid == "" {
conn.Close()
log.Error("devid required")
return
}
client := &Client{
br: br,
devid: devid,
ws: conn,
closeChan: make(chan byte),
outMessage: make(chan *wsOutMessage, 1000),
}
if isDev {
desc := r.URL.Query().Get("description")
sessions := make(map[uint8]string)
dev := &Device{client, desc, time.Now().Unix(), sessions}
if keepalive > 0 {
go dev.keepAlive(int64(keepalive))
}
go dev.readAlway()
br.connecting <- dev
} else {
user := &User{client, ""}
go user.readAlway()
br.logining <- user
}
go client.writePump()
}
+77 -54
View File
@@ -3,116 +3,139 @@ package main
import (
"fmt"
"io"
"io/ioutil"
"net/http"
"sync"
"time"
jsoniter "github.com/json-iterator/go"
log "github.com/sirupsen/logrus"
"github.com/gorilla/websocket"
)
const CommandTimeout = time.Second * 30
const (
RTTY_CMD_ERR_INVALID = 1001
RTTY_CMD_ERR_OFFLINE = 1002
RTTY_CMD_ERR_BUSY = 1003
RTTY_CMD_ERR_TIMEOUT = 1004
RTTY_CMD_ERR_PENDING = 1005
RTTY_CMD_ERR_INVALID_TOKEN = 1006
RttyCmdErrInvalid = 1001
RttyCmdErrOffline = 1002
RttyCmdErrBusy = 1003
RttyCmdErrTimeout = 1004
RttyCmdErrPending = 1005
RttyCmdErrInvalidToken = 1006
)
var cmdErrMsg = map[int]string{
RTTY_CMD_ERR_INVALID: "invalid format",
RTTY_CMD_ERR_OFFLINE: "device offline",
RTTY_CMD_ERR_BUSY: "server is busy",
RTTY_CMD_ERR_TIMEOUT: "timeout",
RTTY_CMD_ERR_PENDING: "pending",
RTTY_CMD_ERR_INVALID_TOKEN: "invalid token",
RttyCmdErrInvalid: "invalid format",
RttyCmdErrOffline: "device offline",
RttyCmdErrBusy: "server is busy",
RttyCmdErrTimeout: "timeout",
RttyCmdErrPending: "pending",
RttyCmdErrInvalidToken: "invalid token",
}
type commandStatus struct {
type CommandStatus struct {
ts time.Time
token string
resp string
t *time.Timer
tmr *time.Timer
}
type CommandInfo struct {
Devid string `json:"devid"`
Cmd string `json:"cmd"`
Sid string `json:"sid"`
Username string `json:"username"`
Password string `json:"password"`
}
var commands sync.Map
type CommandReq struct {
done chan struct{}
token string
content []byte
w http.ResponseWriter
}
func handleCmdResp(data []byte) {
func handleCmdResp(br *Broker, data []byte) {
token := jsoniter.Get(data, "token").ToString()
if cmd, ok := commands.Load(token); ok {
cmd := cmd.(*commandStatus)
if cmd, ok := br.commands[token]; ok {
cmd.resp = jsoniter.Get(data, "attrs").ToString()
}
}
func cmdErrReply(err int, w http.ResponseWriter) {
fmt.Fprintf(w, `{"err": %d, "msg":"%s"}`, err, cmdErrMsg[err])
func cmdErrReply(err int, req *CommandReq) {
fmt.Fprintf(req.w, `{"err": %d, "msg":"%s"}`, err, cmdErrMsg[err])
close(req.done)
}
func serveCmd(br *Broker, w http.ResponseWriter, r *http.Request) {
token := r.URL.Query().Get("token")
func handleCmdReq(br *Broker, req *CommandReq) {
token := req.token
if token != "" {
cmd, ok := commands.Load(token)
if ok {
cmd := cmd.(*commandStatus)
if cmd, ok := br.commands[token]; ok {
if len(cmd.resp) == 0 {
cmdErrReply(RTTY_CMD_ERR_PENDING, w)
if time.Now().Sub(cmd.ts) > CommandTimeout {
cmdErrReply(RttyCmdErrTimeout, req)
} else {
cmdErrReply(RttyCmdErrPending, req)
}
} else {
commands.Delete(token)
io.WriteString(w, cmd.resp)
cmd.t.Stop()
io.WriteString(req.w, cmd.resp)
close(req.done)
br.clearCmd <- token
}
} else {
cmdErrReply(RTTY_CMD_ERR_INVALID_TOKEN, w)
cmdErrReply(RttyCmdErrInvalidToken, req)
}
return
}
body, err := ioutil.ReadAll(r.Body)
if err != nil {
log.Error(err)
cmdErrReply(RTTY_CMD_ERR_INVALID, w)
return
}
cmdInfo := CommandInfo{}
err = jsoniter.Unmarshal(body, &cmdInfo)
if _, ok := httpSessions.Get(cmdInfo.Sid); err != nil || cmdInfo.Cmd == "" || cmdInfo.Devid == "" || ok == false {
cmdErrReply(RTTY_CMD_ERR_INVALID, w)
err := jsoniter.Unmarshal(req.content, &cmdInfo)
if err != nil || cmdInfo.Cmd == "" || cmdInfo.Devid == "" {
cmdErrReply(RttyCmdErrInvalid, req)
return
}
dev, ok := br.devices[cmdInfo.Devid]
if !ok {
cmdErrReply(RTTY_CMD_ERR_OFFLINE, w)
cmdErrReply(RttyCmdErrOffline, req)
return
}
token = genUniqueID("cmd")
cmd := &commandStatus{
cmd := &CommandStatus{
ts: time.Now(),
token: token,
t: time.AfterFunc(30*time.Second, func() {
commands.Delete(token)
tmr: time.AfterFunc(CommandTimeout+time.Second*2, func() {
br.clearCmd <- token
}),
}
commands.Store(token, cmd)
br.commands[token] = cmd
msg := fmt.Sprintf(`{"type":"cmd","token":"%s","attrs":%s}`, token, body)
dev.wsWrite(websocket.TextMessage, []byte(msg))
username := jsoniter.Get(req.content, "username").ToString()
password := jsoniter.Get(req.content, "password").ToString()
cmdName := jsoniter.Get(req.content, "cmd").ToString()
params := jsoniter.Get(req.content, "params")
fmt.Fprintf(w, `{"token":"%s"}`, token)
var data []byte
data = append(data, username...)
data = append(data, 0)
data = append(data, password...)
data = append(data, 0)
data = append(data, cmdName...)
data = append(data, 0)
data = append(data, token...)
data = append(data, 0)
data = append(data, byte(params.Size()))
for i := 0; i < params.Size(); i++ {
data = append(data, params.Get(i).ToString()...)
data = append(data, 0)
}
dev.writeMsg(MsgTypeCmd, data)
fmt.Fprintf(req.w, `{"token":"%s"}`, token)
close(req.done)
}
+74
View File
@@ -0,0 +1,74 @@
package main
import (
"flag"
"github.com/kylelemons/go-gypsy/yaml"
log "github.com/sirupsen/logrus"
"os"
)
type RttysConfig struct {
addrDev string
addrUser string
sslCert string
sslKey string
username string
password string
token string
baseURL string
}
func setConfigOpt(yamlCfg *yaml.File, name string, opt *string) {
val, err := yamlCfg.Get(name)
if err != nil {
return
}
*opt = val
}
func parseConfig() *RttysConfig {
cfg := &RttysConfig{}
flag.StringVar(&cfg.addrDev, "addr-dev", ":5912", "address to listen device")
flag.StringVar(&cfg.addrUser, "addr-user", ":5913", "address to listen user")
flag.StringVar(&cfg.sslCert, "ssl-cert", "./rttys.crt", "certFile Path")
flag.StringVar(&cfg.sslKey, "ssl-key", "./rttys.key", "keyFile Path")
flag.StringVar(&cfg.token, "token", "", "token to use")
flag.StringVar(&cfg.baseURL, "base-url", "/", "base url to serve on")
conf := flag.String("conf", "./rttys.conf", "config file to load")
genToken := flag.Bool("gen-token", false, "generate token")
flag.Parse()
if *genToken {
genTokenAndExit()
}
yamlCfg, err := yaml.ReadFile(*conf)
if err == nil {
setConfigOpt(yamlCfg, "addr-dev", &cfg.addrDev)
setConfigOpt(yamlCfg, "addr-user", &cfg.addrUser)
setConfigOpt(yamlCfg, "ssl-cert", &cfg.sslCert)
setConfigOpt(yamlCfg, "ssl-key", &cfg.sslKey)
setConfigOpt(yamlCfg, "username", &cfg.username)
setConfigOpt(yamlCfg, "password", &cfg.password)
setConfigOpt(yamlCfg, "token", &cfg.token)
setConfigOpt(yamlCfg, "base-url", &cfg.baseURL)
}
if cfg.sslCert != "" && cfg.sslKey != "" {
_, err := os.Lstat(cfg.sslCert)
if err != nil {
log.Error(err)
cfg.sslCert = ""
}
_, err = os.Lstat(cfg.sslKey)
if err != nil {
log.Error(err)
cfg.sslKey = ""
}
}
return cfg
}
+227 -72
View File
@@ -1,104 +1,259 @@
package main
import (
"github.com/gorilla/websocket"
"bufio"
"crypto/rand"
"crypto/tls"
log "github.com/sirupsen/logrus"
"io"
"net"
"strings"
"sync"
"time"
)
const (
/* Max session id for each device */
RTTY_MAX_SESSION_ID_DEV = 5
MsgTypeRegister = 0x00
MsgTypeLogin = 0x01
MsgTypeLogout = 0x02
MsgTypeTermData = 0x03
MsgTypeWinsize = 0x04
MsgTypeCmd = 0x05
MsgTypeHeartbeat = 0x06
MsgTypeFile = 0x07
)
type Device struct {
*Client
desc string /* description of the device */
timestamp int64 /* Connection time */
sessions map[uint8]string /* sessions of each device */
}
const HeartbeatInterval = time.Second * 5
type DeviceInfo struct {
ID string `json:"id"`
Uptime int64 `json:"uptime"`
Description string `json:"description"`
type Device struct {
br *Broker
id string
desc string /* description of the device */
timestamp int64 /* Connection time */
token string
conn net.Conn
loginMutex sync.Mutex
user *User /* User who is wait login */
active time.Time
closeMutex sync.Mutex
closed bool
closeCh chan struct{}
}
type DevMessage struct {
msgType int
data []byte
dev *Device
devid string
sid uint8
data []byte
isFileMsg bool
}
func (dev *Device) Close() {
dev.Client.Close()
dev.br.disconnecting <- dev
}
func (dev *Device) login(user *User) bool {
defer dev.loginMutex.Unlock()
func (dev *Device) getFreeSid() uint8 {
for sid := uint8(1); sid <= RTTY_MAX_SESSION_ID_DEV; sid++ {
if _, ok := dev.sessions[sid]; !ok {
return sid
}
dev.loginMutex.Lock()
if dev.user != nil {
return false
}
return uint8(0)
dev.user = user
dev.writeMsg(MsgTypeLogin, []byte{})
return true
}
/*
* If the Server does not receive a PING Packet from the Client within one and
* a half times the Keep Alive time period, the server will disconnect the
* Connection
*/
func (dev *Device) keepAlive(keepalive int64) {
defer dev.Close()
func (dev *Device) handleLogin(code byte, sid byte) {
defer dev.loginMutex.Unlock()
ticker := time.NewTicker(time.Second * time.Duration(keepalive))
dev.loginMutex.Lock()
if dev.user == nil {
return
}
user := dev.user
dev.user = nil
if code == 1 {
log.Errorf("login fail, device busy")
user.loginAck(LoginErrorBusy)
return
}
dev.br.newSession <- &Session{dev, user, sid}
}
func (dev *Device) logout(sid byte) {
dev.writeMsg(MsgTypeLogout, []byte{sid})
}
func (dev *Device) keepAlive() {
ticker := time.NewTicker(time.Second)
defer ticker.Stop()
last := time.Now().Unix()
keepalive = keepalive * 3
/* Get the current ping handler */
pingHandler := dev.ws.PingHandler()
dev.ws.SetPingHandler(func(appData string) error {
last = time.Now().Unix()
return pingHandler(appData)
})
for {
select {
case <-dev.closeChan:
return
case <-ticker.C:
now := time.Now().Unix()
if now-last > keepalive {
log.Error("Inactive device in long time, now kill it: %s", dev.devid)
if time.Now().Sub(dev.active) > HeartbeatInterval*3/2 {
log.Errorf("Inactive device in long time, now kill it: %s %s", dev.id, time.Now())
dev.close()
return
}
}
}
}
func (dev *Device) readAlway() {
defer dev.Close()
for {
msgType, data, err := dev.ws.ReadMessage()
if err != nil {
if websocket.IsUnexpectedCloseError(err, websocket.CloseGoingAway, websocket.CloseAbnormalClosure) {
log.Errorf("error: %v", err)
}
break
}
msg := &DevMessage{msgType, data, dev}
select {
case dev.br.inDevMessage <- msg:
case <-dev.closeChan:
log.Error("closeChan from readAlway")
case <-dev.closeCh:
return
}
}
}
func (dev *Device) close() {
defer dev.closeMutex.Unlock()
dev.closeMutex.Lock()
if !dev.closed {
dev.closed = true
time.AfterFunc(time.Second, func() {
dev.br.unregister <- dev
close(dev.closeCh)
dev.conn.Close()
log.Infof("Device '%s' closed", dev.id)
})
}
}
func (dev *Device) writeMsg(typ byte, data []byte) {
b := []byte{typ}
b = append(b, intToBytes(len(data), 2)...)
b = append(b, data...)
dev.conn.Write(b)
}
func parseDeviceInfo(br *bufio.Reader) (string, string, string) {
id, _ := br.ReadString(0)
desc, _ := br.ReadString(0)
token, _ := br.ReadString(0)
id = strings.Trim(id, "\000")
desc = strings.Trim(desc, "\000")
token = strings.Trim(token, "\000")
return id, desc, token
}
func (dev *Device) readLoop() {
defer dev.close()
br := bufio.NewReaderSize(dev.conn, 4096+100)
msgLen := -1
for {
if msgLen < 0 {
b, err := br.Peek(3)
if err != nil {
if err != io.EOF && !strings.Contains(err.Error(), "use of closed network connection") {
log.Error(err)
}
return
}
msgLen = bytesToIntU(b[1:])
}
_, err := br.Peek(msgLen + 3)
if err != nil {
log.Error(err)
return
}
typ, _ := br.ReadByte()
br.Discard(2)
dev.active = time.Now()
switch typ {
case MsgTypeRegister:
id, desc, token := parseDeviceInfo(br)
dev.id = id
dev.token = token
dev.desc = desc
dev.br.register <- dev
case MsgTypeLogin:
code, _ := br.ReadByte()
sid := byte(0)
if code == 0 {
sid, _ = br.ReadByte()
}
dev.handleLogin(code, sid)
case MsgTypeLogout:
sid, _ := br.ReadByte()
dev.br.logout <- dev.id + string(sid+'0')
case MsgTypeTermData:
fallthrough
case MsgTypeFile:
sid, _ := br.ReadByte()
data := make([]byte, msgLen-1)
br.Read(data)
dev.br.devMessage <- &DevMessage{dev.id, sid, data, typ == MsgTypeFile}
case MsgTypeCmd:
data := make([]byte, msgLen)
br.Read(data)
dev.br.cmdMessage <- data
case MsgTypeHeartbeat:
dev.writeMsg(MsgTypeHeartbeat, []byte{})
default:
log.Error("invalid msg type")
br.Discard(msgLen)
}
msgLen = -1
}
}
func listenDevice(br *Broker, cfg *RttysConfig) {
ln, err := net.Listen("tcp", cfg.addrDev)
if err != nil {
log.Fatal(err)
}
defer ln.Close()
if cfg.sslCert != "" && cfg.sslKey != "" {
crt, err := tls.LoadX509KeyPair(cfg.sslCert, cfg.sslKey)
if err != nil {
log.Fatalln(err)
}
tlsConfig := &tls.Config{}
tlsConfig.Certificates = []tls.Certificate{crt}
tlsConfig.Time = time.Now
tlsConfig.Rand = rand.Reader
ln = tls.NewListener(ln, tlsConfig)
log.Info("Listen device on: ", cfg.addrDev, " SSL on")
} else {
log.Info("Listen device on: ", cfg.addrDev, " SSL off")
}
for {
conn, err := ln.Accept()
if err != nil {
log.Error(err)
continue
}
dev := &Device{
br: br,
conn: conn,
closeCh: make(chan struct{}),
active: time.Now(),
timestamp: time.Now().Unix(),
}
go dev.readLoop()
go dev.keepAlive()
}
}
+8 -3
View File
@@ -2698,7 +2698,8 @@
"version": "4.6.0",
"resolved": "https://registry.npmjs.org/co/-/co-4.6.0.tgz",
"integrity": "sha1-bqa989hTrlTMuOR7+gvz+QMfsYQ=",
"dev": true
"dev": true,
"optional": true
},
"coa": {
"version": "2.0.2",
@@ -6603,6 +6604,7 @@
"resolved": "https://registry.npmjs.org/levn/-/levn-0.3.0.tgz",
"integrity": "sha1-OwmSTt+fCDwEkP3UwLxEIeBHZO4=",
"dev": true,
"optional": true,
"requires": {
"prelude-ls": "~1.1.2",
"type-check": "~0.3.2"
@@ -8424,7 +8426,8 @@
"version": "1.1.2",
"resolved": "https://registry.npmjs.org/prelude-ls/-/prelude-ls-1.1.2.tgz",
"integrity": "sha1-IZMqVJ9eUv/ZqCf1cOBL5iqX2lQ=",
"dev": true
"dev": true,
"optional": true
},
"prepend-http": {
"version": "2.0.0",
@@ -9054,7 +9057,8 @@
"version": "4.0.8",
"resolved": "https://registry.npmjs.org/rx-lite/-/rx-lite-4.0.8.tgz",
"integrity": "sha1-Cx4Rr4vESDbwSmQH6S2kJGe3lEQ=",
"dev": true
"dev": true,
"optional": true
},
"rx-lite-aggregates": {
"version": "4.0.8",
@@ -10274,6 +10278,7 @@
"resolved": "https://registry.npmjs.org/type-check/-/type-check-0.3.2.tgz",
"integrity": "sha1-WITKtRLPHTVeP7eE8wgEsrUg23I=",
"dev": true,
"optional": true,
"requires": {
"prelude-ls": "~1.1.2"
}
+90 -203
View File
@@ -1,223 +1,110 @@
const blk_size = 8912; /* 8KB */
const MsgTypeFileStartDownload = 0x00;
const MsgTypeFileInfo = 0x01;
const MsgTypeFileData = 0x02;
const MsgTypeFileCanceled = 0x03;
const blk_size = 4096; /* 4KB */
function RttyFile(ws, term, opt) {
this.state = '';
this.ws = ws;
this.term = term;
this.cache = [];
function RttyFile(wsSend, on_download) {
this.buffer = [];
this.canceled = false;
this.to_term = function(octets) {
this.term.write(Buffer.from(octets).toString());
}
this.recvFileMsg = function(data) {
let type = data[0];
this.detect = function(input) {
let type = '';
data = data.slice(1);
if (input.byteLength < 3)
return '';
switch (type) {
case MsgTypeFileStartDownload:
on_download();
break;
input = new Uint8Array(input);
case MsgTypeFileInfo:
this.name = data.toString();
this.buffer = [];
break;
let pos = input.indexOf(0xB6);
if (pos < 0)
return '';
if (pos > input.length - 3)
return '';
if (input[pos + 1] != 0xBC)
return '';
type = String.fromCharCode(input[pos + 2]);
if (type != 's' && type != 'r')
return '';
if (pos > 0)
this.to_term(input.slice(0, pos));
if (type == 's')
this.cache = Array.prototype.slice.call(input.slice(pos + 3));
else
this.to_term(input.slice(pos + 3));
this.start_ts = new Date().getTime() / 1000;
return type;
}
this.consume = function(input) {
if (this.state == '') {
let t = this.detect(input);
if (t == '') {
this.to_term(Buffer.from(input));
return;
}
if (t == 'r') {
this.state = 'send_pending';
opt.on_detect('r');
} else if (t == 's') {
this.state = 'recving';
opt.on_detect('s');
}
} else if (this.state == 'recving' || this.state == 'abort_recv') {
this.recvFile(input);
} else if (this.state == 'sending') {
input = new Uint8Array(input);
if (input.length == 3) {
if (input[0] == 0xB6 && input[1] == 0xBC && String.fromCharCode(input[2]) == 'e') {
this.state = 'abort';
return;
}
}
this.to_term(Buffer.from(input));
} else {
this.to_term(Buffer.from(input));
}
}
this.readFile = function(offset, size) {
let blob = this.file.slice(offset, offset + size);
this.fr.readAsArrayBuffer(blob);
}
this.sendEof = function() {
let b = new Uint8Array([0x03]);
this.ws.send(b);
this.state = '';
}
this.abort = function() {
this.state = 'abort';
}
this.abortRecv = function() {
let b = new Uint8Array([0x03]);
this.ws.send(b);
this.state = 'abort_recv';
}
this.sendInfo = function(file) {
let b = Buffer.alloc(6 + file.name.length);
b[0] = 0x01; /* packet type: file info */
b[1] = file.name.length;
b.write(file.name, 2);
b.writeUInt32BE(file.size, 2 + file.name.length);
this.ws.send(b);
this.file = file;
}
this.sendData = function(data) {
let b = Buffer.alloc(3);
let piece = new Uint8Array(data);
b[0] = 0x02; /* packet type: file data */
b.writeUInt16BE(piece.length, 1);
this.ws.send(b);
this.ws.send(piece);
}
this.sendFile = function(file) {
this.fr = new FileReader();
this.state = 'sending';
let offset = 0;
this.sendInfo(file);
this.fr.onload = (e) => {
this.sendData(e.target.result);
offset += e.loaded;
if (this.state != 'abort' && offset < file.size) {
this.readFile(offset, blk_size);
return;
}
this.sendEof();
};
this.readFile(offset, blk_size);
}
this.recvFile = function(input) {
input = Array.prototype.slice.call(new Uint8Array(input));
this.cache.push.apply(this.cache, input);
while (this.cache.length > 0) {
let type = this.cache[0];
switch (type) {
case 0x01: /* file info */ {
if (this.cache.length < 2)
return;
let nl = this.cache[1];
if (this.cache.length < nl + 2)
return;
this.cache.splice(0, 2);
this.name = Buffer.from(this.cache.splice(0, nl)).toString();
this.size = Buffer.from(this.cache.splice(0, 4)).readUInt32BE(0);
this.offset = 0;
this.buffer = [];
break;
}
case 0x02: /* file data */ {
if (this.cache.length < 3)
return;
let dl = Buffer.from(this.cache.slice(1,3)).readUInt16BE(0);
if (this.cache.length < dl + 3)
return;
this.cache.splice(0, 3);
this.buffer.push(new Uint8Array(this.cache.splice(0, dl)));
this.offset += dl;
let now_ts = new Date().getTime() / 1000;
let unit = 'K';
let offset = this.offset / 1024;
if (offset / 1024 > 0) {
offset /= 1024;
unit = 'M';
}
this.term.write(' %d%% %.2f %sB %.3fs\r'.format(this.offset / this.size * 100, offset, unit, now_ts - this.start_ts));
break;
}
case 0x03: /* file eof */ {
this.cache = [];
this.term.write('\n');
if (this.state == 'abort_recv') {
this.state = '';
return;
}
this.state = '';
case MsgTypeFileData:
if (data.length === 0) {
let blob = new Blob(this.buffer);
let url = URL.createObjectURL(blob);
let el = document.createElement("a");
el.style.display = "none";
let el = document.createElement('a');
el.style.display = 'none';
el.href = url;
el.download = this.name;
document.body.appendChild(el);
el.click();
document.body.removeChild(el);
break;
this.buffer = [];
} else {
this.buffer.push(data);
}
default:
// console.error('invalid type:' + type);
break;
case MsgTypeFileCanceled:
this.buffer = [];
this.canceled = true;
break;
}
return true;
}
this.cancel = function() {
let b = Buffer.from([1, MsgTypeFileCanceled]);
wsSend(b);
}
this.sendInfo = function() {
let buf = [];
buf.push(Buffer.from([1, MsgTypeFileInfo]));
let sb = new Buffer(4);
sb.writeUInt32BE(this.file.size, 0);
buf.push(sb);
buf.push(Buffer.from(this.file.name));
wsSend(Buffer.concat(buf));
}
this.sendData = function(data) {
let b = Buffer.concat([Buffer.from([1, MsgTypeFileData]), Buffer.from(data)]);
wsSend(b);
}
this.readFile = function(offset, size) {
let blob = this.file.slice(offset, offset + size);
this.fr.readAsArrayBuffer(blob);
}
this.sendFile = function(file) {
this.fr = new FileReader();
this.file = file;
this.canceled = false;
this.sendInfo();
let offset = 0;
this.fr.onload = e => {
this.sendData(e.target.result);
offset += e.loaded;
if (this.canceled)
return;
if (offset < file.size) {
this.readFile(offset, blk_size);
return;
}
}
this.sendData([]);
};
this.readFile(offset, blk_size);
}
}
+13 -35
View File
@@ -3,9 +3,7 @@
<Button style="margin-right: 4px;" type="primary" shape="circle" icon="md-refresh" @click="handleRefresh" :disabled="loading">{{$t('Refresh List')}}</Button>
<Input style="margin-right: 4px;width:200px" v-model="filterString" icon="search" size="large" @on-change="handleSearch" :placeholder="$t('Please enter the filter key...')" />
<Button style="margin-right: 4px;" @click="showCmdForm" type="primary" :disabled="cmdStatus.execing > 0"><Icon type="search"/>{{$t('executive command')}}</Button>
<div class="counter">
{{ $t('device-count', {count: devlists.length}) }}
</div>
<div class="counter">{{ $t('device-count', {count: devlists.length}) }}</div>
<Table :height="tableHeight" :loading="loading" :columns="devlistTitle" :data="filtered" style="margin-top: 10px; width: 100%" :no-data-text="$t('No devices connected')" @on-selection-change='handleSelection'>
<template slot-scope="{ row }" slot="uptime">
<span>{{ '%t'.format(row.uptime) }}</span>
@@ -17,22 +15,18 @@
<Modal v-model="cmdModal" :title="$t('executive command')">
<Form :model="cmdData" ref="cmdForm" :rules="cmdRuleValidate" :label-width="80">
<FormItem :label="$t('Username')" prop="username">
<Input v-model="cmdData.username"></Input>
<Input v-model="cmdData.username"/>
</FormItem>
<FormItem :label="$t('Password')" prop="password">
<Input type="password" v-model="cmdData.password"></Input>
<Input type="password" v-model="cmdData.password"/>
</FormItem>
<FormItem :label="$t('Command')" prop="cmd">
<Input v-model="cmdData.cmd"></Input>
<Input v-model="cmdData.cmd"/>
</FormItem>
<FormItem :label="$t('Parameter')" prop="params">
<Tag v-for="(item, index) in cmdData.params" :key="item + index" closable @on-close="handleDelCmdParam(index)" :fade="false">{{ item }}</Tag>
<Input v-model="cmdData.currentParam" icon="md-add-circle" :placeholder="$t('Please enter a single parameter')" @on-click="handleAddCmdParam" @on-keyup.enter="handleAddCmdParam" />
</FormItem>
<FormItem :label="$t('Environment variable')" prop="env">
<Tag v-for="(v, k) in cmdData.env" :key="v + k" closable @on-close="handleDelCmdEnv(k)" :fade="false">{{ k + '=' + v }}</Tag>
<Input v-model="cmdData.currentEnv" icon="md-add-circle" :placeholder="$t('Please enter a single environment')" @on-click="handleAddCmdEnv" @on-keyup.enter="handleAddCmdEnv" />
</FormItem>
</Form>
<div slot="footer">
<Button type="primary" @click="doCmd">{{$t('OK')}}</Button>
@@ -40,16 +34,16 @@
</div>
</Modal>
<Modal v-model="cmdStatus.modal" :title="$t('status of executive command')" :closable="false" :mask-closable="false">
<Progress :percent="cmdStatusPercent" status="active"></Progress>
<Progress :percent="cmdStatusPercent" status="active"/>
<p>{{ $t('cmd-status-total', {count: cmdStatus.total}) }}</p>
<p>{{ $t('cmd-status-fail', {count: cmdStatus.fail}) }}</p>
<div slot="footer">
<Button type="primary" size="large" :disabled="cmdStatus.execing > 0" @click="showCmdResp">{{$t('OK')}}</Button>
<Button type="error" size="large" :disabled="cmdStatus.execing == 0" @click="ignoreCmdResp">{{$t('Ignore')}}</Button>
<Button type="error" size="large" :disabled="cmdStatus.execing === 0" @click="ignoreCmdResp">{{$t('Ignore')}}</Button>
</div>
</Modal>
<Modal v-model="cmdStatus.respModal" :title="$t('Response of executive command')" :width="1000">
<Table :columns="cmdStatus.response.columns" :data="cmdStatus.response.data" height="300" :no-data-text="$t('No Response')"></Table>
<Table :columns="cmdStatus.response.columns" :data="cmdStatus.response.data" height="300" :no-data-text="$t('No Response')"/>
<div slot="footer"></div>
</Modal>
</div>
@@ -149,9 +143,7 @@ export default {
password: '',
cmd: '',
params: [],
currentParam: '',
env: {},
currentEnv: ''
currentParam: ''
},
cmdRuleValidate: {
username: [
@@ -212,12 +204,12 @@ export default {
this.$axios.get(process.env.BASE_URL + 'cmd?token=' + token).then((response) => {
let resp = response.data;
if (resp.err == 1005) {
if (resp.err === 1005) {
item.querying = false;
return;
}
if (resp.err && resp.err != 0)
if (resp.err && resp.err !== 0)
this.cmdStatus.fail++;
this.cmdStatus.execing--;
@@ -249,23 +241,11 @@ export default {
},
handleAddCmdParam() {
this.cmdData.currentParam = this.cmdData.currentParam.trim();
if (this.cmdData.currentParam != '') {
if (this.cmdData.currentParam !== '') {
this.cmdData.params.push(this.cmdData.currentParam);
this.cmdData.currentParam = '';
}
},
handleDelCmdEnv(key) {
this.$delete(this.cmdData.env, key);
},
handleAddCmdEnv() {
this.cmdData.currentEnv = this.cmdData.currentEnv.trim();
if (this.cmdData.currentEnv != '') {
let e = this.cmdData.currentEnv.split('=');
if (e.length == 2)
this.$set(this.cmdData.env, [e[0]], e[1]);
this.cmdData.currentEnv = '';
}
},
doCmd() {
this.$refs['cmdForm'].validate((valid) => {
if (valid) {
@@ -284,8 +264,7 @@ export default {
password: this.cmdData.password,
sid: sessionStorage.getItem('rtty-sid'),
cmd: this.cmdData.cmd.trim(),
params: this.cmdData.params,
env: this.cmdData.env
params: this.cmdData.params
};
this.$axios.post(process.env.BASE_URL + 'cmd', data).then((response) => {
@@ -330,8 +309,7 @@ export default {
},
computed: {
cmdStatusPercent() {
let percent = (this.cmdStatus.total - this.cmdStatus.execing) / this.cmdStatus.total * 100;
return parseInt(percent);
return (this.cmdStatus.total - this.cmdStatus.execing) / this.cmdStatus.total * 100;
}
},
mounted() {
+35 -46
View File
@@ -25,6 +25,9 @@ import RttyFile from '../plugins/rtty-file'
Terminal.applyAddon(fit);
Terminal.applyAddon(overlay);
const LoginErrorOffline = 0x01;
const LoginErrorBusy = 0x02;
export default {
name: 'Rtty',
data() {
@@ -62,7 +65,7 @@ export default {
},
cancelUpfile() {
this.term.focus();
this.rf.sendEof();
this.rf.cancel();
},
doUpload() {
if (!this.upfile.file) {
@@ -73,6 +76,9 @@ export default {
this.upfile.modal = false;
this.term.focus();
this.rf.sendFile(this.upfile.file);
},
wsSendData(type, data) {
this.ws.send(Buffer.concat([Buffer.from([type]), Buffer.from(data)]));
}
},
mounted() {
@@ -107,90 +113,73 @@ export default {
term.on('resize', (size) => {
setTimeout(() => {
let msg = {type: "winsize", sid: this.sid, cols: size.cols, rows: size.rows};
let msg = {type: "winsize", cols: size.cols, rows: size.rows};
ws.send(JSON.stringify(msg));
term.showOverlay(size.cols + 'x' + size.rows);
}, 500);
});
this.term = term;
this.rf = new RttyFile(ws, term, {
on_detect: (t) => {
if (t == 'r')
this.upfile.modal = true;
else if (t == 's')
;
}
this.rf = new RttyFile(data => {
this.ws.send(data);
}, () => {
this.upfile.modal = true;
});
};
ws.onmessage = (ev) => {
let term = this.term;
if (typeof ev.data == 'string') {
if (typeof ev.data === 'string') {
let msg = JSON.parse(ev.data);
if (msg.type == "login") {
if (msg.err == 1) {
if (msg.type === "login") {
if (msg.err === LoginErrorOffline) {
this.$Message.error(this.$t('Device offline'));
this.logout();
return;
} else if (msg.err == 2) {
} else if (msg.err === LoginErrorBusy) {
this.$Message.error(this.$t('Sessions is full'));
this.logout();
return;
}
this.sid = msg.sid;
msg = {type: 'winsize', sid: this.sid, cols: term.cols, rows: term.rows};
msg = {type: 'winsize', cols: term.cols, rows: term.rows};
ws.send(JSON.stringify(msg));
term.on('data', (data) => {
if (this.rf.state != '') {
if (data.length == 1) {
let key = data.charCodeAt(0);
/* Ctrl + C, Esc */
if (key == 3 || key == 27) {
if (this.rf.state == 'recving') {
this.rf.abortRecv();
} else {
this.upfile.modal = false;
if (this.rf.state == 'send_pending')
this.rf.sendEof();
else
this.rf.abort();
}
}
}
return;
}
this.ws.send(Buffer.from(data));
this.wsSendData(0, data);
});
} else if (msg.type == 'logout') {
} else if (msg.type === 'logout') {
this.logout();
}
} else {
let data = Buffer.from(ev.data);
let isFileMsg = data[0] === 1;
if (isFileMsg) {
this.rf.recvFileMsg(data.slice(1));
return;
}
data = data.slice(1).toString();
if (!this.recvTTYCnt)
this.recvTTYCnt = 0;
this.recvTTYCnt++;
if (this.recvTTYCnt < 4) {
let data = Buffer.from(ev.data).toString();
if (data.match('login:') && this.username && this.username != '') {
ws.send(Buffer.from(this.username + '\n'));
if (data.match('login:') && this.username && this.username !== '') {
this.wsSendData(0, this.username + '\n');
return;
}
if (data.match('Password:') && this.password && this.password != '') {
ws.send(Buffer.from(this.password + '\n'));
if (data.match('Password:') && this.password && this.password !== '') {
this.wsSendData(0, this.password + '\n');
return;
}
}
this.rf.consume(ev.data);
term.write(data);
}
};
+8 -4
View File
@@ -3,13 +3,17 @@ module.exports = {
devServer: {
proxy: {
'/devs': {
target: 'http://127.0.0.1:5912'
target: 'http://127.0.0.1:5913'
},
'/login': {
target: 'http://127.0.0.1:5912'
'/signin': {
target: 'http://127.0.0.1:5913'
},
'/cmd': {
target: 'http://127.0.0.1:5912'
target: 'http://127.0.0.1:5913'
},
'/ws': {
ws: true,
target: 'http://127.0.0.1:5913'
}
}
}
+54 -31
View File
@@ -2,17 +2,16 @@ package main
import (
"fmt"
"net/http"
"os"
"strconv"
"time"
jsoniter "github.com/json-iterator/go"
"github.com/rakyll/statik/fs"
log "github.com/sirupsen/logrus"
"github.com/zhaojh329/rttys/cache"
"github.com/zhaojh329/rttys/pwauth"
_ "github.com/zhaojh329/rttys/statik"
"io/ioutil"
"net/http"
"strconv"
"time"
)
type Credentials struct {
@@ -82,22 +81,54 @@ func httpStart(br *Broker, cfg *RttysConfig) {
}
http.HandleFunc(cfg.baseURL+"/ws", func(w http.ResponseWriter, r *http.Request) {
serveWs(br, w, r, cfg)
if _, ok := httpSessions.Get(r.URL.Query().Get("sid")); !ok {
http.Error(w, "Invalid sid", http.StatusForbidden)
return
}
serveUser(br, w, r)
})
http.HandleFunc(cfg.baseURL+"/cmd", func(w http.ResponseWriter, r *http.Request) {
allowOrigin(w)
serveCmd(br, w, r)
done := make(chan struct{})
req := &CommandReq{
done: done,
w: w,
}
if r.Method == "GET" {
req.token = r.URL.Query().Get("token")
} else if r.Method == "POST" {
content, err := ioutil.ReadAll(r.Body)
if err != nil {
log.Error(err)
return
}
sid := jsoniter.Get(content, "sid").ToString()
if _, ok := httpSessions.Get(sid); !ok {
http.Error(w, "Forbidden", http.StatusForbidden)
return
}
req.content = content
} else {
http.Error(w, "MethodNotAllowed", http.StatusMethodNotAllowed)
return
}
br.cmdReq <- req
<-done
})
http.HandleFunc(cfg.baseURL+"/signin", func(w http.ResponseWriter, r *http.Request) {
var creds Credentials
// Get the JSON body and decode into credentials
err := jsoniter.NewDecoder(r.Body).Decode(&creds)
if err != nil {
// If the structure of the body is wrong, return an HTTP error
w.WriteHeader(http.StatusBadRequest)
http.Error(w, "Bad Request", http.StatusBadRequest)
return
}
@@ -118,11 +149,17 @@ func httpStart(br *Broker, cfg *RttysConfig) {
})
http.HandleFunc(cfg.baseURL+"/devs", func(w http.ResponseWriter, r *http.Request) {
type DeviceInfo struct {
ID string `json:"id"`
Uptime int64 `json:"uptime"`
Description string `json:"description"`
}
if !httpAuth(w, r) {
return
}
devs := []DeviceInfo{}
devs := make([]DeviceInfo, 0)
for id, dev := range br.devices {
dev := DeviceInfo{id, time.Now().Unix() - dev.timestamp, dev.desc}
@@ -138,11 +175,11 @@ func httpStart(br *Broker, cfg *RttysConfig) {
hfunc := http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path == "/" {
t := r.URL.Query().Get("t")
t := r.URL.Query().Get("tmr")
id := r.URL.Query().Get("id")
if t == "" && id == "" {
http.Redirect(w, r, cfg.baseURL+"?t="+strconv.FormatInt(time.Now().Unix(), 10), http.StatusFound)
http.Redirect(w, r, cfg.baseURL+"?tmr="+strconv.FormatInt(time.Now().Unix(), 10), http.StatusFound)
return
}
}
@@ -157,24 +194,10 @@ func httpStart(br *Broker, cfg *RttysConfig) {
}
if cfg.sslCert != "" && cfg.sslKey != "" {
_, err := os.Lstat(cfg.sslCert)
if err != nil {
log.Error(err)
cfg.sslCert = ""
}
_, err = os.Lstat(cfg.sslKey)
if err != nil {
log.Error(err)
cfg.sslKey = ""
}
}
if cfg.sslCert != "" && cfg.sslKey != "" {
log.Info("Listen on: ", cfg.addr, " SSL on")
log.Fatal(http.ListenAndServeTLS(cfg.addr, cfg.sslCert, cfg.sslKey, nil))
log.Info("Listen user on: ", cfg.addrUser, " SSL on")
log.Fatal(http.ListenAndServeTLS(cfg.addrUser, cfg.sslCert, cfg.sslKey, nil))
} else {
log.Info("Listen on: ", cfg.addr, " SSL off")
log.Fatal(http.ListenAndServe(cfg.addr, nil))
log.Info("Listen user on: ", cfg.addrUser, " SSL off")
log.Fatal(http.ListenAndServe(cfg.addrUser, nil))
}
}
+7 -89
View File
@@ -1,37 +1,16 @@
package main
import (
"crypto/md5"
"crypto/rand"
"encoding/binary"
"encoding/hex"
"flag"
"fmt"
"io"
"os"
"runtime"
"time"
"golang.org/x/crypto/ssh/terminal"
"github.com/kylelemons/go-gypsy/yaml"
"github.com/zhaojh329/rttys/version"
"github.com/howeyc/gopass"
"github.com/rifflock/lfshook"
log "github.com/sirupsen/logrus"
"github.com/zhaojh329/rttys/version"
"golang.org/x/crypto/ssh/terminal"
)
type RttysConfig struct {
addr string
sslCert string
sslKey string
username string
password string
token string
baseURL string
}
func init() {
if terminal.IsTerminal(int(os.Stdout.Fd())) {
return
@@ -42,7 +21,6 @@ func init() {
func main() {
cfg := parseConfig()
log.Info("Go Version: ", runtime.Version())
log.Info("Go OS/Arch: ", runtime.GOOS, "/", runtime.GOARCH)
@@ -50,73 +28,13 @@ func main() {
log.Info("Git Commit: ", version.GitCommit())
log.Info("Build Time: ", version.BuildTime())
br := newBroker()
br := newBroker(cfg.token)
go br.run()
httpStart(br, cfg)
}
go listenDevice(br, cfg)
go httpStart(br, cfg)
func genUniqueID(extra string) string {
buf := make([]byte, 20)
binary.BigEndian.PutUint32(buf, uint32(time.Now().Unix()))
io.ReadFull(rand.Reader, buf[4:])
h := md5.New()
h.Write(buf)
h.Write([]byte(extra))
return hex.EncodeToString(h.Sum(nil))
}
func setConfigOpt(yamlCfg *yaml.File, name string, opt *string) {
val, err := yamlCfg.Get(name)
if err != nil {
return
for {
time.Sleep(time.Second)
}
*opt = val
}
func parseConfig() *RttysConfig {
cfg := &RttysConfig{}
flag.StringVar(&cfg.addr, "addr", ":5912", "address to listen")
flag.StringVar(&cfg.sslCert, "ssl-cert", "./rttys.crt", "certFile Path")
flag.StringVar(&cfg.sslKey, "ssl-key", "./rttys.key", "keyFile Path")
flag.StringVar(&cfg.token, "token", "", "token to use")
flag.StringVar(&cfg.baseURL, "base-url", "/", "base url to serve on")
conf := flag.String("conf", "./rttys.conf", "config file to load")
genToken := flag.Bool("gen-token", false, "generate token")
flag.Parse()
if *genToken {
genTokenAndExit()
}
yamlCfg, err := yaml.ReadFile(*conf)
if err == nil {
setConfigOpt(yamlCfg, "addr", &cfg.addr)
setConfigOpt(yamlCfg, "ssl-cert", &cfg.sslCert)
setConfigOpt(yamlCfg, "ssl-key", &cfg.sslKey)
setConfigOpt(yamlCfg, "username", &cfg.username)
setConfigOpt(yamlCfg, "password", &cfg.password)
setConfigOpt(yamlCfg, "token", &cfg.token)
setConfigOpt(yamlCfg, "base-url", &cfg.baseURL)
}
return cfg
}
func genTokenAndExit() {
password, err := gopass.GetPasswdPrompt("Please set a password:", true, os.Stdin, os.Stdout)
if err != nil {
log.Fatal(err)
}
token := genUniqueID(string(password))
fmt.Println("Your token is:", token)
os.Exit(0)
}
+2 -1
View File
@@ -1,4 +1,5 @@
#addr: :5912
#addr-dev: :5912
#addr-user: :5913
# default from system
#username: rttys
+1 -1
View File
File diff suppressed because one or more lines are too long
+79 -16
View File
@@ -1,36 +1,99 @@
package main
import (
"fmt"
"github.com/gorilla/websocket"
log "github.com/sirupsen/logrus"
"sync"
"net/http"
)
const (
LoginErrorNone = 0x00
LoginErrorOffline = 0x01
LoginErrorBusy = 0x02
)
type User struct {
*Client
sid string
br *Broker
sid string
devid string
conn *websocket.Conn
closeMutex sync.Mutex
closed bool
}
type UsrMessage struct {
sid string
msgType int
data []byte
user *User
}
func (user *User) Close() {
user.Client.Close()
user.br.logouting <- user
var upgrader = websocket.Upgrader{
CheckOrigin: func(r *http.Request) bool {
return true
},
}
func (user *User) readAlway() {
defer user.Close()
func (u *User) writeMessage(messageType int, data []byte) {
u.conn.WriteMessage(messageType, data)
}
func (u *User) close() {
defer u.closeMutex.Unlock()
u.closeMutex.Lock()
if !u.closed {
u.closed = true
u.conn.Close()
u.br.logout <- u.sid
}
}
func (u *User) loginAck(code int) {
msg := fmt.Sprintf(`{"type":"login","err":%d}`, code)
u.writeMessage(websocket.TextMessage, []byte(msg))
}
func (u *User) readLoop() {
defer u.close()
for {
msgType, data, err := user.ws.ReadMessage()
msgType, data, err := u.conn.ReadMessage()
if err != nil {
if websocket.IsUnexpectedCloseError(err, websocket.CloseGoingAway, websocket.CloseAbnormalClosure) {
log.Error(err)
}
break
}
msg := &UsrMessage{msgType, data, user}
select {
case user.br.inUsrMessage <- msg:
case <-user.closeChan:
return
}
u.br.userMessage <- &UsrMessage{u.sid, msgType, data}
}
}
func serveUser(br *Broker, w http.ResponseWriter, r *http.Request) {
devid := r.URL.Query().Get("devid")
if devid == "" {
http.Error(w, "devid required", http.StatusForbidden)
return
}
conn, err := upgrader.Upgrade(w, r, nil)
if err != nil {
http.Error(w, "Upgrade fail", http.StatusBadRequest)
log.Error(err)
return
}
user := &User{
br: br,
conn: conn,
devid: devid,
}
go user.readLoop()
br.login <- user
}
+85
View File
@@ -0,0 +1,85 @@
package main
import (
"bytes"
"crypto/md5"
"crypto/rand"
"encoding/binary"
"encoding/hex"
"fmt"
"github.com/howeyc/gopass"
log "github.com/sirupsen/logrus"
"io"
"os"
"time"
)
func bytesToIntU(b []byte) int {
if len(b) == 3 {
b = append([]byte{0}, b...)
}
bytesBuffer := bytes.NewBuffer(b)
switch len(b) {
case 1:
var tmp uint8
binary.Read(bytesBuffer, binary.BigEndian, &tmp)
return int(tmp)
case 2:
var tmp uint16
binary.Read(bytesBuffer, binary.BigEndian, &tmp)
return int(tmp)
case 4:
var tmp uint32
binary.Read(bytesBuffer, binary.BigEndian, &tmp)
return int(tmp)
default:
return 0
}
}
func intToBytes(n int, b byte) []byte {
switch b {
case 1:
tmp := int8(n)
bytesBuffer := bytes.NewBuffer([]byte{})
binary.Write(bytesBuffer, binary.BigEndian, &tmp)
return bytesBuffer.Bytes()
case 2:
tmp := int16(n)
bytesBuffer := bytes.NewBuffer([]byte{})
binary.Write(bytesBuffer, binary.BigEndian, &tmp)
return bytesBuffer.Bytes()
case 3, 4:
tmp := int32(n)
bytesBuffer := bytes.NewBuffer([]byte{})
binary.Write(bytesBuffer, binary.BigEndian, &tmp)
return bytesBuffer.Bytes()
}
return []byte{}
}
func genUniqueID(extra string) string {
buf := make([]byte, 20)
binary.BigEndian.PutUint32(buf, uint32(time.Now().Unix()))
io.ReadFull(rand.Reader, buf[4:])
h := md5.New()
h.Write(buf)
h.Write([]byte(extra))
return hex.EncodeToString(h.Sum(nil))
}
func genTokenAndExit() {
password, err := gopass.GetPasswdPrompt("Please set a password:", true, os.Stdin, os.Stdout)
if err != nil {
log.Fatal(err)
}
token := genUniqueID(string(password))
fmt.Println("Your token is:", token)
os.Exit(0)
}
+1 -1
View File
@@ -1,6 +1,6 @@
package version
const version = "2.10.3"
const version = "3.0.0"
var (
gitCommit = ""