Add HTTP upload configuration support
This commit is contained in:
@@ -1,10 +1,24 @@
|
||||
package config
|
||||
|
||||
type Config struct {
|
||||
Server Server `yaml:"server"`
|
||||
Secrets Secrets `yaml:"secrets"`
|
||||
Pipelines []Pipeline `yaml:"pipelines"`
|
||||
}
|
||||
|
||||
type Server struct {
|
||||
HTTP HTTPServer `yaml:"http"`
|
||||
}
|
||||
|
||||
type HTTPServer struct {
|
||||
Bind string `yaml:"bind"`
|
||||
StagingRoot string `yaml:"staging_root"`
|
||||
MaxUploadSize *ByteSize `yaml:"max_upload_size"`
|
||||
QueueSize int `yaml:"queue_size"`
|
||||
MaxConcurrency int `yaml:"max_concurrency"`
|
||||
Retention *Duration `yaml:"retention"`
|
||||
}
|
||||
|
||||
type Secrets struct {
|
||||
Directory string `yaml:"directory"`
|
||||
}
|
||||
@@ -50,6 +64,13 @@ type Backend struct {
|
||||
ForcePath *bool `yaml:"force_path_style"`
|
||||
Creds Credentials `yaml:"credentials"`
|
||||
SSH SSH `yaml:",inline"`
|
||||
Upload HTTPUpload `yaml:",inline"`
|
||||
}
|
||||
|
||||
type HTTPUpload struct {
|
||||
TokenEnv string `yaml:"token_env"`
|
||||
StagingPath string `yaml:"staging_path"`
|
||||
MaxUploadSize *ByteSize `yaml:"max_upload_size"`
|
||||
}
|
||||
|
||||
type SSH struct {
|
||||
|
||||
@@ -1,13 +1,19 @@
|
||||
package config
|
||||
|
||||
import "gitea.maximumdirect.net/eric/distributor/internal/transform"
|
||||
import (
|
||||
"path/filepath"
|
||||
"time"
|
||||
|
||||
"gitea.maximumdirect.net/eric/distributor/internal/transform"
|
||||
)
|
||||
|
||||
const DefaultConfigPath = "/usr/local/etc/distributor/config.yml"
|
||||
|
||||
const (
|
||||
BackendLocal = "local"
|
||||
BackendSSH = "ssh"
|
||||
BackendS3 = "s3"
|
||||
BackendLocal = "local"
|
||||
BackendSSH = "ssh"
|
||||
BackendS3 = "s3"
|
||||
BackendHTTPUpload = "http_upload"
|
||||
)
|
||||
|
||||
const (
|
||||
@@ -38,10 +44,23 @@ const (
|
||||
|
||||
const DefaultS3Region = "us-east-1"
|
||||
|
||||
const (
|
||||
DefaultHTTPBind = "127.0.0.1:8080"
|
||||
DefaultHTTPStagingRoot = "/var/spool/distributor"
|
||||
DefaultHTTPMaxUploadSize = ByteSize(20 * 1024 * 1024)
|
||||
DefaultHTTPQueueSize = 16
|
||||
DefaultHTTPMaxConcurrency = 1
|
||||
DefaultHTTPRetention = Duration(24 * time.Hour)
|
||||
)
|
||||
|
||||
func ApplyDefaults(cfg *Config) {
|
||||
applyHTTPServerDefaults(&cfg.Server.HTTP)
|
||||
for pipelineIndex := range cfg.Pipelines {
|
||||
pipeline := &cfg.Pipelines[pipelineIndex]
|
||||
applyBackendDefaults(&pipeline.Source)
|
||||
if pipeline.Source.Backend == BackendHTTPUpload {
|
||||
applyHTTPUploadDefaults(&pipeline.Source.Upload, pipeline.ID, cfg.Server.HTTP)
|
||||
}
|
||||
if pipeline.Validation.OnDigestMismatch == "" {
|
||||
pipeline.Validation.OnDigestMismatch = ValidationActionFail
|
||||
}
|
||||
@@ -76,6 +95,44 @@ func ApplyDefaults(cfg *Config) {
|
||||
}
|
||||
}
|
||||
|
||||
func applyHTTPServerDefaults(server *HTTPServer) {
|
||||
if server.Bind == "" {
|
||||
server.Bind = DefaultHTTPBind
|
||||
}
|
||||
if server.StagingRoot == "" {
|
||||
server.StagingRoot = DefaultHTTPStagingRoot
|
||||
}
|
||||
if server.MaxUploadSize == nil {
|
||||
server.MaxUploadSize = byteSize(DefaultHTTPMaxUploadSize)
|
||||
}
|
||||
if server.QueueSize == 0 {
|
||||
server.QueueSize = DefaultHTTPQueueSize
|
||||
}
|
||||
if server.MaxConcurrency == 0 {
|
||||
server.MaxConcurrency = DefaultHTTPMaxConcurrency
|
||||
}
|
||||
if server.Retention == nil {
|
||||
server.Retention = duration(DefaultHTTPRetention)
|
||||
}
|
||||
}
|
||||
|
||||
func applyHTTPUploadDefaults(upload *HTTPUpload, pipelineID string, server HTTPServer) {
|
||||
if upload.StagingPath == "" && pipelineID != "" {
|
||||
upload.StagingPath = filepath.Join(server.StagingRoot, pipelineID)
|
||||
}
|
||||
if upload.MaxUploadSize == nil && server.MaxUploadSize != nil {
|
||||
upload.MaxUploadSize = byteSize(*server.MaxUploadSize)
|
||||
}
|
||||
}
|
||||
|
||||
func byteSize(value ByteSize) *ByteSize {
|
||||
return &value
|
||||
}
|
||||
|
||||
func duration(value Duration) *Duration {
|
||||
return &value
|
||||
}
|
||||
|
||||
func applyBackendDefaults(backend *Backend) {
|
||||
if backend.Backend == BackendSSH {
|
||||
if backend.Port == 0 {
|
||||
|
||||
@@ -208,6 +208,123 @@ pipelines:
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadFileDefaultsHTTPServerConfig(t *testing.T) {
|
||||
cfg := loadConfig(t, `
|
||||
pipelines:
|
||||
- id: reports
|
||||
source:
|
||||
backend: local
|
||||
path: /source
|
||||
destinations:
|
||||
- id: archive
|
||||
backend: local
|
||||
path: /archive
|
||||
`)
|
||||
|
||||
server := cfg.Server.HTTP
|
||||
if got, want := server.Bind, DefaultHTTPBind; got != want {
|
||||
t.Fatalf("server.http.bind = %q, want %q", got, want)
|
||||
}
|
||||
if got, want := server.StagingRoot, DefaultHTTPStagingRoot; got != want {
|
||||
t.Fatalf("server.http.staging_root = %q, want %q", got, want)
|
||||
}
|
||||
if got, want := *server.MaxUploadSize, DefaultHTTPMaxUploadSize; got != want {
|
||||
t.Fatalf("server.http.max_upload_size = %s, want %s", got, want)
|
||||
}
|
||||
if got, want := server.QueueSize, DefaultHTTPQueueSize; got != want {
|
||||
t.Fatalf("server.http.queue_size = %d, want %d", got, want)
|
||||
}
|
||||
if got, want := server.MaxConcurrency, DefaultHTTPMaxConcurrency; got != want {
|
||||
t.Fatalf("server.http.max_concurrency = %d, want %d", got, want)
|
||||
}
|
||||
if got, want := *server.Retention, DefaultHTTPRetention; got != want {
|
||||
t.Fatalf("server.http.retention = %s, want %s", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadFileAcceptsHTTPUploadSourceConfig(t *testing.T) {
|
||||
cfg := loadConfig(t, `
|
||||
server:
|
||||
http:
|
||||
bind: 127.0.0.1:9090
|
||||
staging_root: /srv/distributor/staging
|
||||
max_upload_size: 64MB
|
||||
queue_size: 32
|
||||
max_concurrency: 2
|
||||
retention: 48h
|
||||
pipelines:
|
||||
- id: weather-daily
|
||||
source:
|
||||
backend: http_upload
|
||||
token_env: WEATHER_DAILY_UPLOAD_TOKEN
|
||||
staging_path: /srv/distributor/staging/weather-daily
|
||||
max_upload_size: 32MB
|
||||
destinations:
|
||||
- id: archive
|
||||
backend: local
|
||||
path: /archive
|
||||
`)
|
||||
|
||||
server := cfg.Server.HTTP
|
||||
if got, want := server.Bind, "127.0.0.1:9090"; got != want {
|
||||
t.Fatalf("server.http.bind = %q, want %q", got, want)
|
||||
}
|
||||
if got, want := server.StagingRoot, "/srv/distributor/staging"; got != want {
|
||||
t.Fatalf("server.http.staging_root = %q, want %q", got, want)
|
||||
}
|
||||
if got, want := *server.MaxUploadSize, ByteSize(64*1024*1024); got != want {
|
||||
t.Fatalf("server.http.max_upload_size = %s, want %s", got, want)
|
||||
}
|
||||
if got, want := server.QueueSize, 32; got != want {
|
||||
t.Fatalf("server.http.queue_size = %d, want %d", got, want)
|
||||
}
|
||||
if got, want := server.MaxConcurrency, 2; got != want {
|
||||
t.Fatalf("server.http.max_concurrency = %d, want %d", got, want)
|
||||
}
|
||||
if got, want := server.Retention.String(), "48h0m0s"; got != want {
|
||||
t.Fatalf("server.http.retention = %s, want %s", got, want)
|
||||
}
|
||||
|
||||
source := cfg.Pipelines[0].Source
|
||||
if got, want := source.Backend, BackendHTTPUpload; got != want {
|
||||
t.Fatalf("source.backend = %q, want %q", got, want)
|
||||
}
|
||||
if got, want := source.Upload.TokenEnv, "WEATHER_DAILY_UPLOAD_TOKEN"; got != want {
|
||||
t.Fatalf("source.token_env = %q, want %q", got, want)
|
||||
}
|
||||
if got, want := source.Upload.StagingPath, "/srv/distributor/staging/weather-daily"; got != want {
|
||||
t.Fatalf("source.staging_path = %q, want %q", got, want)
|
||||
}
|
||||
if got, want := *source.Upload.MaxUploadSize, ByteSize(32*1024*1024); got != want {
|
||||
t.Fatalf("source.max_upload_size = %s, want %s", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadFileDefaultsHTTPUploadSourceConfig(t *testing.T) {
|
||||
cfg := loadConfig(t, `
|
||||
server:
|
||||
http:
|
||||
max_upload_size: 12MB
|
||||
pipelines:
|
||||
- id: weather-daily
|
||||
source:
|
||||
backend: http_upload
|
||||
token_env: WEATHER_DAILY_UPLOAD_TOKEN
|
||||
destinations:
|
||||
- id: archive
|
||||
backend: local
|
||||
path: /archive
|
||||
`)
|
||||
|
||||
source := cfg.Pipelines[0].Source
|
||||
if got, want := source.Upload.StagingPath, "/var/spool/distributor/weather-daily"; got != want {
|
||||
t.Fatalf("source.staging_path = %q, want %q", got, want)
|
||||
}
|
||||
if got, want := *source.Upload.MaxUploadSize, ByteSize(12*1024*1024); got != want {
|
||||
t.Fatalf("source.max_upload_size = %s, want %s", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadFileValidBackendConfigs(t *testing.T) {
|
||||
tests := map[string]string{
|
||||
"local": `
|
||||
@@ -396,6 +513,26 @@ func TestLoadFileRejectsInvalidS3Config(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadFileRejectsInvalidHTTPUploadConfig(t *testing.T) {
|
||||
tests := map[string]string{
|
||||
"server size": `server: {http: {max_upload_size: 20XB}}`,
|
||||
"source size": `pipelines: [{id: reports, source: {backend: http_upload, token_env: UPLOAD_TOKEN, max_upload_size: 20XB}, destinations: [{id: archive, backend: local, path: /archive}]}]`,
|
||||
"zero source size": `pipelines: [{id: reports, source: {backend: http_upload, token_env: UPLOAD_TOKEN, max_upload_size: 0B}, destinations: [{id: archive, backend: local, path: /archive}]}]`,
|
||||
"server duration": `server: {http: {retention: forever}}`,
|
||||
"zero server duration": `server: {http: {retention: 0s}}`,
|
||||
"missing token env": `pipelines: [{id: reports, source: {backend: http_upload}, destinations: [{id: archive, backend: local, path: /archive}]}]`,
|
||||
"destination http upload": `pipelines: [{id: reports, source: {backend: local, path: /source}, destinations: [{id: ingest, backend: http_upload}]}]`,
|
||||
"literal token": `pipelines: [{id: reports, source: {backend: http_upload, token: secret, token_env: UPLOAD_TOKEN}, destinations: [{id: archive, backend: local, path: /archive}]}]`,
|
||||
"unknown server field": `server: {http: {surprise: true}}`,
|
||||
"unknown source field": `pipelines: [{id: reports, source: {backend: http_upload, token_env: UPLOAD_TOKEN, surprise: true}, destinations: [{id: archive, backend: local, path: /archive}]}]`,
|
||||
}
|
||||
for name, body := range tests {
|
||||
t.Run(name, func(t *testing.T) {
|
||||
assertLoadError(t, body, "")
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestLoadFileDefaultsSSHConfig(t *testing.T) {
|
||||
cfg := loadConfig(t, `
|
||||
pipelines:
|
||||
|
||||
117
internal/config/quantity.go
Normal file
117
internal/config/quantity.go
Normal file
@@ -0,0 +1,117 @@
|
||||
package config
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"gopkg.in/yaml.v3"
|
||||
)
|
||||
|
||||
type ByteSize int64
|
||||
|
||||
type Duration time.Duration
|
||||
|
||||
func (size *ByteSize) UnmarshalYAML(value *yaml.Node) error {
|
||||
if value.Kind != yaml.ScalarNode || value.Tag != "!!str" {
|
||||
return fmt.Errorf("size must be a string with B, KB, MB, or GB suffix")
|
||||
}
|
||||
|
||||
var raw string
|
||||
if err := value.Decode(&raw); err != nil {
|
||||
return err
|
||||
}
|
||||
parsed, err := ParseByteSize(raw)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
*size = parsed
|
||||
return nil
|
||||
}
|
||||
|
||||
func (size ByteSize) String() string {
|
||||
value := int64(size)
|
||||
if value == 0 {
|
||||
return "0B"
|
||||
}
|
||||
units := []struct {
|
||||
suffix string
|
||||
multiplier int64
|
||||
}{
|
||||
{suffix: "GB", multiplier: 1024 * 1024 * 1024},
|
||||
{suffix: "MB", multiplier: 1024 * 1024},
|
||||
{suffix: "KB", multiplier: 1024},
|
||||
{suffix: "B", multiplier: 1},
|
||||
}
|
||||
for _, unit := range units {
|
||||
if value%unit.multiplier == 0 {
|
||||
return strconv.FormatInt(value/unit.multiplier, 10) + unit.suffix
|
||||
}
|
||||
}
|
||||
return strconv.FormatInt(value, 10) + "B"
|
||||
}
|
||||
|
||||
func ParseByteSize(raw string) (ByteSize, error) {
|
||||
value := strings.TrimSpace(raw)
|
||||
if value == "" {
|
||||
return 0, fmt.Errorf("size is required")
|
||||
}
|
||||
|
||||
units := []struct {
|
||||
suffix string
|
||||
multiplier int64
|
||||
}{
|
||||
{suffix: "GB", multiplier: 1024 * 1024 * 1024},
|
||||
{suffix: "MB", multiplier: 1024 * 1024},
|
||||
{suffix: "KB", multiplier: 1024},
|
||||
{suffix: "B", multiplier: 1},
|
||||
}
|
||||
for _, unit := range units {
|
||||
number, ok := strings.CutSuffix(value, unit.suffix)
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
if strings.TrimSpace(number) != number || number == "" {
|
||||
return 0, fmt.Errorf("size must be an integer followed by B, KB, MB, or GB")
|
||||
}
|
||||
parsed, err := strconv.ParseInt(number, 10, 64)
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("size must be an integer followed by B, KB, MB, or GB")
|
||||
}
|
||||
if parsed < 0 {
|
||||
return 0, fmt.Errorf("size must be non-negative")
|
||||
}
|
||||
const maxInt64 = int64(1<<63 - 1)
|
||||
if parsed > 0 && parsed > maxInt64/unit.multiplier {
|
||||
return 0, fmt.Errorf("size is too large")
|
||||
}
|
||||
return ByteSize(parsed * unit.multiplier), nil
|
||||
}
|
||||
return 0, fmt.Errorf("size must use B, KB, MB, or GB suffix")
|
||||
}
|
||||
|
||||
func (duration *Duration) UnmarshalYAML(value *yaml.Node) error {
|
||||
if value.Kind != yaml.ScalarNode || value.Tag != "!!str" {
|
||||
return fmt.Errorf("duration must be a string duration")
|
||||
}
|
||||
|
||||
var raw string
|
||||
if err := value.Decode(&raw); err != nil {
|
||||
return err
|
||||
}
|
||||
parsed, err := time.ParseDuration(raw)
|
||||
if err != nil {
|
||||
return fmt.Errorf("duration must be a valid duration: %w", err)
|
||||
}
|
||||
*duration = Duration(parsed)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (duration Duration) String() string {
|
||||
return time.Duration(duration).String()
|
||||
}
|
||||
|
||||
func (duration Duration) AsDuration() time.Duration {
|
||||
return time.Duration(duration)
|
||||
}
|
||||
@@ -22,6 +22,8 @@ func (e ValidationErrors) Error() string {
|
||||
func Validate(cfg Config) error {
|
||||
var errs ValidationErrors
|
||||
|
||||
errs = validateHTTPServer(errs, "server.http", cfg.Server.HTTP)
|
||||
|
||||
if len(cfg.Pipelines) == 0 {
|
||||
errs = append(errs, "pipelines is required")
|
||||
}
|
||||
@@ -72,14 +74,56 @@ func Validate(cfg Config) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func validateHTTPServer(errs ValidationErrors, context string, server HTTPServer) ValidationErrors {
|
||||
if server.Bind == "" {
|
||||
errs = append(errs, context+".bind is required")
|
||||
}
|
||||
if server.StagingRoot == "" {
|
||||
errs = append(errs, context+".staging_root is required")
|
||||
}
|
||||
if server.MaxUploadSize == nil || *server.MaxUploadSize <= 0 {
|
||||
errs = append(errs, context+".max_upload_size must be greater than zero")
|
||||
}
|
||||
if server.QueueSize <= 0 {
|
||||
errs = append(errs, context+".queue_size must be greater than zero")
|
||||
}
|
||||
if server.MaxConcurrency <= 0 {
|
||||
errs = append(errs, context+".max_concurrency must be greater than zero")
|
||||
}
|
||||
if server.Retention == nil || *server.Retention <= 0 {
|
||||
errs = append(errs, context+".retention must be greater than zero")
|
||||
}
|
||||
return errs
|
||||
}
|
||||
|
||||
func validateSourceBackend(errs ValidationErrors, context string, backend Backend) ValidationErrors {
|
||||
if backend.Backend == BackendHTTPUpload {
|
||||
return validateHTTPUploadSource(errs, context, backend.Upload)
|
||||
}
|
||||
return validateBackend(errs, context, backend.Backend, backend.Host, backend.Port, backend.Path, backend.Endpoint, backend.Bucket, backend.Prefix, backend.SSH.HostKeyPolicy, backend.Creds)
|
||||
}
|
||||
|
||||
func validateDestinationBackend(errs ValidationErrors, context string, destination Destination) ValidationErrors {
|
||||
if destination.Backend == BackendHTTPUpload {
|
||||
errs = append(errs, context+".backend "+BackendHTTPUpload+" is only supported for sources")
|
||||
return errs
|
||||
}
|
||||
return validateBackend(errs, context, destination.Backend, destination.Host, destination.Port, destination.Path, destination.Endpoint, destination.Bucket, destination.Prefix, destination.SSH.HostKeyPolicy, destination.Creds)
|
||||
}
|
||||
|
||||
func validateHTTPUploadSource(errs ValidationErrors, context string, upload HTTPUpload) ValidationErrors {
|
||||
if upload.TokenEnv == "" {
|
||||
errs = append(errs, context+".token_env is required for http_upload backend")
|
||||
}
|
||||
if upload.StagingPath == "" {
|
||||
errs = append(errs, context+".staging_path is required for http_upload backend")
|
||||
}
|
||||
if upload.MaxUploadSize == nil || *upload.MaxUploadSize <= 0 {
|
||||
errs = append(errs, context+".max_upload_size must be greater than zero")
|
||||
}
|
||||
return errs
|
||||
}
|
||||
|
||||
func validateBackend(errs ValidationErrors, context, backend, host string, port int, path, endpoint, bucket, prefix string, hostKeyPolicy HostKeyPolicy, creds Credentials) ValidationErrors {
|
||||
switch backend {
|
||||
case "":
|
||||
|
||||
Reference in New Issue
Block a user