mirror of
https://github.com/wavetermdev/backup.git
synced 2026-08-05 13:57:07 -07:00
port to electron (#33)
This commit is contained in:
+217
@@ -0,0 +1,217 @@
|
||||
// Copyright 2024, Command Line Inc.
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package web
|
||||
|
||||
import (
|
||||
"encoding/base64"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"io/fs"
|
||||
"log"
|
||||
"net/http"
|
||||
"runtime/debug"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/gorilla/mux"
|
||||
"github.com/wavetermdev/thenextwave/pkg/filestore"
|
||||
"github.com/wavetermdev/thenextwave/pkg/service"
|
||||
"github.com/wavetermdev/thenextwave/pkg/wavebase"
|
||||
)
|
||||
|
||||
type WebFnType = func(http.ResponseWriter, *http.Request)
|
||||
|
||||
// Header constants
|
||||
const (
|
||||
CacheControlHeaderKey = "Cache-Control"
|
||||
CacheControlHeaderNoCache = "no-cache"
|
||||
|
||||
ContentTypeHeaderKey = "Content-Type"
|
||||
ContentTypeJson = "application/json"
|
||||
ContentTypeBinary = "application/octet-stream"
|
||||
|
||||
ContentLengthHeaderKey = "Content-Length"
|
||||
LastModifiedHeaderKey = "Last-Modified"
|
||||
|
||||
WaveZoneFileInfoHeaderKey = "X-ZoneFileInfo"
|
||||
)
|
||||
|
||||
const HttpReadTimeout = 5 * time.Second
|
||||
const HttpWriteTimeout = 21 * time.Second
|
||||
const HttpMaxHeaderBytes = 60000
|
||||
const HttpTimeoutDuration = 21 * time.Second
|
||||
|
||||
const MainServerAddr = "127.0.0.1:1719" // wavesrv, P=16+1, S=19, PS=1719
|
||||
const WebSocketServerAddr = "127.0.0.1:1723" // wavesrv:websocket, P=16+1, W=23, PW=1723
|
||||
const MainServerDevAddr = "127.0.0.1:8190"
|
||||
const WebSocketServerDevAddr = "127.0.0.1:8191"
|
||||
const WSStateReconnectTime = 30 * time.Second
|
||||
const WSStatePacketChSize = 20
|
||||
|
||||
type WebFnOpts struct {
|
||||
AllowCaching bool
|
||||
JsonErrors bool
|
||||
}
|
||||
|
||||
func handleService(w http.ResponseWriter, r *http.Request) {
|
||||
bodyData, err := io.ReadAll(r.Body)
|
||||
if err != nil {
|
||||
http.Error(w, "Unable to read request body", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
defer r.Body.Close()
|
||||
if r.Method != http.MethodPost {
|
||||
http.Error(w, "Invalid request method", http.StatusMethodNotAllowed)
|
||||
return
|
||||
}
|
||||
var webCall service.WebCallType
|
||||
err = json.Unmarshal(bodyData, &webCall)
|
||||
if err != nil {
|
||||
http.Error(w, fmt.Sprintf("invalid request body: %v", err), http.StatusBadRequest)
|
||||
}
|
||||
|
||||
rtn := service.CallService(r.Context(), webCall)
|
||||
jsonRtn, err := json.Marshal(rtn)
|
||||
if err != nil {
|
||||
http.Error(w, fmt.Sprintf("error serializing response: %v", err), http.StatusInternalServerError)
|
||||
}
|
||||
w.Header().Set(ContentTypeHeaderKey, ContentTypeJson)
|
||||
w.Header().Set(ContentLengthHeaderKey, fmt.Sprintf("%d", len(jsonRtn)))
|
||||
w.WriteHeader(http.StatusOK)
|
||||
w.Write(jsonRtn)
|
||||
}
|
||||
|
||||
func marshalReturnValue(data any, err error) []byte {
|
||||
var mapRtn = make(map[string]any)
|
||||
if err != nil {
|
||||
mapRtn["error"] = err.Error()
|
||||
} else {
|
||||
mapRtn["success"] = true
|
||||
mapRtn["data"] = data
|
||||
}
|
||||
rtn, err := json.Marshal(mapRtn)
|
||||
if err != nil {
|
||||
return marshalReturnValue(nil, fmt.Errorf("error serializing response: %v", err))
|
||||
}
|
||||
return rtn
|
||||
}
|
||||
|
||||
func handleWaveFile(w http.ResponseWriter, r *http.Request) {
|
||||
zoneId := r.URL.Query().Get("zoneid")
|
||||
name := r.URL.Query().Get("name")
|
||||
if _, err := uuid.Parse(zoneId); err != nil {
|
||||
http.Error(w, fmt.Sprintf("invalid zoneid: %v", err), http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
if name == "" {
|
||||
http.Error(w, "name is required", http.StatusBadRequest)
|
||||
return
|
||||
|
||||
}
|
||||
file, err := filestore.WFS.Stat(r.Context(), zoneId, name)
|
||||
if err == fs.ErrNotExist {
|
||||
http.NotFound(w, r)
|
||||
return
|
||||
}
|
||||
if err != nil {
|
||||
http.Error(w, fmt.Sprintf("error getting file info: %v", err), http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
jsonFileBArr, err := json.Marshal(file)
|
||||
if err != nil {
|
||||
http.Error(w, fmt.Sprintf("error serializing file info: %v", err), http.StatusInternalServerError)
|
||||
}
|
||||
// can make more efficient by checking modtime + If-Modified-Since headers to allow caching
|
||||
w.Header().Set(ContentTypeHeaderKey, ContentTypeBinary)
|
||||
w.Header().Set(ContentLengthHeaderKey, fmt.Sprintf("%d", file.Size))
|
||||
w.Header().Set(WaveZoneFileInfoHeaderKey, base64.StdEncoding.EncodeToString(jsonFileBArr))
|
||||
w.Header().Set(LastModifiedHeaderKey, time.UnixMilli(file.ModTs).UTC().Format(http.TimeFormat))
|
||||
if file.Size == 0 {
|
||||
w.WriteHeader(http.StatusOK)
|
||||
return
|
||||
}
|
||||
for offset := file.DataStartIdx(); offset < file.Size; offset += filestore.DefaultPartDataSize {
|
||||
_, data, err := filestore.WFS.ReadAt(r.Context(), zoneId, name, offset, filestore.DefaultPartDataSize)
|
||||
if err != nil {
|
||||
if offset == 0 {
|
||||
http.Error(w, fmt.Sprintf("error reading file: %v", err), http.StatusInternalServerError)
|
||||
} else {
|
||||
// nothing to do, the headers have already been sent
|
||||
log.Printf("error reading file %s/%s @ %d: %v\n", zoneId, name, offset, err)
|
||||
}
|
||||
return
|
||||
}
|
||||
w.Write(data)
|
||||
}
|
||||
}
|
||||
|
||||
func handleStreamFile(w http.ResponseWriter, r *http.Request) {
|
||||
fileName := r.URL.Query().Get("path")
|
||||
fileName = wavebase.ExpandHomeDir(fileName)
|
||||
http.ServeFile(w, r, fileName)
|
||||
}
|
||||
|
||||
func WebFnWrap(opts WebFnOpts, fn WebFnType) WebFnType {
|
||||
return func(w http.ResponseWriter, r *http.Request) {
|
||||
defer func() {
|
||||
recErr := recover()
|
||||
if recErr == nil {
|
||||
return
|
||||
}
|
||||
panicStr := fmt.Sprintf("panic: %v", recErr)
|
||||
log.Printf("panic: %v\n", recErr)
|
||||
debug.PrintStack()
|
||||
if opts.JsonErrors {
|
||||
jsonRtn := marshalReturnValue(nil, fmt.Errorf(panicStr))
|
||||
w.Header().Set(ContentTypeHeaderKey, ContentTypeJson)
|
||||
w.Header().Set(ContentLengthHeaderKey, fmt.Sprintf("%d", len(jsonRtn)))
|
||||
w.WriteHeader(http.StatusOK)
|
||||
w.Write(jsonRtn)
|
||||
} else {
|
||||
http.Error(w, panicStr, http.StatusInternalServerError)
|
||||
}
|
||||
}()
|
||||
if !opts.AllowCaching {
|
||||
w.Header().Set(CacheControlHeaderKey, CacheControlHeaderNoCache)
|
||||
}
|
||||
// reqAuthKey := r.Header.Get("X-AuthKey")
|
||||
// if reqAuthKey == "" {
|
||||
// w.WriteHeader(http.StatusInternalServerError)
|
||||
// w.Write([]byte("no x-authkey header"))
|
||||
// return
|
||||
// }
|
||||
// if reqAuthKey != scbase.WaveAuthKey {
|
||||
// w.WriteHeader(http.StatusInternalServerError)
|
||||
// w.Write([]byte("x-authkey header is invalid"))
|
||||
// return
|
||||
// }
|
||||
fn(w, r)
|
||||
}
|
||||
}
|
||||
|
||||
// blocking
|
||||
// TODO: create listener separately and use http.Serve, so we can signal SIGUSR1 in a better way
|
||||
func RunWebServer() {
|
||||
gr := mux.NewRouter()
|
||||
gr.HandleFunc("/wave/stream-file", WebFnWrap(WebFnOpts{AllowCaching: true}, handleStreamFile))
|
||||
gr.HandleFunc("/wave/file", WebFnWrap(WebFnOpts{AllowCaching: false}, handleWaveFile))
|
||||
gr.HandleFunc("/wave/service", WebFnWrap(WebFnOpts{JsonErrors: true}, handleService))
|
||||
serverAddr := MainServerAddr
|
||||
if wavebase.IsDevMode() {
|
||||
serverAddr = MainServerDevAddr
|
||||
}
|
||||
server := &http.Server{
|
||||
Addr: serverAddr,
|
||||
ReadTimeout: HttpReadTimeout,
|
||||
WriteTimeout: HttpWriteTimeout,
|
||||
MaxHeaderBytes: HttpMaxHeaderBytes,
|
||||
Handler: http.TimeoutHandler(gr, HttpTimeoutDuration, "Timeout"),
|
||||
}
|
||||
log.Printf("Running main server on %s\n", serverAddr)
|
||||
err := server.ListenAndServe()
|
||||
if err != nil {
|
||||
log.Printf("ERROR: %v\n", err)
|
||||
}
|
||||
}
|
||||
+215
@@ -0,0 +1,215 @@
|
||||
// Copyright 2024, Command Line Inc.
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package web
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"log"
|
||||
"net/http"
|
||||
"runtime/debug"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/gorilla/mux"
|
||||
"github.com/gorilla/websocket"
|
||||
"github.com/wavetermdev/thenextwave/pkg/eventbus"
|
||||
)
|
||||
|
||||
const wsReadWaitTimeout = 15 * time.Second
|
||||
const wsWriteWaitTimeout = 10 * time.Second
|
||||
const wsPingPeriodTickTime = 10 * time.Second
|
||||
const wsInitialPingTime = 1 * time.Second
|
||||
|
||||
func RunWebSocketServer() {
|
||||
gr := mux.NewRouter()
|
||||
gr.HandleFunc("/ws", HandleWs)
|
||||
serverAddr := WebSocketServerDevAddr
|
||||
server := &http.Server{
|
||||
Addr: serverAddr,
|
||||
ReadTimeout: HttpReadTimeout,
|
||||
WriteTimeout: HttpWriteTimeout,
|
||||
MaxHeaderBytes: HttpMaxHeaderBytes,
|
||||
Handler: gr,
|
||||
}
|
||||
server.SetKeepAlivesEnabled(false)
|
||||
log.Printf("Running websocket server on %s\n", serverAddr)
|
||||
err := server.ListenAndServe()
|
||||
if err != nil {
|
||||
log.Printf("[error] trying to run websocket server: %v\n", err)
|
||||
}
|
||||
}
|
||||
|
||||
var WebSocketUpgrader = websocket.Upgrader{
|
||||
ReadBufferSize: 4 * 1024,
|
||||
WriteBufferSize: 32 * 1024,
|
||||
HandshakeTimeout: 1 * time.Second,
|
||||
CheckOrigin: func(r *http.Request) bool { return true },
|
||||
}
|
||||
|
||||
func HandleWs(w http.ResponseWriter, r *http.Request) {
|
||||
err := HandleWsInternal(w, r)
|
||||
if err != nil {
|
||||
http.Error(w, err.Error(), http.StatusInternalServerError)
|
||||
}
|
||||
}
|
||||
|
||||
func getMessageType(jmsg map[string]any) string {
|
||||
if str, ok := jmsg["type"].(string); ok {
|
||||
return str
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
func getStringFromMap(jmsg map[string]any, key string) string {
|
||||
if str, ok := jmsg[key].(string); ok {
|
||||
return str
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
func processMessage(jmsg map[string]any, outputCh chan any) {
|
||||
msgType := getMessageType(jmsg)
|
||||
if msgType != "rpc" {
|
||||
return
|
||||
}
|
||||
reqId := getStringFromMap(jmsg, "reqid")
|
||||
var rtnErr error
|
||||
defer func() {
|
||||
r := recover()
|
||||
if r != nil {
|
||||
rtnErr = fmt.Errorf("panic: %v", r)
|
||||
log.Printf("panic in processMessage: %v\n", r)
|
||||
debug.PrintStack()
|
||||
}
|
||||
if rtnErr == nil {
|
||||
return
|
||||
}
|
||||
rtn := map[string]any{"type": "rpcresp", "reqid": reqId, "error": rtnErr.Error()}
|
||||
outputCh <- rtn
|
||||
}()
|
||||
method := getStringFromMap(jmsg, "method")
|
||||
rtnErr = fmt.Errorf("unknown method %q", method)
|
||||
}
|
||||
|
||||
func ReadLoop(conn *websocket.Conn, outputCh chan any, closeCh chan any) {
|
||||
readWait := wsReadWaitTimeout
|
||||
conn.SetReadLimit(64 * 1024)
|
||||
conn.SetReadDeadline(time.Now().Add(readWait))
|
||||
defer close(closeCh)
|
||||
for {
|
||||
_, message, err := conn.ReadMessage()
|
||||
if err != nil {
|
||||
log.Printf("ReadPump error: %v\n", err)
|
||||
break
|
||||
}
|
||||
jmsg := map[string]any{}
|
||||
err = json.Unmarshal(message, &jmsg)
|
||||
if err != nil {
|
||||
log.Printf("Error unmarshalling json: %v\n", err)
|
||||
break
|
||||
}
|
||||
conn.SetReadDeadline(time.Now().Add(readWait))
|
||||
msgType := getMessageType(jmsg)
|
||||
if msgType == "pong" {
|
||||
// nothing
|
||||
continue
|
||||
}
|
||||
if msgType == "ping" {
|
||||
now := time.Now()
|
||||
pongMessage := map[string]interface{}{"type": "pong", "stime": now.UnixMilli()}
|
||||
outputCh <- pongMessage
|
||||
continue
|
||||
}
|
||||
go processMessage(jmsg, outputCh)
|
||||
}
|
||||
}
|
||||
|
||||
func WritePing(conn *websocket.Conn) error {
|
||||
now := time.Now()
|
||||
pingMessage := map[string]interface{}{"type": "ping", "stime": now.UnixMilli()}
|
||||
jsonVal, _ := json.Marshal(pingMessage)
|
||||
_ = conn.SetWriteDeadline(time.Now().Add(wsWriteWaitTimeout)) // no error
|
||||
err := conn.WriteMessage(websocket.TextMessage, jsonVal)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func WriteLoop(conn *websocket.Conn, outputCh chan any, closeCh chan any) {
|
||||
ticker := time.NewTicker(wsInitialPingTime)
|
||||
defer ticker.Stop()
|
||||
initialPing := true
|
||||
for {
|
||||
select {
|
||||
case msg := <-outputCh:
|
||||
var barr []byte
|
||||
var err error
|
||||
if _, ok := msg.([]byte); ok {
|
||||
barr = msg.([]byte)
|
||||
} else {
|
||||
barr, err = json.Marshal(msg)
|
||||
if err != nil {
|
||||
log.Printf("cannot marshal websocket message: %v\n", err)
|
||||
// just loop again
|
||||
break
|
||||
}
|
||||
}
|
||||
err = conn.WriteMessage(websocket.TextMessage, barr)
|
||||
if err != nil {
|
||||
conn.Close()
|
||||
log.Printf("WritePump error: %v\n", err)
|
||||
return
|
||||
}
|
||||
|
||||
case <-ticker.C:
|
||||
err := WritePing(conn)
|
||||
if err != nil {
|
||||
log.Printf("WritePump error: %v\n", err)
|
||||
return
|
||||
}
|
||||
if initialPing {
|
||||
initialPing = false
|
||||
ticker.Reset(wsPingPeriodTickTime)
|
||||
}
|
||||
|
||||
case <-closeCh:
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func HandleWsInternal(w http.ResponseWriter, r *http.Request) error {
|
||||
windowId := r.URL.Query().Get("windowid")
|
||||
if windowId == "" {
|
||||
return fmt.Errorf("windowid is required")
|
||||
}
|
||||
conn, err := WebSocketUpgrader.Upgrade(w, r, nil)
|
||||
if err != nil {
|
||||
return fmt.Errorf("WebSocket Upgrade Failed: %v", err)
|
||||
}
|
||||
defer conn.Close()
|
||||
wsConnId := uuid.New().String()
|
||||
log.Printf("New websocket connection: windowid:%s connid:%s\n", windowId, wsConnId)
|
||||
outputCh := make(chan any, 100)
|
||||
closeCh := make(chan any)
|
||||
eventbus.RegisterWSChannel(wsConnId, windowId, outputCh)
|
||||
defer eventbus.UnregisterWSChannel(wsConnId)
|
||||
wg := &sync.WaitGroup{}
|
||||
wg.Add(2)
|
||||
go func() {
|
||||
// read loop
|
||||
defer wg.Done()
|
||||
ReadLoop(conn, outputCh, closeCh)
|
||||
}()
|
||||
go func() {
|
||||
// write loop
|
||||
defer wg.Done()
|
||||
WriteLoop(conn, outputCh, closeCh)
|
||||
}()
|
||||
wg.Wait()
|
||||
return nil
|
||||
}
|
||||
Reference in New Issue
Block a user