801 lines
25 KiB
Go
801 lines
25 KiB
Go
// Package workspaceops implements the closed host-side workspace preprocessing contract.
|
|
package workspaceops
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"crypto/sha256"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"os"
|
|
"path/filepath"
|
|
"regexp"
|
|
"sort"
|
|
"strings"
|
|
|
|
"github.com/aritmolab/thothii/tools/thothctl/internal/compose"
|
|
"github.com/aritmolab/thothii/tools/thothctl/internal/config"
|
|
"github.com/aritmolab/thothii/tools/thothctl/internal/safeio"
|
|
)
|
|
|
|
var (
|
|
workspacePattern = regexp.MustCompile(`^[a-z][a-z0-9-]{2,62}$`)
|
|
runIDPattern = regexp.MustCompile(`^[0-9a-f]{32}$`)
|
|
reviewedCandidatesDigest = regexp.MustCompile(`^sha256:[0-9a-f]{64}$`)
|
|
)
|
|
|
|
const (
|
|
maxFromSQLFiles = 32
|
|
maxAssumptions = 256
|
|
)
|
|
|
|
type Runner interface {
|
|
Run(context.Context, []string, io.Reader) (compose.Result, error)
|
|
}
|
|
|
|
type Request interface {
|
|
workspaceRequest()
|
|
workspaceID() string
|
|
JSONMode() bool
|
|
operatorCommand() string
|
|
stdinEnvelope() (requestEnvelope, error)
|
|
}
|
|
|
|
type baseRequest struct {
|
|
Workspace string
|
|
JSON bool
|
|
}
|
|
|
|
func (b baseRequest) workspaceID() string { return b.Workspace }
|
|
func (b baseRequest) JSONMode() bool { return b.JSON }
|
|
|
|
type InspectRequest struct{ baseRequest }
|
|
|
|
type DwhRequest struct {
|
|
baseRequest
|
|
Resume string
|
|
}
|
|
|
|
type SuggestFksRequest struct {
|
|
baseRequest
|
|
FromSQL []string
|
|
Assume []string
|
|
Output string
|
|
}
|
|
|
|
type CheckSchemaRequest struct {
|
|
baseRequest
|
|
Annotations string
|
|
ReviewedCandidates string
|
|
}
|
|
|
|
type IndexSchemaRequest struct{ baseRequest }
|
|
|
|
type EvidenceRequest struct {
|
|
baseRequest
|
|
DryRun bool
|
|
Resume string
|
|
}
|
|
|
|
type RunRequest struct {
|
|
baseRequest
|
|
Resume string
|
|
}
|
|
|
|
func (InspectRequest) workspaceRequest() {}
|
|
func (DwhRequest) workspaceRequest() {}
|
|
func (SuggestFksRequest) workspaceRequest() {}
|
|
func (CheckSchemaRequest) workspaceRequest() {}
|
|
func (IndexSchemaRequest) workspaceRequest() {}
|
|
func (EvidenceRequest) workspaceRequest() {}
|
|
func (RunRequest) workspaceRequest() {}
|
|
|
|
func (InspectRequest) operatorCommand() string { return "inspect" }
|
|
func (DwhRequest) operatorCommand() string { return "preprocess-dwh" }
|
|
func (SuggestFksRequest) operatorCommand() string { return "schema-suggest-fks" }
|
|
func (CheckSchemaRequest) operatorCommand() string {
|
|
return "schema-check"
|
|
}
|
|
func (IndexSchemaRequest) operatorCommand() string { return "index-schema" }
|
|
func (EvidenceRequest) operatorCommand() string { return "preprocess-evidence" }
|
|
func (RunRequest) operatorCommand() string { return "preprocess-run" }
|
|
|
|
func (r InspectRequest) stdinEnvelope() (requestEnvelope, error) {
|
|
return requestEnvelope{SchemaVersion: 1, WorkspaceID: r.Workspace}, nil
|
|
}
|
|
|
|
func (r DwhRequest) stdinEnvelope() (requestEnvelope, error) {
|
|
return requestEnvelope{SchemaVersion: 1, WorkspaceID: r.Workspace, Resume: r.Resume}, nil
|
|
}
|
|
|
|
func (r SuggestFksRequest) stdinEnvelope() (requestEnvelope, error) {
|
|
envelope := requestEnvelope{SchemaVersion: 1, WorkspaceID: r.Workspace, Assume: append([]string(nil), r.Assume...)}
|
|
totalBytes := 0
|
|
for _, path := range r.FromSQL {
|
|
contents, err := safeio.ReadCanonicalUTF8(path, 1<<20)
|
|
if err != nil {
|
|
return requestEnvelope{}, errors.New("SQL input could not be read safely")
|
|
}
|
|
totalBytes += len(contents)
|
|
if totalBytes > 16<<20 {
|
|
return requestEnvelope{}, errors.New("SQL input total exceeds 16 MiB")
|
|
}
|
|
envelope.SQLFiles = append(envelope.SQLFiles, inputFile{Name: filepath.Base(path), SQL: contents})
|
|
}
|
|
return envelope, nil
|
|
}
|
|
|
|
func (r CheckSchemaRequest) stdinEnvelope() (requestEnvelope, error) {
|
|
annotations, err := safeio.ReadCanonicalUTF8(r.Annotations, 16<<20)
|
|
if err != nil {
|
|
return requestEnvelope{}, errors.New("annotation file could not be read safely")
|
|
}
|
|
return requestEnvelope{
|
|
SchemaVersion: 1,
|
|
WorkspaceID: r.Workspace,
|
|
Annotations: annotations,
|
|
ReviewedCandidates: r.ReviewedCandidates,
|
|
}, nil
|
|
}
|
|
|
|
func (r IndexSchemaRequest) stdinEnvelope() (requestEnvelope, error) {
|
|
return requestEnvelope{SchemaVersion: 1, WorkspaceID: r.Workspace}, nil
|
|
}
|
|
|
|
func (r EvidenceRequest) stdinEnvelope() (requestEnvelope, error) {
|
|
return requestEnvelope{SchemaVersion: 1, WorkspaceID: r.Workspace, Resume: r.Resume, DryRun: r.DryRun}, nil
|
|
}
|
|
|
|
func (r RunRequest) stdinEnvelope() (requestEnvelope, error) {
|
|
return requestEnvelope{SchemaVersion: 1, WorkspaceID: r.Workspace, Resume: r.Resume}, nil
|
|
}
|
|
|
|
type requestEnvelope struct {
|
|
SchemaVersion int `json:"schemaVersion"`
|
|
WorkspaceID string `json:"workspaceId"`
|
|
Resume string `json:"resumeRunId,omitempty"`
|
|
DryRun bool `json:"dryRun,omitempty"`
|
|
Assume []string `json:"assume,omitempty"`
|
|
SQLFiles []inputFile `json:"fromSql,omitempty"`
|
|
Annotations string `json:"annotationsYaml,omitempty"`
|
|
ReviewedCandidates string `json:"reviewedCandidatesDigest,omitempty"`
|
|
}
|
|
|
|
type inputFile struct {
|
|
Name string `json:"name"`
|
|
SQL string `json:"sql"`
|
|
}
|
|
|
|
type Result struct {
|
|
SchemaVersion int `json:"schemaVersion"`
|
|
Status string `json:"status"`
|
|
Code string `json:"code"`
|
|
WorkspaceID string `json:"workspaceId"`
|
|
WorkspaceRevision string `json:"workspaceRevision"`
|
|
DescriptorBlob string `json:"descriptorBlob"`
|
|
Operation string `json:"operation"`
|
|
RunID string `json:"runId,omitempty"`
|
|
ChildRuns map[string]string `json:"childRuns,omitempty"`
|
|
CompletedStages []string `json:"completedStages"`
|
|
Counts map[string]int `json:"counts,omitempty"`
|
|
ArtifactIdentities []ArtifactIdentity `json:"artifactIdentities,omitempty"`
|
|
SuggestedFksYAML string `json:"suggestedFksYaml,omitempty"`
|
|
EffectiveConfigIdentity string `json:"effectiveConfigIdentity,omitempty"`
|
|
ConfigFingerprint string `json:"configFingerprint,omitempty"`
|
|
InputFingerprint string `json:"inputFingerprint,omitempty"`
|
|
Warnings []string `json:"warnings,omitempty"`
|
|
}
|
|
|
|
type ArtifactIdentity struct {
|
|
Kind string `json:"kind"`
|
|
Digest string `json:"digest"`
|
|
}
|
|
|
|
type operationResponse struct {
|
|
Result
|
|
SuggestedFksYAML string `json:"suggestedFksYaml,omitempty"`
|
|
}
|
|
|
|
type Stage string
|
|
|
|
type ExitClass string
|
|
|
|
const (
|
|
StageRenderedConfig Stage = "rendered-config"
|
|
StageImageInspect Stage = "image-inspect"
|
|
StageComposeRun Stage = "compose-run"
|
|
|
|
ExitClassNonzero ExitClass = "nonzero-exit"
|
|
ExitClassUnavailable ExitClass = "unavailable"
|
|
ExitClassTimeout ExitClass = "timeout"
|
|
ExitClassInvocation ExitClass = "invocation-failure"
|
|
)
|
|
|
|
type OperationError struct {
|
|
stage Stage
|
|
class ExitClass
|
|
detail string
|
|
}
|
|
|
|
func (e *OperationError) Error() string {
|
|
detail := strings.TrimSpace(e.detail)
|
|
if len(detail) > 0 {
|
|
// Bounded, sanitized operator/container detail so operators can diagnose failures
|
|
// without leaking secrets; the full renderer sanitizes further before output.
|
|
if len(detail) > 2048 {
|
|
detail = detail[:2048]
|
|
}
|
|
return fmt.Sprintf("stage=%s class=%s: %s", e.stage, e.class, detail)
|
|
}
|
|
return fmt.Sprintf("stage=%s class=%s", e.stage, e.class)
|
|
}
|
|
|
|
func (e *OperationError) Stage() Stage { return e.stage }
|
|
func (e *OperationError) Class() ExitClass { return e.class }
|
|
func (e *OperationError) Detail() string { return e.detail }
|
|
|
|
func Parse(args []string) (Request, error) {
|
|
if len(args) == 0 {
|
|
return nil, errors.New("workspace requires a subcommand")
|
|
}
|
|
switch args[0] {
|
|
case "inspect":
|
|
parsed, err := parseInspect(args[1:])
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return parsed, nil
|
|
case "preprocess":
|
|
return parsePreprocess(args[1:])
|
|
case "schema":
|
|
return parseSchema(args[1:])
|
|
case "index-schema":
|
|
parsed, err := parseIndexSchema(args[1:])
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return parsed, nil
|
|
default:
|
|
return nil, fmt.Errorf("unknown workspace command %q", args[0])
|
|
}
|
|
}
|
|
|
|
func Execute(ctx context.Context, installation config.Installation, runner Runner, request Request) (Result, error) {
|
|
envelope, err := request.stdinEnvelope()
|
|
if err != nil {
|
|
return Result{}, err
|
|
}
|
|
rendered, err := runDocker(ctx, runner, StageRenderedConfig, installation.ComposeArgs("config", "--format", "json"))
|
|
if err != nil {
|
|
return Result{}, err
|
|
}
|
|
imageReference, err := selectedCoreImage(rendered.Stdout)
|
|
if err != nil {
|
|
return Result{}, err
|
|
}
|
|
imageID, err := immutableImageID(ctx, runner, imageReference)
|
|
if err != nil {
|
|
return Result{}, err
|
|
}
|
|
override, cleanup, err := maintenanceOverride(installation, imageID)
|
|
if err != nil {
|
|
return Result{}, err
|
|
}
|
|
defer cleanup()
|
|
stdin, err := encodeEnvelope(envelope)
|
|
if err != nil {
|
|
return Result{}, err
|
|
}
|
|
args, err := installation.ComposeArgsWithFinalOverride(
|
|
override,
|
|
"run", "--rm", "--no-deps", "--no-TTY", "--name", ownedContainerName(installation, request), "workspace-maintenance", request.operatorCommand(),
|
|
)
|
|
if err != nil {
|
|
return Result{}, err
|
|
}
|
|
result, err := runDocker(ctx, runner, StageComposeRun, args, bytes.NewReader(stdin))
|
|
// The operator emits one authoritative JSON result on stdout and encodes its status in the
|
|
// exit code (0 success, 3 operator checkpoint/block, 1 operational failure). A nonzero
|
|
// exit is therefore still a valid machine result whenever stdout parses; only a missing or
|
|
// malformed payload becomes an error.
|
|
response, parseErr := parseResponse(result.Stdout)
|
|
if parseErr != nil {
|
|
if err != nil {
|
|
return Result{}, err
|
|
}
|
|
return Result{}, parseErr
|
|
}
|
|
if suggest, ok := request.(SuggestFksRequest); ok && suggest.Output != "" {
|
|
if response.SuggestedFksYAML == "" {
|
|
return Result{}, errors.New("workspace maintenance did not return the requested FK artifact")
|
|
}
|
|
if digest := suggestedArtifactDigest(response); digest != "" {
|
|
sum := sha256.Sum256([]byte(response.SuggestedFksYAML))
|
|
if digest != "sha256:"+fmt.Sprintf("%x", sum[:]) {
|
|
return Result{}, errors.New("workspace maintenance returned an FK artifact with a mismatched digest")
|
|
}
|
|
}
|
|
if err := safeio.WriteCanonicalNewFile(suggest.Output, []byte(response.SuggestedFksYAML), 0o600); err != nil {
|
|
return Result{}, errors.New("workspace FK output file could not be created safely")
|
|
}
|
|
}
|
|
response.Result.SuggestedFksYAML = response.SuggestedFksYAML
|
|
return response.Result, nil
|
|
}
|
|
|
|
func encodeEnvelope(envelope requestEnvelope) ([]byte, error) {
|
|
encoded, err := json.Marshal(envelope)
|
|
if err != nil {
|
|
return nil, errors.New("workspace request could not be encoded")
|
|
}
|
|
if len(encoded) > 1<<20 {
|
|
return nil, errors.New("workspace request exceeds the bounded stdin contract")
|
|
}
|
|
return append(encoded, '\n'), nil
|
|
}
|
|
|
|
func parseResponse(document string) (operationResponse, error) {
|
|
decoder := json.NewDecoder(strings.NewReader(document))
|
|
decoder.DisallowUnknownFields()
|
|
var response operationResponse
|
|
if err := decoder.Decode(&response); err != nil {
|
|
return operationResponse{}, errors.New("workspace maintenance returned invalid JSON")
|
|
}
|
|
var extra any
|
|
if err := decoder.Decode(&extra); !errors.Is(err, io.EOF) {
|
|
return operationResponse{}, errors.New("workspace maintenance returned trailing output")
|
|
}
|
|
if err := validateResult(response.Result); err != nil {
|
|
return operationResponse{}, err
|
|
}
|
|
return response, nil
|
|
}
|
|
|
|
func validateResult(result Result) error {
|
|
if result.SchemaVersion != 1 {
|
|
return errors.New("workspace maintenance returned an unsupported schema version")
|
|
}
|
|
if !workspacePattern.MatchString(result.WorkspaceID) {
|
|
return errors.New("workspace maintenance returned an invalid workspace identity")
|
|
}
|
|
if result.Status != "failed" {
|
|
if len(result.WorkspaceRevision) != 40 || !isLowerHex(result.WorkspaceRevision) {
|
|
return errors.New("workspace maintenance returned an invalid workspace revision")
|
|
}
|
|
if !strings.HasPrefix(result.DescriptorBlob, "sha256:") || len(result.DescriptorBlob) != len("sha256:")+64 || !isLowerHex(strings.TrimPrefix(result.DescriptorBlob, "sha256:")) {
|
|
return errors.New("workspace maintenance returned an invalid descriptor digest")
|
|
}
|
|
}
|
|
validStatuses := map[string]struct{}{"succeeded": {}, "unchanged": {}, "dry_run": {}, "blocked": {}, "failed": {}}
|
|
if _, ok := validStatuses[result.Status]; !ok {
|
|
return errors.New("workspace maintenance returned an invalid status")
|
|
}
|
|
if strings.TrimSpace(result.Code) == "" || strings.TrimSpace(result.Operation) == "" || result.CompletedStages == nil {
|
|
return errors.New("workspace maintenance omitted required fields")
|
|
}
|
|
if result.Status != "failed" {
|
|
for _, digest := range result.ArtifactIdentities {
|
|
if strings.TrimSpace(digest.Kind) == "" || !strings.HasPrefix(digest.Digest, "sha256:") {
|
|
return errors.New("workspace maintenance returned an invalid artifact identity")
|
|
}
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func isLowerHex(value string) bool {
|
|
for _, r := range value {
|
|
if !(r >= '0' && r <= '9' || r >= 'a' && r <= 'f') {
|
|
return false
|
|
}
|
|
}
|
|
return value != ""
|
|
}
|
|
|
|
func parseInspect(args []string) (InspectRequest, error) {
|
|
base, err := parseBaseFlags(args, false)
|
|
if err != nil {
|
|
return InspectRequest{}, err
|
|
}
|
|
return InspectRequest{baseRequest: base}, nil
|
|
}
|
|
|
|
func parsePreprocess(args []string) (Request, error) {
|
|
if len(args) == 0 {
|
|
return nil, errors.New("workspace preprocess requires dwh, evidence, or run")
|
|
}
|
|
switch args[0] {
|
|
case "dwh":
|
|
base, resume, dryRun, err := parseResumeFlags(args[1:], false)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if dryRun {
|
|
return nil, errors.New("workspace preprocess dwh does not accept --dry-run")
|
|
}
|
|
return DwhRequest{baseRequest: base, Resume: resume}, nil
|
|
case "evidence":
|
|
base, resume, dryRun, err := parseResumeFlags(args[1:], true)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return EvidenceRequest{baseRequest: base, Resume: resume, DryRun: dryRun}, nil
|
|
case "run":
|
|
base, resume, dryRun, err := parseResumeFlags(args[1:], false)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if dryRun {
|
|
return nil, errors.New("workspace preprocess run does not accept --dry-run")
|
|
}
|
|
return RunRequest{baseRequest: base, Resume: resume}, nil
|
|
default:
|
|
return nil, fmt.Errorf("unknown workspace preprocess command %q", args[0])
|
|
}
|
|
}
|
|
|
|
func parseSchema(args []string) (Request, error) {
|
|
if len(args) == 0 {
|
|
return nil, errors.New("workspace schema requires suggest-fks or check")
|
|
}
|
|
switch args[0] {
|
|
case "suggest-fks":
|
|
return parseSuggestFks(args[1:])
|
|
case "check":
|
|
return parseSchemaCheck(args[1:])
|
|
default:
|
|
return nil, fmt.Errorf("unknown workspace schema command %q", args[0])
|
|
}
|
|
}
|
|
|
|
func parseIndexSchema(args []string) (IndexSchemaRequest, error) {
|
|
base, err := parseBaseFlags(args, false)
|
|
if err != nil {
|
|
return IndexSchemaRequest{}, err
|
|
}
|
|
return IndexSchemaRequest{baseRequest: base}, nil
|
|
}
|
|
|
|
func parseResumeFlags(args []string, allowDryRun bool) (baseRequest, string, bool, error) {
|
|
var resume string
|
|
var dryRun bool
|
|
base, seen, err := parseSharedFlags(args, map[string]func(string) error{
|
|
"--resume": func(value string) error {
|
|
if resume != "" {
|
|
return errors.New("--resume may be supplied once")
|
|
}
|
|
if !runIDPattern.MatchString(value) {
|
|
return errors.New("--resume must be 32 lowercase hex characters")
|
|
}
|
|
resume = value
|
|
return nil
|
|
},
|
|
}, map[string]func() error{
|
|
"--dry-run": func() error {
|
|
if !allowDryRun {
|
|
return errors.New("--dry-run is not accepted here")
|
|
}
|
|
if dryRun {
|
|
return errors.New("--dry-run may be supplied once")
|
|
}
|
|
dryRun = true
|
|
return nil
|
|
},
|
|
})
|
|
if err != nil {
|
|
return baseRequest{}, "", false, err
|
|
}
|
|
if !seen.workspace {
|
|
return baseRequest{}, "", false, errors.New("--workspace is required")
|
|
}
|
|
return base, resume, dryRun, nil
|
|
}
|
|
|
|
func parseSuggestFks(args []string) (SuggestFksRequest, error) {
|
|
request := SuggestFksRequest{}
|
|
base, seen, err := parseSharedFlags(args, map[string]func(string) error{
|
|
"--from-sql": func(value string) error {
|
|
if len(request.FromSQL) >= maxFromSQLFiles {
|
|
return fmt.Errorf("--from-sql may be supplied at most %d times", maxFromSQLFiles)
|
|
}
|
|
request.FromSQL = append(request.FromSQL, value)
|
|
return nil
|
|
},
|
|
"--assume": func(value string) error {
|
|
if len(request.Assume) >= maxAssumptions {
|
|
return fmt.Errorf("--assume may be supplied at most %d times", maxAssumptions)
|
|
}
|
|
if len(value) > 256 || !strings.Contains(value, "=") {
|
|
return errors.New("--assume values must be column=table entries up to 256 bytes")
|
|
}
|
|
left, right, _ := strings.Cut(value, "=")
|
|
if strings.TrimSpace(left) == "" || strings.TrimSpace(right) == "" {
|
|
return errors.New("--assume values must be column=table entries up to 256 bytes")
|
|
}
|
|
request.Assume = append(request.Assume, value)
|
|
return nil
|
|
},
|
|
"--output": func(value string) error {
|
|
if request.Output != "" {
|
|
return errors.New("--output may be supplied once")
|
|
}
|
|
request.Output = value
|
|
return nil
|
|
},
|
|
}, nil)
|
|
if err != nil {
|
|
return SuggestFksRequest{}, err
|
|
}
|
|
if !seen.workspace {
|
|
return SuggestFksRequest{}, errors.New("--workspace is required")
|
|
}
|
|
request.baseRequest = base
|
|
return request, nil
|
|
}
|
|
|
|
func parseSchemaCheck(args []string) (CheckSchemaRequest, error) {
|
|
request := CheckSchemaRequest{}
|
|
base, seen, err := parseSharedFlags(args, map[string]func(string) error{
|
|
"--annotations": func(value string) error {
|
|
if request.Annotations != "" {
|
|
return errors.New("--annotations may be supplied once")
|
|
}
|
|
request.Annotations = value
|
|
return nil
|
|
},
|
|
"--reviewed-candidates": func(value string) error {
|
|
if request.ReviewedCandidates != "" {
|
|
return errors.New("--reviewed-candidates may be supplied once")
|
|
}
|
|
if !reviewedCandidatesDigest.MatchString(value) {
|
|
return errors.New("--reviewed-candidates must be sha256:<64 lowercase hex>")
|
|
}
|
|
request.ReviewedCandidates = value
|
|
return nil
|
|
},
|
|
}, nil)
|
|
if err != nil {
|
|
return CheckSchemaRequest{}, err
|
|
}
|
|
if !seen.workspace {
|
|
return CheckSchemaRequest{}, errors.New("--workspace is required")
|
|
}
|
|
if (request.Annotations == "") != (request.ReviewedCandidates == "") {
|
|
return CheckSchemaRequest{}, errors.New("--annotations and --reviewed-candidates must be supplied together")
|
|
}
|
|
request.baseRequest = base
|
|
return request, nil
|
|
}
|
|
|
|
func parseBaseFlags(args []string, allowDryRun bool) (baseRequest, error) {
|
|
base, seen, err := parseSharedFlags(args, nil, nil)
|
|
if err != nil {
|
|
return baseRequest{}, err
|
|
}
|
|
if !seen.workspace {
|
|
return baseRequest{}, errors.New("--workspace is required")
|
|
}
|
|
return base, nil
|
|
}
|
|
|
|
type seenFlags struct {
|
|
workspace bool
|
|
json bool
|
|
}
|
|
|
|
func parseSharedFlags(args []string, valueHandlers map[string]func(string) error, boolHandlers map[string]func() error) (baseRequest, seenFlags, error) {
|
|
request := baseRequest{}
|
|
seen := seenFlags{}
|
|
valueHandlers = cloneValueHandlers(valueHandlers)
|
|
boolHandlers = cloneBoolHandlers(boolHandlers)
|
|
for len(args) > 0 {
|
|
flag := args[0]
|
|
if flag == "--" {
|
|
return baseRequest{}, seenFlags{}, errors.New("passthrough separators are not supported")
|
|
}
|
|
switch flag {
|
|
case "--workspace":
|
|
if len(args) < 2 {
|
|
return baseRequest{}, seenFlags{}, errors.New("--workspace requires a value")
|
|
}
|
|
if seen.workspace {
|
|
return baseRequest{}, seenFlags{}, errors.New("--workspace must be supplied exactly once")
|
|
}
|
|
workspace := args[1]
|
|
if !workspacePattern.MatchString(workspace) {
|
|
return baseRequest{}, seenFlags{}, errors.New("--workspace must match [a-z][a-z0-9-]{2,62}")
|
|
}
|
|
request.Workspace, seen.workspace, args = workspace, true, args[2:]
|
|
case "--json":
|
|
if seen.json {
|
|
return baseRequest{}, seenFlags{}, errors.New("--json may be supplied once")
|
|
}
|
|
request.JSON, seen.json, args = true, true, args[1:]
|
|
default:
|
|
if handler, ok := boolHandlers[flag]; ok {
|
|
if err := handler(); err != nil {
|
|
return baseRequest{}, seenFlags{}, err
|
|
}
|
|
args = args[1:]
|
|
continue
|
|
}
|
|
handler, ok := valueHandlers[flag]
|
|
if !ok {
|
|
return baseRequest{}, seenFlags{}, fmt.Errorf("unknown workspace option %q", flag)
|
|
}
|
|
if len(args) < 2 {
|
|
return baseRequest{}, seenFlags{}, fmt.Errorf("%s requires a value", flag)
|
|
}
|
|
if err := handler(args[1]); err != nil {
|
|
return baseRequest{}, seenFlags{}, err
|
|
}
|
|
args = args[2:]
|
|
}
|
|
}
|
|
return request, seen, nil
|
|
}
|
|
|
|
func cloneValueHandlers(source map[string]func(string) error) map[string]func(string) error {
|
|
if len(source) == 0 {
|
|
return map[string]func(string) error{}
|
|
}
|
|
clone := make(map[string]func(string) error, len(source))
|
|
for key, handler := range source {
|
|
clone[key] = handler
|
|
}
|
|
return clone
|
|
}
|
|
|
|
func cloneBoolHandlers(source map[string]func() error) map[string]func() error {
|
|
if len(source) == 0 {
|
|
return map[string]func() error{}
|
|
}
|
|
clone := make(map[string]func() error, len(source))
|
|
for key, handler := range source {
|
|
clone[key] = handler
|
|
}
|
|
return clone
|
|
}
|
|
|
|
func selectedCoreImage(document string) (string, error) {
|
|
var rendered struct {
|
|
Services map[string]struct {
|
|
Image string `json:"image"`
|
|
} `json:"services"`
|
|
}
|
|
if err := json.Unmarshal([]byte(document), &rendered); err != nil {
|
|
return "", errors.New("rendered Compose configuration is invalid")
|
|
}
|
|
core, exists := rendered.Services["core"]
|
|
if !exists || strings.TrimSpace(core.Image) == "" {
|
|
return "", errors.New("selected core image is unavailable")
|
|
}
|
|
return core.Image, nil
|
|
}
|
|
|
|
func immutableImageID(ctx context.Context, runner Runner, reference string) (string, error) {
|
|
result, err := runDocker(ctx, runner, StageImageInspect, []string{"image", "inspect", "--format", "{{.Id}}", reference})
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
id := strings.TrimSpace(result.Stdout)
|
|
if !strings.HasPrefix(id, "sha256:") || len(id) != len("sha256:")+64 || !isLowerHex(strings.TrimPrefix(id, "sha256:")) {
|
|
return "", errors.New("selected core image did not resolve to an immutable sha256 image id")
|
|
}
|
|
return id, nil
|
|
}
|
|
|
|
func maintenanceOverride(installation config.Installation, imageID string) (string, func(), error) {
|
|
control := installation.ControlDirectory()
|
|
if err := os.MkdirAll(control, 0o700); err != nil {
|
|
return "", func() {}, errors.New("workspace maintenance control directory could not be created")
|
|
}
|
|
info, err := os.Lstat(control)
|
|
if err != nil || !info.IsDir() || info.Mode()&os.ModeSymlink != 0 {
|
|
return "", func() {}, errors.New("workspace maintenance control directory is unsafe")
|
|
}
|
|
directory, err := os.MkdirTemp(control, "workspace-maintenance-")
|
|
if err != nil {
|
|
return "", func() {}, errors.New("workspace maintenance override directory could not be created")
|
|
}
|
|
cleanup := func() {
|
|
_ = os.Remove(filepath.Join(directory, "override.yaml"))
|
|
_ = os.Remove(directory)
|
|
}
|
|
path := filepath.Join(directory, "override.yaml")
|
|
contents := []string{
|
|
"services:",
|
|
" core:",
|
|
" image: " + strconvQuote(imageID),
|
|
" pull_policy: never",
|
|
" workspace-maintenance:",
|
|
" image: " + strconvQuote(imageID),
|
|
" pull_policy: never",
|
|
"",
|
|
}
|
|
if err := os.WriteFile(path, []byte(strings.Join(contents, "\n")), 0o600); err != nil {
|
|
cleanup()
|
|
return "", func() {}, errors.New("workspace maintenance override could not be written")
|
|
}
|
|
return path, cleanup, nil
|
|
}
|
|
|
|
func strconvQuote(value string) string {
|
|
encoded, _ := json.Marshal(value)
|
|
return string(encoded)
|
|
}
|
|
|
|
func ownedContainerName(installation config.Installation, request Request) string {
|
|
parts := []string{installation.ProjectName(), request.workspaceID(), request.operatorCommand()}
|
|
for index, value := range parts {
|
|
parts[index] = strings.NewReplacer("/", "-", ":", "-", "@", "-", "_", "-").Replace(value)
|
|
}
|
|
return strings.Join(parts, "-")
|
|
}
|
|
|
|
func isExpectedOperatorExit(result compose.Result, err error) bool {
|
|
if err == nil || result.ExitCode != 3 {
|
|
return false
|
|
}
|
|
var operationErr *OperationError
|
|
return errors.As(err, &operationErr) && operationErr.class == ExitClassNonzero
|
|
}
|
|
|
|
func runDocker(ctx context.Context, runner Runner, stage Stage, args []string, stdin ...io.Reader) (compose.Result, error) {
|
|
var input io.Reader
|
|
if len(stdin) > 0 {
|
|
input = stdin[0]
|
|
}
|
|
result, err := runner.Run(ctx, args, input)
|
|
if err != nil {
|
|
class := ExitClassInvocation
|
|
switch {
|
|
case errors.Is(ctx.Err(), context.DeadlineExceeded):
|
|
class = ExitClassTimeout
|
|
case result.ExitCode == 127:
|
|
class = ExitClassUnavailable
|
|
case result.ExitCode != 0:
|
|
class = ExitClassNonzero
|
|
}
|
|
detail := result.Stderr
|
|
if strings.TrimSpace(detail) == "" {
|
|
detail = err.Error()
|
|
}
|
|
return result, &OperationError{stage: stage, class: class, detail: detail}
|
|
}
|
|
return result, nil
|
|
}
|
|
|
|
func suggestedArtifactDigest(response operationResponse) string {
|
|
for _, artifact := range response.ArtifactIdentities {
|
|
if strings.HasPrefix(artifact.Digest, "sha256:") {
|
|
return artifact.Digest
|
|
}
|
|
}
|
|
return ""
|
|
}
|
|
|
|
func Human(result Result) string {
|
|
lines := []string{
|
|
fmt.Sprintf("workspace: %s", result.WorkspaceID),
|
|
fmt.Sprintf("operation: %s", result.Operation),
|
|
fmt.Sprintf("status: %s", result.Status),
|
|
fmt.Sprintf("code: %s", result.Code),
|
|
fmt.Sprintf("revision: %s", result.WorkspaceRevision),
|
|
}
|
|
if result.RunID != "" {
|
|
lines = append(lines, fmt.Sprintf("run: %s", result.RunID))
|
|
}
|
|
if len(result.CompletedStages) > 0 {
|
|
stages := append([]string(nil), result.CompletedStages...)
|
|
sort.Strings(stages)
|
|
lines = append(lines, fmt.Sprintf("completed: %s", strings.Join(stages, ", ")))
|
|
}
|
|
for _, warning := range result.Warnings {
|
|
lines = append(lines, fmt.Sprintf("warning: %s", warning))
|
|
}
|
|
return strings.Join(lines, "\n") + "\n"
|
|
}
|