trying to improve tenzir deployement

This commit is contained in:
Hari Krishna
2024-04-26 15:59:50 +00:00
parent e05e487db9
commit f825306a85
+558 -140
View File
@@ -38,6 +38,7 @@ import (
"github.com/docker/docker/api/types/mount" "github.com/docker/docker/api/types/mount"
"github.com/docker/docker/api/types/network" "github.com/docker/docker/api/types/network"
"github.com/docker/docker/api/types/swarm" "github.com/docker/docker/api/types/swarm"
"github.com/docker/go-connections/nat"
//"github.com/docker/docker/api/types/filters" //"github.com/docker/docker/api/types/filters"
dockerclient "github.com/docker/docker/client" dockerclient "github.com/docker/docker/client"
@@ -1495,6 +1496,13 @@ func main() {
initializeImages() initializeImages()
go func() {
if err := deployTenzirNode(); err != nil {
// Handle the error here
}
}()
workerImage := fmt.Sprintf("%s/%s/shuffle-worker:%s", baseimageregistry, baseimagename, workerVersion) workerImage := fmt.Sprintf("%s/%s/shuffle-worker:%s", baseimageregistry, baseimagename, workerVersion)
if len(newWorkerImage) > 0 { if len(newWorkerImage) > 0 {
workerImage = newWorkerImage workerImage = newWorkerImage
@@ -1666,14 +1674,14 @@ func main() {
newrequests := []shuffle.ExecutionRequest{} newrequests := []shuffle.ExecutionRequest{}
for _, incRequest := range executionRequests.Data { for _, incRequest := range executionRequests.Data {
// Looking for specific jobs // Looking for specific jobs
if incRequest.Type == "PIPELINE_CREATE" || incRequest.Type == "PIPELINE_UPDATE" || incRequest.Type == "PIPELINE_DELETE" { if incRequest.Type == "PIPELINE_CREATE" || incRequest.Type == "PIPELINE_STOP" || incRequest.Type == "PIPELINE_DELETE" {
err := handlePipeline(incRequest) err := handlePipeline(incRequest)
if err != nil { if err != nil {
log.Printf("[ERROR] Failed handling pipeline: %s", err) log.Printf("[ERROR] Failed handling pipeline: %s", err)
} else {
toBeRemoved.Data = append(toBeRemoved.Data, incRequest)
} }
toBeRemoved.Data = append(toBeRemoved.Data, incRequest)
} else if incRequest.Type == "DOCKER_IMAGE_DOWNLOAD" { } else if incRequest.Type == "DOCKER_IMAGE_DOWNLOAD" {
log.Printf("[INFO] Should delete -> download new image %#v", incRequest.ExecutionArgument) log.Printf("[INFO] Should delete -> download new image %#v", incRequest.ExecutionArgument)
@@ -1868,154 +1876,156 @@ func main() {
} }
func deployPipeline(image, identifier, command string) error { // func deployPipeline(image, identifier, command string) error {
if isKubernetes == "true" { // if isKubernetes == "true" {
return errors.New("Kubernetes not implemented") // return errors.New("Kubernetes not implemented")
} // }
ctx := context.Background() // ctx := context.Background()
hostConfig := &container.HostConfig{ // hostConfig := &container.HostConfig{
LogConfig: container.LogConfig{ // LogConfig: container.LogConfig{
Type: "json-file", // Type: "json-file",
Config: map[string]string{ // Config: map[string]string{
"max-size": "10m", // "max-size": "10m",
}, // },
}, // },
Resources: container.Resources{}, // Resources: container.Resources{},
} // }
hostConfig.NetworkMode = container.NetworkMode(fmt.Sprintf("container:%s", containerId)) // hostConfig.NetworkMode = container.NetworkMode(fmt.Sprintf("container:%s", containerId))
if strings.ToLower(cleanupEnv) != "false" { // if strings.ToLower(cleanupEnv) != "false" {
hostConfig.AutoRemove = true // hostConfig.AutoRemove = true
} // }
envVariables := []string{ // envVariables := []string{
} // }
// Add volume binds for storage // // Add volume binds for storage
// Want read/write with full access for the container // // Want read/write with full access for the container
//sourceFolder := "/Users/frikky/git/shuffle/shuffle-database" // //sourceFolder := "/Users/frikky/git/shuffle/shuffle-database"
//destinationFolder := "/tmp/storage" // //destinationFolder := "/tmp/storage"
//hostConfig.Mounts = append(hostConfig.Mounts, mount.Mount{ // //hostConfig.Mounts = append(hostConfig.Mounts, mount.Mount{
// Type: mount.TypeBind, // // Type: mount.TypeBind,
// Source: sourceFolder, // // Source: sourceFolder,
// Target: destinationFolder, // // Target: destinationFolder,
//}) // //})
// FIXME: Is using sigma "automatically" here good? // // FIXME: Is using sigma "automatically" here good?
// Or is it better to run it as a separate workflow? // // Or is it better to run it as a separate workflow?
if strings.Contains(command, "sigma") { // if strings.Contains(command, "sigma") {
log.Printf("[DEBUG] Should LOAD sigma from backend in realtime and dump it in a folder inside the container") // log.Printf("[DEBUG] Should LOAD sigma from backend in realtime and dump it in a folder inside the container")
//sourceFolder := "/tmp/tenzir/sigma" // //sourceFolder := "/tmp/tenzir/sigma"
//sigmaFolder := "/tmp/tenzir/sigma" // //sigmaFolder := "/tmp/tenzir/sigma"
//hostConfig.Mounts = append(hostConfig.Mounts, mount.Mount{ // //hostConfig.Mounts = append(hostConfig.Mounts, mount.Mount{
// Type: mount.TypeBind, // // Type: mount.TypeBind,
// Source: sigmaFolder, // // Source: sigmaFolder,
// Target: sigmaFolder, // // Target: sigmaFolder,
//} // //}
} // }
config := &container.Config{ // config := &container.Config{
Image: image, // Image: image,
Env: envVariables, // Env: envVariables,
Cmd: []string{ // Cmd: []string{
command, // command,
}, // },
} // }
// Add label to container in case of zombies // // Add label to container in case of zombies
config.Labels = map[string]string{ // config.Labels = map[string]string{
"name": identifier, // "name": identifier,
"shuffle": "shuffle", // "shuffle": "shuffle",
} // }
cont, err := dockercli.ContainerCreate( // cont, err := dockercli.ContainerCreate(
ctx, // ctx,
config, // config,
hostConfig, // hostConfig,
nil, // nil,
nil, // nil,
identifier, // identifier,
) // )
if err != nil { // if err != nil {
if strings.Contains(fmt.Sprintf("%s", err), "Conflict. The container name ") { // if strings.Contains(fmt.Sprintf("%s", err), "Conflict. The container name ") {
log.Printf("[DEBUG] Pipeline Container %s already exists, removing it", identifier) // log.Printf("[DEBUG] Pipeline Container %s already exists, removing it", identifier)
} else { // } else {
log.Printf("[ERROR] Failed to create pipeline container %s: %s", identifier, err) // log.Printf("[ERROR] Failed to create pipeline container %s: %s", identifier, err)
return err // return err
} // }
} // }
containerStartOptions := container.StartOptions{} // containerStartOptions := container.StartOptions{}
err = dockercli.ContainerStart( // err = dockercli.ContainerStart(
ctx, // ctx,
cont.ID, // cont.ID,
containerStartOptions, // containerStartOptions,
) // )
if err != nil { // if err != nil {
if strings.Contains(fmt.Sprintf("%s", err), "cannot join network") || strings.Contains(fmt.Sprintf("%s", err), "No such container") { // if strings.Contains(fmt.Sprintf("%s", err), "cannot join network") || strings.Contains(fmt.Sprintf("%s", err), "No such container") {
hostConfig.NetworkMode = "" // hostConfig.NetworkMode = ""
cont, err = dockercli.ContainerCreate( // cont, err = dockercli.ContainerCreate(
ctx, // ctx,
config, // config,
hostConfig, // hostConfig,
nil, // nil,
nil, // nil,
identifier+"-2", // identifier+"-2",
) // )
if err != nil { // if err != nil {
log.Printf("[ERROR] Failed to CREATE pipeline container (2): %s", err) // log.Printf("[ERROR] Failed to CREATE pipeline container (2): %s", err)
} // }
err = dockercli.ContainerStart( // err = dockercli.ContainerStart(
ctx, // ctx,
cont.ID, // cont.ID,
containerStartOptions, // containerStartOptions,
) // )
if err != nil { // if err != nil {
log.Printf("[ERROR] Failed to start pipeline container (2): %s", err) // log.Printf("[ERROR] Failed to start pipeline container (2): %s", err)
return err // return err
} // }
} else { // } else {
log.Printf("[ERROR] Failed initial pipeline container start. Quitting as this is NOT a simple network issue. Err: %s", err) // log.Printf("[ERROR] Failed initial pipeline container start. Quitting as this is NOT a simple network issue. Err: %s", err)
} // }
if err != nil { // if err != nil {
log.Printf("[ERROR] Failed to start pipeline container in environment %s: %s", environment, err) // log.Printf("[ERROR] Failed to start pipeline container in environment %s: %s", environment, err)
return err // return err
} else { // } else {
log.Printf("[INFO] Pipeline Container created (1). Environment %s: docker logs %s", environment, cont.ID) // log.Printf("[INFO] Pipeline Container created (1). Environment %s: docker logs %s", environment, cont.ID)
} // }
stats, err := dockercli.ContainerInspect(ctx, cont.ID) // stats, err := dockercli.ContainerInspect(ctx, cont.ID)
if err != nil { // if err != nil {
log.Printf("[ERROR] Failed checking pipeline with containername '%s'", cont.ID) // log.Printf("[ERROR] Failed checking pipeline with containername '%s'", cont.ID)
return nil // return nil
} // }
containerStatus := stats.ContainerJSONBase.State.Status // containerStatus := stats.ContainerJSONBase.State.Status
log.Printf("[DEBUG] Status of pipeline '%s' is %s. Should be running. Will reset", containerName, containerStatus) // log.Printf("[DEBUG] Status of pipeline '%s' is %s. Should be running. Will reset", containerName, containerStatus)
} // }
// // Wait for the container to finish
// /*
// statusCh, errCh := dockercli.ContainerWait(ctx, cont.ID, container.WaitConditionNotRunning)
// select {
// case err := <-errCh:
// if err != nil {
// log.Printf("[ERROR] Failed to wait for container: %s", err)
// }
// case <-statusCh:
// log.Printf("[INFO] Container finished")
// }
// */
// return nil
// }
// Wait for the container to finish
/*
statusCh, errCh := dockercli.ContainerWait(ctx, cont.ID, container.WaitConditionNotRunning)
select {
case err := <-errCh:
if err != nil {
log.Printf("[ERROR] Failed to wait for container: %s", err)
}
case <-statusCh:
log.Printf("[INFO] Container finished")
}
*/
return nil
}
// Tenzir command samples // Tenzir command samples
// docker pull ghcr.io/dominiklohmann/tenzir-arm64:latest // docker pull ghcr.io/dominiklohmann/tenzir-arm64:latest
@@ -2024,34 +2034,442 @@ func deployPipeline(image, identifier, command string) error {
// Read from Cache and send it to a webhook // Read from Cache and send it to a webhook
// docker run tenzir/tenzir:latest 'from http://192.168.86.44:5002/api/v1/orgs/7e9b9007-5df2-4b47-bca5-c4d267ef2943/cache/CIDR%20ranges?type=text&authorization=cec9d01f-09b2-4419-8a0a-76c6046e3fef read lines | to http://192.168.86.44:5002/api/v1/hooks/webhook_665ace5f-f27b-496a-a365-6e07eb61078c write lines' // docker run tenzir/tenzir:latest 'from http://192.168.86.44:5002/api/v1/orgs/7e9b9007-5df2-4b47-bca5-c4d267ef2943/cache/CIDR%20ranges?type=text&authorization=cec9d01f-09b2-4419-8a0a-76c6046e3fef read lines | to http://192.168.86.44:5002/api/v1/hooks/webhook_665ace5f-f27b-496a-a365-6e07eb61078c write lines'
func handlePipeline(incRequest shuffle.ExecutionRequest) error { func handlePipeline(incRequest shuffle.ExecutionRequest) error {
if len(incRequest.ExecutionArgument) == 0 { if incRequest.Type != "PIPELINE_STOP" && len(incRequest.ExecutionArgument) == 0 {
log.Printf("[ERROR] No execution argument found for pipeline create. Skipping") log.Printf("[ERROR] No execution argument found for pipeline create. Skipping")
return errors.New("No execution argument found for pipeline create. Skipping") return errors.New("no execution argument found for pipeline create. Skipping")
} }
image := "tenzir/tenzir:latest" //image := "tenzir/tenzir:latest"
identifier := fmt.Sprintf("shuffle-%s", strings.ToLower(strings.ReplaceAll(incRequest.ExecutionSource, " ", "-"))) identifier := fmt.Sprintf("shuffle-%s", strings.ToLower(strings.ReplaceAll(incRequest.ExecutionSource, " ", "-")))
command := incRequest.ExecutionArgument command := incRequest.ExecutionArgument
if incRequest.Type == "PIPELINE_CREATE" { if incRequest.Type == "PIPELINE_CREATE" {
log.Printf("[INFO] Should delete -> recreate new pipeline %#v. Name: %#v", incRequest.ExecutionArgument, identifier) log.Printf("[INFO] Should delete -> recreate new pipeline %#v. Name: %#v", incRequest.ExecutionArgument, identifier)
err := deployPipeline(image, identifier, command) //err := deployPipeline(image, identifier, command)
pipelineId, err := createPipeline(command, identifier)
if err != nil { if err != nil {
log.Printf("[ERROR] Failed to deploy pipeline: %s", err) log.Printf("[ERROR] Failed to deploy pipeline: %s", err)
return err return err
} else { } else {
log.Printf("[INFO] Pipeline deployed successfully") log.Printf("[INFO] Pipeline deployed successfully with Id: %s", pipelineId)
newErr := savePipelineData(pipelineId, identifier, "running")
if newErr != nil {
log.Printf("[DEBUG] failed to save the pipeline data: %s", newErr)
} else {
log.Printf("[INFO] succesfully saved the pipeline info ")
}
} }
} else if incRequest.Type == "PIPELINE_DELETE" { } else if incRequest.Type == "PIPELINE_DELETE" {
log.Printf("[INFO] Should delete pipeline %#v", incRequest.ExecutionArgument) log.Printf("[INFO] Should delete pipeline %#v", incRequest.ExecutionArgument)
} else if incRequest.Type == "PIPELINE_UPDATE" { pipelineId, err := searchPipeline(identifier)
log.Printf("[INFO] Should update pipeline %#v", incRequest.ExecutionArgument) if err != nil {
log.Printf("[ERROR] Failed searching for Pipeline with name %s reason:%s ", identifier, err)
return err
}
err = deletePipeline(pipelineId)
if err != nil {
log.Printf("[ERROR] Failed Deleting Pipeline %s", err)
return err
}
} else if incRequest.Type == "PIPELINE_STOP" {
log.Printf("[INFO] Should stop the pipeline %#v", identifier)
pipelineId, err := searchPipeline(identifier)
if err != nil {
log.Printf("[ERROR] Failed searching for Pipeline with name %s reason:%s ", identifier, err)
return err
}
state, err := updatePipelineState(pipelineId, "stop")
if err != nil {
log.Printf("[ERROR] Failed to stop Pipeline: %s reason:%s ", pipelineId, err)
return err
} else {
log.Printf("[INFO] successfully stopped the Pipeline: %s", pipelineId)
}
err = savePipelineData(pipelineId, identifier, state)
if err != nil {
log.Printf("[DEBUG] failed to save the pipeline data: %s", err)
} else {
log.Printf("[INFO] succesfully saved the pipeline info ")
}
} else { } else {
log.Printf("[ERROR] Unknown type for pipeline: %s", incRequest.Type) log.Printf("[ERROR] Unknown type for pipeline: %s", incRequest.Type)
return errors.New("Unknown type for pipeline") return errors.New("unknown type for pipeline")
} }
return nil
}
func deployTenzirNode() error {
if isKubernetes == "true" {
return errors.New("kubernetes not implemented")
}
ctx := context.Background()
imageName := "tenzir/tenzir"
containerName := "tenzir-node"
healthconfig := &container.HealthConfig{
Test: []string{"tenzir --connection-timeout=30s --connection-retry-delay=1s 'api /ping'"},
Interval: 30 * time.Second,
Retries: 1,
}
config := &container.Config{
Cmd: []string{"--commands=web server --mode=dev --bind=0.0.0.0"},
Image: imageName,
Healthcheck: healthconfig,
ExposedPorts: nat.PortSet{"5160/tcp": struct{}{}},
Entrypoint: []string{containerName},
}
hostConfig := &container.HostConfig{
PortBindings: nat.PortMap{
"5160/tcp": []nat.PortBinding{{HostPort: "5160"}},
},
Mounts: []mount.Mount{
{
Type: mount.TypeVolume,
Source: containerName,
Target: "/var/lib/tenzir/",
},
},
LogConfig: container.LogConfig{
Type: "json-file",
Config: map[string]string{
"max-size": "10m",
},
},
VolumeDriver: "local",
}
// do we need to pull manually ??
pullOptions := types.ImagePullOptions{}
out, err := dockercli.ImagePull(ctx, imageName, pullOptions)
if err != nil {
log.Printf("[ERROR] Failed to pull the tenzir image %s",err)
}
defer out.Close()
containerStartOptions := container.StartOptions{}
_, err = dockercli.ContainerCreate(ctx, config, hostConfig, nil, nil, containerName)
if err != nil {
if strings.Contains(fmt.Sprintf("%s", err), "Conflict. The container name ") {
log.Printf("[DEBUG] Tenzir Node Container already exists, starting it")
err = dockercli.ContainerStart(ctx, containerName, containerStartOptions)
if err != nil {
log.Printf("[ERROR] Failed to start existing Tenzir Node container: %v", err)
return err
}
log.Printf("[INFO] Existing Tenzir Node container started successfully")
return nil
}
return fmt.Errorf("failed to create Tenzir Node container: %v", err)
}
log.Printf("[INFO] New Tenzir Node container created successfully")
err = dockercli.ContainerStart(ctx, containerName, containerStartOptions)
if err != nil {
log.Printf("[ERROR] Failed to start new Tenzir Node container: %v", err)
return err
}
log.Printf("[INFO] New Tenzir Node container started successfully")
return nil
}
func createPipeline(command, identifier string) (string, error) {
toBeDeleted := false
pipelineId, err := searchPipeline(identifier)
url := fmt.Sprintf("%s/api/v0/pipelines/create", tenzirUrl)
forwardMethod := "POST"
if err != nil {
if strings.Contains(fmt.Sprintf("%s", err), "no existing pipeline found") {
log.Printf("[INFO] No existing pipeline found with name: %s. Creating a new one!", identifier)
} else {
log.Printf("[ERROR] Failed to search for existing pipeline but continuing anyway : %s", err)
}
} else {
log.Printf("[INFO] an existing pipeline found with ID: %s. it will be deleted", pipelineId)
toBeDeleted = true
}
requestBody := map[string]interface{}{
"definition": command,
"name": identifier,
"hidden": false,
"autostart": map[string]bool{
"created": true,
"completed": false,
"failed": false,
},
"autodelete": map[string]bool{
"completed": false,
"failed": true,
"stopped": false,
},
}
requestBodyJSON, err := json.Marshal(requestBody)
if err != nil {
log.Printf("[ERROR] failed marshalling body: %s", err)
return "", err
}
forwardData := bytes.NewBuffer(requestBodyJSON)
req, err := http.NewRequest(
forwardMethod,
url,
forwardData,
)
if err != nil {
log.Printf("[ERROR] Failed to create HTTP request: %s", err)
return "", err
}
req.Header.Set("Content-Type", "application/json")
client := &http.Client{Timeout: 10 * time.Second}
resp, err := client.Do(req)
if err != nil {
log.Printf("[ERROR] Failed to send HTTP request: %s", err)
return "", err
}
defer resp.Body.Close()
if resp.StatusCode != 200 {
log.Printf("[DEBUG] status code is %d instead of 200", resp.StatusCode)
return "", fmt.Errorf("got the status code %d instead of 200", resp.StatusCode)
}
type PipelineResponse struct {
ID string `json:"id"`
}
var response PipelineResponse
if err := json.NewDecoder(resp.Body).Decode(&response); err != nil {
log.Printf("[ERROR] decoding response: %s", err)
return "", err
}
if response.ID == "" {
log.Println("[DEBUG] ID not found or empty in response")
return "", errors.New("pipeline ID not found or empty in the response")
}
id := response.ID
if toBeDeleted {
go deletePipeline(pipelineId)
}
return id, nil
}
func updatePipelineState(pipelineId, action string) (string, error) {
url := fmt.Sprintf("%s/api/v0/pipeline/update", tenzirUrl)
forwardMethod := "POST"
requestBody := map[string]interface{}{
"id": pipelineId,
"action": action,
"autostart": map[string]bool{
"created": true,
"completed": false,
"failed": false,
},
"autodelete": map[string]bool{
"completed": false,
"failed": true,
"stopped": false,
},
}
requestBodyJSON, err := json.Marshal(requestBody)
if err != nil {
return "", err
}
forwardData := bytes.NewBuffer(requestBodyJSON)
req, err := http.NewRequest(
forwardMethod,
url,
forwardData,
)
if err != nil {
log.Printf("[ERROR] Failed to create HTTP request: %s", err)
return "", err
}
req.Header.Set("Content-Type", "application/json")
client := &http.Client{Timeout: 10 * time.Second}
resp, err := client.Do(req)
if err != nil {
log.Printf("[ERROR] Failed to send HTTP request: %s", err)
return "", err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
return "", fmt.Errorf("got the status code %d instead of 200", resp.StatusCode)
}
body, err := ioutil.ReadAll(resp.Body)
if err != nil {
return "", err
}
var responseData struct {
Pipeline struct {
State string `json:"state"`
} `json:"pipeline"`
}
if err := json.Unmarshal(body, &responseData); err != nil {
return "", err
}
return responseData.Pipeline.State, nil
}
func deletePipeline(pipelineId string) error {
requestBody := map[string]string{
"id": pipelineId,
}
url := fmt.Sprintf("%s/api/v0/pipeline/delete", tenzirUrl)
forwardMethod := "POST"
requestBodyJSON, err := json.Marshal(requestBody)
if err != nil {
log.Println("[ERROR] failed marshalling request body:", err)
return err
}
forwardData := bytes.NewBuffer(requestBodyJSON)
req, err := http.NewRequest(
forwardMethod,
url,
forwardData,
)
if err != nil {
log.Printf("[ERROR] Failed to create HTTP request: %s", err)
return err
}
req.Header.Set("Content-Type", "application/json")
client := &http.Client{Timeout: 10 * time.Second}
resp, err := client.Do(req)
if err != nil {
log.Printf("[ERROR] Failed to send HTTP request: %s", err)
return err
}
defer resp.Body.Close()
if resp.StatusCode != 200 {
log.Printf("[DEBUG] The deletion of pipeline with ID: %s is unsucessful as status code is NOT 200 !!!", pipelineId)
return fmt.Errorf("got the status code %d instead of 200", resp.StatusCode)
}
log.Printf("[INFO] pipeline with ID: %s deleted successfully", pipelineId)
return nil
}
func searchPipeline(identifier string) (string, error) {
type pipelineInfo struct {
ID string `json:"id"`
Name string `json:"name"`
}
var reqBody []byte
url := fmt.Sprintf("%s/api/v0/pipeline/list", tenzirUrl)
resp, err := http.Post(url, "application/json", bytes.NewBuffer(reqBody))
if err != nil {
return "", err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
return "", fmt.Errorf("got the status code %d instead of 200", resp.StatusCode)
}
body, err := ioutil.ReadAll(resp.Body)
if err != nil {
return "", err
}
var responseData struct {
Pipelines []pipelineInfo `json:"pipelines"`
}
if err := json.Unmarshal(body, &responseData); err != nil {
return "", err
}
for _, pipeline := range responseData.Pipelines {
if pipeline.Name == identifier {
return pipeline.ID, nil
}
}
return "", errors.New("no existing pipeline found with name")
}
func savePipelineData(pipelineId, identifier, status string) error {
url := fmt.Sprintf("%s/api/v1/triggers/pipeline/save", tenzirUrl)
forwardMethod := "PUT"
payload := map[string]interface{}{
"pipeline_id": pipelineId,
"trigger_id": identifier,
"status": status,
}
payloadBytes, err := json.Marshal(payload)
if err != nil {
log.Printf("[ERROR] Failed to marshal payload: %s", err)
return err
}
forwardData := bytes.NewBuffer(payloadBytes)
req, err := http.NewRequest(
forwardMethod,
url,
forwardData,
)
if err != nil {
log.Printf("[ERROR] Failed to create HTTP request: %s", err)
return err
}
req.Header.Set("Content-Type", "application/json")
client := &http.Client{Timeout: 10 * time.Second}
resp, err := client.Do(req)
if err != nil {
log.Printf("[ERROR] Failed to send HTTP request: %s", err)
return err
}
defer resp.Body.Close()
if resp.StatusCode != 200 {
log.Printf("[ERROR] Received non-successful HTTP status code: %d", resp.StatusCode)
return fmt.Errorf("unexpected HTTP status code: %d", resp.StatusCode)
}
return nil return nil
} }