Added demo UI for realtime data

This commit is contained in:
frikky
2022-01-20 10:15:44 +01:00
parent 6a941df7d4
commit ad8618fce3
2 changed files with 122 additions and 6 deletions
+82
View File
@@ -0,0 +1,82 @@
import React, { useState, useEffect, useLayoutEffect } from "react";
const KeepAlive = (defaultprops) => {
const { globalUrl, isLoggedIn, isLoaded, userdata } = defaultprops;
const [data, setData] = React.useState([]);
const [update, setUpdate] = React.useState(0);
const onChunkedResponseComplete = (result) => {
console.log('all done!', result)
}
const onChunkedResponseError = (err) => {
console.error(err)
}
const processChunkedResponse = async (response) => {
var text = '';
var reader = response.body.getReader()
var decoder = new TextDecoder();
const appendChunks = (result) => {
var chunk = decoder.decode(result.value || new Uint8Array, {stream: !result.done});
data.push(chunk)
setData(data)
setUpdate(Math.random());
console.log('got chunk of', chunk.length, 'bytes. Value: ', chunk)
text += chunk;
//console.log('text so far is', text.length, 'bytes\n');
if (result.done) {
console.log('returning')
return text;
} else {
return readChunk()
}
}
const readChunk = () => {
return reader.read().then(appendChunks);
}
return readChunk();
}
const getWorkflowStream = async () => {
await fetch(globalUrl + "/api/v1/workflows/a843fe90-585a-4693-8b7c-0b4dbce3347d/stream", {
method: "GET",
headers: {
"Content-Type": "application/json",
Accept: "application/json",
},
credentials: "include",
})
.then(processChunkedResponse)
.then(onChunkedResponseComplete)
.catch(onChunkedResponseError)
}
useEffect(() => {
console.log("FIRST REQUEST")
getWorkflowStream()
}, [])
return (
<div>
DONE
{data.map((innerdata, index) => {
return (
<div key={index}>
{innerdata}
</div>
)
})}
</div>
)
}
export default KeepAlive;
+40 -6
View File
@@ -9,6 +9,10 @@ package main
// frikky@debian:~/git/shuffle/functions/onprem/worker$ docker service create --replicas 5 --name shuffle-workers --env SHUFFLE_SWARM_CONFIG=run --publish published=33333,target=33333 ghcr.io/frikky/shuffle-worker:nightly
// Potential issues:
// Default network could be same as on the host
// Ingress network may not exist (default)
import (
"github.com/shuffle/shuffle-shared"
@@ -200,13 +204,39 @@ func deployServiceWorkers(image string) {
networkName = swarmNetworkName
}
ingressOptions := types.NetworkCreate{
Driver: "overlay",
Attachable: false,
Ingress: true,
IPAM: &network.IPAM{
Driver: "default",
Config: []network.IPAMConfig{
network.IPAMConfig{
Subnet: "10.225.225.0/24",
Gateway: "10.225.225.1",
},
},
},
}
_, err := dockercli.NetworkCreate(
ctx,
"ingress",
ingressOptions,
)
if err != nil {
log.Printf("[WARNING] Ingress network may already exist: %s", err)
}
//docker network create --driver=overlay workers
// Specific subnet?
networkCreateOptions := types.NetworkCreate{
Driver: "overlay",
Attachable: true,
Ingress: false,
IPAM: &network.IPAM{
Driver: "overlay",
Driver: "default",
Config: []network.IPAMConfig{
network.IPAMConfig{
Subnet: "10.224.224.0/24",
@@ -215,7 +245,7 @@ func deployServiceWorkers(image string) {
},
},
}
_, err := dockercli.NetworkCreate(
_, err = dockercli.NetworkCreate(
ctx,
networkName,
networkCreateOptions,
@@ -435,7 +465,9 @@ func deployWorker(image string, identifier string, env []string, executionReques
//var swarmConfig = os.Getenv("SHUFFLE_SWARM_CONFIG")
parsedUuid := uuid.NewV4()
if swarmConfig == "run" || swarmConfig == "swarm" {
// Stopping goroutine, as it just becomes too fast on startup
// FIXME: Should we handle replies properly?
// In certain cases, a workflow may e.g. be aborted already. If it's aborted, that returns
// a 401 from the worker, which returns an error here
err := sendWorkerRequest(executionRequest)
if err != nil {
log.Printf("[ERROR] Failed worker request for %s: %s", executionRequest.ExecutionId, err)
@@ -785,6 +817,8 @@ func main() {
checkSwarmService(ctx)
log.Printf("[DEBUG] Deploying worker image %s to swarm", workerImage)
deployServiceWorkers(workerImage)
log.Printf("[DEBUG] Waiting 30 seconds to ensure workers are deployed. Run: \"docker service ls\" for more info")
time.Sleep(time.Duration(30) * time.Second)
//deployServiceWorkers(workerImage)
}
@@ -847,7 +881,7 @@ func main() {
// FIXME - add check for StatusCode
if newresp.StatusCode != 200 {
if hasStarted {
log.Printf("[WARNING] Bad statuscode: %d", newresp.StatusCode)
log.Printf("[WARNING] Bad statuscode from backend: %d", newresp.StatusCode)
}
} else {
if !hasStarted {
@@ -1258,8 +1292,8 @@ func sendWorkerRequest(workflowExecution shuffle.ExecutionRequest) error {
}
if newresp.StatusCode != 200 {
log.Printf("[ERROR] Error running request - status code is %d, not 200. Body: %s", newresp.StatusCode, string(body))
return errors.New(fmt.Sprintf("Bad statuscode: %d - expecting 200", newresp.StatusCode))
log.Printf("[ERROR] Error running worker request - status code is %d, not 200. Body: %s", newresp.StatusCode, string(body))
return errors.New(fmt.Sprintf("Bad statuscode from worker: %d - expecting 200", newresp.StatusCode))
}
_ = body