diff --git a/pkg/tunnel/replay_part_upload.go b/pkg/tunnel/replay_part_upload.go index 110663d..f224076 100644 --- a/pkg/tunnel/replay_part_upload.go +++ b/pkg/tunnel/replay_part_upload.go @@ -263,23 +263,33 @@ func (p *PartUploader) GetStorage() storage.ReplayStorage { 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) } // 上传到存储 @@ -331,6 +341,10 @@ func (p *PartUploader) uploadToStorage(uploadPath string) { } +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) diff --git a/pkg/tunnel/replay_part_upload_test.go b/pkg/tunnel/replay_part_upload_test.go new file mode 100644 index 0000000..9e1e8b9 --- /dev/null +++ b/pkg/tunnel/replay_part_upload_test.go @@ -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) + } +}