Add distributor upload adapter

This commit is contained in:
2026-06-07 23:34:21 +00:00
parent a2ba6f5382
commit 64cae8c4d9
5 changed files with 525 additions and 1 deletions

View File

@@ -0,0 +1,229 @@
// Package distributor adapts weatherreporter report artifacts to distributor uploads.
package distributor
import (
"context"
"errors"
"fmt"
"net/http"
"os"
"strings"
"time"
distributorbundle "gitea.maximumdirect.net/eric/distributor/pkg/bundle"
distributorupload "gitea.maximumdirect.net/eric/distributor/pkg/upload"
"gitea.maximumdirect.net/eric/weatherreporter/internal/config"
)
type Client struct {
Endpoint string
TokenEnv string
Timeout time.Duration
newUploadClient uploadClientFactory
}
type UploadRequest struct {
BundleID string
IdempotencyKey string
SourcePath string
BundlePath string
}
type UploadResult struct {
RunID string
Status string
}
type IdempotencyConflictError struct {
Err error
}
func (e *IdempotencyConflictError) Error() string {
if e == nil || e.Err == nil {
return "distributor idempotency conflict"
}
return e.Err.Error()
}
func (e *IdempotencyConflictError) Unwrap() error {
if e == nil {
return nil
}
return e.Err
}
type uploadClientFactory func(endpoint, token string, timeout time.Duration) (uploadClient, error)
type uploadClient interface {
UploadFiles(ctx context.Context, opts uploadFilesOptions) (uploadFilesResult, error)
}
type uploadFilesOptions struct {
BundleID string
IdempotencyKey string
SourcePath string
BundlePath string
}
type uploadFilesResult struct {
RunID string
Status string
}
func New(cfg config.DistributorNotifyConfig) *Client {
return newClient(cfg, newDistributorUploadClient)
}
func newClient(cfg config.DistributorNotifyConfig, factory uploadClientFactory) *Client {
if factory == nil {
factory = newDistributorUploadClient
}
return &Client{
Endpoint: cfg.Endpoint,
TokenEnv: cfg.TokenEnv,
Timeout: cfg.Timeout,
newUploadClient: factory,
}
}
func (c *Client) Upload(ctx context.Context, req UploadRequest) (UploadResult, error) {
if c == nil {
return UploadResult{}, fmt.Errorf("distributor client is nil")
}
if c.Endpoint == "" {
return UploadResult{}, fmt.Errorf("distributor endpoint is required")
}
if c.TokenEnv == "" {
return UploadResult{}, fmt.Errorf("distributor token environment variable is required")
}
if req.BundleID == "" {
return UploadResult{}, fmt.Errorf("distributor bundle id is required")
}
if req.IdempotencyKey == "" {
return UploadResult{}, fmt.Errorf("distributor idempotency key is required for bundle %q", req.BundleID)
}
if req.SourcePath == "" {
return UploadResult{}, fmt.Errorf("distributor source path is required for bundle %q", req.BundleID)
}
if req.BundlePath == "" {
return UploadResult{}, fmt.Errorf("distributor bundle path is required for bundle %q", req.BundleID)
}
if c.newUploadClient == nil {
return UploadResult{}, fmt.Errorf("distributor upload client factory is required for endpoint %q", c.Endpoint)
}
token := os.Getenv(c.TokenEnv)
if token == "" {
return UploadResult{}, fmt.Errorf("distributor token environment variable %q is not set", c.TokenEnv)
}
uploadClient, err := c.newUploadClient(c.Endpoint, token, c.Timeout)
if err != nil {
return UploadResult{}, fmt.Errorf("create distributor upload client for endpoint %q: %w", c.Endpoint, redactToken(err, token))
}
runCtx := ctx
if runCtx == nil {
runCtx = context.Background()
}
cancel := func() {}
if c.Timeout > 0 {
runCtx, cancel = context.WithTimeout(runCtx, c.Timeout)
}
defer cancel()
result, err := uploadClient.UploadFiles(runCtx, uploadFilesOptions{
BundleID: req.BundleID,
IdempotencyKey: req.IdempotencyKey,
SourcePath: req.SourcePath,
BundlePath: req.BundlePath,
})
if err != nil {
return UploadResult{}, wrapUploadError(err, uploadErrorContext{
Endpoint: c.Endpoint,
BundleID: req.BundleID,
IdempotencyKey: req.IdempotencyKey,
SourcePath: req.SourcePath,
BundlePath: req.BundlePath,
Token: token,
})
}
return UploadResult{
RunID: result.RunID,
Status: result.Status,
}, nil
}
type distributorUploadClient struct {
client *distributorupload.Client
}
func newDistributorUploadClient(endpoint, token string, timeout time.Duration) (uploadClient, error) {
httpClient := (*http.Client)(nil)
if timeout > 0 {
httpClient = &http.Client{Timeout: timeout}
}
client, err := distributorupload.NewClient(distributorupload.ClientOptions{
Endpoint: endpoint,
Token: token,
HTTPClient: httpClient,
})
if err != nil {
return nil, err
}
return distributorUploadClient{client: client}, nil
}
func (c distributorUploadClient) UploadFiles(ctx context.Context, opts uploadFilesOptions) (uploadFilesResult, error) {
result, err := c.client.UploadFiles(ctx, distributorupload.UploadFilesOptions{
ID: opts.BundleID,
IdempotencyKey: opts.IdempotencyKey,
Files: []distributorbundle.BundleFile{
{SourcePath: opts.SourcePath, Path: opts.BundlePath},
},
})
if err != nil {
return uploadFilesResult{}, err
}
return uploadFilesResult{
RunID: result.RunID,
Status: result.Status,
}, nil
}
type uploadErrorContext struct {
Endpoint string
BundleID string
IdempotencyKey string
SourcePath string
BundlePath string
Token string
}
func wrapUploadError(err error, ctx uploadErrorContext) error {
var conflict *distributorupload.IdempotencyConflictError
isConflict := errors.As(err, &conflict)
err = redactToken(err, ctx.Token)
if isConflict {
return &IdempotencyConflictError{
Err: fmt.Errorf("upload distributor bundle %q to endpoint %q with idempotency key %q from source %q as bundle path %q: idempotency conflict: %w", ctx.BundleID, ctx.Endpoint, ctx.IdempotencyKey, ctx.SourcePath, ctx.BundlePath, err),
}
}
return fmt.Errorf("upload distributor bundle %q to endpoint %q with idempotency key %q from source %q as bundle path %q: %w", ctx.BundleID, ctx.Endpoint, ctx.IdempotencyKey, ctx.SourcePath, ctx.BundlePath, err)
}
func redactToken(err error, token string) error {
if err == nil || token == "" {
return err
}
return errors.New(redactTokenString(err.Error(), token))
}
func redactTokenString(value, token string) string {
if token == "" {
return value
}
return strings.ReplaceAll(value, token, "[redacted]")
}

View File

@@ -0,0 +1,243 @@
package distributor
import (
"context"
"errors"
"fmt"
"strings"
"testing"
"time"
distributorupload "gitea.maximumdirect.net/eric/distributor/pkg/upload"
"gitea.maximumdirect.net/eric/weatherreporter/internal/config"
)
func TestUploadUsesConfiguredClientAndSingleFile(t *testing.T) {
cfg := config.Defaults().Notify.Distributor
cfg.Endpoint = "https://distributor.example.test"
cfg.TokenEnv = "DISTRIBUTOR_UPLOAD_TOKEN"
cfg.Timeout = 15 * time.Second
t.Setenv(cfg.TokenEnv, "secret-token")
factory := &fakeUploadFactory{
client: &fakeUploadClient{
result: uploadFilesResult{RunID: "run-123", Status: "accepted"},
},
}
client := newClient(cfg, factory.newClient)
result, err := client.Upload(context.Background(), UploadRequest{
BundleID: "weatherreporter.home.daily.run",
IdempotencyKey: "weatherreporter.home.daily.run",
SourcePath: "/tmp/report.md",
BundlePath: "daily.md",
})
if err != nil {
t.Fatalf("Upload() error = %v", err)
}
if result.RunID != "run-123" || result.Status != "accepted" {
t.Fatalf("result = %#v, want accepted run", result)
}
if factory.endpoint != cfg.Endpoint {
t.Fatalf("factory endpoint = %q, want %q", factory.endpoint, cfg.Endpoint)
}
if factory.token != "secret-token" {
t.Fatalf("factory token = %q, want secret-token", factory.token)
}
if factory.timeout != 15*time.Second {
t.Fatalf("factory timeout = %s, want 15s", factory.timeout)
}
got := factory.client.opts
if got.BundleID != "weatherreporter.home.daily.run" {
t.Fatalf("BundleID = %q, want weatherreporter.home.daily.run", got.BundleID)
}
if got.IdempotencyKey != "weatherreporter.home.daily.run" {
t.Fatalf("IdempotencyKey = %q, want weatherreporter.home.daily.run", got.IdempotencyKey)
}
if got.SourcePath != "/tmp/report.md" {
t.Fatalf("SourcePath = %q, want /tmp/report.md", got.SourcePath)
}
if got.BundlePath != "daily.md" {
t.Fatalf("BundlePath = %q, want daily.md", got.BundlePath)
}
}
func TestUploadRejectsMissingInputs(t *testing.T) {
cfg := config.Defaults().Notify.Distributor
t.Setenv(cfg.TokenEnv, "secret-token")
tests := []struct {
name string
mutate func(*Client, *UploadRequest)
wantErr string
}{
{
name: "Token",
mutate: func(c *Client, req *UploadRequest) {
t.Setenv(c.TokenEnv, "")
},
wantErr: "token environment variable",
},
{
name: "SourcePath",
mutate: func(c *Client, req *UploadRequest) {
req.SourcePath = ""
},
wantErr: "source path is required",
},
{
name: "BundlePath",
mutate: func(c *Client, req *UploadRequest) {
req.BundlePath = ""
},
wantErr: "bundle path is required",
},
{
name: "UploadClientFactory",
mutate: func(c *Client, req *UploadRequest) {
c.newUploadClient = nil
},
wantErr: "upload client factory is required",
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
t.Setenv(cfg.TokenEnv, "secret-token")
client := newClient(cfg, (&fakeUploadFactory{client: &fakeUploadClient{}}).newClient)
req := validUploadRequest()
tt.mutate(client, &req)
_, err := client.Upload(context.Background(), req)
if err == nil {
t.Fatal("Upload() error = nil, want error")
}
if !strings.Contains(err.Error(), tt.wantErr) {
t.Fatalf("error = %q, want %q", err.Error(), tt.wantErr)
}
if strings.Contains(err.Error(), "secret-token") {
t.Fatalf("error = %q, want no token value", err.Error())
}
})
}
}
func TestUploadWrapsFactoryErrorWithoutToken(t *testing.T) {
cfg := config.Defaults().Notify.Distributor
cfg.Endpoint = "https://distributor.example.test"
t.Setenv(cfg.TokenEnv, "secret-token")
factory := &fakeUploadFactory{
err: fmt.Errorf("factory failed with secret-token"),
}
client := newClient(cfg, factory.newClient)
_, err := client.Upload(context.Background(), validUploadRequest())
if err == nil {
t.Fatal("Upload() error = nil, want error")
}
if strings.Contains(err.Error(), "secret-token") {
t.Fatalf("error = %q, want no token value", err.Error())
}
if !strings.Contains(err.Error(), cfg.Endpoint) {
t.Fatalf("error = %q, want endpoint context", err.Error())
}
}
func TestUploadWrapsUploadFailureWithContextWithoutToken(t *testing.T) {
cfg := config.Defaults().Notify.Distributor
cfg.Endpoint = "https://distributor.example.test"
t.Setenv(cfg.TokenEnv, "secret-token")
factory := &fakeUploadFactory{
client: &fakeUploadClient{err: fmt.Errorf("server rejected secret-token")},
}
client := newClient(cfg, factory.newClient)
req := validUploadRequest()
_, err := client.Upload(context.Background(), req)
if err == nil {
t.Fatal("Upload() error = nil, want error")
}
for _, want := range []string{cfg.Endpoint, req.BundleID, req.IdempotencyKey, req.SourcePath, req.BundlePath} {
if !strings.Contains(err.Error(), want) {
t.Fatalf("error = %q, want context %q", err.Error(), want)
}
}
if strings.Contains(err.Error(), "secret-token") {
t.Fatalf("error = %q, want no token value", err.Error())
}
}
func TestUploadPreservesIdempotencyConflictDiagnosis(t *testing.T) {
cfg := config.Defaults().Notify.Distributor
cfg.Endpoint = "https://distributor.example.test"
t.Setenv(cfg.TokenEnv, "secret-token")
factory := &fakeUploadFactory{
client: &fakeUploadClient{
err: &distributorupload.IdempotencyConflictError{
HTTPError: distributorupload.HTTPError{
StatusCode: 409,
Status: "409 Conflict",
Message: "conflicting upload for secret-token",
},
},
},
}
client := newClient(cfg, factory.newClient)
_, err := client.Upload(context.Background(), validUploadRequest())
if err == nil {
t.Fatal("Upload() error = nil, want error")
}
var conflict *IdempotencyConflictError
if !errors.As(err, &conflict) {
t.Fatalf("Upload() error = %T %v, want IdempotencyConflictError", err, err)
}
if !strings.Contains(err.Error(), "idempotency conflict") {
t.Fatalf("error = %q, want idempotency conflict diagnosis", err.Error())
}
if strings.Contains(err.Error(), "secret-token") {
t.Fatalf("error = %q, want no token value", err.Error())
}
}
func validUploadRequest() UploadRequest {
return UploadRequest{
BundleID: "weatherreporter.home.daily.run",
IdempotencyKey: "weatherreporter.home.daily.run",
SourcePath: "/tmp/report.md",
BundlePath: "daily.md",
}
}
type fakeUploadFactory struct {
endpoint string
token string
timeout time.Duration
client *fakeUploadClient
err error
}
func (f *fakeUploadFactory) newClient(endpoint, token string, timeout time.Duration) (uploadClient, error) {
f.endpoint = endpoint
f.token = token
f.timeout = timeout
if f.err != nil {
return nil, f.err
}
return f.client, nil
}
type fakeUploadClient struct {
opts uploadFilesOptions
result uploadFilesResult
err error
}
func (c *fakeUploadClient) UploadFiles(ctx context.Context, opts uploadFilesOptions) (uploadFilesResult, error) {
c.opts = opts
if c.err != nil {
return uploadFilesResult{}, c.err
}
return c.result, nil
}