Merge pull request #7271 from dokku/7268-k3s-nginx-logs
Add support for non-local nginx implementations
This commit is contained in:
@@ -2583,6 +2583,25 @@ DOKKU_SCHEDULER="$1"; APP="$2";
|
||||
# TODO
|
||||
```
|
||||
|
||||
### `scheduler-proxy-logs`
|
||||
|
||||
> [!WARNING]
|
||||
> The scheduler plugin trigger apis are under development and may change
|
||||
> between minor releases until the 1.0 release.
|
||||
- Description: Allows you to run scheduler commands when retrieving failed container logs
|
||||
- Invoked by: `dokku nginx:access-logs` and `dokku nginx:error-logs`
|
||||
- Arguments: `$DOKKU_SCHEDULER $APP $PROXY_TYPE $LOG_TYPE $TAIL $NUM_LINES`
|
||||
- Example:
|
||||
|
||||
```shell
|
||||
#!/usr/bin/env bash
|
||||
|
||||
set -eo pipefail; [[ $DOKKU_TRACE ]] && set -x
|
||||
DOKKU_SCHEDULER="$1"; APP="$2"; $PROXY_TYPE="$3"; LOG_TYPE="$4"; TAIL="$5"; NUM_LINES="$6"
|
||||
|
||||
# TODO
|
||||
```
|
||||
|
||||
### `scheduler-pre-restore`
|
||||
|
||||
> [!WARNING]
|
||||
|
||||
@@ -74,10 +74,23 @@ nginx_logs() {
|
||||
dokku_log_fail "$NGINX_LOGS_TYPE logs are disabled for this app"
|
||||
fi
|
||||
|
||||
local tail=false
|
||||
local num=0
|
||||
if [[ $3 == "-t" ]]; then
|
||||
tail=true
|
||||
else
|
||||
num=20
|
||||
fi
|
||||
|
||||
local DOKKU_SCHEDULER=$(get_app_scheduler "$APP")
|
||||
if [[ "$DOKKU_SCHEDULER" != "docker-local" ]]; then
|
||||
plugn trigger scheduler-nginx-logs "$DOKKU_SCHEDULER" "$APP" "nginx" "$NGINX_LOGS_TYPE" "$tail" "$num"
|
||||
fi
|
||||
|
||||
if [[ "$tail" == "true" ]]; then
|
||||
local NGINX_LOGS_ARGS="-F"
|
||||
else
|
||||
local NGINX_LOGS_ARGS="-n 20"
|
||||
local NGINX_LOGS_ARGS="-n $num"
|
||||
fi
|
||||
|
||||
tail "$NGINX_LOGS_ARGS" "$NGINX_LOGS_PATH"
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
SUBCOMMANDS = subcommands/annotations:set subcommands/autoscaling-auth:set subcommands/autoscaling-auth:report subcommands/cluster-add subcommands/cluster-list subcommands/cluster-remove subcommands/initialize subcommands/labels:set subcommands/report subcommands/set subcommands/show-kubeconfig subcommands/uninstall
|
||||
TRIGGERS = triggers/install triggers/post-app-clone-setup triggers/post-app-rename-setup triggers/post-delete triggers/report triggers/scheduler-app-status triggers/scheduler-deploy triggers/scheduler-enter triggers/scheduler-logs triggers/scheduler-post-delete triggers/scheduler-run triggers/scheduler-run-list triggers/scheduler-stop
|
||||
TRIGGERS = triggers/install triggers/post-app-clone-setup triggers/post-app-rename-setup triggers/post-delete triggers/report triggers/scheduler-app-status triggers/scheduler-deploy triggers/scheduler-enter triggers/scheduler-logs triggers/scheduler-nginx-logs triggers/scheduler-post-delete triggers/scheduler-run triggers/scheduler-run-list triggers/scheduler-stop
|
||||
BUILD = commands subcommands triggers
|
||||
PLUGIN_NAME = scheduler-k3s
|
||||
|
||||
|
||||
@@ -28,9 +28,6 @@ import (
|
||||
"k8s.io/apimachinery/pkg/api/resource"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/apimachinery/pkg/util/wait"
|
||||
corev1client "k8s.io/client-go/kubernetes/typed/core/v1"
|
||||
"k8s.io/client-go/tools/remotecommand"
|
||||
"k8s.io/kubectl/pkg/util/term"
|
||||
"k8s.io/kubernetes/pkg/client/conditions"
|
||||
"k8s.io/utils/ptr"
|
||||
"mvdan.cc/sh/v3/shell"
|
||||
@@ -371,11 +368,6 @@ func createKubernetesNamespace(ctx context.Context, namespaceName string) error
|
||||
}
|
||||
|
||||
func enterPod(ctx context.Context, input EnterPodInput) error {
|
||||
coreclient, err := corev1client.NewForConfig(&input.Clientset.RestConfig)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Error creating corev1 client: %w", err)
|
||||
}
|
||||
|
||||
labelSelector := []string{}
|
||||
for k, v := range input.SelectedPod.Labels {
|
||||
labelSelector = append(labelSelector, fmt.Sprintf("%s=%s", k, v))
|
||||
@@ -385,7 +377,7 @@ func enterPod(ctx context.Context, input EnterPodInput) error {
|
||||
input.WaitTimeout = 5
|
||||
}
|
||||
|
||||
err = waitForPodBySelectorRunning(ctx, WaitForPodBySelectorRunningInput{
|
||||
err := waitForPodBySelectorRunning(ctx, WaitForPodBySelectorRunningInput{
|
||||
Clientset: input.Clientset,
|
||||
Namespace: input.SelectedPod.Namespace,
|
||||
LabelSelector: strings.Join(labelSelector, ","),
|
||||
@@ -405,46 +397,12 @@ func enterPod(ctx context.Context, input EnterPodInput) error {
|
||||
return fmt.Errorf("No container specified and no default container found")
|
||||
}
|
||||
|
||||
req := coreclient.RESTClient().Post().
|
||||
Resource("pods").
|
||||
Namespace(input.SelectedPod.Namespace).
|
||||
Name(input.SelectedPod.Name).
|
||||
SubResource("exec")
|
||||
|
||||
req.Param("container", input.SelectedContainerName)
|
||||
req.Param("stdin", "true")
|
||||
req.Param("stdout", "true")
|
||||
req.Param("stderr", "true")
|
||||
req.Param("tty", "true")
|
||||
|
||||
if input.Entrypoint != "" {
|
||||
req.Param("command", input.Entrypoint)
|
||||
}
|
||||
for _, cmd := range input.Command {
|
||||
req.Param("command", cmd)
|
||||
}
|
||||
|
||||
t := term.TTY{
|
||||
In: os.Stdin,
|
||||
Out: os.Stdout,
|
||||
Raw: true,
|
||||
}
|
||||
size := t.GetSize()
|
||||
sizeQueue := t.MonitorSize(size)
|
||||
|
||||
return t.Safe(func() error {
|
||||
exec, err := remotecommand.NewSPDYExecutor(&input.Clientset.RestConfig, "POST", req.URL())
|
||||
if err != nil {
|
||||
return fmt.Errorf("Error creating executor: %w", err)
|
||||
}
|
||||
|
||||
return exec.StreamWithContext(ctx, remotecommand.StreamOptions{
|
||||
Stdin: os.Stdin,
|
||||
Stdout: os.Stdout,
|
||||
Stderr: os.Stderr,
|
||||
Tty: true,
|
||||
TerminalSizeQueue: sizeQueue,
|
||||
})
|
||||
return input.Clientset.ExecCommand(ctx, ExecCommandInput{
|
||||
Command: input.Command,
|
||||
ContainerName: input.SelectedContainerName,
|
||||
Entrypoint: input.Entrypoint,
|
||||
Name: input.SelectedPod.Name,
|
||||
Namespace: input.SelectedPod.Namespace,
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
@@ -1,11 +1,19 @@
|
||||
package scheduler_k3s
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"os/signal"
|
||||
"strconv"
|
||||
"strings"
|
||||
"syscall"
|
||||
|
||||
"github.com/dokku/dokku/plugins/common"
|
||||
"github.com/fatih/color"
|
||||
"github.com/go-openapi/jsonpointer"
|
||||
kedav1alpha1 "github.com/kedacore/keda/v2/apis/keda/v1alpha1"
|
||||
appsv1 "k8s.io/api/apps/v1"
|
||||
@@ -19,9 +27,12 @@ import (
|
||||
"k8s.io/apimachinery/pkg/types"
|
||||
"k8s.io/client-go/dynamic"
|
||||
"k8s.io/client-go/kubernetes"
|
||||
corev1client "k8s.io/client-go/kubernetes/typed/core/v1"
|
||||
"k8s.io/client-go/rest"
|
||||
"k8s.io/client-go/tools/clientcmd"
|
||||
clientcmdapi "k8s.io/client-go/tools/clientcmd/api"
|
||||
"k8s.io/client-go/tools/remotecommand"
|
||||
"k8s.io/kubectl/pkg/util/term"
|
||||
"k8s.io/utils/ptr"
|
||||
)
|
||||
|
||||
@@ -285,6 +296,74 @@ func (k KubernetesClient) DeleteSecret(ctx context.Context, input DeleteSecretIn
|
||||
return k.Client.CoreV1().Secrets(input.Namespace).Delete(ctx, input.Name, metav1.DeleteOptions{})
|
||||
}
|
||||
|
||||
// ExecCommandInput contains all the information needed to execute a command in a Kubernetes pod
|
||||
type ExecCommandInput struct {
|
||||
// Command is the command to execute
|
||||
Command []string
|
||||
|
||||
// ContainerName is the Kubernetes container name
|
||||
ContainerName string
|
||||
|
||||
// Entrypoint is the command entrypoint
|
||||
Entrypoint string
|
||||
|
||||
// Name is the Kubernetes pod name
|
||||
Name string
|
||||
|
||||
// Namespace is the Kubernetes namespace
|
||||
Namespace string
|
||||
}
|
||||
|
||||
// ExecCommand executes a command in a Kubernetes pod
|
||||
func (k KubernetesClient) ExecCommand(ctx context.Context, input ExecCommandInput) error {
|
||||
coreclient, err := corev1client.NewForConfig(&k.RestConfig)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Error creating corev1 client: %w", err)
|
||||
}
|
||||
|
||||
req := coreclient.RESTClient().Post().
|
||||
Resource("pods").
|
||||
Namespace(input.Namespace).
|
||||
Name(input.Name).
|
||||
SubResource("exec")
|
||||
|
||||
req.Param("container", input.ContainerName)
|
||||
req.Param("stdin", "true")
|
||||
req.Param("stdout", "true")
|
||||
req.Param("stderr", "true")
|
||||
req.Param("tty", "true")
|
||||
|
||||
if input.Entrypoint != "" {
|
||||
req.Param("command", input.Entrypoint)
|
||||
}
|
||||
for _, cmd := range input.Command {
|
||||
req.Param("command", cmd)
|
||||
}
|
||||
|
||||
t := term.TTY{
|
||||
In: os.Stdin,
|
||||
Out: os.Stdout,
|
||||
Raw: true,
|
||||
}
|
||||
size := t.GetSize()
|
||||
sizeQueue := t.MonitorSize(size)
|
||||
|
||||
return t.Safe(func() error {
|
||||
exec, err := remotecommand.NewSPDYExecutor(&k.RestConfig, "POST", req.URL())
|
||||
if err != nil {
|
||||
return fmt.Errorf("Error creating executor: %w", err)
|
||||
}
|
||||
|
||||
return exec.StreamWithContext(ctx, remotecommand.StreamOptions{
|
||||
Stdin: os.Stdin,
|
||||
Stdout: os.Stdout,
|
||||
Stderr: os.Stderr,
|
||||
Tty: true,
|
||||
TerminalSizeQueue: sizeQueue,
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
// GetNodeInput contains all the information needed to get a Kubernetes node
|
||||
type GetNodeInput struct {
|
||||
// Name is the Kubernetes node name
|
||||
@@ -606,3 +685,136 @@ func (k KubernetesClient) ScaleDeployment(ctx context.Context, input ScaleDeploy
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
type StreamLogsInput struct {
|
||||
// Namespace is the Kubernetes namespace
|
||||
Namespace string
|
||||
|
||||
// DeploymentName is the Kubernetes deployment name
|
||||
DeploymentName string
|
||||
|
||||
// ContainerName is the Kubernetes container name
|
||||
ContainerName string
|
||||
|
||||
// Follow is whether to follow the logs
|
||||
Follow bool
|
||||
|
||||
// TailLines is the number of lines to tail
|
||||
TailLines int64
|
||||
|
||||
// Quiet is whether to suppress output
|
||||
Quiet bool
|
||||
}
|
||||
|
||||
func (k KubernetesClient) StreamLogs(ctx context.Context, input StreamLogsInput) error {
|
||||
ctx, cancel := context.WithCancel(ctx)
|
||||
signals := make(chan os.Signal, 1)
|
||||
signal.Notify(signals, os.Interrupt, syscall.SIGHUP,
|
||||
syscall.SIGINT,
|
||||
syscall.SIGQUIT,
|
||||
syscall.SIGTERM)
|
||||
go func() {
|
||||
<-signals
|
||||
cancel()
|
||||
}()
|
||||
|
||||
if err := k.Ping(); err != nil {
|
||||
return fmt.Errorf("kubernetes api not available: %w", err)
|
||||
}
|
||||
|
||||
labelSelector := []string{fmt.Sprintf("app.kubernetes.io/part-of=%s", input.DeploymentName)}
|
||||
processIndex := 0
|
||||
if input.ContainerName != "" {
|
||||
parts := strings.SplitN(input.ContainerName, ".", 2)
|
||||
if len(parts) == 2 {
|
||||
var err error
|
||||
input.ContainerName = parts[0]
|
||||
processIndex, err = strconv.Atoi(parts[1])
|
||||
if err != nil {
|
||||
return fmt.Errorf("Error parsing process index: %w", err)
|
||||
}
|
||||
}
|
||||
labelSelector = append(labelSelector, fmt.Sprintf("app.kubernetes.io/name=%s", input.ContainerName))
|
||||
}
|
||||
|
||||
pods, err := k.ListPods(ctx, ListPodsInput{
|
||||
Namespace: input.Namespace,
|
||||
LabelSelector: strings.Join(labelSelector, ","),
|
||||
})
|
||||
if err != nil {
|
||||
return fmt.Errorf("Error listing pods: %w", err)
|
||||
}
|
||||
if len(pods) == 0 {
|
||||
return fmt.Errorf("No pods found for app %s", input.DeploymentName)
|
||||
}
|
||||
|
||||
ch := make(chan bool)
|
||||
|
||||
if os.Getenv("FORCE_TTY") == "1" {
|
||||
color.NoColor = false
|
||||
}
|
||||
|
||||
colors := []color.Attribute{
|
||||
color.FgRed,
|
||||
color.FgYellow,
|
||||
color.FgGreen,
|
||||
color.FgCyan,
|
||||
color.FgBlue,
|
||||
color.FgMagenta,
|
||||
}
|
||||
// colorIndex := 0
|
||||
for i := 0; i < len(pods); i++ {
|
||||
if processIndex > 0 && i != (processIndex-1) {
|
||||
continue
|
||||
}
|
||||
|
||||
logOptions := v1.PodLogOptions{
|
||||
Follow: input.Follow,
|
||||
}
|
||||
if input.TailLines > 0 {
|
||||
logOptions.TailLines = ptr.To(input.TailLines)
|
||||
}
|
||||
|
||||
podColor := colors[i%len(colors)]
|
||||
dynoText := color.New(podColor).SprintFunc()
|
||||
podName := pods[i].Name
|
||||
podLogs, err := k.Client.CoreV1().Pods(input.Namespace).GetLogs(podName, &logOptions).Stream(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
buffer := bufio.NewReader(podLogs)
|
||||
go func(ctx context.Context, buffer *bufio.Reader, prettyText func(a ...interface{}) string, ch chan bool) {
|
||||
defer func() {
|
||||
ch <- true
|
||||
}()
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done(): // if cancel() execute
|
||||
ch <- true
|
||||
return
|
||||
default:
|
||||
str, readErr := buffer.ReadString('\n')
|
||||
if readErr == io.EOF {
|
||||
break
|
||||
}
|
||||
|
||||
if str == "" {
|
||||
continue
|
||||
}
|
||||
|
||||
if !input.Quiet {
|
||||
str = fmt.Sprintf("%s %s", dynoText(fmt.Sprintf("app[%s]:", podName)), str)
|
||||
}
|
||||
|
||||
_, err := fmt.Print(str)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
}(ctx, buffer, dynoText, ch)
|
||||
}
|
||||
<-ch
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -84,6 +84,22 @@ func main() {
|
||||
}
|
||||
|
||||
err = scheduler_k3s.TriggerSchedulerLogs(scheduler, appName, processType, tail, quiet, numLines)
|
||||
case "scheduler-proxy-logs":
|
||||
var tail bool
|
||||
var numLines int64
|
||||
scheduler := flag.Arg(0)
|
||||
appName := flag.Arg(1)
|
||||
proxyType := flag.Arg(2)
|
||||
logType := flag.Arg(3)
|
||||
tail, err = strconv.ParseBool(flag.Arg(4))
|
||||
if err != nil {
|
||||
tail = false
|
||||
}
|
||||
numLines, err = strconv.ParseInt(flag.Arg(5), 10, 64)
|
||||
if err != nil {
|
||||
numLines = 0
|
||||
}
|
||||
err = scheduler_k3s.TriggerSchedulerProxyLogs(scheduler, appName, proxyType, logType, tail, numLines)
|
||||
case "scheduler-stop":
|
||||
scheduler := flag.Arg(0)
|
||||
appName := flag.Arg(1)
|
||||
|
||||
@@ -1,14 +1,12 @@
|
||||
package scheduler_k3s
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"context"
|
||||
"crypto/rand"
|
||||
"encoding/base64"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"os/signal"
|
||||
"path/filepath"
|
||||
@@ -22,6 +20,7 @@ import (
|
||||
"github.com/dokku/dokku/plugins/common"
|
||||
"github.com/dokku/dokku/plugins/config"
|
||||
"github.com/dokku/dokku/plugins/cron"
|
||||
nginxvhosts "github.com/dokku/dokku/plugins/nginx-vhosts"
|
||||
"github.com/fatih/color"
|
||||
"github.com/gosimple/slug"
|
||||
"github.com/kballard/go-shellquote"
|
||||
@@ -29,7 +28,6 @@ import (
|
||||
corev1 "k8s.io/api/core/v1"
|
||||
v1 "k8s.io/api/core/v1"
|
||||
"k8s.io/kubernetes/pkg/client/conditions"
|
||||
"k8s.io/utils/ptr"
|
||||
)
|
||||
|
||||
// TriggerInstall runs the install step for the scheduler-k3s plugin
|
||||
@@ -828,121 +826,68 @@ func TriggerSchedulerLogs(scheduler string, appName string, processType string,
|
||||
return nil
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
signals := make(chan os.Signal, 1)
|
||||
signal.Notify(signals, os.Interrupt, syscall.SIGHUP,
|
||||
syscall.SIGINT,
|
||||
syscall.SIGQUIT,
|
||||
syscall.SIGTERM)
|
||||
go func() {
|
||||
<-signals
|
||||
cancel()
|
||||
}()
|
||||
clientset, err := NewKubernetesClient()
|
||||
if err != nil {
|
||||
return fmt.Errorf("Error creating kubernetes client: %w", err)
|
||||
}
|
||||
|
||||
return clientset.StreamLogs(context.Background(), StreamLogsInput{
|
||||
DeploymentName: appName,
|
||||
Namespace: getComputedNamespace(appName),
|
||||
ContainerName: processType,
|
||||
TailLines: numLines,
|
||||
Follow: tail,
|
||||
Quiet: quiet,
|
||||
})
|
||||
}
|
||||
|
||||
// TriggerSchedulerProxyLogs displays nginx logs for a given application
|
||||
func TriggerSchedulerProxyLogs(scheduler string, appName string, proxyType string, logType string, tail bool, numLines int64) error {
|
||||
if scheduler != "k3s" || proxyType != "nginx" {
|
||||
return nil
|
||||
}
|
||||
|
||||
clientset, err := NewKubernetesClient()
|
||||
if err != nil {
|
||||
return fmt.Errorf("Error creating kubernetes client: %w", err)
|
||||
}
|
||||
|
||||
if err := clientset.Ping(); err != nil {
|
||||
return fmt.Errorf("kubernetes api not available: %w", err)
|
||||
filename := ""
|
||||
if logType == "access" {
|
||||
filename = nginxvhosts.ComputedAccessLogPath(appName)
|
||||
} else if logType == "error" {
|
||||
filename = nginxvhosts.ComputedErrorLogPath(appName)
|
||||
} else {
|
||||
return errors.New("Invalid log type")
|
||||
}
|
||||
|
||||
labelSelector := []string{fmt.Sprintf("app.kubernetes.io/part-of=%s", appName)}
|
||||
processIndex := 0
|
||||
if processType != "" {
|
||||
parts := strings.SplitN(processType, ".", 2)
|
||||
if len(parts) == 2 {
|
||||
processType = parts[0]
|
||||
processIndex, err = strconv.Atoi(parts[1])
|
||||
if err != nil {
|
||||
return fmt.Errorf("Error parsing process index: %w", err)
|
||||
}
|
||||
}
|
||||
labelSelector = append(labelSelector, fmt.Sprintf("app.kubernetes.io/name=%s", processType))
|
||||
command := []string{"tail"}
|
||||
if tail {
|
||||
command = append(command, "-F")
|
||||
}
|
||||
if numLines > 0 {
|
||||
command = append(command, fmt.Sprintf("-n %d", numLines))
|
||||
}
|
||||
command = append(command, filename)
|
||||
|
||||
namespace := getComputedNamespace(appName)
|
||||
pods, err := clientset.ListPods(ctx, ListPodsInput{
|
||||
Namespace: namespace,
|
||||
LabelSelector: strings.Join(labelSelector, ","),
|
||||
pods, err := clientset.ListPods(context.Background(), ListPodsInput{
|
||||
Namespace: "ingress-nginx",
|
||||
LabelSelector: "app.kubernetes.io/name=ingress-nginx",
|
||||
})
|
||||
if err != nil {
|
||||
return fmt.Errorf("Error listing pods: %w", err)
|
||||
}
|
||||
|
||||
if len(pods) == 0 {
|
||||
return fmt.Errorf("No pods found for app %s", appName)
|
||||
return errors.New("No pods found for ingress-nginx")
|
||||
}
|
||||
|
||||
ch := make(chan bool)
|
||||
|
||||
if os.Getenv("FORCE_TTY") == "1" {
|
||||
color.NoColor = false
|
||||
}
|
||||
|
||||
colors := []color.Attribute{
|
||||
color.FgRed,
|
||||
color.FgYellow,
|
||||
color.FgGreen,
|
||||
color.FgCyan,
|
||||
color.FgBlue,
|
||||
color.FgMagenta,
|
||||
}
|
||||
// colorIndex := 0
|
||||
for i := 0; i < len(pods); i++ {
|
||||
if processIndex > 0 && i != (processIndex-1) {
|
||||
continue
|
||||
}
|
||||
|
||||
logOptions := v1.PodLogOptions{
|
||||
Follow: tail,
|
||||
}
|
||||
if numLines > 0 {
|
||||
logOptions.TailLines = ptr.To(numLines)
|
||||
}
|
||||
|
||||
podColor := colors[i%len(colors)]
|
||||
dynoText := color.New(podColor).SprintFunc()
|
||||
podName := pods[i].Name
|
||||
podLogs, err := clientset.Client.CoreV1().Pods(namespace).GetLogs(podName, &logOptions).Stream(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
buffer := bufio.NewReader(podLogs)
|
||||
go func(ctx context.Context, buffer *bufio.Reader, prettyText func(a ...interface{}) string, ch chan bool) {
|
||||
defer func() {
|
||||
ch <- true
|
||||
}()
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done(): // if cancel() execute
|
||||
ch <- true
|
||||
return
|
||||
default:
|
||||
str, readErr := buffer.ReadString('\n')
|
||||
if readErr == io.EOF {
|
||||
break
|
||||
}
|
||||
|
||||
if str == "" {
|
||||
continue
|
||||
}
|
||||
|
||||
if !quiet {
|
||||
str = fmt.Sprintf("%s %s", dynoText(fmt.Sprintf("app[%s]:", podName)), str)
|
||||
}
|
||||
|
||||
_, err := fmt.Print(str)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
}(ctx, buffer, dynoText, ch)
|
||||
}
|
||||
<-ch
|
||||
|
||||
return nil
|
||||
return clientset.ExecCommand(context.Background(), ExecCommandInput{
|
||||
Command: command,
|
||||
ContainerName: "ingress-nginx",
|
||||
Name: pods[0].Name,
|
||||
Namespace: "ingress-nginx",
|
||||
})
|
||||
}
|
||||
|
||||
// TriggerSchedulerRun runs a command in an ephemeral container
|
||||
|
||||
Reference in New Issue
Block a user