Files
ThothII/tools/thothctl/cmd/thothctl/main.go
T

948 lines
32 KiB
Go

// thothctl is the host-side operator command for a local ThothII installation.
package main
import (
"bufio"
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"io"
"os"
"os/exec"
"path/filepath"
"strconv"
"strings"
"unicode"
"github.com/aritmolab/thothii/tools/thothctl/internal/compose"
"github.com/aritmolab/thothii/tools/thothctl/internal/config"
"github.com/aritmolab/thothii/tools/thothctl/internal/output"
"github.com/aritmolab/thothii/tools/thothctl/internal/pi"
"github.com/aritmolab/thothii/tools/thothctl/internal/serverops"
"github.com/aritmolab/thothii/tools/thothctl/internal/workspaceops"
)
const (
maxPublicStdoutBytes = 1 << 20
maxPublicStderrBytes = 64 << 10
)
const usage = `Usage: thothctl --installation <absolute-path>/thothii-installation.yaml <command>
Commands:
status Show the Compose service state.
doctor Validate Docker, Compose, rendered configuration, line endings, volumes, and health.
logs Show the latest 200 sanitized service log lines (bounded; no follow mode).
start Start the installation in the background.
stop Stop the installation.
update --check-only Validate the current installation without changing containers.
sessions migrate --yes
Run only the server session migrator and verify pending=[] and drifted=[].
remove Display exact stopped app container IDs without mutation.
remove --yes ID... Remove only the stopped IDs copied from the preceding display.
pi status Show the Pi version embedded in core.
pi doctor Check Pi preconditions without changing the installation.
pi test Run the temporary Pi/core smoke checks.
pi check Alias for pi test.
pi configure [--provider P --model M --thinking low|medium|high]
Select closed backend defaults interactively on a TTY; all flags are required otherwise.
pi update --version V --source build --yes [--drain]
Rebuild a pinned Pi version and recreate only core.
pi update --version V --source pull --image IMAGE@sha256:DIGEST --yes [--drain]
Pull an immutable candidate and recreate only core.
pi rollback --yes Restore the image recorded by the latest Pi update.
pi maintenance status
Show the durable core admission-gate state.
pi maintenance recover --yes
Verify a terminal installation, remove stale lifecycle files, and clear maintenance.
pi logs Show the latest 200 sanitized core log lines (bounded; no follow mode).
workspace inspect --workspace ID [--json]
workspace preprocess dwh|evidence|run --workspace ID [options]
workspace schema suggest-fks|check --workspace ID [options]
workspace index-schema --workspace ID [--json]
`
func main() {
os.Exit(run(context.Background(), os.Args[1:], os.Stdout, os.Stderr))
}
func run(ctx context.Context, args []string, stdout, stderr io.Writer) int {
isWorkspaceCommand := len(args) > 2 && args[2] == "workspace"
// Workspace results are untrusted child output, so only that dispatch receives
// the public bounds. Legacy renderers intentionally retain their established
// behavior and must not silently truncate successful output.
if len(args) > 2 && args[2] == "workspace" {
stdout = &boundedWriter{dst: stdout, maximum: maxPublicStdoutBytes}
stderr = &boundedWriter{dst: stderr, maximum: maxPublicStderrBytes}
}
if len(args) == 1 && (args[0] == "--help" || args[0] == "-h") {
fmt.Fprint(stdout, usage)
return 0
}
installationPath, command, commandArgs, err := parseArgs(args)
if err != nil {
if isWorkspaceCommand {
return writeWorkspaceError(stderr, err, nil, 2)
}
fmt.Fprintf(stderr, "thothctl: %s\n\n%s", err, usage)
return 2
}
installation, err := config.Load(installationPath)
if err != nil {
if isWorkspaceCommand {
return writeWorkspaceError(stderr, err, nil, 2)
}
fmt.Fprintf(stderr, "thothctl: %s\n", output.Sanitize(err.Error(), nil))
return 2
}
secretFiles, err := installation.SecretFiles()
if err != nil {
if isWorkspaceCommand {
return writeWorkspaceError(stderr, errors.New("installation secret declarations could not be read"), nil, 2)
}
fmt.Fprintln(stderr, "thothctl: installation secret declarations could not be read")
return 2
}
secretValues, err := output.SecretValuesFromFiles(secretFiles)
if err != nil {
if isWorkspaceCommand {
return writeWorkspaceError(stderr, errors.New("declared secret file could not be read"), nil, 2)
}
fmt.Fprintln(stderr, "thothctl: declared secret file could not be read")
return 2
}
runner := compose.NewRunner("")
if command == "workspace" {
workspaceCommand, parseErr := workspaceops.ParseWorkspaceCommand(append([]string{"workspace"}, commandArgs...))
if parseErr != nil {
return writeWorkspaceError(stderr, parseErr, secretValues, 2)
}
jsonMode := true
switch c := workspaceCommand.(type) {
case workspaceops.InspectCommand:
jsonMode = c.JSON
case workspaceops.DwhRequest:
jsonMode = c.JSON
case workspaceops.SuggestFksRequest:
jsonMode = c.JSON
case workspaceops.CheckSchemaRequest:
jsonMode = c.JSON
case workspaceops.IndexSchemaRequest:
jsonMode = c.JSON
case workspaceops.EvidenceRequest:
jsonMode = c.JSON
case workspaceops.RunRequest:
jsonMode = c.JSON
}
var finalOutput []byte
result, operationErr := workspaceops.RunWithProjectorAndSecrets(ctx, installation, runner, workspaceCommand, nil, func(result workspaceops.Result) (workspaceops.Result, error) {
return projectWorkspaceResult(result, secretValues)
}, secretValues, func(result workspaceops.Result) error {
var err error
if jsonMode {
finalOutput, err = encodeWorkspaceJSON(result)
} else {
finalOutput, err = encodeWorkspaceHuman(result)
}
if err != nil {
return errors.New("workspace output could not be encoded")
}
if len(finalOutput) > maxPublicStdoutBytes {
return errors.New("workspace result exceeds output limit")
}
// This is the exact once-encoded byte slice that will be published.
// Scan it after encoding so JSON keys and human renderer chrome are
// inside the same no-secret boundary as typed result values.
if containsDeclaredSecretBytes(finalOutput, secretValues) {
return errors.New("workspace output contains a declared secret")
}
return nil
})
if operationErr != nil {
if workspaceUsageError(operationErr) {
return writeWorkspaceError(stderr, operationErr, secretValues, 2)
}
return writeWorkspaceError(stderr, operationErr, secretValues, 1)
}
if _, err := stdout.Write(finalOutput); err != nil {
// The candidate publication precedes stdout. A physical stdout
// failure therefore requires reconciliation before retrying; it is
// not an output-limit rejection.
return writeWorkspaceError(stderr, errors.New("workspace output write failed; reconcile any committed candidate before retrying"), secretValues, 1)
}
if result.Status == "blocked" {
return 3
}
if result.Status == "failed" {
return 1
}
return 0
}
var result compose.Result
switch command {
case "status":
if len(commandArgs) != 0 {
return commandUsageError(stderr, "status does not accept arguments")
}
result, err = runner.Run(ctx, installation.ComposeArgs("ps", "--format", "json"), nil)
case "logs":
logArgs, argumentError := logsArgs(commandArgs)
if argumentError != nil {
return commandUsageError(stderr, argumentError.Error())
}
result, err = runner.Run(ctx, installation.ComposeArgs(logArgs...), nil)
case "start":
if len(commandArgs) != 0 {
return commandUsageError(stderr, "start does not accept arguments")
}
result, err = runner.Run(ctx, installation.ComposeArgs("up", "--detach", "--remove-orphans"), nil)
case "stop":
if len(commandArgs) != 0 {
return commandUsageError(stderr, "stop does not accept arguments")
}
result, err = runner.Run(ctx, installation.ComposeArgs("stop"), nil)
case "update":
if len(commandArgs) != 1 || commandArgs[0] != "--check-only" {
return commandUsageError(stderr, "update currently requires --check-only")
}
result, err = runner.Run(ctx, installation.ComposeArgs("config", "--quiet"), nil)
case "doctor":
if len(commandArgs) != 0 {
return commandUsageError(stderr, "doctor does not accept arguments")
}
return doctor(ctx, installation, runner, secretValues, stdout, stderr)
case "pi":
return piCommand(ctx, installation, runner, commandArgs, secretValues, stdout, stderr)
case "sessions":
if len(commandArgs) != 2 || commandArgs[0] != "migrate" || commandArgs[1] != "--yes" {
return commandUsageError(stderr, "sessions migrate requires --yes")
}
status, operationErr := serverops.MigrateSessions(ctx, installation, runner, true)
if operationErr != nil {
return serverOperationFailure(stderr, operationErr, secretValues)
}
if encodeErr := json.NewEncoder(stdout).Encode(status); encodeErr != nil {
fmt.Fprintln(stderr, "thothctl: migration status could not be written")
return 1
}
return 0
case "remove":
var confirmedIDs []string
if len(commandArgs) > 0 {
if commandArgs[0] != "--yes" || len(commandArgs) < 2 {
return commandUsageError(stderr, "remove requires either no arguments or --yes followed by every displayed container ID")
}
confirmedIDs = commandArgs[1:]
}
removal, operationErr := serverops.Remove(ctx, installation, runner, confirmedIDs)
writeRemovalTargets(stdout, installation.ProjectName(), removal.Targets)
if errors.Is(operationErr, serverops.ErrConfirmationRequired) {
fmt.Fprint(stderr, "thothctl: inspect the exact targets above, then re-run with remove --yes")
for _, target := range removal.Targets {
fmt.Fprintf(stderr, " %s", target.ID)
}
fmt.Fprintln(stderr)
return 2
}
if operationErr != nil {
return serverOperationFailure(stderr, operationErr, secretValues)
}
fmt.Fprintf(stdout, "Removed %d stopped app containers; verified %d preserved paths.\n", len(removal.Targets), removal.Preserved)
return 0
default:
return commandUsageError(stderr, fmt.Sprintf("unknown command %q", command))
}
return writeResult(result, err, secretValues, stdout, stderr)
}
func workspaceUsageError(err error) bool {
if err == nil || errors.Is(err, compose.ErrOutputLimit) {
return false
}
message := err.Error()
// These are host-side grammar/local-file failures. Child envelope/result and
// bounded-stream failures are operational and deliberately remain exit 1.
return strings.Contains(message, "unsafe") || strings.Contains(message, "request exceeds") ||
strings.Contains(message, "SQL input exceeds") || strings.Contains(message, "annotation input exceeds")
}
type boundedWriter struct {
dst io.Writer
maximum int64
written int64
}
func (w *boundedWriter) Write(p []byte) (int, error) {
remaining := w.maximum - w.written
if remaining <= 0 {
return 0, io.ErrShortWrite
}
if int64(len(p)) > remaining {
n, err := w.dst.Write(p[:int(remaining)])
w.written += int64(n)
if err != nil {
return n, err
}
return n, io.ErrShortWrite
}
n, err := w.dst.Write(p)
w.written += int64(n)
return n, err
}
func containsDeclaredSecretBytes(contents []byte, secrets []string) bool {
for _, secret := range secrets {
if secret != "" && bytes.Contains(contents, []byte(secret)) {
return true
}
}
return false
}
// writeWorkspaceError is the sole workspace stderr path. It validates the
// complete rendered line, including its fixed prefix, before writing; if no
// deterministic safe line exists it emits empty stderr rather than risk a
// declared-secret collision.
func writeWorkspaceError(stderr io.Writer, err error, secrets []string, code int) int {
message := "workspace operation failed"
if err != nil {
message = output.Sanitize(err.Error(), secrets)
}
candidates := [][]byte{[]byte("workspace operation failed\n")}
if secrets != nil {
candidates = append([][]byte{[]byte("thothctl: " + message + "\n")}, candidates...)
}
for _, candidate := range candidates {
if !containsDeclaredSecretBytes(candidate, secrets) {
_, _ = stderr.Write(candidate)
return code
}
}
return code
}
func encodeWorkspaceJSON(result workspaceops.Result) ([]byte, error) {
var encoded bytes.Buffer
encoder := json.NewEncoder(&encoded)
// Keep public output compact and avoid HTML-escape amplification of warnings.
encoder.SetEscapeHTML(false)
if err := encoder.Encode(result); err != nil {
return nil, err
}
return encoded.Bytes(), nil
}
func writeWorkspaceJSON(w io.Writer, result workspaceops.Result) error {
encoded, err := encodeWorkspaceJSON(result)
if err != nil || len(encoded) > maxPublicStdoutBytes {
return io.ErrShortWrite
}
_, err = w.Write(encoded)
return err
}
func encodeWorkspaceHuman(result workspaceops.Result) ([]byte, error) {
var encoded bytes.Buffer
if err := renderWorkspaceHuman(&encoded, result); err != nil {
return nil, err
}
return encoded.Bytes(), nil
}
func projectWorkspaceResult(result workspaceops.Result, secretValues []string) (workspaceops.Result, error) {
// Public fields are either closed identities (which must never be rewritten)
// or human-rendered strings (which are redacted and checked for spoofing).
project := func(value string) (string, error) {
value = output.Sanitize(value, secretValues)
for _, r := range value {
if unicode.IsControl(r) || unicode.In(r, unicode.Cf, unicode.Zl, unicode.Zp) {
return "", errors.New("unsafe character in workspace result")
}
}
return value, nil
}
closed := func(value string) (string, error) {
public, err := project(value)
if err != nil || public != value {
return "", errors.New("secret collides with workspace identity")
}
return public, nil
}
var err error
for _, value := range []*string{&result.Status, &result.Code, &result.WorkspaceID, &result.WorkspaceRevision, &result.DescriptorBlob, &result.Operation} {
if *value, err = closed(*value); err != nil {
return workspaceops.Result{}, err
}
}
if result.RunID, err = closed(result.RunID); err != nil {
return workspaceops.Result{}, err
}
if result.ChildRuns != nil {
childRuns := make(map[string]string, len(result.ChildRuns))
for key, value := range result.ChildRuns {
publicKey, keyErr := project(key)
if keyErr != nil {
return workspaceops.Result{}, keyErr
}
publicValue, valueErr := closed(value)
if valueErr != nil {
return workspaceops.Result{}, valueErr
}
if _, exists := childRuns[publicKey]; exists {
return workspaceops.Result{}, errors.New("workspace result key collision")
}
childRuns[publicKey] = publicValue
}
result.ChildRuns = childRuns
}
if result.Counts != nil {
counts := make(map[string]int, len(result.Counts))
for key, value := range result.Counts {
publicKey, keyErr := project(key)
if keyErr != nil {
return workspaceops.Result{}, keyErr
}
if _, exists := counts[publicKey]; exists {
return workspaceops.Result{}, errors.New("workspace result key collision")
}
counts[publicKey] = value
}
result.Counts = counts
}
for i := range result.CompletedStages {
if result.CompletedStages[i], err = project(result.CompletedStages[i]); err != nil {
return workspaceops.Result{}, err
}
}
for i := range result.Warnings {
if result.Warnings[i], err = project(result.Warnings[i]); err != nil {
return workspaceops.Result{}, err
}
}
for i := range result.ArtifactIdentities {
if result.ArtifactIdentities[i].Kind, err = project(result.ArtifactIdentities[i].Kind); err != nil {
return workspaceops.Result{}, err
}
if result.ArtifactIdentities[i].Digest, err = closed(result.ArtifactIdentities[i].Digest); err != nil {
return workspaceops.Result{}, err
}
}
if publicResultContainsSecret(result, secretValues) {
return workspaceops.Result{}, errors.New("workspace result contains a declared secret")
}
return result, nil
}
func publicResultContainsSecret(result workspaceops.Result, secrets []string) bool {
values := []string{result.Status, result.Code, result.WorkspaceID, result.WorkspaceRevision, result.DescriptorBlob, result.Operation, result.RunID}
values = append(values, result.CompletedStages...)
values = append(values, result.Warnings...)
for key, value := range result.ChildRuns {
values = append(values, key, value)
}
for key := range result.Counts {
values = append(values, key)
}
for _, artifact := range result.ArtifactIdentities {
values = append(values, artifact.Kind, artifact.Digest)
}
for _, value := range values {
for _, secret := range secrets {
if secret != "" && strings.Contains(value, secret) {
return true
}
}
}
return false
}
func renderWorkspaceHuman(w io.Writer, result workspaceops.Result) error {
if result.Code == workspaceops.CodeRegistryBootstrapRecoveryConflict {
_, err := fmt.Fprintln(w, "Bootstrap recovery is ambiguous or corrupt; inspect the installation registry jobs.")
return err
}
if _, err := fmt.Fprintf(w, "Workspace %s: %s (%s)\n", result.WorkspaceID, result.Status, result.Code); err != nil {
return err
}
if result.RunID != "" {
if _, err := fmt.Fprintf(w, "Run: %s\n", result.RunID); err != nil {
return err
}
}
if len(result.CompletedStages) > 0 {
_, err := fmt.Fprintf(w, "Completed stages: %s\n", strings.Join(result.CompletedStages, ", "))
return err
}
return nil
}
func writeRemovalTargets(outputWriter io.Writer, project string, targets []serverops.Container) {
fmt.Fprintf(outputWriter, "Removal targets for installation project %s:\n", project)
if len(targets) == 0 {
fmt.Fprintln(outputWriter, " (none)")
return
}
for _, target := range targets {
fmt.Fprintf(outputWriter, " service=%s name=%s id=%s state=%s\n", target.Service, target.Name, target.ID, target.State)
}
}
func serverOperationFailure(stderr io.Writer, err error, secretValues []string) int {
message := output.Sanitize(err.Error(), secretValues)
var operationErr *serverops.OperationError
if errors.As(err, &operationErr) && operationErr.Detail() != "" {
detail := output.SanitizeDetail(operationErr.Detail(), secretValues)
fmt.Fprintf(stderr, "thothctl: %s: %s\n", message, detail)
} else {
fmt.Fprintf(stderr, "thothctl: %s\n", message)
}
if errors.Is(err, serverops.ErrConfirmationRequired) || errors.Is(err, serverops.ErrUnsafeState) {
return 2
}
return 1
}
// installationRunner transforms only Compose invocations into the installation's validated,
// profile-specific argument list. Direct Docker image commands remain host-side and use arguments.
type installationRunner struct {
installation config.Installation
runner compose.Runner
}
func (r installationRunner) SessionInventoryScope() string {
if r.installation.Profile == "local" {
return "mine"
}
return "all"
}
func (r installationRunner) Run(ctx context.Context, args []string, stdin io.Reader) (compose.Result, error) {
if len(args) > 0 && args[0] == "compose" {
return r.runner.Run(ctx, r.installation.ComposeArgs(args[1:]...), stdin)
}
return r.runner.Run(ctx, args, stdin)
}
func piCommand(ctx context.Context, installation config.Installation, runner compose.Runner, args []string, secretValues []string, stdout, stderr io.Writer) int {
if len(args) == 0 {
return commandUsageError(stderr, "pi requires a subcommand")
}
controlled := installationRunner{installation: installation, runner: runner}
switch args[0] {
case "status":
if len(args) != 1 {
return commandUsageError(stderr, "pi status does not accept arguments")
}
version, err := pi.Status(ctx, controlled)
if err != nil {
return piFailure(stderr, err, secretValues)
}
fmt.Fprintf(stdout, "Pi version: %s\n", output.Sanitize(version, secretValues))
return 0
case "doctor":
if len(args) != 1 {
return commandUsageError(stderr, "pi doctor does not accept arguments")
}
if err := pi.Doctor(ctx, controlled); err != nil {
return piFailure(stderr, err, secretValues)
}
fmt.Fprintln(stdout, "Pi preflight checks passed.")
return 0
case "test", "check":
if len(args) != 1 {
return commandUsageError(stderr, "pi test does not accept arguments")
}
if err := pi.Test(ctx, controlled); err != nil {
return piFailure(stderr, err, secretValues)
}
fmt.Fprintln(stdout, "Pi/core smoke checks passed.")
return 0
case "logs":
if len(args) != 1 {
return commandUsageError(stderr, "pi logs does not support --follow; use bounded snapshots")
}
logArgs := []string{"logs", "--tail", "200", "core"}
result, err := controlled.Run(ctx, append([]string{"compose"}, logArgs...), nil)
return writeResult(result, err, secretValues, stdout, stderr)
case "configure":
authFile, authErr := installation.EnvironmentValue("PI_AUTH_FILE")
if authErr != nil || strings.TrimSpace(authFile) == "" {
return commandUsageError(stderr, "PI_AUTH_FILE must name the actual protected host credential file")
}
defaults, err := resolvePiConfigure(ctx, controlled, args[1:], os.Stdin, stdout, stdinIsTTY(os.Stdin))
if err != nil {
return commandUsageError(stderr, err.Error())
}
if err := pi.Configure(ctx, controlled, defaults); err != nil {
return piFailure(stderr, err, secretValues)
}
fmt.Fprintf(stdout, "Pi defaults applied and read back. Provider credentials remain only in the host file %s (mode 0600). Never pass credentials to thothctl.\n", authFile)
return 0
case "update":
request, err := parsePiUpdateArgs(args[1:], installation.UpdateStatePath())
if err != nil {
return commandUsageError(stderr, err.Error())
}
result, err := pi.Update(ctx, controlled, request)
if err != nil {
return piFailure(stderr, err, secretValues)
}
if result.Phase == pi.PhaseNoop {
fmt.Fprintf(stdout, "Pi already runs requested version %s; no container was recreated.\n", request.Version)
return 0
}
fmt.Fprintf(stdout, "Pi update verified. Recovery metadata: %s\n", result.StatePath)
return 0
case "rollback":
if len(args) != 2 || args[1] != "--yes" {
return commandUsageError(stderr, "pi rollback requires --yes")
}
result, err := pi.Rollback(ctx, controlled, installation.UpdateStatePath(), true)
if err != nil {
return piFailure(stderr, err, secretValues)
}
fmt.Fprintf(stdout, "Pi rollback restored the recorded core image. Recovery metadata: %s\n", result.StatePath)
return 0
case "maintenance":
if len(args) == 2 && args[1] == "status" {
status, err := pi.MaintenanceStatus(ctx, controlled)
if err != nil {
return piFailure(stderr, err, secretValues)
}
fmt.Fprintf(stdout, "Pi maintenance active: %t (admissions: %d)\n", status.Active, status.Admissions)
return 0
}
if len(args) == 3 && args[1] == "recover" && args[2] == "--yes" {
statePath := installation.UpdateStatePath()
if err := pi.RecoverMaintenance(ctx, controlled, statePath, true); err != nil {
return piFailure(stderr, err, secretValues)
}
fmt.Fprintln(stdout, "Pi maintenance recovery verified; stale lifecycle files were removed and admissions are open.")
return 0
}
return commandUsageError(stderr, "pi maintenance requires status or recover --yes")
default:
return commandUsageError(stderr, fmt.Sprintf("unknown pi command %q", args[0]))
}
}
func resolvePiConfigure(
ctx context.Context,
runner pi.Runner,
args []string,
input io.Reader,
prompt io.Writer,
isTTY bool,
) (pi.Defaults, error) {
if len(args) > 0 {
return parsePiConfigureArgs(args)
}
if !isTTY {
return pi.Defaults{}, errors.New("non-interactive pi configure requires --provider --model --thinking")
}
options, err := pi.ConfigurationOptions(ctx, runner)
if err != nil {
return pi.Defaults{}, err
}
providers := uniqueProviders(options)
scanner := bufio.NewScanner(input)
provider, err := numberedChoice(scanner, prompt, "provider", providers)
if err != nil {
return pi.Defaults{}, err
}
models := make([]string, 0)
for _, option := range options {
if option.Provider == provider {
models = append(models, option.ID)
}
}
model, err := numberedChoice(scanner, prompt, "model", models)
if err != nil {
return pi.Defaults{}, err
}
thinking, err := numberedChoice(scanner, prompt, "thinking level", []string{"low", "medium", "high"})
if err != nil {
return pi.Defaults{}, err
}
return pi.Defaults{Provider: provider, Model: model, Thinking: thinking}, nil
}
func uniqueProviders(options []pi.ModelOption) []string {
seen := make(map[string]bool)
providers := make([]string, 0)
for _, option := range options {
if !seen[option.Provider] {
seen[option.Provider] = true
providers = append(providers, option.Provider)
}
}
return providers
}
func numberedChoice(scanner *bufio.Scanner, output io.Writer, label string, choices []string) (string, error) {
if len(choices) == 0 {
return "", fmt.Errorf("Pi returned no %s choices", label)
}
fmt.Fprintf(output, "Select %s:\n", label)
for index, choice := range choices {
fmt.Fprintf(output, " %d) %s\n", index+1, choice)
}
for {
fmt.Fprintf(output, "Choice [1-%d]: ", len(choices))
if !scanner.Scan() {
return "", fmt.Errorf("interactive %s selection ended before a choice was entered", label)
}
selected, err := strconv.Atoi(strings.TrimSpace(scanner.Text()))
if err == nil && selected >= 1 && selected <= len(choices) {
return choices[selected-1], nil
}
fmt.Fprintln(output, "Enter one of the listed numbers.")
}
}
func stdinIsTTY(input *os.File) bool {
info, err := input.Stat()
return err == nil && info.Mode()&os.ModeCharDevice != 0
}
func parsePiConfigureArgs(args []string) (pi.Defaults, error) {
var value pi.Defaults
for len(args) > 0 {
if len(args) < 2 {
return pi.Defaults{}, errors.New("configure options require values")
}
key, v := args[0], args[1]
args = args[2:]
switch key {
case "--provider":
value.Provider = v
case "--model":
value.Model = v
case "--thinking":
value.Thinking = v
default:
return pi.Defaults{}, fmt.Errorf("unknown pi configure option %q", key)
}
}
if value.Provider == "" || value.Model == "" || value.Thinking == "" {
return pi.Defaults{}, errors.New("pi configure requires --provider --model --thinking; THT_LLM_URL stays Compose-managed")
}
return value, nil
}
func parsePiUpdateArgs(args []string, statePath string) (pi.Request, error) {
request := pi.Request{StatePath: statePath}
for len(args) > 0 {
switch args[0] {
case "--version":
if len(args) < 2 || request.Version != "" {
return pi.Request{}, errors.New("pi update requires one --version <pinned-version>")
}
request.Version, args = args[1], args[2:]
case "--source":
if len(args) < 2 {
return pi.Request{}, errors.New("--source requires build or pull")
}
request.Source, args = pi.Source(args[1]), args[2:]
case "--image":
if len(args) < 2 || request.Image != "" {
return pi.Request{}, errors.New("--image requires one digest-pinned image reference")
}
request.Image, args = args[1], args[2:]
case "--yes":
if request.Confirm {
return pi.Request{}, errors.New("--yes may be supplied once")
}
request.Confirm, args = true, args[1:]
case "--drain":
if request.Drain {
return pi.Request{}, errors.New("--drain may be supplied once")
}
request.Drain, args = true, args[1:]
default:
return pi.Request{}, fmt.Errorf("unknown pi update option %q", args[0])
}
}
if request.Version == "" {
return pi.Request{}, errors.New("pi update requires --version <pinned-version>")
}
if request.Source == "" {
return pi.Request{}, errors.New("pi update requires explicit --source build or pull")
}
if request.Source != pi.BuildSource && request.Source != pi.PullSource {
return pi.Request{}, errors.New("--source requires build or pull")
}
if request.Source == pi.PullSource && request.Image == "" {
return pi.Request{}, errors.New("--source pull requires --image <digest-reference>")
}
if request.Source == pi.BuildSource && request.Image != "" {
return pi.Request{}, errors.New("--image is valid only with --source pull")
}
return request, nil
}
func piFailure(stderr io.Writer, err error, secretValues []string) int {
code := 1
if errors.Is(err, pi.ErrConfirmationRequired) || errors.Is(err, pi.ErrInvalidRequest) || errors.Is(err, pi.ErrActiveSessions) || errors.Is(err, pi.ErrInterruptedUpdate) {
code = 2
}
var childExit interface{ ExitCode() int }
if errors.As(err, &childExit) && childExit.ExitCode() != 0 {
code = childExit.ExitCode()
}
fmt.Fprintf(stderr, "thothctl: %s\n", output.Sanitize(err.Error(), secretValues))
return code
}
func parseArgs(args []string) (string, string, []string, error) {
if len(args) < 3 || args[0] != "--installation" {
return "", "", nil, errors.New("--installation <absolute-path> is required before the command")
}
if !filepath.IsAbs(args[1]) {
return "", "", nil, errors.New("--installation must be an absolute path")
}
return args[1], args[2], args[3:], nil
}
func logsArgs(args []string) ([]string, error) {
if len(args) == 0 {
return []string{"logs", "--tail", "200"}, nil
}
return nil, errors.New("logs does not accept arguments; use bounded snapshots")
}
func commandUsageError(stderr io.Writer, message string) int {
fmt.Fprintf(stderr, "thothctl: %s\n", message)
return 2
}
func writeResult(result compose.Result, err error, secretValues []string, stdout, stderr io.Writer) int {
if result.Stdout != "" {
fmt.Fprint(stdout, output.Sanitize(result.Stdout, secretValues))
}
if result.Stderr != "" {
fmt.Fprint(stderr, output.Sanitize(result.Stderr, secretValues))
}
if err == nil {
return 0
}
if errors.Is(err, exec.ErrNotFound) {
fmt.Fprintln(stderr, "thothctl: Docker is not installed or is not on PATH")
}
if result.ExitCode != 0 {
return result.ExitCode
}
return 1
}
func doctor(ctx context.Context, installation config.Installation, runner compose.Runner, secretValues []string, stdout, stderr io.Writer) int {
checks := [][]string{
{"version", "--format", "{{.Client.Version}}"},
{"compose", "version", "--short"},
installation.ComposeArgs("config", "--quiet"),
installation.ComposeArgs("config", "--format", "json"),
installation.ComposeArgs("ps", "--format", "json"),
}
var renderedConfig, status string
for index, args := range checks {
result, err := runner.Run(ctx, args, nil)
if err != nil {
return writeResult(result, err, secretValues, stdout, stderr)
}
if index == 3 {
renderedConfig = result.Stdout
}
if index == 4 {
status = result.Stdout
}
}
if err := requireLF(installation.ProjectDirectory); err != nil {
fmt.Fprintf(stderr, "thothctl: %s\n", err)
return 1
}
if err := requireVolumes(renderedConfig); err != nil {
fmt.Fprintf(stderr, "thothctl: %s\n", err)
return 1
}
if err := requireHealthyServices(status); err != nil {
fmt.Fprintf(stderr, "thothctl: %s\n", err)
return 1
}
fmt.Fprintln(stdout, "Doctor checks passed.")
return 0
}
func requireLF(root string) error {
return filepath.WalkDir(root, func(path string, entry os.DirEntry, walkErr error) error {
if walkErr != nil {
return walkErr
}
if entry.IsDir() || !requiresLF(entry.Name()) {
return nil
}
contents, err := os.ReadFile(path)
if err != nil {
return err
}
if strings.Contains(string(contents), "\r\n") {
return fmt.Errorf("CRLF line endings found in %s", filepath.Base(path))
}
return nil
})
}
func requiresLF(name string) bool {
if name == "Dockerfile" || strings.HasPrefix(name, "Dockerfile.") || strings.HasSuffix(name, ".Dockerfile") {
return true
}
for _, suffix := range []string{".sh", ".yml", ".yaml"} {
if strings.HasSuffix(name, suffix) {
return true
}
}
return false
}
func requireVolumes(renderedConfig string) error {
var document struct {
Volumes map[string]json.RawMessage `json:"volumes"`
}
if err := json.Unmarshal([]byte(renderedConfig), &document); err != nil {
return fmt.Errorf("Compose returned invalid rendered configuration")
}
if len(document.Volumes) == 0 {
return errors.New("rendered Compose configuration declares no volumes")
}
return nil
}
func requireHealthyServices(status string) error {
var services []struct {
Service string `json:"Service"`
State string `json:"State"`
Health string `json:"Health"`
}
if err := json.Unmarshal([]byte(status), &services); err != nil {
return errors.New("Compose returned invalid service status")
}
seen := map[string]bool{}
for _, service := range services {
if service.Service != "core" && service.Service != "frontend" {
continue
}
if service.State != "running" || service.Health != "healthy" {
return fmt.Errorf("%s is not healthy", service.Service)
}
seen[service.Service] = true
}
for _, service := range []string{"core", "frontend"} {
if !seen[service] {
return fmt.Errorf("%s service is not running", service)
}
}
return nil
}