// Package distributor adapts weatherreporter report artifacts to distributor uploads. package distributor import ( "context" "encoding/json" "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 CreatedAt time.Time } type UploadResult struct { RunID string Status string UploadStatus string StatusError string RunStatus *RunStatus } type RunStatus struct { RunID string PipelineID string Status string AcceptedAt time.Time StartedAt *time.Time FinishedAt *time.Time Report json.RawMessage Error 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) Status(ctx context.Context, runID string) (runStatus, error) } type uploadFilesOptions struct { BundleID string IdempotencyKey string SourcePath string BundlePath string CreatedAt time.Time } type uploadFilesResult struct { RunID string Status string } type runStatus struct { RunID string PipelineID string Status string AcceptedAt time.Time StartedAt *time.Time FinishedAt *time.Time Report json.RawMessage Error string } const statusPollInterval = 250 * time.Millisecond 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, CreatedAt: req.CreatedAt, }) 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, }) } uploadResult := UploadResult{ RunID: result.RunID, Status: result.Status, UploadStatus: result.Status, } status, statusErr := waitForRunStatus(runCtx, uploadClient, result.RunID, c.Timeout > 0) if status.RunID != "" || status.Status != "" { uploadResult.RunStatus = &RunStatus{ RunID: status.RunID, PipelineID: status.PipelineID, Status: status.Status, AcceptedAt: status.AcceptedAt, StartedAt: status.StartedAt, FinishedAt: status.FinishedAt, Report: append(json.RawMessage(nil), status.Report...), Error: redactTokenString(status.Error, token), } if status.Status != "" { uploadResult.Status = status.Status } } if statusErr != nil { uploadResult.StatusError = redactTokenString(statusErr.Error(), token) return uploadResult, nil } if status.Status == "failed" { return uploadResult, fmt.Errorf("distributor run %q failed: %s", status.RunID, uploadResult.RunStatus.Error) } return uploadResult, nil } func waitForRunStatus(ctx context.Context, client uploadClient, runID string, poll bool) (runStatus, error) { status, err := client.Status(ctx, runID) if err != nil || terminalRunStatus(status.Status) || !poll { return status, err } for { timer := time.NewTimer(statusPollInterval) select { case <-ctx.Done(): timer.Stop() return status, fmt.Errorf("distributor run %q did not reach terminal status before timeout: %w", runID, ctx.Err()) case <-timer.C: } next, err := client.Status(ctx, runID) if err != nil { return status, err } status = next if terminalRunStatus(status.Status) { return status, nil } } } func terminalRunStatus(status string) bool { return status == "succeeded" || status == "failed" } 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, Created: opts.CreatedAt, 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 } func (c distributorUploadClient) Status(ctx context.Context, runID string) (runStatus, error) { status, err := c.client.Status(ctx, runID) if err != nil { return runStatus{}, err } return runStatus{ RunID: status.RunID, PipelineID: status.PipelineID, Status: status.Status, AcceptedAt: status.AcceptedAt, StartedAt: status.StartedAt, FinishedAt: status.FinishedAt, Report: append(json.RawMessage(nil), status.Report...), Error: status.Error, }, 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]") }