Files
Codex 82e2c91f42
Publish documentation / publish (push) Successful in 1m27s
feat: implement memory and evidence administration with guided repairs
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.
2026-09-10 10:31:34 +02:00

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"
}