Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
44 changes: 29 additions & 15 deletions pkg/tunnel/replay_part_upload.go
Original file line number Diff line number Diff line change
Expand Up @@ -263,23 +263,33 @@
const recordDirTimeFormat = "2006-01-02"

func (p *PartUploader) uploadToStorage(uploadPath string) {
// check whether to use ENABLE_VIDEO_WORKER
if videoWorkerClient := NewWorkerClient(*config.GlobalConfig); videoWorkerClient != nil {
taskCfg := videoworker.TaskConfig{
Width: p.Info.OptimalScreenWidth,
Height: p.Info.OptimalScreenHeight,
Bitrate: 1,
}
taskId, err := videoWorkerClient.CreateReplaySessionTask(p.SessionId, uploadPath, &taskCfg)
if err == nil {
logger.Infof("Create replay session VideoWorker task success, task id: %s", taskId)
if err = os.RemoveAll(p.RootPath); err != nil {
logger.Errorf("PartUploader %s remove root path %s error: %v", p.SessionId, p.RootPath, err)
if p.TermCfg == nil {
// Without Core's storage configuration there is no safe fallback target.
// Leave local replay files intact for a later recovery attempt.
logger.Errorf("PartUploader %s cannot upload replay without terminal config; retaining local files", p.SessionId)
return
}
// An environment switch alone must not route Community Edition recordings
// to Video Worker. Core supplies the license state in terminal configuration.
if shouldUseVideoWorker(config.GlobalConfig.EnableVideoWorker, p.TermCfg) {
videoWorkerClient := NewWorkerClient(*config.GlobalConfig)
if videoWorkerClient != nil {
taskCfg := videoworker.TaskConfig{
Width: p.Info.OptimalScreenWidth,
Height: p.Info.OptimalScreenHeight,
Bitrate: 1,
}
return
taskId, err := videoWorkerClient.CreateReplaySessionTask(p.SessionId, uploadPath, &taskCfg)
if err == nil {
logger.Infof("Create replay session VideoWorker task success, task id: %s", taskId)
if err = os.RemoveAll(p.RootPath); err != nil {
logger.Errorf("PartUploader %s remove root path %s error: %v", p.SessionId, p.RootPath, err)
}
return
}
// videoWorkerClient failed then try to use self storage to upload
logger.Errorf("Create replay session task error: %v, try to use self storage", err)
}
// videoWorkerClient failed then try to use self storage to upload
logger.Errorf("Create replay session task error: %v, try to use self storage", err)
}

// 上传到存储
Expand Down Expand Up @@ -331,6 +341,10 @@

}

func shouldUseVideoWorker(enabled bool, terminalCfg *model.TerminalConfig) bool {
return enabled && terminalCfg != nil && terminalCfg.LicenseIsValid
}

func (p *PartUploader) RecordLifecycleLog(event model.LifecycleEvent, logObj model.SessionLifecycleLog) {
if err := p.ApiClient.RecordSessionLifecycleLog(p.SessionId, event, logObj); err != nil {
logger.Errorf("Record session %s lifecycle %s log err: %s", p.SessionId, event, err)
Expand Down Expand Up @@ -377,7 +391,7 @@
if err != nil {
return scan, err
}
defer fd.Close()

Check failure on line 394 in pkg/tunnel/replay_part_upload.go

View workflow job for this annotation

GitHub Actions / lint

Error return value of `fd.Close` is not checked (errcheck)
info, err := fd.Stat()
if err != nil {
return scan, err
Expand Down
38 changes: 38 additions & 0 deletions pkg/tunnel/replay_part_upload_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,38 @@
package tunnel

import (
"os"
"testing"

"github.com/jumpserver-dev/sdk-go/model"
)

func TestVideoWorkerRequiresValidEnterpriseLicense(t *testing.T) {
tests := []struct {
name string
enabled bool
terminal *model.TerminalConfig
want bool
}{
{name: "switch off", terminal: &model.TerminalConfig{LicenseIsValid: true}},
{name: "community edition", enabled: true, terminal: &model.TerminalConfig{}},
{name: "missing terminal config", enabled: true},
{name: "licensed and enabled", enabled: true, terminal: &model.TerminalConfig{LicenseIsValid: true}, want: true},
}
for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
if got := shouldUseVideoWorker(test.enabled, test.terminal); got != test.want {
t.Fatalf("shouldUseVideoWorker() = %t, want %t", got, test.want)
}
})
}
}

func TestMissingTerminalConfigRetainsReplay(t *testing.T) {
root := t.TempDir()
uploader := PartUploader{SessionId: "session-id", RootPath: root}
uploader.uploadToStorage(root)
if _, err := os.Stat(root); err != nil {
t.Fatalf("source directory was not retained for recovery: %v", err)
}
}
Loading