From 467529fae2639c471ae67fdf441bf7f4064ca530 Mon Sep 17 00:00:00 2001 From: Amp Date: Wed, 5 Aug 2026 11:24:53 +0000 Subject: [PATCH 1/6] Add service container commands Amp-Thread-ID: https://ampcode.com/threads/T-019fd17c-2d3f-74ce-ac88-8726a8936cdb Co-authored-by: Arjun Komath --- agent/internal/agent/handlers.go | 17 + agent/internal/agent/workqueue.go | 32 +- agent/internal/container/runtime.go | 84 +++++ agent/internal/container/runtime_test.go | 44 +++ agent/internal/http/client.go | 12 +- .../services/[serviceId]/commands/page.tsx | 310 ++++++++++++++++++ web/app/api/inngest/route.ts | 2 + web/app/api/projects/[id]/services/route.ts | 6 +- web/app/api/services/[id]/commands/route.ts | 177 ++++++++++ .../service/service-layout-client.tsx | 4 +- web/db/schema.ts | 39 +++ web/db/types.ts | 4 +- web/lib/inngest/functions/crons.ts | 11 + web/lib/inngest/functions/index.ts | 1 + web/lib/navigation.ts | 6 + web/lib/scheduler.ts | 44 ++- web/lib/service-command-retention.ts | 18 + web/lib/work-queue.ts | 55 +++- web/tests/inngest-route.test.ts | 1 + web/tests/navigation.test.ts | 3 +- web/tests/service-command-retention.test.ts | 52 +++ web/tests/service-commands-route.test.ts | 181 ++++++++++ web/tests/work-queue.test.ts | 95 +++++- 23 files changed, 1174 insertions(+), 24 deletions(-) create mode 100644 web/app/(dashboard)/dashboard/projects/[slug]/[env]/services/[serviceId]/commands/page.tsx create mode 100644 web/app/api/services/[id]/commands/route.ts create mode 100644 web/lib/service-command-retention.ts create mode 100644 web/tests/service-command-retention.test.ts create mode 100644 web/tests/service-commands-route.test.ts diff --git a/agent/internal/agent/handlers.go b/agent/internal/agent/handlers.go index 0b8aeb29..3686144d 100644 --- a/agent/internal/agent/handlers.go +++ b/agent/internal/agent/handlers.go @@ -10,6 +10,7 @@ import ( "os/exec" "path/filepath" "time" + "unicode/utf8" "techulus/cloud-agent/internal/build" "techulus/cloud-agent/internal/container" @@ -19,6 +20,22 @@ import ( "techulus/cloud-agent/internal/registryauth" ) +func (a *Agent) ProcessCommand(item agenthttp.WorkQueueItem) (container.CommandResult, error) { + var payload struct { + CommandRunID string `json:"commandRunId"` + DeploymentID string `json:"deploymentId"` + ContainerID string `json:"containerId"` + Command string `json:"command"` + } + if err := json.Unmarshal([]byte(item.Payload), &payload); err != nil { + return container.CommandResult{}, fmt.Errorf("failed to parse command payload: %w", err) + } + if payload.CommandRunID != item.ID || payload.DeploymentID == "" || payload.ContainerID == "" || payload.Command == "" || utf8.RuneCountInString(payload.Command) > 4096 { + return container.CommandResult{}, fmt.Errorf("invalid command payload") + } + return container.ExecCommand(payload.ContainerID, payload.Command) +} + func (a *Agent) ProcessRestart(item agenthttp.WorkQueueItem) error { var payload struct { DeploymentID string `json:"deploymentId"` diff --git a/agent/internal/agent/workqueue.go b/agent/internal/agent/workqueue.go index 5791a21f..2f6411c5 100644 --- a/agent/internal/agent/workqueue.go +++ b/agent/internal/agent/workqueue.go @@ -7,6 +7,7 @@ import ( "os" "time" + "techulus/cloud-agent/internal/container" agenthttp "techulus/cloud-agent/internal/http" ) @@ -93,7 +94,19 @@ func (a *Agent) processLeasedWorkItem(item agenthttp.WorkQueueItem) { status := "completed" errorMsg := "" restartAfterReport := false - if err := a.ProcessWorkItem(item); err != nil { + var commandResult *container.CommandResult + var processErr error + if item.Type == "command" { + result, err := a.ProcessCommand(item) + processErr = err + if err == nil { + commandResult = &result + } + } else { + processErr = a.ProcessWorkItem(item) + } + if processErr != nil { + err := processErr if errors.Is(err, errAgentUpgradeRestartNeeded) { restartAfterReport = true } else { @@ -109,12 +122,25 @@ func (a *Agent) processLeasedWorkItem(item agenthttp.WorkQueueItem) { if !restartAfterReport && a.activeWorkItem != nil && a.activeWorkItem.ID == item.ID && a.activeWorkItem.Attempt == item.Attempt { a.activeWorkItem = nil } - a.pendingWorkResults = append(a.pendingWorkResults, agenthttp.CompletedWorkItem{ + completed := agenthttp.CompletedWorkItem{ ID: item.ID, Attempt: item.Attempt, Status: status, Error: errorMsg, - }) + } + if commandResult != nil { + completed.Output = commandResult.Output + completed.ExitCode = &commandResult.ExitCode + completed.OutputTruncated = commandResult.Truncated + if commandResult.TimedOut { + completed.Status = "failed" + completed.Error = "command timed out after 60 seconds" + completed.TimedOut = true + } else if commandResult.ExitCode != 0 { + completed.Status = "failed" + } + } + a.pendingWorkResults = append(a.pendingWorkResults, completed) a.workMutex.Unlock() a.RequestStatusReport("work item " + status) diff --git a/agent/internal/container/runtime.go b/agent/internal/container/runtime.go index ee4d08fb..e1212fb0 100644 --- a/agent/internal/container/runtime.go +++ b/agent/internal/container/runtime.go @@ -1,6 +1,7 @@ package container import ( + "bytes" "context" "encoding/json" "errors" @@ -9,12 +10,95 @@ import ( "os" "os/exec" "strings" + "sync" "time" + "unicode/utf8" "techulus/cloud-agent/internal/retry" "techulus/cloud-agent/internal/wireguard" ) +const CommandOutputLimit = 64 * 1024 + +var commandTimeout = 60 * time.Second + +type CommandResult struct { + Output string + ExitCode int + Truncated bool + TimedOut bool +} + +type limitedBuffer struct { + mutex sync.Mutex + buffer bytes.Buffer + truncated bool +} + +func (b *limitedBuffer) snapshot() (string, bool) { + b.mutex.Lock() + defer b.mutex.Unlock() + return b.buffer.String(), b.truncated +} + +func (b *limitedBuffer) Write(p []byte) (int, error) { + b.mutex.Lock() + defer b.mutex.Unlock() + + n := len(p) + remaining := CommandOutputLimit - b.buffer.Len() + if remaining > 0 { + writeLength := min(remaining, len(p)) + _, _ = b.buffer.Write(p[:writeLength]) + if writeLength < len(p) { + b.truncated = true + } + } else if len(p) > 0 { + b.truncated = true + } + return n, nil +} + +func ExecCommand(containerID, command string) (CommandResult, error) { + running, err := IsContainerRunning(containerID) + if err != nil { + return CommandResult{}, err + } + if !running { + return CommandResult{}, fmt.Errorf("container is not running") + } + ctx, cancel := context.WithTimeout(context.Background(), commandTimeout) + defer cancel() + cmd := exec.CommandContext(ctx, "podman", "exec", containerID, "/bin/sh", "-c", command) + var output limitedBuffer + cmd.Stdout, cmd.Stderr = &output, &output + err = cmd.Run() + outputText, truncated := output.snapshot() + outputText = strings.ToValidUTF8(outputText, "�") + if len(outputText) > CommandOutputLimit { + outputText = outputText[:CommandOutputLimit] + for !utf8.ValidString(outputText) { + outputText = outputText[:len(outputText)-1] + } + truncated = true + } + result := CommandResult{Output: outputText, Truncated: truncated} + if ctx.Err() == context.DeadlineExceeded { + result.TimedOut = true + result.ExitCode = 124 + return result, nil + } + if err == nil { + return result, nil + } + var exitErr *exec.ExitError + if errors.As(err, &exitErr) { + result.ExitCode = exitErr.ExitCode() + return result, nil + } + return result, fmt.Errorf("failed to execute command: %w", err) +} + func ContainerExists(containerID string) (bool, error) { cmd := exec.Command("podman", "inspect", "--format", "json", containerID) output, err := cmd.CombinedOutput() diff --git a/agent/internal/container/runtime_test.go b/agent/internal/container/runtime_test.go index f983aba5..7055c154 100644 --- a/agent/internal/container/runtime_test.go +++ b/agent/internal/container/runtime_test.go @@ -8,8 +8,52 @@ import ( "slices" "strings" "testing" + "time" ) +func TestExecCommand(t *testing.T) { + dir := t.TempDir() + script := `#!/bin/sh +if [ "$1" = inspect ]; then echo '[{"State":{"Running":true}}]'; exit 0; fi +case "$5" in + success) printf 'hello';; + failure) printf 'bad'; exit 7;; + exact) head -c 65536 /dev/zero | tr '\0' x;; + large) head -c 70000 /dev/zero | tr '\0' x;; + timeout) sleep 1;; +esac` + if err := os.WriteFile(filepath.Join(dir, "podman"), []byte(script), 0o755); err != nil { + t.Fatal(err) + } + t.Setenv("PATH", dir+":"+os.Getenv("PATH")) + tests := []struct { + name string + exit int + truncated, timedOut bool + }{ + {"success", 0, false, false}, {"failure", 7, false, false}, {"exact", 0, false, false}, {"large", 0, true, false}, {"timeout", 124, false, true}, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + old := commandTimeout + if tt.timedOut { + commandTimeout = 10 * time.Millisecond + } + defer func() { commandTimeout = old }() + result, err := ExecCommand("container", tt.name) + if err != nil { + t.Fatal(err) + } + if result.ExitCode != tt.exit || result.Truncated != tt.truncated || result.TimedOut != tt.timedOut { + t.Fatalf("unexpected result: %+v", result) + } + if len(result.Output) > CommandOutputLimit { + t.Fatalf("output exceeded limit: %d", len(result.Output)) + } + }) + } +} + func TestBuildPodmanPullArgs(t *testing.T) { tests := []struct { name string diff --git a/agent/internal/http/client.go b/agent/internal/http/client.go index 5cf21ede..b825e720 100644 --- a/agent/internal/http/client.go +++ b/agent/internal/http/client.go @@ -358,10 +358,14 @@ type StatusReport struct { } type CompletedWorkItem struct { - ID string `json:"id"` - Attempt int `json:"attempt"` - Status string `json:"status"` - Error string `json:"error,omitempty"` + ID string `json:"id"` + Attempt int `json:"attempt"` + Status string `json:"status"` + Error string `json:"error,omitempty"` + Output string `json:"output,omitempty"` + ExitCode *int `json:"exitCode,omitempty"` + OutputTruncated bool `json:"outputTruncated,omitempty"` + TimedOut bool `json:"timedOut,omitempty"` } type ActiveWorkItem struct { diff --git a/web/app/(dashboard)/dashboard/projects/[slug]/[env]/services/[serviceId]/commands/page.tsx b/web/app/(dashboard)/dashboard/projects/[slug]/[env]/services/[serviceId]/commands/page.tsx new file mode 100644 index 00000000..52e2a7aa --- /dev/null +++ b/web/app/(dashboard)/dashboard/projects/[slug]/[env]/services/[serviceId]/commands/page.tsx @@ -0,0 +1,310 @@ +"use client"; + +import { Loader2, Terminal, XCircle } from "lucide-react"; +import { useMemo, useState } from "react"; +import useSWRInfinite from "swr/infinite"; +import { useService } from "@/components/service/service-layout-client"; +import { Badge } from "@/components/ui/badge"; +import { Button } from "@/components/ui/button"; +import { Card, CardContent, CardHeader, CardTitle } from "@/components/ui/card"; +import { + Empty, + EmptyDescription, + EmptyMedia, + EmptyTitle, +} from "@/components/ui/empty"; +import { + NativeSelect, + NativeSelectOption, +} from "@/components/ui/native-select"; +import { Textarea } from "@/components/ui/textarea"; +import { formatDateTime, formatRelativeTime } from "@/lib/date"; +import { isObservedReady } from "@/lib/deployment-status"; +import { fetcher } from "@/lib/fetcher"; + +type CommandStatus = + | "pending" + | "running" + | "succeeded" + | "failed" + | "timed_out"; + +type CommandRun = { + id: string; + command: string; + status: CommandStatus; + output: string | null; + exitCode: number | null; + outputTruncated: boolean; + errorMessage: string | null; + actor: { name: string }; + serverName: string; + containerId: string; + createdAt: string; + startedAt: string | null; + completedAt: string | null; +}; + +type CommandHistory = { + commands: CommandRun[]; + nextCursor: string | null; +}; + +const STATUS_LABELS: Record = { + pending: "Queued", + running: "Running", + succeeded: "Succeeded", + failed: "Failed", + timed_out: "Timed out", +}; + +function CommandStatusBadge({ command }: { command: CommandRun }) { + const failed = command.status === "failed" || command.status === "timed_out"; + return ( + + {STATUS_LABELS[command.status]} + {command.exitCode !== null ? ` · exit ${command.exitCode}` : null} + + ); +} + +export default function CommandsPage() { + const { service } = useService(); + const targets = service.deployments.filter( + (deployment) => + deployment.containerId && + deployment.runtimeDesiredState === "running" && + isObservedReady(deployment.observedPhase) && + deployment.server?.status === "online", + ); + const [deploymentId, setDeploymentId] = useState(targets[0]?.id ?? ""); + const selectedDeploymentId = targets.some( + (target) => target.id === deploymentId, + ) + ? deploymentId + : (targets[0]?.id ?? ""); + const [command, setCommand] = useState(""); + const [submitting, setSubmitting] = useState(false); + const [submitError, setSubmitError] = useState(null); + + const { + data, + error: historyError, + isLoading, + isValidating, + mutate, + size, + setSize, + } = useSWRInfinite( + (pageIndex, previousPage) => { + if (previousPage && !previousPage.nextCursor) return null; + const cursor = pageIndex === 0 ? null : previousPage?.nextCursor; + return `/api/services/${service.id}/commands${cursor ? `?cursor=${encodeURIComponent(cursor)}` : ""}`; + }, + fetcher, + { + refreshInterval: (pages) => + pages?.some((page) => + page.commands.some( + (item) => item.status === "pending" || item.status === "running", + ), + ) + ? 2000 + : 0, + revalidateOnFocus: true, + }, + ); + + const history = useMemo(() => { + const byId = new Map(); + for (const page of data ?? []) { + for (const item of page.commands) byId.set(item.id, item); + } + return [...byId.values()]; + }, [data]); + const hasMore = data?.[data.length - 1]?.nextCursor != null; + const isLoadingMore = isValidating && Boolean(data?.[size - 1] === undefined); + + async function runCommand() { + setSubmitting(true); + setSubmitError(null); + try { + const response = await fetch(`/api/services/${service.id}/commands`, { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ + deploymentId: selectedDeploymentId, + command, + }), + }); + if (!response.ok) { + const body = (await response.json()) as { + error?: string; + message?: string; + }; + throw new Error( + body.error ?? body.message ?? "Command could not be queued", + ); + } + setCommand(""); + await mutate(); + } catch (error) { + setSubmitError( + error instanceof Error ? error.message : "Command could not be queued", + ); + } finally { + setSubmitting(false); + } + } + + return ( +
+ + + Run command + + + {targets.length === 0 ? ( +

+ No ready, running container is available on an online server. +

+ ) : ( + setDeploymentId(event.target.value)} + > + {targets.map((target) => ( + + {target.server?.name} · {target.containerId?.slice(0, 12)} + + ))} + + )} +