mirror of
https://github.com/wavetermdev/backup.git
synced 2026-08-05 13:57:07 -07:00
working on wave OSC escapes, modes for the terminal (#46)
This commit is contained in:
@@ -25,6 +25,7 @@ var CommandToTypeMap = map[string]reflect.Type{
|
||||
BlockCommand_Input: reflect.TypeOf(BlockInputCommand{}),
|
||||
BlockCommand_SetView: reflect.TypeOf(BlockSetViewCommand{}),
|
||||
BlockCommand_SetMeta: reflect.TypeOf(BlockSetMetaCommand{}),
|
||||
BlockCommand_Message: reflect.TypeOf(BlockMessageCommand{}),
|
||||
}
|
||||
|
||||
func CommandTypeUnionMeta() tsgenmeta.TypeUnionMeta {
|
||||
@@ -35,6 +36,7 @@ func CommandTypeUnionMeta() tsgenmeta.TypeUnionMeta {
|
||||
reflect.TypeOf(BlockInputCommand{}),
|
||||
reflect.TypeOf(BlockSetViewCommand{}),
|
||||
reflect.TypeOf(BlockSetMetaCommand{}),
|
||||
reflect.TypeOf(BlockMessageCommand{}),
|
||||
},
|
||||
}
|
||||
}
|
||||
@@ -96,3 +98,12 @@ type BlockSetMetaCommand struct {
|
||||
func (smc *BlockSetMetaCommand) GetCommand() string {
|
||||
return BlockCommand_SetMeta
|
||||
}
|
||||
|
||||
type BlockMessageCommand struct {
|
||||
Command string `json:"command" tstype:"\"message\""`
|
||||
Message string `json:"message"`
|
||||
}
|
||||
|
||||
func (bmc *BlockMessageCommand) GetCommand() string {
|
||||
return BlockCommand_Message
|
||||
}
|
||||
|
||||
@@ -11,6 +11,7 @@ import (
|
||||
"fmt"
|
||||
"io"
|
||||
"log"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
@@ -39,6 +40,7 @@ type BlockController struct {
|
||||
InputCh chan BlockCommand
|
||||
Status string
|
||||
|
||||
PtyBuffer *PtyBuffer
|
||||
ShellProc *shellexec.ShellProc
|
||||
ShellInputCh chan *BlockInputCommand
|
||||
}
|
||||
@@ -90,7 +92,7 @@ func (bc *BlockController) Close() {
|
||||
|
||||
const DefaultTermMaxFileSize = 256 * 1024
|
||||
|
||||
func (bc *BlockController) handleShellProcData(data []byte, seqNum int) error {
|
||||
func (bc *BlockController) handleShellProcData(data []byte) error {
|
||||
ctx, cancelFn := context.WithTimeout(context.Background(), DefaultTimeout)
|
||||
defer cancelFn()
|
||||
err := filestore.WFS.AppendData(ctx, bc.BlockId, "main", data)
|
||||
@@ -104,7 +106,6 @@ func (bc *BlockController) handleShellProcData(data []byte, seqNum int) error {
|
||||
"blockid": bc.BlockId,
|
||||
"blockfile": "main",
|
||||
"ptydata": base64.StdEncoding.EncodeToString(data),
|
||||
"seqnum": seqNum,
|
||||
},
|
||||
})
|
||||
return nil
|
||||
@@ -159,15 +160,13 @@ func (bc *BlockController) DoRunShellCommand(rc *RunShellOpts) error {
|
||||
bc.ShellProc = nil
|
||||
bc.ShellInputCh = nil
|
||||
}()
|
||||
seqNum := 0
|
||||
buf := make([]byte, 4096)
|
||||
for {
|
||||
nr, err := bc.ShellProc.Pty.Read(buf)
|
||||
seqNum++
|
||||
if nr > 0 {
|
||||
handleDataErr := bc.handleShellProcData(buf[:nr], seqNum)
|
||||
if handleDataErr != nil {
|
||||
log.Printf("error handling shell data: %v\n", handleDataErr)
|
||||
bc.PtyBuffer.AppendData(buf[:nr])
|
||||
if bc.PtyBuffer.Err != nil {
|
||||
log.Printf("error processing pty data: %v\n", bc.PtyBuffer.Err)
|
||||
break
|
||||
}
|
||||
}
|
||||
@@ -269,6 +268,15 @@ func StartBlockController(ctx context.Context, blockId string) error {
|
||||
Status: "init",
|
||||
InputCh: make(chan BlockCommand),
|
||||
}
|
||||
ptyBuffer := MakePtyBuffer(bc.handleShellProcData, func(cmd BlockCommand) error {
|
||||
if strings.HasPrefix(cmd.GetCommand(), "controller:") {
|
||||
bc.InputCh <- cmd
|
||||
} else {
|
||||
ProcessStaticCommand(blockId, cmd)
|
||||
}
|
||||
return nil
|
||||
})
|
||||
bc.PtyBuffer = ptyBuffer
|
||||
blockControllerMap[blockId] = bc
|
||||
go bc.Run(blockData)
|
||||
return nil
|
||||
@@ -304,6 +312,21 @@ func ProcessStaticCommand(blockId string, cmdGen BlockCommand) error {
|
||||
if err != nil {
|
||||
return fmt.Errorf("error updating block: %w", err)
|
||||
}
|
||||
// send a waveobj:update event
|
||||
updatedBlock, err := wstore.DBGet[*wstore.Block](ctx, blockId)
|
||||
if err != nil {
|
||||
return fmt.Errorf("error getting block: %w", err)
|
||||
}
|
||||
eventbus.SendEvent(eventbus.WSEventType{
|
||||
EventType: "waveobj:update",
|
||||
ORef: waveobj.MakeORef(wstore.OType_Block, blockId).String(),
|
||||
Data: wstore.WaveObjUpdate{
|
||||
UpdateType: wstore.UpdateType_Update,
|
||||
OType: wstore.OType_Block,
|
||||
OID: blockId,
|
||||
Obj: updatedBlock,
|
||||
},
|
||||
})
|
||||
return nil
|
||||
case *BlockSetMetaCommand:
|
||||
log.Printf("SETMETA: %s | %v\n", blockId, cmd.Meta)
|
||||
@@ -314,6 +337,9 @@ func ProcessStaticCommand(blockId string, cmdGen BlockCommand) error {
|
||||
if block == nil {
|
||||
return nil
|
||||
}
|
||||
if block.Meta == nil {
|
||||
block.Meta = make(map[string]any)
|
||||
}
|
||||
for k, v := range cmd.Meta {
|
||||
if v == nil {
|
||||
delete(block.Meta, k)
|
||||
@@ -325,6 +351,24 @@ func ProcessStaticCommand(blockId string, cmdGen BlockCommand) error {
|
||||
if err != nil {
|
||||
return fmt.Errorf("error updating block: %w", err)
|
||||
}
|
||||
// send a waveobj:update event
|
||||
updatedBlock, err := wstore.DBGet[*wstore.Block](ctx, blockId)
|
||||
if err != nil {
|
||||
return fmt.Errorf("error getting block: %w", err)
|
||||
}
|
||||
eventbus.SendEvent(eventbus.WSEventType{
|
||||
EventType: "waveobj:update",
|
||||
ORef: waveobj.MakeORef(wstore.OType_Block, blockId).String(),
|
||||
Data: wstore.WaveObjUpdate{
|
||||
UpdateType: wstore.UpdateType_Update,
|
||||
OType: wstore.OType_Block,
|
||||
OID: blockId,
|
||||
Obj: updatedBlock,
|
||||
},
|
||||
})
|
||||
return nil
|
||||
case *BlockMessageCommand:
|
||||
log.Printf("MESSAGE: %s | %q\n", blockId, cmd.Message)
|
||||
return nil
|
||||
default:
|
||||
return fmt.Errorf("unknown command type %T", cmdGen)
|
||||
|
||||
@@ -0,0 +1,119 @@
|
||||
// Copyright 2024, Command Line Inc.
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package blockcontroller
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
|
||||
"github.com/wavetermdev/thenextwave/pkg/wshutil"
|
||||
)
|
||||
|
||||
const (
|
||||
Mode_Normal = "normal"
|
||||
Mode_Esc = "esc"
|
||||
Mode_WaveEsc = "waveesc"
|
||||
)
|
||||
|
||||
type PtyBuffer struct {
|
||||
Mode string
|
||||
EscSeqBuf []byte
|
||||
DataOutputFn func([]byte) error
|
||||
CommandOutputFn func(BlockCommand) error
|
||||
Err error
|
||||
}
|
||||
|
||||
func MakePtyBuffer(dataOutputFn func([]byte) error, commandOutputFn func(BlockCommand) error) *PtyBuffer {
|
||||
return &PtyBuffer{
|
||||
Mode: Mode_Normal,
|
||||
DataOutputFn: dataOutputFn,
|
||||
CommandOutputFn: commandOutputFn,
|
||||
}
|
||||
}
|
||||
|
||||
func (b *PtyBuffer) setErr(err error) {
|
||||
if b.Err == nil {
|
||||
b.Err = err
|
||||
}
|
||||
}
|
||||
|
||||
func (b *PtyBuffer) processWaveEscSeq(escSeq []byte) {
|
||||
jmsg := make(map[string]any)
|
||||
err := json.Unmarshal(escSeq, &jmsg)
|
||||
if err != nil {
|
||||
b.setErr(fmt.Errorf("error unmarshalling Wave OSC sequence data: %w", err))
|
||||
return
|
||||
}
|
||||
cmd, err := ParseCmdMap(jmsg)
|
||||
if err != nil {
|
||||
b.setErr(fmt.Errorf("error parsing Wave OSC command: %w", err))
|
||||
return
|
||||
}
|
||||
err = b.CommandOutputFn(cmd)
|
||||
if err != nil {
|
||||
b.setErr(fmt.Errorf("error processing Wave OSC command: %w", err))
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
func (b *PtyBuffer) AppendData(data []byte) {
|
||||
outputBuf := make([]byte, 0, len(data))
|
||||
for _, ch := range data {
|
||||
if b.Mode == Mode_WaveEsc {
|
||||
if ch == wshutil.ESC {
|
||||
// terminates the escape sequence (and the rest was invalid)
|
||||
b.Mode = Mode_Normal
|
||||
outputBuf = append(outputBuf, b.EscSeqBuf...)
|
||||
outputBuf = append(outputBuf, ch)
|
||||
b.EscSeqBuf = nil
|
||||
} else if ch == wshutil.BEL || ch == wshutil.ST {
|
||||
// terminates the escpae sequence (is a valid Wave OSC command)
|
||||
b.Mode = Mode_Normal
|
||||
waveEscSeq := b.EscSeqBuf[len(wshutil.WaveOSCPrefix):]
|
||||
b.EscSeqBuf = nil
|
||||
b.processWaveEscSeq(waveEscSeq)
|
||||
} else {
|
||||
b.EscSeqBuf = append(b.EscSeqBuf, ch)
|
||||
}
|
||||
continue
|
||||
}
|
||||
if b.Mode == Mode_Esc {
|
||||
if ch == wshutil.ESC || ch == wshutil.BEL || ch == wshutil.ST {
|
||||
// these all terminate the escape sequence (invalid, not a Wave OSC)
|
||||
b.Mode = Mode_Normal
|
||||
outputBuf = append(outputBuf, b.EscSeqBuf...)
|
||||
outputBuf = append(outputBuf, ch)
|
||||
} else {
|
||||
if ch == wshutil.WaveOSCPrefixBytes[len(b.EscSeqBuf)] {
|
||||
// we're still building what could be a Wave OSC sequence
|
||||
b.EscSeqBuf = append(b.EscSeqBuf, ch)
|
||||
} else {
|
||||
// this is not a Wave OSC sequence, just an escape sequence
|
||||
b.Mode = Mode_Normal
|
||||
outputBuf = append(outputBuf, b.EscSeqBuf...)
|
||||
outputBuf = append(outputBuf, ch)
|
||||
continue
|
||||
}
|
||||
// check to see if we have a full Wave OSC prefix
|
||||
if len(b.EscSeqBuf) == len(wshutil.WaveOSCPrefixBytes) {
|
||||
b.Mode = Mode_WaveEsc
|
||||
}
|
||||
}
|
||||
continue
|
||||
}
|
||||
// Mode_Normal
|
||||
if ch == wshutil.ESC {
|
||||
b.Mode = Mode_Esc
|
||||
b.EscSeqBuf = []byte{ch}
|
||||
continue
|
||||
}
|
||||
outputBuf = append(outputBuf, ch)
|
||||
}
|
||||
if len(outputBuf) > 0 {
|
||||
err := b.DataOutputFn(outputBuf)
|
||||
if err != nil {
|
||||
b.setErr(fmt.Errorf("error processing data output: %w", err))
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user