More delay and network management when issues occur

This commit is contained in:
Frikky
2025-08-22 00:02:11 +02:00
parent d3519cd164
commit 065dc160b3
6 changed files with 146 additions and 24 deletions
+1 -1
View File
@@ -10,7 +10,7 @@ require (
github.com/docker/docker v28.2.2+incompatible
github.com/docker/go-connections v0.5.0
github.com/satori/go.uuid v1.2.0
github.com/shuffle/shuffle-shared v0.8.99
github.com/shuffle/shuffle-shared v0.9.2
k8s.io/api v0.33.1
k8s.io/apimachinery v0.33.1
)
+4 -4
View File
@@ -142,8 +142,8 @@ github.com/felixge/httpsnoop v1.0.4 h1:NFTV2Zj1bL4mc9sqWACXbQFVBBg2W3GPvqp8/ESS2
github.com/felixge/httpsnoop v1.0.4/go.mod h1:m8KPJKqk1gH5J9DgRY2ASl2lWCfGKXixSwevea8zH2U=
github.com/frikky/kin-openapi v0.42.0 h1:d5Z6vnuQ6RnCCPIxZaDL+TH2ODLxT8abytOt+Zh+Kd0=
github.com/frikky/kin-openapi v0.42.0/go.mod h1:ev9OZAw7Bv5p0w93j91++6a1ElPzGcCofst+kmrWsj4=
github.com/frikky/schemaless v0.0.16 h1:4d2ZktB9xGsAusbbKliOI8TuriSrdIMzD/6ToY3wkz8=
github.com/frikky/schemaless v0.0.16/go.mod h1:jT48kTcmr1q3o8i+8qe7g+eCsbwaz2Q9CjOJevQQzQs=
github.com/frikky/schemaless v0.0.17 h1:Tlg7td64r/EGAfQAL8LIlB8mOEnSRAgfJ8TERxzJSvM=
github.com/frikky/schemaless v0.0.17/go.mod h1:jT48kTcmr1q3o8i+8qe7g+eCsbwaz2Q9CjOJevQQzQs=
github.com/fxamacker/cbor/v2 v2.7.0 h1:iM5WgngdRBanHcxugY4JySA0nk1wZorNOpTgCMedv5E=
github.com/fxamacker/cbor/v2 v2.7.0/go.mod h1:pxXPTn3joSm21Gbwsv0w9OSA2y1HFR9qXEeXQVeNoDQ=
github.com/ghodss/yaml v1.0.0 h1:wQHKEahhL6wmXdzwWG11gIVCkOv05bNOh+Rxn0yngAk=
@@ -324,8 +324,8 @@ github.com/sendgrid/sendgrid-go v3.16.1+incompatible h1:zWhTmB0Y8XCDzeWIm2/BIt1G
github.com/sendgrid/sendgrid-go v3.16.1+incompatible/go.mod h1:QRQt+LX/NmgVEvmdRw0VT/QgUn499+iza2FnDca9fg8=
github.com/sergi/go-diff v1.3.2-0.20230802210424-5b0b94c5c0d3 h1:n661drycOFuPLCN3Uc8sB6B/s6Z4t2xvBgU1htSHuq8=
github.com/sergi/go-diff v1.3.2-0.20230802210424-5b0b94c5c0d3/go.mod h1:A0bzQcvG0E7Rwjx0REVgAGH58e96+X0MeOfepqsbeW4=
github.com/shuffle/shuffle-shared v0.8.84 h1:ElIMQYjKBVOiadbiGkSzt/lPU5xqaQwRxQvk9wx/xYM=
github.com/shuffle/shuffle-shared v0.8.84/go.mod h1:RdfNxqCPI+zU4jQKy3E/p4Io2injm7LpSKQUCDHNtLk=
github.com/shuffle/shuffle-shared v0.9.2 h1:xl/dTNKWh9mol2XPaGXCUlUTqLuwrwLxlEuv7AINwkk=
github.com/shuffle/shuffle-shared v0.9.2/go.mod h1:PFLG4eaxz8zQc9Lbc0Ok/nG8rNsaBkZ0TNar+VU+J/g=
github.com/sirupsen/logrus v1.7.0/go.mod h1:yWOB1SBYBC5VeMP7gHvWumXLIWorT60ONWic61uBYv0=
github.com/sirupsen/logrus v1.9.3 h1:dueUQJ1C2q9oE3F7wvmSGAaVtTmUizReu6fjN8uqzbQ=
github.com/sirupsen/logrus v1.9.3/go.mod h1:naHLuLoDiP4jHNo9R0sCBMtWGeIprob74mVsIT4qYEQ=
+1 -1
View File
@@ -90,7 +90,7 @@ var baseimageregistry = os.Getenv("SHUFFLE_BASE_IMAGE_REGISTRY")
//var baseimagetagsuffix = os.Getenv("SHUFFLE_BASE_IMAGE_TAG_SUFFIX")
// Used for cloud with auth
// Used for cloud with auth. Onprem in certain cases too.
var auth = os.Getenv("AUTH")
var org = os.Getenv("ORG")
+1 -1
View File
@@ -10,7 +10,7 @@ require (
github.com/docker/docker v28.2.2+incompatible
github.com/gorilla/mux v1.8.1
github.com/satori/go.uuid v1.2.0
github.com/shuffle/shuffle-shared v0.8.99
github.com/shuffle/shuffle-shared v0.9.2
k8s.io/api v0.33.1
k8s.io/apimachinery v0.33.1
k8s.io/client-go v0.33.1
+4 -4
View File
@@ -142,8 +142,8 @@ github.com/felixge/httpsnoop v1.0.4 h1:NFTV2Zj1bL4mc9sqWACXbQFVBBg2W3GPvqp8/ESS2
github.com/felixge/httpsnoop v1.0.4/go.mod h1:m8KPJKqk1gH5J9DgRY2ASl2lWCfGKXixSwevea8zH2U=
github.com/frikky/kin-openapi v0.42.0 h1:d5Z6vnuQ6RnCCPIxZaDL+TH2ODLxT8abytOt+Zh+Kd0=
github.com/frikky/kin-openapi v0.42.0/go.mod h1:ev9OZAw7Bv5p0w93j91++6a1ElPzGcCofst+kmrWsj4=
github.com/frikky/schemaless v0.0.16 h1:4d2ZktB9xGsAusbbKliOI8TuriSrdIMzD/6ToY3wkz8=
github.com/frikky/schemaless v0.0.16/go.mod h1:jT48kTcmr1q3o8i+8qe7g+eCsbwaz2Q9CjOJevQQzQs=
github.com/frikky/schemaless v0.0.17 h1:Tlg7td64r/EGAfQAL8LIlB8mOEnSRAgfJ8TERxzJSvM=
github.com/frikky/schemaless v0.0.17/go.mod h1:jT48kTcmr1q3o8i+8qe7g+eCsbwaz2Q9CjOJevQQzQs=
github.com/fxamacker/cbor/v2 v2.7.0 h1:iM5WgngdRBanHcxugY4JySA0nk1wZorNOpTgCMedv5E=
github.com/fxamacker/cbor/v2 v2.7.0/go.mod h1:pxXPTn3joSm21Gbwsv0w9OSA2y1HFR9qXEeXQVeNoDQ=
github.com/ghodss/yaml v1.0.0 h1:wQHKEahhL6wmXdzwWG11gIVCkOv05bNOh+Rxn0yngAk=
@@ -326,8 +326,8 @@ github.com/sendgrid/sendgrid-go v3.16.1+incompatible h1:zWhTmB0Y8XCDzeWIm2/BIt1G
github.com/sendgrid/sendgrid-go v3.16.1+incompatible/go.mod h1:QRQt+LX/NmgVEvmdRw0VT/QgUn499+iza2FnDca9fg8=
github.com/sergi/go-diff v1.3.2-0.20230802210424-5b0b94c5c0d3 h1:n661drycOFuPLCN3Uc8sB6B/s6Z4t2xvBgU1htSHuq8=
github.com/sergi/go-diff v1.3.2-0.20230802210424-5b0b94c5c0d3/go.mod h1:A0bzQcvG0E7Rwjx0REVgAGH58e96+X0MeOfepqsbeW4=
github.com/shuffle/shuffle-shared v0.8.84 h1:ElIMQYjKBVOiadbiGkSzt/lPU5xqaQwRxQvk9wx/xYM=
github.com/shuffle/shuffle-shared v0.8.84/go.mod h1:RdfNxqCPI+zU4jQKy3E/p4Io2injm7LpSKQUCDHNtLk=
github.com/shuffle/shuffle-shared v0.9.2 h1:xl/dTNKWh9mol2XPaGXCUlUTqLuwrwLxlEuv7AINwkk=
github.com/shuffle/shuffle-shared v0.9.2/go.mod h1:PFLG4eaxz8zQc9Lbc0Ok/nG8rNsaBkZ0TNar+VU+J/g=
github.com/sirupsen/logrus v1.7.0/go.mod h1:yWOB1SBYBC5VeMP7gHvWumXLIWorT60ONWic61uBYv0=
github.com/sirupsen/logrus v1.9.3 h1:dueUQJ1C2q9oE3F7wvmSGAaVtTmUizReu6fjN8uqzbQ=
github.com/sirupsen/logrus v1.9.3/go.mod h1:naHLuLoDiP4jHNo9R0sCBMtWGeIprob74mVsIT4qYEQ=
+135 -13
View File
@@ -683,6 +683,10 @@ func deployk8sApp(image string, identifier string, env []string) error {
return err
}
// Giving the service time to start before we contineu anything
log.Printf("[DEBUG] Waiting 20 seconds before moving on to let app '%s' start properly. Service: %s (k8s)", name, image)
time.Sleep(20 * time.Second)
return nil
}
@@ -3135,7 +3139,7 @@ func findActiveSwarmNodes(dockercli *dockerclient.Client) (int64, error) {
}
/*** STARTREMOVE ***/
func deploySwarmService(dockercli *dockerclient.Client, name, image string, deployport int) error {
func deploySwarmService(dockercli *dockerclient.Client, name, image string, deployport int, retry bool) error {
log.Printf("[DEBUG] Deploying service for %s to swarm on port %d", name, deployport)
//containerName := fmt.Sprintf("shuffle-worker-%s", parsedUuid)
@@ -3296,6 +3300,21 @@ func deploySwarmService(dockercli *dockerclient.Client, name, image string, depl
_ = service
if err != nil {
if strings.Contains(fmt.Sprintf("%s", err), "network") && strings.Contains(fmt.Sprintf("%s", err), "not found") {
log.Printf("[DEBUG] Network %s not found. Trying to initialize it.", networkName)
networkErr := initSwarmNetwork()
if networkErr != nil {
log.Printf("[ERROR] Failed initializing swarm network: %s", err)
//return err
}
// Retry deploying the service
if !retry {
return deploySwarmService(dockercli, name, image, deployport, true)
}
}
log.Printf("[DEBUG] Failed deploying %s with image %s: %s", name, image, err)
return err
}
@@ -3351,13 +3370,108 @@ func findAppInfoKubernetes(image, name string, env []string) error {
return err
}
func findAppInfo(image, name string) (int, error) {
// Backups in case networks are removed
func initSwarmNetwork() error {
ctx := context.Background()
dockercli, err := dockerclient.NewEnvClient()
if err != nil {
log.Printf("[ERROR] Unable to create docker client (2): %s", err)
return -1, err
return err
}
// Create the network options with the specified MTU
options := make(map[string]string)
mtu := 1500
options["com.docker.network.driver.mtu"] = fmt.Sprintf("%d", mtu)
ingressOptions := network.CreateOptions{
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?
networkName := "shuffle_swarm_executions"
if len(swarmNetworkName) > 0 {
networkName = swarmNetworkName
}
networkCreateOptions := network.CreateOptions{
Driver: "overlay",
Options: options,
Attachable: true,
Ingress: false,
IPAM: &network.IPAM{
Driver: "default",
Config: []network.IPAMConfig{
network.IPAMConfig{
Subnet: "10.224.224.0/24",
Gateway: "10.224.224.1",
},
},
},
}
_, err = dockercli.NetworkCreate(
ctx,
networkName,
networkCreateOptions,
)
if err != nil {
log.Printf("[WARNING] Swarm Executions network may already exist: %s", err)
}
networkName = "shuffle-executions"
networkCreateOptions = network.CreateOptions{
Driver: "overlay",
Options: options,
Attachable: true,
Ingress: false,
IPAM: &network.IPAM{
Driver: "default",
Config: []network.IPAMConfig{
network.IPAMConfig{
Subnet: "10.223.223.0/24",
Gateway: "10.223.223.1",
},
},
},
}
_, err = dockercli.NetworkCreate(
ctx,
networkName,
networkCreateOptions,
)
if err != nil {
log.Printf("[WARNING] Swarm Executions network may already exist: %s", err)
}
return nil
}
func findAppInfo(image, name string) (int, error) {
highest := baseport
exposedPort := -1
@@ -3379,6 +3493,12 @@ func findAppInfo(image, name string) (int, error) {
//Filters:
if exposedPort == -1 {
dockercli, err := dockerclient.NewEnvClient()
if err != nil {
log.Printf("[ERROR] Unable to create docker client (2): %s", err)
return -1, err
}
serviceListOptions := types.ServiceListOptions{}
services, err := dockercli.ServiceList(
context.Background(),
@@ -3443,26 +3563,29 @@ func findAppInfo(image, name string) (int, error) {
if exposedPort >= 0 {
//log.Printf("[INFO] Found service %s on port %d - no need to deploy another", name, exposedPort)
} else {
dockercli, err := dockerclient.NewEnvClient()
if err != nil {
log.Printf("[ERROR] Unable to create docker client (2): %s", err)
return -1, err
}
// Increment by 1 for highest port
if highest <= baseport {
highest = baseport
}
highest += 1
err = deploySwarmService(dockercli, name, image, highest)
err = deploySwarmService(dockercli, name, image, highest, false)
if err != nil {
log.Printf("[WARNING] NOT Found service: %s. error: %s", name, err)
return highest, err
} else {
log.Printf("[DEBUG] Deployed app with name %s", name)
log.Printf("[DEBUG] Waiting 20 seconds before moving on to let app '%s' start properly. Service: %s", name, image)
time.Sleep(time.Duration(20) * time.Second)
}
exposedPort = highest
if appsInitialized {
log.Printf("[DEBUG] Waiting 30 seconds before moving on to let app start")
time.Sleep(time.Duration(30) * time.Second)
}
//return exposedPort, errors.New("Deployed app %s")
}
return exposedPort, nil
@@ -3652,7 +3775,7 @@ func sendAppRequest(ctx context.Context, incomingUrl, appName string, port int,
func baseDeploy() {
var cli *dockerclient.Client
var err error
//var err error
if isKubernetes != "true" {
cli, err := dockerclient.NewEnvClient()
@@ -3714,8 +3837,7 @@ func baseDeploy() {
//deployApp(cli, value, identifier, env, workflowExecution, action)
log.Printf("[DEBUG] Deploying app with identifier %s to ensure basic apps are available from the get-go", identifier)
err = deployApp(cli, value, identifier, env, workflowExecution, action)
_ = err
go deployApp(cli, value, identifier, env, workflowExecution, action)
//err := deployApp(cli, value, identifier, env, workflowExecution, action)
//if err != nil {
// log.Printf("[DEBUG] Failed deploying app %s: %s", value, err)