diff --git a/.dockerignore b/.dockerignore new file mode 100644 index 0000000..93b5832 --- /dev/null +++ b/.dockerignore @@ -0,0 +1,5 @@ +test_input/ +test_output/ +.git/ +.gitignore +Readme.md \ No newline at end of file diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..95f8f7d --- /dev/null +++ b/.gitignore @@ -0,0 +1,5 @@ +test_input/ +test_output/ +*.mp4 + +/compressor \ No newline at end of file diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..e4cd758 --- /dev/null +++ b/Dockerfile @@ -0,0 +1,38 @@ +# Build stage +FROM golang:1.21-alpine AS builder + +WORKDIR /app + +# Copy go mod +COPY go.mod go.sum ./ +RUN go mod download + +# Copy source +COPY . . + +# Build +RUN CGO_ENABLED=0 GOOS=linux go build -a -installsuffix cgo -o compressor ./cmd/compressor + +# Runtime stage with CUDA support +FROM nvidia/cuda:11.8-runtime-ubuntu20.04 + +# Install ffmpeg and ca-certificates +RUN apt-get update && apt-get install -y ffmpeg ca-certificates && rm -rf /var/lib/apt/lists/* + +# Create app user +RUN groupadd -r appgroup && useradd -r -g appgroup appuser + +# Copy binary +COPY --from=builder /app/compressor /usr/local/bin/compressor + +# Create directories +RUN mkdir -p /input /output && chown -R appuser:appgroup /input /output + +USER appuser + +# Health check +HEALTHCHECK --interval=30s --timeout=10s --start-period=5s --retries=3 \ + CMD curl -f http://localhost:8080/status || exit 1 + +# Default command +CMD ["/usr/local/bin/compressor"] \ No newline at end of file diff --git a/Readme.md b/Readme.md index 9daeafb..c945ca5 100644 --- a/Readme.md +++ b/Readme.md @@ -1 +1,78 @@ -test +Compressor is a lightweight folder watcher that transcodes videos with `ffmpeg`. + +## Features + +- Watches a single level input directory and schedules new files for compression. +- Renames files with a `.processing` suffix to coordinate multiple replicas. +- Runs the configured `ffmpeg` command ( GPU friendly defaults supplied ). +- Writes results into an output directory and optionally deletes sources. +- Exposes a `/status` endpoint that returns HTTP 200 for health checks. + +## Configuration + +Environment variables drive the runtime configuration. Defaults are shown in parentheses. + +| Variable | Description | +| ------------------------------------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | +| `INPUT_DIR` (`/input`) | Directory to scan and watch for new videos. | +| `OUTPUT_DIR` (`/output`) | Directory where encoded files are written. | +| `VIDEO_EXTENSIONS` (`.mp4,.mkv,.mov,.avi,.flv,.wmv,.m4v,.webm,.ts`) | Comma separated list of extensions that should be processed. | +| `FFMPEG_BIN` (`ffmpeg`) | Binary to invoke. | +| `FFMPEG_COMMAND` | Arguments passed to `ffmpeg`. Must include `{{input}}` and `{{output}}` placeholders. Default: `-y -hwaccel cuda -hwaccel_device 0 -i {{input}} -c:v hevc_nvenc -qp 25 -preset p6 -gpu 0 -b_qfactor 1.1 -b_ref_mode middle -bf 3 -g 250 -i_qfactor 0.75 -max_muxing_queue_size 1024 -multipass 1 -rc vbr -rc-lookahead 20 -temporal-aq 1 -tune hq -c:a aac -af volume=2.0 {{output}}` | +| `FFMPEG_COMMAND_CPU` | CPU fallback arguments if GPU not detected. Default: `-y -i {{input}} -c:v libx264 -preset slow -crf 22 -c:a aac {{output}}` | +| `OUTPUT_EXTENSION` (`.mp4`) | Extension applied to the output file name. | +| `DELETE_SOURCE` (`false`) | When `true`, removes the processed input file instead of restoring it. | +| `PROCESSING_SUFFIX` (`.processing`) | Suffix appended while a file is in flight. | +| `MAX_CONCURRENT` (`1`) | Number of concurrent transcodes. Consider GPU capacity when raising. | +| `QUEUE_SIZE` (`128`) | Work queue buffer length. | +| `FILE_STABILITY_DURATION` (`3s`) | How long a file size must remain unchanged before processing. | +| `RESCAN_INTERVAL` (`30s`) | Periodic full directory rescan interval. | +| `PORT` (`8080`) | Port for the HTTP `/status` endpoint. | + +Placeholders are shell escaped before the command line is parsed, so paths containing spaces are handled safely. + +## Running Locally + +```bash +go build -o compressor ./cmd/compressor +INPUT_DIR=$(pwd)/test_input OUTPUT_DIR=$(pwd)/test_output ./compressor +``` + +Place a video file in `test_input/` and watch it get processed to `test_output/`. The test directories are gitignored. + +## Container Notes + +- Mount the hot folder to `/input` and the destination to `/output`. +- Expose the health endpoint through your orchestrator: `http://:8080/status`. +- Provide GPU-capable `ffmpeg` binaries in the image (for example via CUDA base images). + +## Docker + +Build and run locally with GPU: + +```bash +docker build -t compressor . +docker run --gpus all -v $(pwd)/test_input:/input -v $(pwd)/test_output:/output -p 8080:8080 compressor +``` + +## Docker Compose + +```bash +docker-compose up --build +``` + +## Kubernetes + +Apply the YAMLs: + +```bash +kubectl apply -f k8s/ +``` + +This creates PVCs for input/output volumes. Mount your persistent volumes accordingly. The deployment requests 1 GPU. + +## Health Endpoint + +- `GET /status` → `200 OK` with body `ok` + +The watcher logs every successful encode with source and destination paths. diff --git a/build.ps1 b/build.ps1 new file mode 100644 index 0000000..200f8f7 --- /dev/null +++ b/build.ps1 @@ -0,0 +1,19 @@ +# Build and push multi-arch Docker images +param( + [string]$Registry = "meisterlala/compressor" +) + +# Enable experimental features for buildx +docker buildx create --use --name multiarch 2>$null || docker buildx use multiarch + +# Build and push latest tag for linux/amd64 and linux/arm64 +$fullTag = "$Registry`:latest" +Write-Host "Building and pushing $fullTag for linux/amd64 and linux/arm64..." + +docker buildx build --platform linux/amd64, linux/arm64 ` + --tag $fullTag ` + --push . + +Write-Host "Build and push complete!" + +Write-Host "Build and push complete!" \ No newline at end of file diff --git a/cmd/compressor/config.go b/cmd/compressor/config.go new file mode 100644 index 0000000..22b4b22 --- /dev/null +++ b/cmd/compressor/config.go @@ -0,0 +1,163 @@ +package main + +import ( + "errors" + "fmt" + "log" + "os" + "os/exec" + "path/filepath" + "strings" + "time" +) + +const ( + defaultInputDir = "/input" + defaultOutputDir = "/output" + defaultProcessingSuffix = ".processing" + defaultOutputExtension = ".mp4" + defaultQueueSize = 128 + defaultMaxConcurrent = 1 + defaultHTTPPort = "8080" + defaultRescanInterval = 30 * time.Second + defaultStabilityDuration = 3 * time.Second +) + +const defaultFFMPEGCommand = "-y -hwaccel cuda -hwaccel_device 0 -i {{input}} -c:v hevc_nvenc -qp 25 -preset p6 -gpu 0 -b_qfactor 1.1 -b_ref_mode middle -bf 3 -g 250 -i_qfactor 0.75 -max_muxing_queue_size 1024 -multipass 1 -rc vbr -rc-lookahead 20 -temporal-aq 1 -tune hq -c:a aac -af volume=2.0 {{output}}" + +const defaultFFMPEGCommandCPU = "-y -i {{input}} -c:v libx264 -preset slow -crf 22 -c:a aac {{output}}" + +var defaultExtensions = []string{".mp4", ".mkv", ".mov", ".avi", ".flv", ".wmv", ".m4v", ".webm", ".ts"} + +type config struct { + inputDir string + outputDir string + ffmpegBinary string + ffmpegCommand string + deleteSource bool + processingSuffix string + outputExtension string + httpPort string + rescanInterval time.Duration + stabilityWindow time.Duration + queueSize int + maxConcurrent int + extensions map[string]struct{} +} + +func loadConfig() (config, error) { + cfg := config{ + inputDir: getEnv("INPUT_DIR", defaultInputDir), + outputDir: getEnv("OUTPUT_DIR", defaultOutputDir), + ffmpegBinary: getEnv("FFMPEG_BIN", "ffmpeg"), + processingSuffix: getEnv("PROCESSING_SUFFIX", defaultProcessingSuffix), + outputExtension: getEnv("OUTPUT_EXTENSION", defaultOutputExtension), + httpPort: getEnv("PORT", defaultHTTPPort), + queueSize: getEnvInt("QUEUE_SIZE", defaultQueueSize), + maxConcurrent: getEnvInt("MAX_CONCURRENT", defaultMaxConcurrent), + rescanInterval: getEnvDuration("RESCAN_INTERVAL", defaultRescanInterval), + stabilityWindow: getEnvDuration("FILE_STABILITY_DURATION", defaultStabilityDuration), + deleteSource: getEnvBool("DELETE_SOURCE"), + } + + if cfg.maxConcurrent < 1 { + cfg.maxConcurrent = 1 + } + if cfg.queueSize < cfg.maxConcurrent { + cfg.queueSize = cfg.maxConcurrent * 2 + } + if cfg.processingSuffix == "" { + cfg.processingSuffix = defaultProcessingSuffix + } + + extEnv := os.Getenv("VIDEO_EXTENSIONS") + if strings.TrimSpace(extEnv) == "" { + extEnv = strings.Join(defaultExtensions, ",") + } + cfg.extensions = make(map[string]struct{}) + for _, raw := range strings.Split(extEnv, ",") { + trimmed := strings.TrimSpace(raw) + if trimmed == "" { + continue + } + if !strings.HasPrefix(trimmed, ".") { + trimmed = "." + trimmed + } + cfg.extensions[strings.ToLower(trimmed)] = struct{}{} + } + if len(cfg.extensions) == 0 { + return cfg, errors.New("no video extensions configured") + } + + // Detect GPU and set ffmpeg command + gpuAvailable := detectGPU() + if gpuAvailable { + cfg.ffmpegCommand = getEnv("FFMPEG_COMMAND", defaultFFMPEGCommand) + } else { + log.Printf("GPU not detected, falling back to CPU encoding") + cfg.ffmpegCommand = getEnv("FFMPEG_COMMAND_CPU", defaultFFMPEGCommandCPU) + if cfg.ffmpegCommand == "" { + return cfg, errors.New("GPU not available and no CPU ffmpeg command configured") + } + } + + // If inputDir is customized but outputDir is default, assume local testing and set outputDir relative to inputDir + if cfg.outputDir == defaultOutputDir && cfg.inputDir != defaultInputDir { + cfg.outputDir = filepath.Join(filepath.Dir(cfg.inputDir), "test_output") + } + + return cfg, nil +} + +func getEnv(key, fallback string) string { + val := strings.TrimSpace(os.Getenv(key)) + if val == "" { + return fallback + } + return val +} + +func getEnvBool(key string) bool { + val := strings.TrimSpace(os.Getenv(key)) + if val == "" { + return false + } + switch strings.ToLower(val) { + case "1", "true", "yes", "on": + return true + default: + return false + } +} + +func getEnvInt(key string, fallback int) int { + val := strings.TrimSpace(os.Getenv(key)) + if val == "" { + return fallback + } + var parsed int + if _, err := fmt.Sscanf(val, "%d", &parsed); err != nil { + log.Printf("invalid int for %s: %v", key, err) + return fallback + } + return parsed +} + +func getEnvDuration(key string, fallback time.Duration) time.Duration { + val := strings.TrimSpace(os.Getenv(key)) + if val == "" { + return fallback + } + d, err := time.ParseDuration(val) + if err != nil { + log.Printf("invalid duration for %s: %v", key, err) + return fallback + } + return d +} + +func detectGPU() bool { + cmd := exec.Command("nvidia-smi") + err := cmd.Run() + return err == nil +} diff --git a/cmd/compressor/main.go b/cmd/compressor/main.go new file mode 100644 index 0000000..308cab7 --- /dev/null +++ b/cmd/compressor/main.go @@ -0,0 +1,206 @@ +package main + +import ( + "context" + "log" + "os" + "os/signal" + "path/filepath" + "strings" + "sync" + "syscall" + "time" + + "github.com/fsnotify/fsnotify" +) + +var processed sync.Map // track recently processed files to prevent loops + +func main() { + cfg, err := loadConfig() + if err != nil { + log.Fatalf("config error: %v", err) + } + + log.Printf("Configuration loaded:") + log.Printf(" Input Dir: %s", cfg.inputDir) + log.Printf(" Output Dir: %s", cfg.outputDir) + log.Printf(" FFmpeg Binary: %s", cfg.ffmpegBinary) + log.Printf(" FFmpeg Command: %s", cfg.ffmpegCommand) + log.Printf(" Delete Source: %t", cfg.deleteSource) + log.Printf(" Processing Suffix: %s", cfg.processingSuffix) + log.Printf(" Output Extension: %s", cfg.outputExtension) + log.Printf(" HTTP Port: %s", cfg.httpPort) + log.Printf(" Rescan Interval: %v", cfg.rescanInterval) + log.Printf(" Stability Window: %v", cfg.stabilityWindow) + log.Printf(" Queue Size: %d", cfg.queueSize) + log.Printf(" Max Concurrent: %d", cfg.maxConcurrent) + var exts []string + for ext := range cfg.extensions { + exts = append(exts, ext) + } + log.Printf(" Video Extensions: %v", exts) + + if err := os.MkdirAll(cfg.inputDir, 0o755); err != nil { + log.Fatalf("create input dir: %v", err) + } + if err := os.MkdirAll(cfg.outputDir, 0o755); err != nil { + log.Fatalf("create output dir: %v", err) + } + + ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) + defer cancel() + + queue := make(chan string, cfg.queueSize) + var inProgress sync.Map + var wg sync.WaitGroup + + dispatcherCtx, dispatcherCancel := context.WithCancel(ctx) + go func() { + <-dispatcherCtx.Done() + close(queue) + }() + + sem := make(chan struct{}, cfg.maxConcurrent) + + go func() { + for path := range queue { + path := path + sem <- struct{}{} + wg.Add(1) + go func() { + defer func() { + <-sem + inProgress.Delete(path) + wg.Done() + }() + + if err := processFile(ctx, cfg, path); err != nil { + log.Printf("process failed for %s: %v", path, err) + } + }() + } + }() + + enqueue := func(path string) { + if !shouldProcess(cfg, path) { + return + } + // Skip if recently processed to prevent loops + if val, ok := processed.Load(path); ok { + if recent := val.(time.Time); time.Since(recent) < 10*time.Second { + return + } + } + if _, loaded := inProgress.LoadOrStore(path, struct{}{}); loaded { + return + } + select { + case queue <- path: + default: + go func() { queue <- path }() + } + } + + if err := scanAndEnqueue(cfg, enqueue); err != nil { + log.Printf("initial scan failed: %v", err) + } + + watcher, err := fsnotify.NewWatcher() + if err != nil { + log.Fatalf("start watcher: %v", err) + } + defer watcher.Close() + + if err := watcher.Add(cfg.inputDir); err != nil { + log.Fatalf("watch dir: %v", err) + } + + go func() { + for { + select { + case <-ctx.Done(): + return + case event, ok := <-watcher.Events: + if !ok { + return + } + if event.Op&(fsnotify.Create|fsnotify.Rename|fsnotify.Write) != 0 { + enqueue(event.Name) + } + case err, ok := <-watcher.Errors: + if !ok { + return + } + log.Printf("watch error: %v", err) + } + } + }() + + go func() { + ticker := time.NewTicker(cfg.rescanInterval) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + if err := scanAndEnqueue(cfg, enqueue); err != nil { + log.Printf("periodic scan failed: %v", err) + } + } + } + }() + + serverErrs := make(chan error, 1) + go func() { + serverErrs <- runHTTPServer(ctx, cfg.httpPort) + }() + + select { + case <-ctx.Done(): + case err := <-serverErrs: + if err != nil { + log.Printf("http server error: %v", err) + } + } + + dispatcherCancel() + wg.Wait() +} + +func scanAndEnqueue(cfg config, enqueue func(string)) error { + entries, err := os.ReadDir(cfg.inputDir) + if err != nil { + return err + } + for _, entry := range entries { + if entry.IsDir() { + continue + } + fullPath := filepath.Join(cfg.inputDir, entry.Name()) + enqueue(fullPath) + } + return nil +} + +func shouldProcess(cfg config, path string) bool { + if !strings.HasPrefix(path, cfg.inputDir) { + return false + } + info, err := os.Stat(path) + if err != nil { + return false + } + if !info.Mode().IsRegular() { + return false + } + if strings.HasSuffix(path, cfg.processingSuffix) { + return false + } + ext := strings.ToLower(filepath.Ext(path)) + if _, ok := cfg.extensions[ext]; !ok { + return false + } + return true +} diff --git a/cmd/compressor/processor.go b/cmd/compressor/processor.go new file mode 100644 index 0000000..eee22e2 --- /dev/null +++ b/cmd/compressor/processor.go @@ -0,0 +1,165 @@ +package main + +import ( + "context" + "errors" + "fmt" + "log" + "os" + "os/exec" + "path/filepath" + "strings" + "time" + + "al.essio.dev/pkg/shellescape" + "github.com/mattn/go-shellwords" +) + +func processFile(ctx context.Context, cfg config, originalPath string) error { + if err := waitForStability(ctx, originalPath, cfg.stabilityWindow); err != nil { + if errors.Is(err, os.ErrNotExist) { + return nil + } + return fmt.Errorf("stability check: %w", err) + } + + processingPath := originalPath + cfg.processingSuffix + if err := os.Rename(originalPath, processingPath); err != nil { + if errors.Is(err, os.ErrNotExist) { + return nil + } + return fmt.Errorf("rename for processing: %w", err) + } + + success := false + defer func() { + if success { + if cfg.deleteSource { + if err := os.Remove(processingPath); err != nil && !errors.Is(err, os.ErrNotExist) { + log.Printf("cleanup remove failed for %s: %v", processingPath, err) + } + } else { + if err := os.Rename(processingPath, originalPath); err != nil { + log.Printf("restore original failed for %s: %v", processingPath, err) + } + } + } else { + if _, err := os.Stat(processingPath); err == nil { + if err := os.Rename(processingPath, originalPath); err != nil { + log.Printf("restore after failure failed for %s: %v", processingPath, err) + } + } + } + }() + + outputPath, err := buildOutputPath(cfg, originalPath) + if err != nil { + return err + } + + if err := runFFMPEG(ctx, cfg, processingPath, outputPath); err != nil { + if removeErr := os.Remove(outputPath); removeErr != nil && !errors.Is(removeErr, os.ErrNotExist) { + log.Printf("remove partial output %s failed: %v", outputPath, removeErr) + } + return err + } + + success = true + log.Printf("processed %s -> %s", originalPath, outputPath) + processed.Store(originalPath, time.Now()) + return nil +} + +func waitForStability(ctx context.Context, path string, stableFor time.Duration) error { + if stableFor <= 0 { + return nil + } + + const tick = 500 * time.Millisecond + var ( + prevSize int64 = -1 + stableTime time.Time + ) + + ticker := time.NewTicker(tick) + defer ticker.Stop() + + for { + select { + case <-ctx.Done(): + return ctx.Err() + case <-ticker.C: + info, err := os.Stat(path) + if err != nil { + return err + } + if !info.Mode().IsRegular() { + return fmt.Errorf("%s is not a regular file", path) + } + size := info.Size() + if size == prevSize { + if stableTime.IsZero() { + stableTime = time.Now() + } + if time.Since(stableTime) >= stableFor { + return nil + } + } else { + prevSize = size + stableTime = time.Time{} + } + } + } +} + +func buildOutputPath(cfg config, originalPath string) (string, error) { + base := strings.TrimSuffix(filepath.Base(originalPath), filepath.Ext(originalPath)) + ext := cfg.outputExtension + if ext == "" { + ext = filepath.Ext(originalPath) + } + if !strings.HasPrefix(ext, ".") { + ext = "." + ext + } + + candidate := filepath.Join(cfg.outputDir, base+ext) + if _, err := os.Stat(candidate); errors.Is(err, os.ErrNotExist) { + return candidate, nil + } + if err := os.MkdirAll(cfg.outputDir, 0o755); err != nil { + return "", fmt.Errorf("ensure output dir: %w", err) + } + for idx := 1; idx < 10_000; idx++ { + candidate = filepath.Join(cfg.outputDir, fmt.Sprintf("%s_%d%s", base, idx, ext)) + if _, err := os.Stat(candidate); errors.Is(err, os.ErrNotExist) { + return candidate, nil + } + } + return "", fmt.Errorf("unable to find free output name for %s", originalPath) +} + +func runFFMPEG(ctx context.Context, cfg config, inputPath, outputPath string) error { + if err := os.MkdirAll(filepath.Dir(outputPath), 0o755); err != nil { + return fmt.Errorf("prepare output dir: %w", err) + } + + substituted := strings.ReplaceAll(cfg.ffmpegCommand, "{{input}}", shellescape.Quote(inputPath)) + substituted = strings.ReplaceAll(substituted, "{{output}}", shellescape.Quote(outputPath)) + + args, err := shellwords.Parse(substituted) + if err != nil { + return fmt.Errorf("parse ffmpeg args: %w", err) + } + + cmd := exec.CommandContext(ctx, cfg.ffmpegBinary, args...) + cmd.Stdout = os.Stdout + cmd.Stderr = os.Stderr + cmd.Env = os.Environ() + + log.Printf("ffmpeg start: %s -> %s", inputPath, outputPath) + + if err := cmd.Run(); err != nil { + return fmt.Errorf("ffmpeg failed: %w", err) + } + return nil +} diff --git a/cmd/compressor/server.go b/cmd/compressor/server.go new file mode 100644 index 0000000..7ac30db --- /dev/null +++ b/cmd/compressor/server.go @@ -0,0 +1,41 @@ +package main + +import ( + "context" + "errors" + "net/http" + "time" +) + +func runHTTPServer(ctx context.Context, port string) error { + mux := http.NewServeMux() + mux.HandleFunc("/status", func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusOK) + _, _ = w.Write([]byte("ok")) + }) + + server := &http.Server{ + Addr: ":" + port, + Handler: mux, + } + + errCh := make(chan error, 1) + go func() { + if err := server.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) { + errCh <- err + } + close(errCh) + }() + + select { + case <-ctx.Done(): + shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + if err := server.Shutdown(shutdownCtx); err != nil { + return err + } + return nil + case err := <-errCh: + return err + } +} diff --git a/docker-compose.yml b/docker-compose.yml new file mode 100644 index 0000000..2cc7232 --- /dev/null +++ b/docker-compose.yml @@ -0,0 +1,22 @@ +version: "3.8" + +services: + compressor: + build: . + ports: + - "8080:8080" + volumes: + - ./test_input:/input + - ./test_output:/output + environment: + - INPUT_DIR=/input + - OUTPUT_DIR=/output + - DELETE_SOURCE=false + deploy: + resources: + reservations: + devices: + - driver: nvidia + count: 1 + capabilities: [gpu] + restart: unless-stopped diff --git a/go.mod b/go.mod new file mode 100644 index 0000000..b439ee1 --- /dev/null +++ b/go.mod @@ -0,0 +1,11 @@ +module compressor + +go 1.21 + +require ( + al.essio.dev/pkg/shellescape v1.6.0 + github.com/fsnotify/fsnotify v1.9.0 + github.com/mattn/go-shellwords v1.0.12 +) + +require golang.org/x/sys v0.13.0 // indirect diff --git a/go.sum b/go.sum new file mode 100644 index 0000000..231fcd7 --- /dev/null +++ b/go.sum @@ -0,0 +1,10 @@ +al.essio.dev/pkg/shellescape v1.6.0 h1:NxFcEqzFSEVCGN2yq7Huv/9hyCEGVa/TncnOOBBeXHA= +al.essio.dev/pkg/shellescape v1.6.0/go.mod h1:6sIqp7X2P6mThCQ7twERpZTuigpr6KbZWtls1U8I890= +github.com/fsnotify/fsnotify v1.9.0 h1:2Ml+OJNzbYCTzsxtv8vKSFD9PbJjmhYF14k/jKC7S9k= +github.com/fsnotify/fsnotify v1.9.0/go.mod h1:8jBTzvmWwFyi3Pb8djgCCO5IBqzKJ/Jwo8TRcHyHii0= +github.com/google/shlex v0.0.0-20191202100458-e7afc7fbc510 h1:El6M4kTTCOh6aBiKaUGG7oYTSPP8MxqL4YI3kZKwcP4= +github.com/google/shlex v0.0.0-20191202100458-e7afc7fbc510/go.mod h1:pupxD2MaaD3pAXIBCelhxNneeOaAeabZDe5s4K6zSpQ= +github.com/mattn/go-shellwords v1.0.12 h1:M2zGm7EW6UQJvDeQxo4T51eKPurbeFbe8WtebGE2xrk= +github.com/mattn/go-shellwords v1.0.12/go.mod h1:EZzvwXDESEeg03EKmM+RmDnNOPKG4lLtQsUlTZDWQ8Y= +golang.org/x/sys v0.13.0 h1:Af8nKPmuFypiUBjVoU9V20FiaFXOcuZI21p0ycVYYGE= +golang.org/x/sys v0.13.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= diff --git a/k8s/deployment.yaml b/k8s/deployment.yaml new file mode 100644 index 0000000..105723d --- /dev/null +++ b/k8s/deployment.yaml @@ -0,0 +1,52 @@ +apiVersion: apps/v1 +kind: Deployment +metadata: + name: compressor + namespace: compressor + labels: + app: compressor +spec: + replicas: 1 + selector: + matchLabels: + app: compressor + template: + metadata: + labels: + app: compressor + spec: + containers: + - name: compressor + image: compressor:latest + ports: + - containerPort: 8080 + env: + - name: INPUT_DIR + value: "/input" + - name: OUTPUT_DIR + value: "/output" + - name: DELETE_SOURCE + value: "false" + - name: MAX_CONCURRENT + value: "1" + volumeMounts: + - name: input-volume + mountPath: /input + - name: output-volume + mountPath: /output + resources: + requests: + memory: "512Mi" + cpu: "500m" + nvidia.com/gpu: 1 + limits: + memory: "1Gi" + cpu: "1000m" + nvidia.com/gpu: 1 + volumes: + - name: input-volume + persistentVolumeClaim: + claimName: compressor-input + - name: output-volume + persistentVolumeClaim: + claimName: compressor-output diff --git a/k8s/namespace.yaml b/k8s/namespace.yaml new file mode 100644 index 0000000..d89adad --- /dev/null +++ b/k8s/namespace.yaml @@ -0,0 +1,4 @@ +apiVersion: v1 +kind: Namespace +metadata: + name: compressor diff --git a/k8s/pvc.yaml b/k8s/pvc.yaml new file mode 100644 index 0000000..97e1af7 --- /dev/null +++ b/k8s/pvc.yaml @@ -0,0 +1,23 @@ +apiVersion: v1 +kind: PersistentVolumeClaim +metadata: + name: compressor-input + namespace: compressor +spec: + accessModes: + - ReadWriteOnce + resources: + requests: + storage: 10Gi +--- +apiVersion: v1 +kind: PersistentVolumeClaim +metadata: + name: compressor-output + namespace: compressor +spec: + accessModes: + - ReadWriteOnce + resources: + requests: + storage: 50Gi diff --git a/k8s/service.yaml b/k8s/service.yaml new file mode 100644 index 0000000..d7ee9ce --- /dev/null +++ b/k8s/service.yaml @@ -0,0 +1,15 @@ +apiVersion: v1 +kind: Service +metadata: + name: compressor + namespace: compressor + labels: + app: compressor +spec: + type: ClusterIP + ports: + - port: 8080 + targetPort: 8080 + protocol: TCP + selector: + app: compressor