Files
panel/agent/internal/dispatch/files_chunked.go
T
2026-07-14 23:01:33 -07:00

338 lines
9.7 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package dispatch
// Chunked file upload — solves three problems at once:
//
// 1. Cloudflare's request-body cap (~100 MB free, 200 MB pro) silently
// stalls "single-shot" uploads of large mod / scenario archives
// that go through the panel's public hostname.
// 2. The previous single-shot path buffered the entire body in memory
// on both the controller and the agent, peaking at ~3× file size.
// A 3 GB upload would thrash a smaller box.
// 3. Single-shot uploads block the agent's bidi stream for the
// duration of the transfer; chunks are small and interleaved.
//
// Wire protocol (all on the existing AgentEnvelope):
//
// FsWriteChunkRequest { upload_id, instance_id, path, offset, data,
// total_size, is_final }
// FsWriteChunkResult { bytes_received, total_received, finalized,
// error }
//
// The agent owns one `uploadSession` per upload_id. Each session has an
// `*os.File` open in its scratch dir. Chunks append in order. On
// is_final the agent ships the file to the container via a streaming
// CopyTarToContainer call (no full-file in-memory copy), then unlinks.
// Stale sessions are reaped after uploadIdleTTL.
import (
"archive/tar"
"context"
"fmt"
"io"
"os"
"path"
"path/filepath"
"sync"
"time"
"google.golang.org/protobuf/types/known/timestamppb"
panelv1 "github.com/dbledeez/panel/proto/panel/v1"
)
const uploadIdleTTL = 30 * time.Minute
type uploadSession struct {
mu sync.Mutex
uploadID string
instanceID string
path string
tempPath string
f *os.File
bytes int64
total int64
lastTouch time.Time
}
func (d *Dispatcher) uploadScratchDir() string {
dir := filepath.Join(d.dataRoot, ".panel-uploads")
_ = os.MkdirAll(dir, 0o755)
return dir
}
// handleFsWriteChunk processes one chunk in the streaming-upload protocol.
// First chunk creates the session; subsequent chunks append; is_final
// triggers a streaming copy into the container and cleans up.
func (d *Dispatcher) handleFsWriteChunk(corrID string, req *panelv1.FsWriteChunkRequest) {
if req.UploadId == "" {
d.sendChunkResult(corrID, 0, 0, false, "upload_id is required")
return
}
if req.InstanceId == "" || req.Path == "" {
d.sendChunkResult(corrID, 0, 0, false, "instance_id and path are required")
return
}
rec, err := d.lookupRecord(req.InstanceId)
if err != nil {
d.sendChunkResult(corrID, 0, 0, false, err.Error())
return
}
d.mu.Lock()
sess, ok := d.uploads[req.UploadId]
if !ok {
// First chunk for this upload — create session + temp file.
// Use the upload_id in the filename so a recovered scratch dir
// is self-describing for forensic purposes.
base := filepath.Join(d.uploadScratchDir(), "u-"+sanitizeUploadID(req.UploadId))
f, err := os.OpenFile(base, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, 0o600)
if err != nil {
d.mu.Unlock()
d.sendChunkResult(corrID, 0, 0, false, fmt.Errorf("open temp: %w", err).Error())
return
}
sess = &uploadSession{
uploadID: req.UploadId,
instanceID: req.InstanceId,
path: req.Path,
tempPath: base,
f: f,
total: req.TotalSize,
lastTouch: time.Now(),
}
d.uploads[req.UploadId] = sess
}
d.mu.Unlock()
sess.mu.Lock()
defer sess.mu.Unlock()
// Reject stray chunks for a session that's already closed (final
// chunk processed earlier — operator double-clicked, retry, etc.).
if sess.f == nil {
d.sendChunkResult(corrID, 0, sess.bytes, true, "")
return
}
// Allow operator to upload to a renamed path mid-stream? No — pin
// the destination from the first chunk for safety. Subsequent
// chunks with a different path get rejected.
if sess.path != req.Path {
d.sendChunkResult(corrID, 0, sess.bytes, false, "path differs from initial chunk")
return
}
// Append the chunk. The browser is expected to send chunks in
// order; if offset doesn't match the running tail we err so the
// frontend can fall back / retry rather than silently writing a
// hole.
if req.Offset != sess.bytes {
d.sendChunkResult(corrID, 0, sess.bytes, false,
fmt.Sprintf("chunk offset %d != expected %d (out-of-order chunk)", req.Offset, sess.bytes))
return
}
n, err := sess.f.Write(req.Data)
if err != nil {
d.sendChunkResult(corrID, 0, sess.bytes, false, fmt.Errorf("write chunk: %w", err).Error())
return
}
sess.bytes += int64(n)
sess.lastTouch = time.Now()
if !req.IsFinal {
d.sendChunkResult(corrID, int64(n), sess.bytes, false, "")
return
}
// Final chunk — close the temp file, ship it to the container.
if err := sess.f.Sync(); err != nil {
d.sendChunkResult(corrID, int64(n), sess.bytes, false, fmt.Errorf("fsync: %w", err).Error())
return
}
if err := sess.f.Close(); err != nil {
d.sendChunkResult(corrID, int64(n), sess.bytes, false, fmt.Errorf("close temp: %w", err).Error())
return
}
sess.f = nil
if sess.total > 0 && sess.bytes != sess.total {
d.sendChunkResult(corrID, int64(n), sess.bytes, false,
fmt.Sprintf("size mismatch: received %d bytes, expected %d", sess.bytes, sess.total))
return
}
// Ship to the container (or host fallback) without re-reading the
// whole file into memory — we hand Docker an io.Reader that
// streams from disk, wrapped in a tar header on the fly.
if err := d.deliverUploadedFile(rec, sess); err != nil {
d.sendChunkResult(corrID, int64(n), sess.bytes, false, err.Error())
return
}
// Tear down the session — temp file deleted, map slot freed.
_ = os.Remove(sess.tempPath)
d.mu.Lock()
delete(d.uploads, sess.uploadID)
d.mu.Unlock()
d.sendChunkResult(corrID, int64(n), sess.bytes, true, "")
}
// deliverUploadedFile streams the temp file into the destination,
// either via container CopyTarToContainer (running or helper sidecar)
// or via a plain os.Rename for host-fallback instances. Designed so
// the file's bytes never sit in memory in their entirety.
func (d *Dispatcher) deliverUploadedFile(rec *instanceRecord, sess *uploadSession) error {
if d.useContainerOps(rec) {
abs, err := safeJoinAny(rec.BrowseableRoot, rootPaths(rec), sess.path)
if err != nil {
return err
}
ctx, cancel := context.WithTimeout(context.Background(), fsCopyTimeout)
defer cancel()
targetID, err := d.getFsTargetContainerID(ctx, rec)
if err != nil {
return err
}
// mkdir -p the parent so the destination always exists.
_, _, _, _ = d.runtime.ExecCapture(ctx, targetID, []string{"sh", "-c",
"mkdir -p -- " + shellQuote(path.Dir(abs))})
// Open the temp file for read; build a tar wrapper on the fly
// and feed Docker. CopyTarToContainer streams the body without
// materializing it in memory.
fr, err := os.Open(sess.tempPath)
if err != nil {
return err
}
defer fr.Close()
fi, err := fr.Stat()
if err != nil {
return err
}
pr, pw := io.Pipe()
errCh := make(chan error, 1)
go func() {
tw := tar.NewWriter(pw)
err := tw.WriteHeader(&tar.Header{
Name: path.Base(abs),
Mode: 0o644,
Size: fi.Size(),
})
if err == nil {
_, err = io.Copy(tw, fr)
}
if err == nil {
err = tw.Close()
}
_ = pw.CloseWithError(err)
errCh <- err
}()
if err := d.runtime.CopyTarToContainer(ctx, targetID, path.Dir(abs), pr); err != nil {
if encErr := <-errCh; encErr != nil {
return encErr
}
return fmt.Errorf("copy to container: %w", err)
}
if encErr := <-errCh; encErr != nil {
return encErr
}
return nil
}
// Host fallback — just move the temp file into place.
if rec.DataPath == "" {
return fmt.Errorf("no file storage available")
}
abs, err := safeJoinHost(rec.DataPath, sess.path)
if err != nil {
return err
}
if err := os.MkdirAll(filepath.Dir(abs), 0o755); err != nil {
return err
}
// os.Rename is atomic on the same filesystem; if scratch lives on
// a different fs we fall back to copy + delete.
if err := os.Rename(sess.tempPath, abs); err == nil {
return nil
}
src, err := os.Open(sess.tempPath)
if err != nil {
return err
}
defer src.Close()
dst, err := os.OpenFile(abs, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, 0o644)
if err != nil {
return err
}
defer dst.Close()
if _, err := io.Copy(dst, src); err != nil {
return err
}
return nil
}
// uploadReaper closes + removes any session that hasn't been touched
// in uploadIdleTTL. Runs forever; cheap enough to scan once a minute.
func (d *Dispatcher) uploadReaper() {
t := time.NewTicker(time.Minute)
defer t.Stop()
for range t.C {
now := time.Now()
d.mu.Lock()
stale := []*uploadSession{}
for id, s := range d.uploads {
if now.Sub(s.lastTouch) > uploadIdleTTL {
stale = append(stale, s)
delete(d.uploads, id)
}
}
d.mu.Unlock()
for _, s := range stale {
s.mu.Lock()
if s.f != nil {
_ = s.f.Close()
s.f = nil
}
_ = os.Remove(s.tempPath)
s.mu.Unlock()
d.log.Info("upload session reaped", "upload_id", s.uploadID, "bytes", s.bytes)
}
}
}
// sanitizeUploadID strips characters that could escape the scratch dir.
// upload_id is generated by the browser (UUID-like) — defense in depth.
func sanitizeUploadID(s string) string {
out := make([]byte, 0, len(s))
for i := 0; i < len(s); i++ {
c := s[i]
switch {
case (c >= 'a' && c <= 'z'), (c >= 'A' && c <= 'Z'), (c >= '0' && c <= '9'), c == '-', c == '_':
out = append(out, c)
}
if len(out) > 64 {
break
}
}
if len(out) == 0 {
out = []byte("anon")
}
return string(out)
}
func (d *Dispatcher) sendChunkResult(corrID string, n, total int64, finalized bool, errMsg string) {
d.sendEnv(&panelv1.AgentEnvelope{
CorrelationId: corrID,
SentAt: timestamppb.Now(),
Payload: &panelv1.AgentEnvelope_FsWriteChunkResult{
FsWriteChunkResult: &panelv1.FsWriteChunkResult{
BytesReceived: n,
TotalReceived: total,
Finalized: finalized,
Error: errMsg,
},
},
})
}