Files

267 lines
6.9 KiB
Go

package fileedit
import (
"bytes"
"context"
"crypto/sha256"
"crypto/subtle"
"encoding/json"
"errors"
"fmt"
"io"
"net/http"
"strings"
"sync"
"time"
)
// Browse runs inside the read-only Job and reuses Execute for every command.
// Refuse write commands even if the other end violates the protocol.
func Browse(ctx context.Context, worldRoot, url, token string) error {
if token == "" {
return fmt.Errorf("fileedit: browser token is required")
}
client := &http.Client{Timeout: browserIdle + 30*time.Second,
CheckRedirect: func(*http.Request, []*http.Request) error { return http.ErrUseLastResponse }}
answer := BrowseAnswer{}
for {
body, err := json.Marshal(answer)
if err != nil {
return err
}
req, err := http.NewRequestWithContext(ctx, http.MethodPost, url, bytes.NewReader(body))
if err != nil {
return err
}
req.Header.Set("Authorization", "Bearer "+token)
req.Header.Set("Content-Type", "application/json")
resp, err := client.Do(req)
if err != nil {
return err
}
if resp.StatusCode == http.StatusNoContent {
resp.Body.Close()
return nil
}
if resp.StatusCode != http.StatusOK {
resp.Body.Close()
return fmt.Errorf("fileedit: browser exchange returned %s", resp.Status)
}
var command BrowseCommand
err = json.NewDecoder(io.LimitReader(resp.Body, 64<<10)).Decode(&command)
resp.Body.Close()
if err != nil {
return err
}
if command.Op != OpList && command.Op != OpRead {
return fmt.Errorf("fileedit: browser refused operation %q", command.Op)
}
res, err := Execute(worldRoot, Request{Op: command.Op, Path: command.Path})
answer = BrowseAnswer{ID: command.ID}
if err != nil {
answer.Error = err.Error()
continue
}
answer.Result, err = json.Marshal(res)
if err != nil {
return err
}
}
}
const (
BrowserTokenEnv = "FELIS_FILE_BROWSER_TOKEN"
BrowserRoute = "/api/v1/internal/file-browser/"
browserIdle = 45 * time.Second
browserLifetime = 4 * time.Minute
maxBrowsers = 4
)
var errBrowserFull = errors.New("fileedit: browser capacity reached")
// Browser reuses a short-lived, read-only Job per world. Jobs pull commands from
// the existing internal API face: no inbound Pod port or new RBAC is needed.
// Every browser request is still authorized by the normal file API handlers.
type Browser struct {
BaseURL string
mu sync.Mutex
sessions map[string]*browseSession
}
type BrowseCommand struct {
ID string `json:"id"`
Op string `json:"op"`
Path string `json:"path"`
}
type BrowseAnswer struct {
ID string `json:"id"`
Result json.RawMessage `json:"result,omitempty"`
Error string `json:"error,omitempty"`
}
type browseSession struct {
id, key string
tokenHash [sha256.Size]byte
commands chan BrowseCommand
answers chan BrowseAnswer
serial chan struct{}
done chan struct{}
once sync.Once
timer *time.Timer
mu sync.Mutex
connected bool
pending string
}
func (b *Browser) close(s *browseSession) {
b.mu.Lock()
if b.sessions[s.key] == s {
delete(b.sessions, s.key)
}
b.mu.Unlock()
s.once.Do(func() {
if s.timer != nil {
s.timer.Stop()
}
close(s.done)
})
}
func (b *Browser) session(ctx context.Context, p JobParams, start func(context.Context, JobParams) error) (*browseSession, error) {
key := strings.Join([]string{p.Namespace, p.Server, p.WorldPVC, p.NodeName, p.Image}, "\x00")
b.mu.Lock()
defer b.mu.Unlock()
if s := b.sessions[key]; s != nil {
return s, nil
}
if len(b.sessions) >= maxBrowsers {
return nil, errBrowserFull
}
token, err := randomHex(32)
if err != nil {
return nil, err
}
s := &browseSession{id: p.OpID, key: key, tokenHash: sha256.Sum256([]byte(token)),
commands: make(chan BrowseCommand, 1), answers: make(chan BrowseAnswer, 1), serial: make(chan struct{}, 1), done: make(chan struct{})}
if b.sessions == nil {
b.sessions = make(map[string]*browseSession)
}
b.sessions[key] = s
p.BrowserURL, p.BrowserToken = strings.TrimRight(b.BaseURL, "/")+BrowserRoute+s.id, token
p.Deadline = browserLifetime
if err := start(ctx, p); err != nil {
delete(b.sessions, key)
return nil, err
}
if err := ctx.Err(); err != nil {
delete(b.sessions, key)
return nil, err
}
s.timer = time.AfterFunc(browserLifetime, func() { b.close(s) })
return s, nil
}
func (b *Browser) Run(ctx context.Context, p JobParams, start func(context.Context, JobParams) error) ([]byte, error) {
s, err := b.session(ctx, p, start)
if err != nil {
return nil, err
}
select {
case s.serial <- struct{}{}:
case <-ctx.Done():
return nil, ctx.Err()
case <-s.done:
return nil, fmt.Errorf("fileedit: browser closed")
}
defer func() { <-s.serial }()
s.mu.Lock()
s.pending = p.OpID
s.mu.Unlock()
select {
case s.commands <- BrowseCommand{ID: p.OpID, Op: p.Op, Path: p.Path}:
case <-ctx.Done():
b.close(s)
return nil, ctx.Err()
case <-s.done:
return nil, fmt.Errorf("fileedit: browser closed")
}
select {
case answer := <-s.answers:
if answer.Error != "" {
return nil, fmt.Errorf("fileedit: browser read: %s", answer.Error)
}
return answer.Result, nil
case <-ctx.Done():
b.close(s)
return nil, ctx.Err()
case <-s.done:
return nil, fmt.Errorf("fileedit: browser closed before returning a result")
}
}
// ServeHTTP accepts one worker's result and long-polls its next command. The
// random token opens only this world/session, and lives only in the Job's env.
func (b *Browser) ServeHTTP(w http.ResponseWriter, r *http.Request) {
id := strings.TrimPrefix(r.URL.Path, BrowserRoute)
b.mu.Lock()
var s *browseSession
for _, candidate := range b.sessions {
if candidate.id == id {
s = candidate
break
}
}
b.mu.Unlock()
token, bearer := strings.CutPrefix(r.Header.Get("Authorization"), "Bearer ")
sum := sha256.Sum256([]byte(token))
if s == nil || !bearer || subtle.ConstantTimeCompare(sum[:], s.tokenHash[:]) != 1 {
http.Error(w, "unknown browser", http.StatusNotFound)
return
}
var answer BrowseAnswer
decoder := json.NewDecoder(http.MaxBytesReader(w, r.Body, maxLogBytes))
decoder.DisallowUnknownFields()
if err := decoder.Decode(&answer); err != nil {
http.Error(w, "invalid result", http.StatusBadRequest)
return
}
s.mu.Lock()
valid := !s.connected && answer.ID == "" || s.connected && answer.ID != "" && answer.ID == s.pending
if valid {
s.connected = true
if answer.ID != "" {
s.pending = ""
}
}
s.mu.Unlock()
if !valid {
http.Error(w, "unexpected result", http.StatusConflict)
return
}
if answer.ID != "" {
select {
case s.answers <- answer:
case <-s.done:
w.WriteHeader(http.StatusGone)
return
}
}
timer := time.NewTimer(browserIdle)
defer timer.Stop()
select {
case command := <-s.commands:
w.Header().Set("Content-Type", "application/json")
if err := json.NewEncoder(w).Encode(command); err != nil {
b.close(s)
}
case <-timer.C:
b.close(s)
w.WriteHeader(http.StatusNoContent)
case <-s.done:
w.WriteHeader(http.StatusNoContent)
case <-r.Context().Done():
b.close(s)
}
}