fix(auth): harden unified diagnostic execution
This commit is contained in:
@@ -13,9 +13,13 @@ import (
|
||||
"net"
|
||||
"net/url"
|
||||
"os"
|
||||
"os/exec"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"time"
|
||||
"unicode"
|
||||
"unicode/utf16"
|
||||
"unicode/utf8"
|
||||
|
||||
"github.com/aritmolab/thothii/tools/tht/internal/compose"
|
||||
"github.com/aritmolab/thothii/tools/tht/internal/config"
|
||||
@@ -76,9 +80,9 @@ func RunWithRunner(ctx context.Context, installation config.Installation, args [
|
||||
|
||||
// AuthDiagnostic is the closed JSON contract emitted by the backend diagnostic command.
|
||||
type AuthDiagnostic struct {
|
||||
Level string `json:"level"`
|
||||
Code string `json:"code"`
|
||||
Message string `json:"message"`
|
||||
Level string `json:"level"`
|
||||
Code string `json:"code"`
|
||||
Message string `json:"message"`
|
||||
Field *string `json:"field,omitempty"`
|
||||
}
|
||||
|
||||
@@ -89,6 +93,19 @@ type AuthDiagnostics struct {
|
||||
Checks []AuthDiagnostic `json:"checks"`
|
||||
}
|
||||
|
||||
type authDiagnosticWire struct {
|
||||
Level string `json:"level"`
|
||||
Code string `json:"code"`
|
||||
Message string `json:"message"`
|
||||
Field json.RawMessage `json:"field"`
|
||||
}
|
||||
|
||||
type authDiagnosticsWire struct {
|
||||
Ready bool `json:"ready"`
|
||||
Mode string `json:"mode"`
|
||||
Checks []authDiagnosticWire `json:"checks"`
|
||||
}
|
||||
|
||||
func parseCheckArgs(args []string) (jsonMode, interactive bool, err error) {
|
||||
for _, arg := range args {
|
||||
switch arg {
|
||||
@@ -160,20 +177,95 @@ func validAuthDiagnostics(report AuthDiagnostics) bool {
|
||||
if len(report.Checks) == 0 || len(report.Checks) > 129 {
|
||||
return false
|
||||
}
|
||||
seen := make(map[string]struct{}, len(report.Checks))
|
||||
hasError := false
|
||||
for _, check := range report.Checks {
|
||||
if (check.Level != "error" && check.Level != "info") || check.Message == "" || len(check.Message) > 512 {
|
||||
if (check.Level != "error" && check.Level != "info") || !safeDiagnosticText(check.Message) {
|
||||
return false
|
||||
}
|
||||
if _, ok := authDiagnosticCodes[check.Code]; !ok {
|
||||
return false
|
||||
}
|
||||
if check.Field != nil && (*check.Field == "" || len(*check.Field) > 512) {
|
||||
if check.Field != nil && (!safeDiagnosticText(*check.Field) ||
|
||||
check.Code != "oidc_mapped_group_missing" && check.Code != "oidc_mapped_group_ambiguous") {
|
||||
return false
|
||||
}
|
||||
field := ""
|
||||
if check.Field != nil {
|
||||
field = *check.Field
|
||||
}
|
||||
key := check.Code + "\x00" + field
|
||||
if _, duplicate := seen[key]; duplicate {
|
||||
return false
|
||||
}
|
||||
seen[key] = struct{}{}
|
||||
hasError = hasError || check.Level == "error"
|
||||
}
|
||||
if report.Ready {
|
||||
return len(report.Checks) == 1 && report.Checks[0].Level == "info" &&
|
||||
report.Checks[0].Code == "auth_ready" && report.Checks[0].Field == nil
|
||||
}
|
||||
if !hasError {
|
||||
return false
|
||||
}
|
||||
for _, check := range report.Checks {
|
||||
if check.Code == "auth_ready" {
|
||||
return false
|
||||
}
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
func safeDiagnosticText(value string) bool {
|
||||
if value == "" || utf16Length(value) > 512 || strings.TrimSpace(value) != value {
|
||||
return false
|
||||
}
|
||||
for _, character := range value {
|
||||
if unicode.IsControl(character) {
|
||||
return false
|
||||
}
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
func utf16Length(value string) int {
|
||||
length := 0
|
||||
for _, character := range value {
|
||||
length += utf16.RuneLen(character)
|
||||
}
|
||||
return length
|
||||
}
|
||||
|
||||
func decodeAuthDiagnostics(value string) (AuthDiagnostics, error) {
|
||||
if !utf8.ValidString(value) {
|
||||
return AuthDiagnostics{}, errors.New("authentication diagnostic report is invalid")
|
||||
}
|
||||
decoder := json.NewDecoder(strings.NewReader(value))
|
||||
decoder.DisallowUnknownFields()
|
||||
var wire authDiagnosticsWire
|
||||
if err := decoder.Decode(&wire); err != nil || decoder.Decode(&struct{}{}) != io.EOF {
|
||||
return AuthDiagnostics{}, errors.New("authentication diagnostic report is invalid")
|
||||
}
|
||||
report := AuthDiagnostics{Ready: wire.Ready, Mode: wire.Mode, Checks: make([]AuthDiagnostic, 0, len(wire.Checks))}
|
||||
for _, item := range wire.Checks {
|
||||
var field *string
|
||||
if item.Field != nil {
|
||||
var decoded string
|
||||
if string(item.Field) == "null" || json.Unmarshal(item.Field, &decoded) != nil {
|
||||
return AuthDiagnostics{}, errors.New("authentication diagnostic report is invalid")
|
||||
}
|
||||
field = &decoded
|
||||
}
|
||||
report.Checks = append(report.Checks, AuthDiagnostic{
|
||||
Level: item.Level, Code: item.Code, Message: item.Message, Field: field,
|
||||
})
|
||||
}
|
||||
if !validAuthDiagnostics(report) {
|
||||
return AuthDiagnostics{}, errors.New("authentication diagnostic report is invalid")
|
||||
}
|
||||
return report, nil
|
||||
}
|
||||
|
||||
func authenticationSecretValues(installation config.Installation) []string {
|
||||
files, err := installation.SecretFiles()
|
||||
if err != nil {
|
||||
@@ -260,18 +352,29 @@ func runCheck(ctx context.Context, installation config.Installation, runner comp
|
||||
if interactive {
|
||||
command = append(command, "--interactive")
|
||||
}
|
||||
result, err := runner.Run(bounded, installation.ComposeArgs(command...), nil)
|
||||
result, err := compose.RunBounded(runner, bounded, installation.ComposeArgs(command...), nil, compose.CaptureLimits{
|
||||
StdoutBytes: maxAuthDiagnosticOutputBytes,
|
||||
StderrBytes: maxAuthDiagnosticOutputBytes,
|
||||
})
|
||||
secrets := authenticationSecretValues(installation)
|
||||
if err != nil || result.ExitCode != 0 || len(result.Stdout) > maxAuthDiagnosticOutputBytes || len(result.Stderr) > maxAuthDiagnosticOutputBytes {
|
||||
var exitError *exec.ExitError
|
||||
validProcessOutcome := result.ExitCode == 0 && err == nil ||
|
||||
result.ExitCode == 1 && (err == nil || errors.As(err, &exitError))
|
||||
if !validProcessOutcome {
|
||||
return AuthDiagnostics{}, "", errors.New("authentication diagnostic command failed")
|
||||
}
|
||||
decoder := json.NewDecoder(strings.NewReader(result.Stdout))
|
||||
decoder.DisallowUnknownFields()
|
||||
var report AuthDiagnostics
|
||||
if err := decoder.Decode(&report); err != nil || decoder.Decode(&struct{}{}) != io.EOF || !validAuthDiagnostics(report) {
|
||||
report, decodeErr := decodeAuthDiagnostics(result.Stdout)
|
||||
if decodeErr != nil {
|
||||
return AuthDiagnostics{}, "", errors.New("authentication diagnostic report is invalid")
|
||||
}
|
||||
return sanitizeAuthDiagnostics(report, secrets), devicePrompt(result.Stderr, secrets), nil
|
||||
if report.Ready != (result.ExitCode == 0) {
|
||||
return AuthDiagnostics{}, "", errors.New("authentication diagnostic command failed")
|
||||
}
|
||||
safe := sanitizeAuthDiagnostics(report, secrets)
|
||||
if !validAuthDiagnostics(safe) {
|
||||
return AuthDiagnostics{}, "", errors.New("authentication diagnostic report is invalid")
|
||||
}
|
||||
return safe, devicePrompt(result.Stderr, secrets), nil
|
||||
}
|
||||
|
||||
func authFailure(stderr io.Writer, message string) int {
|
||||
|
||||
@@ -49,6 +49,113 @@ func TestAuthCheckRunsOneShotCoreDiagnosticWithPristineJSON(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestAuthCheckEmitsValidFailedReportFromRealExitError(t *testing.T) {
|
||||
installation := authInstallation(newAuthDirectory(t))
|
||||
runner := compose.NewRunner(writeAuthExecutable(t, `#!/bin/sh
|
||||
printf '%s\n' '{"ready":false,"mode":"oidc","checks":[{"level":"error","code":"oidc_secret_missing","message":"A required OIDC or group catalog secret is unavailable."}]}'
|
||||
exit 1
|
||||
`))
|
||||
var stdout, stderr bytes.Buffer
|
||||
|
||||
code := RunWithRunner(context.Background(), installation, []string{"check", "--json"}, strings.NewReader(""), &stdout, &stderr, runner)
|
||||
|
||||
if code != 1 {
|
||||
t.Fatalf("auth check = %d, want diagnostic failure 1; stdout=%q stderr=%q", code, stdout.String(), stderr.String())
|
||||
}
|
||||
var report AuthDiagnostics
|
||||
if err := json.Unmarshal(stdout.Bytes(), &report); err != nil || report.Ready || report.Checks[0].Code != "oidc_secret_missing" {
|
||||
t.Fatalf("auth check stdout is not the pristine failed report: %q: %#v, %v", stdout.String(), report, err)
|
||||
}
|
||||
if stderr.Len() != 0 {
|
||||
t.Fatalf("auth check stderr = %q, want empty", stderr.String())
|
||||
}
|
||||
}
|
||||
|
||||
func TestAuthCheckRejectsExitAndReportSemanticMismatches(t *testing.T) {
|
||||
ready := `{"ready":true,"mode":"oidc","checks":[{"level":"info","code":"auth_ready","message":"Authentication is ready."}]}`
|
||||
failed := `{"ready":false,"mode":"oidc","checks":[{"level":"error","code":"oidc_secret_missing","message":"A required OIDC or group catalog secret is unavailable."}]}`
|
||||
for _, test := range []struct {
|
||||
name string
|
||||
exit int
|
||||
report string
|
||||
}{
|
||||
{name: "zero with failed report", exit: 0, report: failed},
|
||||
{name: "one with ready report", exit: 1, report: ready},
|
||||
{name: "two with failed report", exit: 2, report: failed},
|
||||
} {
|
||||
t.Run(test.name, func(t *testing.T) {
|
||||
runner := runnerFunc(func(_ context.Context, _ []string, _ io.Reader) (compose.Result, error) {
|
||||
var err error
|
||||
if test.exit != 0 {
|
||||
err = errors.New("process exited")
|
||||
}
|
||||
return compose.Result{Stdout: test.report, ExitCode: test.exit}, err
|
||||
})
|
||||
var stdout, stderr bytes.Buffer
|
||||
code := RunWithRunner(context.Background(), authInstallation(newAuthDirectory(t)), []string{"check", "--json"}, strings.NewReader(""), &stdout, &stderr, runner)
|
||||
if code != 1 || stdout.Len() != 0 || stderr.String() != "tht: authentication diagnostics could not be completed\n" {
|
||||
t.Fatalf("mismatch accepted: code=%d stdout=%q stderr=%q", code, stdout.String(), stderr.String())
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestAuthCheckRejectsCancellationEvenIfTheKilledChildReportsExitOne(t *testing.T) {
|
||||
failed := `{"ready":false,"mode":"oidc","checks":[{"level":"error","code":"oidc_secret_missing","message":"A required OIDC or group catalog secret is unavailable."}]}`
|
||||
runner := runnerFunc(func(_ context.Context, _ []string, _ io.Reader) (compose.Result, error) {
|
||||
return compose.Result{Stdout: failed, ExitCode: 1}, context.DeadlineExceeded
|
||||
})
|
||||
var stdout, stderr bytes.Buffer
|
||||
|
||||
code := RunWithRunner(context.Background(), authInstallation(newAuthDirectory(t)), []string{"check", "--json"}, strings.NewReader(""), &stdout, &stderr, runner)
|
||||
|
||||
if code != 1 || stdout.Len() != 0 || stderr.String() != "tht: authentication diagnostics could not be completed\n" {
|
||||
t.Fatalf("cancelled report accepted: code=%d stdout=%q stderr=%q", code, stdout.String(), stderr.String())
|
||||
}
|
||||
}
|
||||
|
||||
func TestAuthDiagnosticsContractRejectsContradictionsDuplicatesAndAttackerFields(t *testing.T) {
|
||||
field := "Configured Group"
|
||||
validFailure := AuthDiagnostics{Ready: false, Mode: "oidc", Checks: []AuthDiagnostic{{
|
||||
Level: "error", Code: "oidc_mapped_group_missing", Message: "A configured authorization group does not exist.", Field: &field,
|
||||
}}}
|
||||
if !validAuthDiagnostics(validFailure) {
|
||||
t.Fatal("valid failed report was rejected")
|
||||
}
|
||||
for _, report := range []AuthDiagnostics{
|
||||
{Ready: true, Mode: "oidc", Checks: []AuthDiagnostic{{Level: "error", Code: "oidc_secret_missing", Message: "failure"}}},
|
||||
{Ready: false, Mode: "oidc", Checks: []AuthDiagnostic{{Level: "info", Code: "auth_ready", Message: "ready"}}},
|
||||
{Ready: false, Mode: "oidc", Checks: []AuthDiagnostic{{Level: "info", Code: "auth_config_invalid", Message: "not an error"}}},
|
||||
{Ready: false, Mode: "oidc", Checks: []AuthDiagnostic{
|
||||
{Level: "error", Code: "oidc_secret_missing", Message: "failure"},
|
||||
{Level: "error", Code: "oidc_secret_missing", Message: "duplicate"},
|
||||
}},
|
||||
{Ready: false, Mode: "oidc", Checks: []AuthDiagnostic{{Level: "error", Code: "oidc_secret_missing", Message: "failure", Field: &field}}},
|
||||
{Ready: false, Mode: "oidc", Checks: []AuthDiagnostic{{Level: "error", Code: "oidc_mapped_group_missing", Message: "failure", Field: stringPointer(" attacker ")}}},
|
||||
{Ready: false, Mode: "oidc", Checks: []AuthDiagnostic{{Level: "error", Code: "oidc_mapped_group_missing", Message: "failure", Field: stringPointer("attacker\u0085field")}}},
|
||||
} {
|
||||
if validAuthDiagnostics(report) {
|
||||
t.Fatalf("invalid authentication report accepted: %#v", report)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestAuthCheckRejectsNullAndUnexpectedDiagnosticFields(t *testing.T) {
|
||||
for _, report := range []string{
|
||||
`{"ready":false,"mode":"oidc","checks":[{"level":"error","code":"oidc_mapped_group_missing","message":"failure","field":null}]}`,
|
||||
`{"ready":false,"mode":"oidc","checks":[{"level":"error","code":"oidc_secret_missing","message":"failure","unexpected":"attacker"}]}`,
|
||||
} {
|
||||
runner := runnerFunc(func(_ context.Context, _ []string, _ io.Reader) (compose.Result, error) {
|
||||
return compose.Result{Stdout: report, ExitCode: 1}, nil
|
||||
})
|
||||
var stdout, stderr bytes.Buffer
|
||||
code := RunWithRunner(context.Background(), authInstallation(newAuthDirectory(t)), []string{"check", "--json"}, strings.NewReader(""), &stdout, &stderr, runner)
|
||||
if code != 1 || stdout.Len() != 0 || stderr.String() != "tht: authentication diagnostics could not be completed\n" {
|
||||
t.Fatalf("hostile field accepted: code=%d stdout=%q stderr=%q", code, stdout.String(), stderr.String())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestAuthCheckInteractiveForwardsOnlyTheValidatedDevicePrompt(t *testing.T) {
|
||||
installation := authInstallation(newAuthDirectory(t))
|
||||
var calls [][]string
|
||||
@@ -96,8 +203,8 @@ func TestAuthCheckRedactsFailedCoreOutputAndRejectsMalformedReports(t *testing.T
|
||||
}
|
||||
runner := runnerFunc(func(_ context.Context, _ []string, _ io.Reader) (compose.Result, error) {
|
||||
return compose.Result{
|
||||
Stdout: "not-json auth-check-secret token-sentinel /private/sentinel $argon2id$hash-sentinel",
|
||||
Stderr: "auth-check-secret cookie-sentinel /private/sentinel",
|
||||
Stdout: "not-json auth-check-secret token-sentinel /private/sentinel $argon2id$hash-sentinel",
|
||||
Stderr: "auth-check-secret cookie-sentinel /private/sentinel",
|
||||
ExitCode: 23,
|
||||
}, errors.New("core failed")
|
||||
})
|
||||
@@ -419,6 +526,17 @@ func writePasswordFile(t *testing.T, password string) string {
|
||||
return path
|
||||
}
|
||||
|
||||
func writeAuthExecutable(t *testing.T, contents string) string {
|
||||
t.Helper()
|
||||
path := filepath.Join(t.TempDir(), "fake-docker")
|
||||
if err := os.WriteFile(path, []byte(contents), 0o700); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return path
|
||||
}
|
||||
|
||||
func stringPointer(value string) *string { return &value }
|
||||
|
||||
func authInstallation(directory string) config.Installation {
|
||||
installation := config.Installation{}
|
||||
installation.Authentication.ConfigDirectory = directory
|
||||
|
||||
@@ -0,0 +1,37 @@
|
||||
//go:build !windows
|
||||
|
||||
package compose
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"os/exec"
|
||||
"syscall"
|
||||
"time"
|
||||
)
|
||||
|
||||
const gracefulTerminationBound = 500 * time.Millisecond
|
||||
const finalTerminationBound = 2 * time.Second
|
||||
|
||||
func configureProcess(command *exec.Cmd) {
|
||||
command.SysProcAttr = &syscall.SysProcAttr{Setpgid: true}
|
||||
}
|
||||
|
||||
func terminateProcess(command *exec.Cmd, done <-chan error) error {
|
||||
if command.Process == nil {
|
||||
return nil
|
||||
}
|
||||
_ = syscall.Kill(-command.Process.Pid, syscall.SIGINT)
|
||||
select {
|
||||
case <-done:
|
||||
return nil
|
||||
case <-time.After(gracefulTerminationBound):
|
||||
}
|
||||
_ = syscall.Kill(-command.Process.Pid, syscall.SIGKILL)
|
||||
_ = command.Process.Kill()
|
||||
select {
|
||||
case <-done:
|
||||
return nil
|
||||
case <-time.After(finalTerminationBound):
|
||||
return errors.New("Docker command could not be reaped")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,30 @@
|
||||
//go:build windows
|
||||
|
||||
package compose
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"os/exec"
|
||||
"syscall"
|
||||
"time"
|
||||
)
|
||||
|
||||
const finalTerminationBound = 2 * time.Second
|
||||
const createNewProcessGroup = 0x00000200
|
||||
|
||||
func configureProcess(command *exec.Cmd) {
|
||||
command.SysProcAttr = &syscall.SysProcAttr{CreationFlags: createNewProcessGroup}
|
||||
}
|
||||
|
||||
func terminateProcess(command *exec.Cmd, done <-chan error) error {
|
||||
if command.Process == nil {
|
||||
return nil
|
||||
}
|
||||
_ = command.Process.Kill()
|
||||
select {
|
||||
case <-done:
|
||||
return nil
|
||||
case <-time.After(finalTerminationBound):
|
||||
return errors.New("Docker command could not be reaped")
|
||||
}
|
||||
}
|
||||
@@ -2,17 +2,29 @@
|
||||
package compose
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"os/exec"
|
||||
"sync"
|
||||
|
||||
"github.com/aritmolab/thothii/tools/tht/internal/config"
|
||||
)
|
||||
|
||||
const defaultCaptureBytes = 4 * 1024 * 1024
|
||||
const maximumCaptureBytes = 64 * 1024 * 1024
|
||||
|
||||
// ErrOutputLimit reports that a child exceeded one of its capture limits.
|
||||
var ErrOutputLimit = errors.New("Docker command output limit exceeded")
|
||||
|
||||
// CaptureLimits bounds each captured stream while the child is running.
|
||||
type CaptureLimits struct {
|
||||
StdoutBytes int
|
||||
StderrBytes int
|
||||
}
|
||||
|
||||
// Result is the captured output and process exit code for one Docker invocation.
|
||||
type Result struct {
|
||||
Stdout string
|
||||
@@ -25,6 +37,10 @@ type Runner interface {
|
||||
Run(context.Context, []string, io.Reader) (Result, error)
|
||||
}
|
||||
|
||||
type boundedRunner interface {
|
||||
RunBounded(context.Context, []string, io.Reader, CaptureLimits) (Result, error)
|
||||
}
|
||||
|
||||
// execRunner executes the Docker CLI. It never invokes a shell.
|
||||
type execRunner struct {
|
||||
binary string
|
||||
@@ -40,21 +56,101 @@ func NewRunner(binary string) Runner {
|
||||
|
||||
// Run invokes Docker with the supplied argument array and optional standard input.
|
||||
func (r execRunner) Run(ctx context.Context, args []string, stdin io.Reader) (Result, error) {
|
||||
command := exec.CommandContext(ctx, r.binary, args...)
|
||||
return r.RunBounded(ctx, args, stdin, CaptureLimits{
|
||||
StdoutBytes: defaultCaptureBytes,
|
||||
StderrBytes: defaultCaptureBytes,
|
||||
})
|
||||
}
|
||||
|
||||
// RunBounded invokes Docker while enforcing both stream limits during capture.
|
||||
func (r execRunner) RunBounded(ctx context.Context, args []string, stdin io.Reader, limits CaptureLimits) (Result, error) {
|
||||
if err := validCaptureLimits(limits); err != nil {
|
||||
return Result{}, err
|
||||
}
|
||||
if err := ctx.Err(); err != nil {
|
||||
return Result{}, err
|
||||
}
|
||||
command := exec.Command(r.binary, args...)
|
||||
configureProcess(command)
|
||||
command.Stdin = stdin
|
||||
var stdout, stderr bytes.Buffer
|
||||
command.Stdout = &stdout
|
||||
command.Stderr = &stderr
|
||||
err := command.Run()
|
||||
overflow := make(chan struct{}, 1)
|
||||
stdout := newCappedBuffer(limits.StdoutBytes, overflow)
|
||||
stderr := newCappedBuffer(limits.StderrBytes, overflow)
|
||||
command.Stdout = stdout
|
||||
command.Stderr = stderr
|
||||
if err := command.Start(); err != nil {
|
||||
return startFailure(err)
|
||||
}
|
||||
done := make(chan error, 1)
|
||||
go func() { done <- command.Wait() }()
|
||||
var err error
|
||||
select {
|
||||
case err = <-done:
|
||||
case <-ctx.Done():
|
||||
_ = terminateProcess(command, done)
|
||||
err = ctx.Err()
|
||||
case <-overflow:
|
||||
_ = terminateProcess(command, done)
|
||||
err = ErrOutputLimit
|
||||
}
|
||||
result := Result{Stdout: stdout.String(), Stderr: stderr.String()}
|
||||
if command.ProcessState != nil {
|
||||
result.ExitCode = command.ProcessState.ExitCode()
|
||||
}
|
||||
if stdout.Overflowed() || stderr.Overflowed() {
|
||||
return result, ErrOutputLimit
|
||||
}
|
||||
if err == nil {
|
||||
return result, nil
|
||||
}
|
||||
if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) || errors.Is(err, ErrOutputLimit) {
|
||||
return result, err
|
||||
}
|
||||
var exitError *exec.ExitError
|
||||
if errors.As(err, &exitError) {
|
||||
result.ExitCode = exitError.ExitCode()
|
||||
return result, err
|
||||
}
|
||||
return result, err
|
||||
}
|
||||
|
||||
// RunBounded uses the production runner's during-capture limits while retaining compatibility
|
||||
// with injected runners, whose already-bounded test results are checked before use.
|
||||
func RunBounded(runner Runner, ctx context.Context, args []string, stdin io.Reader, limits CaptureLimits) (Result, error) {
|
||||
if err := validCaptureLimits(limits); err != nil {
|
||||
return Result{}, err
|
||||
}
|
||||
if bounded, ok := runner.(boundedRunner); ok {
|
||||
return bounded.RunBounded(ctx, args, stdin, limits)
|
||||
}
|
||||
result, err := runner.Run(ctx, args, stdin)
|
||||
if len(result.Stdout) > limits.StdoutBytes || len(result.Stderr) > limits.StderrBytes {
|
||||
return Result{
|
||||
Stdout: boundedString(result.Stdout, limits.StdoutBytes),
|
||||
Stderr: boundedString(result.Stderr, limits.StderrBytes),
|
||||
ExitCode: result.ExitCode,
|
||||
}, ErrOutputLimit
|
||||
}
|
||||
return result, err
|
||||
}
|
||||
|
||||
func validCaptureLimits(limits CaptureLimits) error {
|
||||
if limits.StdoutBytes < 1 || limits.StderrBytes < 1 ||
|
||||
limits.StdoutBytes > maximumCaptureBytes || limits.StderrBytes > maximumCaptureBytes {
|
||||
return errors.New("Docker command capture limits are invalid")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func boundedString(value string, maximum int) string {
|
||||
if len(value) <= maximum {
|
||||
return value
|
||||
}
|
||||
return value[:maximum]
|
||||
}
|
||||
|
||||
func startFailure(err error) (Result, error) {
|
||||
result := Result{}
|
||||
if errors.Is(err, exec.ErrNotFound) || errors.Is(err, os.ErrNotExist) {
|
||||
result.ExitCode = 127
|
||||
return result, fmt.Errorf("%w: %w", exec.ErrNotFound, err)
|
||||
@@ -62,6 +158,55 @@ func (r execRunner) Run(ctx context.Context, args []string, stdin io.Reader) (Re
|
||||
return result, err
|
||||
}
|
||||
|
||||
type cappedBuffer struct {
|
||||
mu sync.Mutex
|
||||
contents []byte
|
||||
maximum int
|
||||
overflow chan<- struct{}
|
||||
exceeded bool
|
||||
}
|
||||
|
||||
func newCappedBuffer(maximum int, overflow chan<- struct{}) *cappedBuffer {
|
||||
capacity := maximum
|
||||
if capacity > 4096 {
|
||||
capacity = 4096
|
||||
}
|
||||
return &cappedBuffer{contents: make([]byte, 0, capacity), maximum: maximum, overflow: overflow}
|
||||
}
|
||||
|
||||
func (b *cappedBuffer) Write(value []byte) (int, error) {
|
||||
b.mu.Lock()
|
||||
defer b.mu.Unlock()
|
||||
remaining := b.maximum - len(b.contents)
|
||||
if remaining > 0 {
|
||||
kept := len(value)
|
||||
if kept > remaining {
|
||||
kept = remaining
|
||||
}
|
||||
b.contents = append(b.contents, value[:kept]...)
|
||||
}
|
||||
if len(value) > remaining && !b.exceeded {
|
||||
b.exceeded = true
|
||||
select {
|
||||
case b.overflow <- struct{}{}:
|
||||
default:
|
||||
}
|
||||
}
|
||||
return len(value), nil
|
||||
}
|
||||
|
||||
func (b *cappedBuffer) String() string {
|
||||
b.mu.Lock()
|
||||
defer b.mu.Unlock()
|
||||
return string(b.contents)
|
||||
}
|
||||
|
||||
func (b *cappedBuffer) Overflowed() bool {
|
||||
b.mu.Lock()
|
||||
defer b.mu.Unlock()
|
||||
return b.exceeded
|
||||
}
|
||||
|
||||
// InstallationRunner applies an installation's validated Compose arguments to commands that
|
||||
// explicitly start with compose. Direct Docker image commands remain host-side.
|
||||
type InstallationRunner struct {
|
||||
@@ -84,3 +229,11 @@ func (r InstallationRunner) Run(ctx context.Context, args []string, stdin io.Rea
|
||||
}
|
||||
return r.Runner.Run(ctx, args, stdin)
|
||||
}
|
||||
|
||||
// RunBounded preserves bounded capture when Compose argument injection is wrapped per installation.
|
||||
func (r InstallationRunner) RunBounded(ctx context.Context, args []string, stdin io.Reader, limits CaptureLimits) (Result, error) {
|
||||
if len(args) > 0 && args[0] == "compose" {
|
||||
return RunBounded(r.Runner, ctx, r.Installation.ComposeArgs(args[1:]...), stdin, limits)
|
||||
}
|
||||
return RunBounded(r.Runner, ctx, args, stdin, limits)
|
||||
}
|
||||
|
||||
@@ -9,6 +9,7 @@ import (
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/aritmolab/thothii/tools/tht/internal/config"
|
||||
)
|
||||
@@ -30,6 +31,45 @@ func TestRunnerPassesEachArgumentWithoutShellSplitting(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestRunnerBoundsFloodingOutputDuringCapture(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
runner := NewRunner(writeExecutable(t, "#!/bin/sh\nwhile :; do printf '0123456789abcdef'; printf 'fedcba9876543210' >&2; done\n"))
|
||||
started := time.Now()
|
||||
result, err := RunBounded(runner, context.Background(), []string{"compose", "run", "--rm", "core"}, nil, CaptureLimits{
|
||||
StdoutBytes: 1024,
|
||||
StderrBytes: 1024,
|
||||
})
|
||||
if !errors.Is(err, ErrOutputLimit) {
|
||||
t.Fatalf("RunBounded() error = %v, want ErrOutputLimit", err)
|
||||
}
|
||||
if len(result.Stdout) > 1024 || len(result.Stderr) > 1024 {
|
||||
t.Fatalf("captured output exceeded limits: stdout=%d stderr=%d", len(result.Stdout), len(result.Stderr))
|
||||
}
|
||||
if elapsed := time.Since(started); elapsed > 5*time.Second {
|
||||
t.Fatalf("overflow teardown took %s", elapsed)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRunnerCancelsAndReapsAHangingChildWithinFinalBound(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
runner := NewRunner(writeExecutable(t, "#!/bin/sh\ntrap '' TERM INT\nwhile :; do sleep 1; done\n"))
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 50*time.Millisecond)
|
||||
defer cancel()
|
||||
started := time.Now()
|
||||
_, err := RunBounded(runner, ctx, []string{"compose", "run", "--rm", "core"}, nil, CaptureLimits{
|
||||
StdoutBytes: 1024,
|
||||
StderrBytes: 1024,
|
||||
})
|
||||
if !errors.Is(err, context.DeadlineExceeded) {
|
||||
t.Fatalf("RunBounded() error = %v, want deadline exceeded", err)
|
||||
}
|
||||
if elapsed := time.Since(started); elapsed > 5*time.Second {
|
||||
t.Fatalf("cancellation teardown took %s", elapsed)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRunnerReturnsTheChildExitCode(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
|
||||
@@ -155,27 +155,36 @@ func RunWithProbe(ctx context.Context, installation config.Installation, runner
|
||||
|
||||
status, statusAvailable, servicesCheck := serviceStatus(ctx, installation, runner, secretValues)
|
||||
coreRunning := false
|
||||
coreHealthy := false
|
||||
if statusAvailable {
|
||||
var err error
|
||||
coreRunning, err = service.CoreRunning(status)
|
||||
if err != nil {
|
||||
var runningErr, healthyErr error
|
||||
coreRunning, runningErr = service.CoreRunning(status)
|
||||
coreHealthy, healthyErr = service.CoreHealthy(status)
|
||||
if runningErr != nil || healthyErr != nil {
|
||||
coreRunning = false
|
||||
coreHealthy = false
|
||||
}
|
||||
}
|
||||
if !configReady || !coreRunning {
|
||||
add("authentication", StatusSkipped, "core is unavailable")
|
||||
} else if !coreHealthy {
|
||||
add("authentication", StatusFailed, "core is running but unhealthy")
|
||||
} else if authenticationCheck(ctx, installation, runner, secretValues) {
|
||||
add("authentication", StatusPassed, "container-local authentication diagnostics passed")
|
||||
} else {
|
||||
add("authentication", StatusFailed, "container-local authentication diagnostics failed")
|
||||
}
|
||||
add(servicesCheck.Name, servicesCheck.Status, servicesCheck.Detail)
|
||||
if !coreRunning {
|
||||
add("core-http", StatusSkipped, "core is not running")
|
||||
add("frontend-http", StatusSkipped, "core is not running")
|
||||
add("workspace-registry", StatusSkipped, "core is not running")
|
||||
add("workflow", StatusSkipped, "core is not running")
|
||||
add("pi", StatusSkipped, "core is not running")
|
||||
if !coreHealthy || servicesCheck.Status != StatusPassed {
|
||||
detail := "required services are not healthy"
|
||||
if !coreRunning {
|
||||
detail = "core is not running"
|
||||
}
|
||||
add("core-http", StatusSkipped, detail)
|
||||
add("frontend-http", StatusSkipped, detail)
|
||||
add("workspace-registry", StatusSkipped, detail)
|
||||
add("workflow", StatusSkipped, detail)
|
||||
add("pi", StatusSkipped, detail)
|
||||
return finalize(report), nil
|
||||
}
|
||||
|
||||
|
||||
@@ -60,6 +60,43 @@ func TestRunSkipsContainerDiagnosticsWhenCoreIsStopped(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestRunFailsAuthenticationWithoutExecWhenCoreIsRunningButUnhealthy(t *testing.T) {
|
||||
installation := doctorInstallation(t, "")
|
||||
runner := &doctorRunner{services: unhealthyCoreServices}
|
||||
|
||||
report, err := Run(context.Background(), installation, runner)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if report.OK || checkStatus(report, "authentication") != StatusFailed {
|
||||
t.Fatalf("Run() report = %#v, want deterministic failed authentication", report)
|
||||
}
|
||||
assertChecklist(t, report, []string{"descriptor", "files", "docker", "compose", "configuration", "authentication", "services", "core-http", "frontend-http", "workspace-registry", "workflow", "pi"})
|
||||
if strings.Contains(strings.Join(runner.calls, "\n"), " exec -T ") {
|
||||
t.Fatalf("Run() invoked exec -T while core was unhealthy: %v", runner.calls)
|
||||
}
|
||||
if detail := checkDetail(report, "authentication"); detail != "core is running but unhealthy" {
|
||||
t.Fatalf("authentication detail = %q, want deterministic unhealthy detail", detail)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRunExecutesOnlyAuthenticationWhenCoreIsHealthyButAnotherServiceIsUnhealthy(t *testing.T) {
|
||||
installation := doctorInstallation(t, "")
|
||||
runner := &doctorRunner{services: unhealthyFrontendServices}
|
||||
|
||||
report, err := Run(context.Background(), installation, runner)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if checkStatus(report, "authentication") != StatusPassed || checkStatus(report, "services") != StatusFailed {
|
||||
t.Fatalf("Run() report = %#v, want auth passed before failed services", report)
|
||||
}
|
||||
calls := strings.Join(runner.calls, "\n")
|
||||
if strings.Count(calls, " exec -T ") != 1 || !strings.Contains(calls, "exec -T core node dist/auth/diagnostic-command.js --json") {
|
||||
t.Fatalf("Run() calls = %s, want only the healthy-core authentication exec", calls)
|
||||
}
|
||||
}
|
||||
|
||||
// Catches host-Python diagnostics or omission of workflow/Pi checks once core is healthy.
|
||||
func TestRunUsesOnlyContainerLocalWorkflowAndPiDiagnosticsWhenCoreRuns(t *testing.T) {
|
||||
installation := doctorInstallation(t, "")
|
||||
@@ -240,6 +277,15 @@ func checkStatus(report Report, name string) string {
|
||||
return ""
|
||||
}
|
||||
|
||||
func checkDetail(report Report, name string) string {
|
||||
for _, check := range report.Checks {
|
||||
if check.Name == name {
|
||||
return check.Detail
|
||||
}
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
func reportText(report Report) string {
|
||||
parts := make([]string, 0, len(report.Checks))
|
||||
for _, check := range report.Checks {
|
||||
@@ -273,3 +319,19 @@ const stoppedServices = `[
|
||||
{"Service":"core","State":"exited","Health":""},
|
||||
{"Service":"frontend","State":"running","Health":"healthy"}
|
||||
]`
|
||||
|
||||
const unhealthyCoreServices = `[
|
||||
{"Service":"core","State":"running","Health":"unhealthy"},
|
||||
{"Service":"frontend","State":"running","Health":"healthy"},
|
||||
{"Service":"qdrant","State":"running","Health":"healthy"},
|
||||
{"Service":"embedding","State":"running","Health":"healthy"},
|
||||
{"Service":"embedding-model-init","State":"exited","ExitCode":0}
|
||||
]`
|
||||
|
||||
const unhealthyFrontendServices = `[
|
||||
{"Service":"core","State":"running","Health":"healthy"},
|
||||
{"Service":"frontend","State":"running","Health":"unhealthy"},
|
||||
{"Service":"qdrant","State":"running","Health":"healthy"},
|
||||
{"Service":"embedding","State":"running","Health":"healthy"},
|
||||
{"Service":"embedding-model-init","State":"exited","ExitCode":0}
|
||||
]`
|
||||
|
||||
@@ -97,6 +97,20 @@ func CoreRunning(value string) (bool, error) {
|
||||
return false, nil
|
||||
}
|
||||
|
||||
// CoreHealthy reports whether core is both running and certified healthy by Compose.
|
||||
func CoreHealthy(value string) (bool, error) {
|
||||
statuses, err := parseStatuses(value)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
for _, item := range statuses {
|
||||
if item.Service == "core" {
|
||||
return strings.EqualFold(item.State, "running") && strings.EqualFold(item.Health, "healthy"), nil
|
||||
}
|
||||
}
|
||||
return false, nil
|
||||
}
|
||||
|
||||
// Healthy verifies the full expected Compose service set, including the one-shot model initializer.
|
||||
func Healthy(value string) error {
|
||||
statuses, err := parseStatuses(value)
|
||||
|
||||
Reference in New Issue
Block a user