From 0d8c15910185eb43d1a5f163b77c99d898266db4 Mon Sep 17 00:00:00 2001 From: Mike Sawka Date: Mon, 19 Aug 2024 14:37:52 -0700 Subject: [PATCH] remote file preview (streaming) working (#248) --- frontend/app/store/wos.ts | 25 +++++-- frontend/app/view/preview/preview.tsx | 12 +++- pkg/web/web.go | 96 ++++++++++++++++++++++++--- pkg/wshrpc/wshremote/wshremote.go | 2 +- 4 files changed, 115 insertions(+), 20 deletions(-) diff --git a/frontend/app/store/wos.ts b/frontend/app/store/wos.ts index d97edb4a..7aa45fc6 100644 --- a/frontend/app/store/wos.ts +++ b/frontend/app/store/wos.ts @@ -61,6 +61,23 @@ function GetObject(oref: string): Promise { return callBackendService("object", "GetObject", [oref], true); } +function debugLogBackendCall(methodName: string, durationStr: string, args: any[]) { + durationStr = "| " + durationStr; + if (methodName == "object.UpdateObject" && args.length > 0) { + console.log("[service] object.UpdateObject", args[0].otype, args[0].oid, durationStr, args[0]); + return; + } + if (methodName == "object.GetObject" && args.length > 0) { + console.log("[service] object.GetObject", args[0], durationStr); + return; + } + if (methodName == "file.StatFile" && args.length >= 2) { + console.log("[service] file.StatFile", args[1], durationStr); + return; + } + console.log("[service]", methodName, durationStr); +} + function callBackendService(service: string, method: string, args: any[], noUIContext?: boolean): Promise { const startTs = Date.now(); let uiContext: UIContext = null; @@ -101,13 +118,7 @@ function callBackendService(service: string, method: string, args: any[], noUICo throw new Error(`call ${methodName} error: ${respData.error}`); } const durationStr = Date.now() - startTs + "ms"; - if (methodName == "object.UpdateObject") { - console.log("Call UpdateObject", args[0].otype, args[0].oid, durationStr, args[0]); - } else if (methodName == "object.GetObject") { - console.log("Call GetObject", args[0], durationStr); - } else { - console.log("Call", methodName, durationStr); - } + debugLogBackendCall(methodName, durationStr, args); return respData.data; }); return prtn; diff --git a/frontend/app/view/preview/preview.tsx b/frontend/app/view/preview/preview.tsx index 951c8cfd..00d1412f 100644 --- a/frontend/app/view/preview/preview.tsx +++ b/frontend/app/view/preview/preview.tsx @@ -374,9 +374,14 @@ function MarkdownPreview({ contentAtom }: { contentAtom: jotai.Atom @@ -516,6 +521,7 @@ function PreviewView({ blockId, model }: { blockId: string; model: PreviewModel const fileName = jotai.useAtomValue(fileNameAtom); const fileInfo = jotai.useAtomValue(statFileAtom); const ceReadOnly = jotai.useAtomValue(ceReadOnlyAtom); + const conn = jotai.useAtomValue(model.connection); let blockIcon = iconForFile(mimeType, fileName); // ensure consistent hook calls @@ -528,7 +534,7 @@ function PreviewView({ blockId, model }: { blockId: string; model: PreviewModel mimeType.startsWith("audio/") || mimeType.startsWith("image/") ) { - view = ; + view = ; } else if (!fileInfo) { view = File Not Found{util.isBlank(fileName) ? null : JSON.stringify(fileName)}; } else if (fileInfo.size > MaxFileSize) { diff --git a/pkg/web/web.go b/pkg/web/web.go index bbb9f724..ffffdaba 100644 --- a/pkg/web/web.go +++ b/pkg/web/web.go @@ -4,6 +4,7 @@ package web import ( + "bytes" "encoding/base64" "encoding/json" "fmt" @@ -24,6 +25,10 @@ import ( "github.com/wavetermdev/thenextwave/pkg/service" "github.com/wavetermdev/thenextwave/pkg/telemetry" "github.com/wavetermdev/thenextwave/pkg/wavebase" + "github.com/wavetermdev/thenextwave/pkg/wshrpc" + "github.com/wavetermdev/thenextwave/pkg/wshrpc/wshclient" + "github.com/wavetermdev/thenextwave/pkg/wshrpc/wshserver" + "github.com/wavetermdev/thenextwave/pkg/wshutil" "github.com/wavetermdev/thenextwave/pkg/wstore" ) @@ -210,15 +215,8 @@ func serveTransparentGIF(w http.ResponseWriter) { w.Write(gifBytes) } -func handleStreamFile(w http.ResponseWriter, r *http.Request) { - fileName := r.URL.Query().Get("path") - if fileName == "" { - http.Error(w, "path is required", http.StatusBadRequest) - return - } - no404 := r.URL.Query().Get("no404") - log.Printf("got no404: %q\n", no404) - if no404 != "" { +func handleLocalStreamFile(w http.ResponseWriter, r *http.Request, fileName string, no404 bool) { + if no404 { log.Printf("streaming file w/no404: %q\n", fileName) // use the custom response writer rw := ¬FoundBlockingResponseWriter{w: w, headers: http.Header{}} @@ -235,6 +233,86 @@ func handleStreamFile(w http.ResponseWriter, r *http.Request) { } } +func handleRemoteStreamFile(w http.ResponseWriter, r *http.Request, conn string, fileName string, no404 bool) error { + client := wshserver.GetMainRpcClient() + streamFileData := wshrpc.CommandRemoteStreamFileData{Path: fileName} + route := wshutil.MakeConnectionRouteId(conn) + rtnCh := wshclient.RemoteStreamFileCommand(client, streamFileData, &wshrpc.RpcOpts{Route: route}) + firstPk := true + var fileInfo *wshrpc.FileInfo + loopDone := false + defer func() { + if loopDone { + return + } + // if loop didn't finish naturally clear it out + go func() { + for range rtnCh { + } + }() + }() + for respUnion := range rtnCh { + if respUnion.Error != nil { + return respUnion.Error + } + if firstPk { + firstPk = false + if len(respUnion.Response.FileInfo) != 1 { + return fmt.Errorf("stream file protocol error, first pk fileinfo len=%d", len(respUnion.Response.FileInfo)) + } + fileInfo = respUnion.Response.FileInfo[0] + if fileInfo.NotFound { + if no404 { + serveTransparentGIF(w) + return nil + } else { + return fmt.Errorf("file not found: %q", fileName) + } + } + if fileInfo.IsDir { + return fmt.Errorf("cannot stream directory: %q", fileName) + } + w.Header().Set(ContentTypeHeaderKey, fileInfo.MimeType) + w.Header().Set(ContentLengthHeaderKey, fmt.Sprintf("%d", fileInfo.Size)) + continue + } + if respUnion.Response.Data64 == "" { + continue + } + decoder := base64.NewDecoder(base64.StdEncoding, bytes.NewReader([]byte(respUnion.Response.Data64))) + _, err := io.Copy(w, decoder) + if err != nil { + log.Printf("error streaming file %q: %v\n", fileName, err) + // not sure what to do here, the headers have already been sent. + // just return + return nil + } + } + loopDone = true + return nil +} + +func handleStreamFile(w http.ResponseWriter, r *http.Request) { + conn := r.URL.Query().Get("connection") + if conn == "" { + conn = wshrpc.LocalConnName + } + fileName := r.URL.Query().Get("path") + if fileName == "" { + http.Error(w, "path is required", http.StatusBadRequest) + return + } + no404 := r.URL.Query().Get("no404") + if conn == wshrpc.LocalConnName { + handleLocalStreamFile(w, r, fileName, no404 != "") + } else { + err := handleRemoteStreamFile(w, r, conn, fileName, no404 != "") + if err != nil { + http.Error(w, fmt.Sprintf("error streaming file: %v", err), http.StatusInternalServerError) + } + } +} + func WriteJsonError(w http.ResponseWriter, errVal error) { w.Header().Set(ContentTypeHeaderKey, ContentTypeJson) w.WriteHeader(http.StatusOK) diff --git a/pkg/wshrpc/wshremote/wshremote.go b/pkg/wshrpc/wshremote/wshremote.go index 1e4dc68d..6783a21a 100644 --- a/pkg/wshrpc/wshremote/wshremote.go +++ b/pkg/wshrpc/wshremote/wshremote.go @@ -159,7 +159,7 @@ func (impl *ServerImpl) remoteStreamFileRegular(ctx context.Context, path string filePos += int64(n) dataCallback(nil, buf[:n]) } - if filePos >= byteRange.End { + if !byteRange.All && filePos >= byteRange.End { break } if errors.Is(err, io.EOF) {