Publish documentation / publish (push) Successful in 1m27s
Add PostgreSQL-backed memory, editable evidence with source review and activation, and human-approved archive repairs across the harness, API, and UI. Include migrations, deployment support, regression coverage, and validation documentation. Refresh permissions from validated session roles so existing administrator logins can access newly deployed archive management features.
531 lines
17 KiB
Go
531 lines
17 KiB
Go
// Package workspaceops implements the closed host-side workspace preprocessing contract.
|
|
package workspaceops
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"os"
|
|
"path/filepath"
|
|
"regexp"
|
|
"sort"
|
|
"strings"
|
|
|
|
"github.com/aritmolab/thothii/tools/tht/internal/compose"
|
|
"github.com/aritmolab/thothii/tools/tht/internal/config"
|
|
)
|
|
|
|
var workspacePattern = regexp.MustCompile(`^[a-z][a-z0-9-]{2,62}$`)
|
|
|
|
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 RunRequest struct{ baseRequest }
|
|
type ClearRequest struct{ baseRequest }
|
|
type EvidenceRequest struct {
|
|
baseRequest
|
|
Action string
|
|
SourceID string
|
|
Revision string
|
|
Decision string
|
|
}
|
|
|
|
func (EvidenceRequest) workspaceRequest() {}
|
|
func (r EvidenceRequest) operatorCommand() string {
|
|
if r.Action == "refresh" {
|
|
return "evidence-refresh"
|
|
}
|
|
if r.Action == "decide" {
|
|
return "evidence-decide"
|
|
}
|
|
return "evidence-consolidate"
|
|
}
|
|
func (r EvidenceRequest) stdinEnvelope() (requestEnvelope, error) {
|
|
return requestEnvelope{SchemaVersion: 1, WorkspaceID: r.Workspace, SourceID: r.SourceID, Revision: r.Revision, Decision: r.Decision}, nil
|
|
}
|
|
|
|
func (InspectRequest) workspaceRequest() {}
|
|
func (RunRequest) workspaceRequest() {}
|
|
func (ClearRequest) workspaceRequest() {}
|
|
|
|
func (InspectRequest) operatorCommand() string { return "inspect" }
|
|
func (RunRequest) operatorCommand() string { return "preprocess-run" }
|
|
func (ClearRequest) operatorCommand() string { return "preprocess-clear" }
|
|
|
|
func (r InspectRequest) stdinEnvelope() (requestEnvelope, error) {
|
|
return requestEnvelope{SchemaVersion: 1, WorkspaceID: r.Workspace}, nil
|
|
}
|
|
|
|
func (r RunRequest) stdinEnvelope() (requestEnvelope, error) {
|
|
return requestEnvelope{SchemaVersion: 1, WorkspaceID: r.Workspace}, nil
|
|
}
|
|
|
|
func (r ClearRequest) stdinEnvelope() (requestEnvelope, error) {
|
|
return requestEnvelope{SchemaVersion: 1, WorkspaceID: r.Workspace}, nil
|
|
}
|
|
|
|
type requestEnvelope struct {
|
|
SchemaVersion int `json:"schemaVersion"`
|
|
WorkspaceID string `json:"workspaceId"`
|
|
SourceID string `json:"sourceId,omitempty"`
|
|
Revision string `json:"revision,omitempty"`
|
|
Decision string `json:"decision,omitempty"`
|
|
}
|
|
|
|
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"`
|
|
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
|
|
}
|
|
|
|
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 "evidence":
|
|
if len(args) < 2 || (args[1] != "consolidate" && args[1] != "refresh" && args[1] != "decide") {
|
|
return nil, errors.New("workspace evidence requires consolidate, refresh or decide")
|
|
}
|
|
request := EvidenceRequest{Action: args[1]}
|
|
baseArgs := []string{}
|
|
seen := map[string]bool{}
|
|
for i := 2; i < len(args); i++ {
|
|
flag := args[i]
|
|
if flag == "--source-id" || flag == "--revision" || flag == "--decision" {
|
|
if request.Action != "decide" || seen[flag] || i+1 >= len(args) {
|
|
return nil, errors.New("invalid Evidence decision flags")
|
|
}
|
|
seen[flag] = true
|
|
i++
|
|
switch flag {
|
|
case "--source-id":
|
|
request.SourceID = args[i]
|
|
case "--revision":
|
|
request.Revision = args[i]
|
|
case "--decision":
|
|
request.Decision = args[i]
|
|
}
|
|
} else {
|
|
baseArgs = append(baseArgs, flag)
|
|
}
|
|
}
|
|
if request.Action == "decide" && (!regexp.MustCompile(`^[a-f0-9]{64}$`).MatchString(request.SourceID) || !regexp.MustCompile(`^[a-f0-9]{64}$`).MatchString(request.Revision) || (request.Decision != "keep" && request.Decision != "replace")) {
|
|
return nil, errors.New("decide requires source-id, revision and decision keep or replace")
|
|
}
|
|
base, err := parseBaseFlags(baseArgs)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
request.baseRequest = base
|
|
return request, nil
|
|
case "inspect":
|
|
parsed, err := parseInspect(args[1:])
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return parsed, nil
|
|
case "preprocess":
|
|
return parsePreprocess(args[1:])
|
|
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
|
|
}
|
|
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)
|
|
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 run or clear")
|
|
}
|
|
switch args[0] {
|
|
case "run":
|
|
base, err := parseBaseFlags(args[1:])
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return RunRequest{baseRequest: base}, nil
|
|
case "clear":
|
|
base, err := parseBaseFlags(args[1:])
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return ClearRequest{baseRequest: base}, nil
|
|
default:
|
|
return nil, fmt.Errorf("unknown workspace preprocess command %q", args[0])
|
|
}
|
|
}
|
|
|
|
func parseBaseFlags(args []string) (baseRequest, error) {
|
|
request := baseRequest{}
|
|
workspaceSeen := false
|
|
jsonSeen := false
|
|
for len(args) > 0 {
|
|
flag := args[0]
|
|
if flag == "--" {
|
|
return baseRequest{}, errors.New("passthrough separators are not supported")
|
|
}
|
|
switch flag {
|
|
case "--workspace":
|
|
if len(args) < 2 {
|
|
return baseRequest{}, errors.New("--workspace requires a value")
|
|
}
|
|
if workspaceSeen {
|
|
return baseRequest{}, errors.New("--workspace must be supplied exactly once")
|
|
}
|
|
workspace := args[1]
|
|
if !workspacePattern.MatchString(workspace) {
|
|
return baseRequest{}, errors.New("--workspace must match [a-z][a-z0-9-]{2,62}")
|
|
}
|
|
request.Workspace, workspaceSeen, args = workspace, true, args[2:]
|
|
case "--json":
|
|
if jsonSeen {
|
|
return baseRequest{}, errors.New("--json may be supplied once")
|
|
}
|
|
request.JSON, jsonSeen, args = true, true, args[1:]
|
|
default:
|
|
return baseRequest{}, fmt.Errorf("unknown workspace option %q", flag)
|
|
}
|
|
}
|
|
if !workspaceSeen {
|
|
return baseRequest{}, errors.New("--workspace is required")
|
|
}
|
|
return request, nil
|
|
}
|
|
|
|
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 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"
|
|
}
|